Archilyzer · Source

archilyzer

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

commit 27ea3025161a1c577b69d03202a32b38a640f564
parent d93cd7ec27ac5431bbe651a4f56dc1689124b856
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sat, 12 Sep 2026 02:06:43 -0400

common: the archive reader moves into the core, with node isolated to one file

`lib/archive/reader.ts` is mcp/src/source.ts's reader, moved: the
`ArchiveReader` interface (mcp's `ShardSource`, all 18 members, same names and
signatures), `RemoteSource`, `PageCache`, `PromiseMap`, and the four tolerant
folds (`buildVideoIndex`, `buildStatsIndex`, `buildDuplicateIndex`,
`parseGroupsManifest`). `reader-fs.ts` holds `LocalSource` and is the ONLY file
here that imports node:fs / node:path; `reader-hub.ts` holds `HubSource`.
reader.ts never value-imports reader-fs.ts — not even a type — which is what
lets the export bundle reach the reader without dragging node:fs behind it.
`next build` is the test of that, not tsc.

Two deliberate non-changes. There is still NO `record(layer, slug, id)`: the
walk is manifest -> shard -> record and a caller that wants one record already
holds the shard, so a per-record fetch would be N reads where the plan needs
one. And every URL now comes from `contract.ts`'s builders rather than a second
interpolation of the same literal — `manifestUrl("subs", slug, base)` is
character-for-character what it replaced, including the thrown error's full URL
and each `recordRead` kind.

`process` appears twice in the browser-safe files and both are guarded: the
page-cache budget knob reads `process.env?.[…]` behind a `typeof process` check
(a dynamic form Next does not inline, which is right for a server-only knob —
a browser reader takes the 48 MB default). `TRANSCRIPT_PLATFORM_LINKS` stayed in
reader-fs.ts, where it was already node-only.

`reader.test.ts` drives the REAL RemoteSource against a Map-backed `fetch`, so
every claim is countable rather than asserted against a stub that agrees with
it: the corpus -> manifest -> slugToPage -> page walk; two concurrent callers
sharing one in-flight read (the fetch is held open, and the count is checked
BEFORE it resolves); a rejection never memoised; the cached 404-as-null for
digests, subs, posts, duplicates and aliases; and the page cache's byte-budget
LRU, including that a touch really does move an entry (the eviction victim is
the least-recently-USED, not the oldest inserted) and that MIN_CACHED_PAGES
keeps a page four times the budget resident rather than thrashing on it.

Gates: tsc --noEmit clean in all six packages; common 1077/1077 (1060 + 17);
mcp 205/205 unchanged.

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

Diffstat:
Acommon/lib/archive/reader-fs.ts | 335+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/archive/reader-hub.ts | 316+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/archive/reader.test.ts | 454+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/archive/reader.ts | 779+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
4 files changed, 1884 insertions(+), 0 deletions(-)

diff --git a/common/lib/archive/reader-fs.ts b/common/lib/archive/reader-fs.ts @@ -0,0 +1,335 @@ +// ─── Local: read composed shards from a directory on disk ─── +// +// THE ONLY FILE UNDER lib/archive/ THAT TOUCHES node:*. reader.ts imports +// nothing from here (not even a type), so a browser bundle that reaches the +// reader cannot drag `node:fs` in behind it. `pnpm --filter export exec next +// build` is what actually proves that — tsc cannot tell a `node:fs` import from +// any other, and the export app is the bundle where it would blow up. +// +// `dir` is a COMPOSED PUBLIC DIR — the output of compose-site.ts, laid out +// exactly as the deployed origin serves it. It is not a corpus on disk, so +// there is no `config.dataDir` to resolve here: the relocation resolver the +// one-core plan asks for belongs to the thing that walks +// `channels/<slug>/data/` (umtool's cues.mjs), not to this. + +import { readFile, readdir } from "node:fs/promises"; +import path from "node:path"; +import type { + ChannelTranscriptsManifest, + ChannelSubsManifest, + Manifest, +} from "../manifest"; +import type { TranscriptDetail, DisplaySummary } from "../transcripts"; +import type { SubsDetail } from "../subs"; +import type { ChannelPostsManifest, Post } from "../posts"; +import { coerceAliasConfig, type SearchAlias } from "../searchAliases"; +import { DUPLICATES_FILENAME, type DuplicateReport } from "../duplicates"; +import type { ChannelDigestsManifest, VideoDigest } from "../digests"; +import type { StatsManifest, VideoStat } from "../stats"; +import { pageFileName } from "./contract"; +import { recordRead } from "./io-stats"; +import { + DEFAULT_PAGE_CONCURRENCY, + EMPTY_GROUPS, + PageCache, + PromiseMap, + buildDuplicateIndex, + buildStatsIndex, + buildVideoIndex, + pageCacheBudgetBytes, + parseGroupsManifest, + type ArchiveReader, + type ChannelGroups, + type ChannelRef, + type DuplicateIndex, + type IndexedVideo, + type SiteCorpusJson, + type VideoAvailability, + type VideoIndex, +} from "./reader"; + +// Read a local JSON file, counting its bytes when instrumentation is on, and +// reporting the raw size so a byte-budgeted cache can account for it. +async function readLocalJsonSized<T>( + file: string, + kind: string, +): Promise<{ value: T; bytes: number }> { + const raw = await readFile(file, "utf8"); + recordRead(kind, raw.length); + return { value: JSON.parse(raw) as T, bytes: raw.length }; +} + +async function readLocalJson<T>(file: string, kind: string): Promise<T> { + return (await readLocalJsonSized<T>(file, kind)).value; +} + +// Opt-out for the composed site's declared origin (see LocalSource.publicOrigin). +// Set TRANSCRIPT_PLATFORM_LINKS=1 to cite platform watch pages instead, which is +// the right answer when a local build's declared site url is not actually +// deployed. +const PREFER_PLATFORM_LINKS = process.env.TRANSCRIPT_PLATFORM_LINKS === "1"; + +// Prefers corpus.json for the channel list (names + counts); falls back to +// listing the transcripts/ subdirectories so it works even pre-Layer-1. +export class LocalSource implements ArchiveReader { + readonly label: string; + private aliases?: SearchAlias[]; + private groups?: ChannelGroups; + private index?: Promise<Map<string, IndexedVideo>>; + private duplicates?: Promise<DuplicateIndex>; + private stats?: Promise<ReadonlyMap<string, VideoStat>>; + // The composed site's own declared origin, learned from corpus.json the first + // time the channel list is read. Undefined = not looked at yet. + private siteOrigin: string | null | undefined; + + constructor(private dir: string) { + this.label = `local:${dir}`; + } + + // The deployed archilyzer viewer these shards were composed for, as declared + // by the dir's own corpus.json (`site.url`). A composed public dir is not an + // anonymous pile of JSON — it names the site it is the build output of — so + // citing that viewer is both possible and the right default: a reader + // following a citation lands in the archive, at the cited second, with the + // transcript around it, rather than on the platform page where the archive's + // whole point (that we still have a copy) is invisible. + // + // Populated by readChannels(), which every read path runs before it renders a + // link. Null when the dir ships no corpus.json (the bare directory-listing + // fallback), or when TRANSCRIPT_PLATFORM_LINKS=1 asks for platform links — + // both fall back to the platform watch page exactly as before. + publicOrigin(): string | null { + return this.siteOrigin ?? null; + } + + // Cached INCLUDING the negative answer, exactly like postsManifests below: + // live-chat search probes every channel in scope, and a video-only channel + // would otherwise cost one failed read per query. + private subsManifests = new PromiseMap<ChannelSubsManifest | null>(); + private digestManifests = new PromiseMap<ChannelDigestsManifest | null>(); + private subsPages = new PageCache<SubsDetail[]>(pageCacheBudgetBytes()); + private postsManifests = new PromiseMap<ChannelPostsManifest | null>(); + private transcriptManifests = new PromiseMap<ChannelTranscriptsManifest>(); + private transcriptPages = new PageCache<TranscriptDetail[]>(pageCacheBudgetBytes()); + + // Drop every cached read. Reached only through listChannels({refresh:true}) — + // the deliberate, explicit staleness escape hatch for a corpus rebuilt under + // a long-lived server. Not a TTL, on purpose. + private resetCaches(): void { + this.subsManifests.clear(); + this.digestManifests.clear(); + this.subsPages.clear(); + this.postsManifests.clear(); + this.transcriptManifests.clear(); + this.transcriptPages.clear(); + this.aliases = undefined; + this.groups = undefined; + this.index = undefined; + this.duplicates = undefined; + this.stats = undefined; + this.siteOrigin = undefined; + } + + subsManifest(ch: ChannelRef): Promise<ChannelSubsManifest | null> { + return this.subsManifests.take(ch.slug, () => + readLocalJson<ChannelSubsManifest>( + path.join(this.dir, "subs", ch.slug, "manifest.json"), + "subsManifest", + ).catch(() => null), // channel ships no subs shards + ); + } + + subsPage(ch: ChannelRef, page: number): Promise<SubsDetail[]> { + return this.subsPages.take(`${ch.slug}:${page}`, () => + readLocalJsonSized<SubsDetail[]>( + path.join(this.dir, "subs", ch.slug, pageFileName(page)), + "subsPage", + ), + ); + } + + // Cached per channel INCLUDING the negative answer: most channels are + // video-only, and a posts-covering search would otherwise re-probe every one + // of them on every query. + postsManifest(ch: ChannelRef): Promise<ChannelPostsManifest | null> { + return this.postsManifests.take(ch.slug, () => + readLocalJson<ChannelPostsManifest>( + path.join(this.dir, "posts", ch.slug, "manifest.json"), + "postsManifest", + ).catch(() => null), // channel ships no posts shards + ); + } + + async postsPage(ch: ChannelRef, page: number): Promise<Post[]> { + return readLocalJson<Post[]>( + path.join(this.dir, "posts", ch.slug, pageFileName(page)), + "postsPage", + ); + } + + async loadAliases(): Promise<SearchAlias[]> { + if (this.aliases) return this.aliases; + try { + this.aliases = coerceAliasConfig( + await readLocalJson(path.join(this.dir, "search-aliases.json"), "aliases"), + ).aliases; + } catch { + this.aliases = []; // no/invalid file — search stays plain + } + return this.aliases; + } + + async loadGroups(): Promise<ChannelGroups> { + if (this.groups) return this.groups; + try { + this.groups = parseGroupsManifest( + await readLocalJson( + path.join(this.dir, "summaries", "manifest.json"), + "summariesManifest", + ), + ); + } catch { + this.groups = EMPTY_GROUPS; // no/invalid manifest — groups off + } + return this.groups; + } + + private channelList?: Promise<ChannelRef[]>; + + listChannels(opts: { refresh?: boolean } = {}): Promise<ChannelRef[]> { + if (opts.refresh) { + this.channelList = undefined; + this.resetCaches(); + } + this.channelList ??= this.readChannels().catch((e: unknown) => { + this.channelList = undefined; // don't memoise a failure + throw e; + }); + return this.channelList; + } + + private async readChannels(): Promise<ChannelRef[]> { + try { + const corpus = await readLocalJson<SiteCorpusJson>( + path.join(this.dir, "corpus.json"), + "corpus", + ); + const declared = corpus.site?.url?.trim(); + this.siteOrigin = + declared && !PREFER_PLATFORM_LINKS ? declared.replace(/\/+$/, "") : null; + if (Array.isArray(corpus.channels) && corpus.channels.length > 0) { + return corpus.channels.map((c) => ({ + key: c.slug, + slug: c.slug, + name: c.name ?? c.slug, + videoCount: c.videoCount, + groupId: c.groupId, + })); + } + } catch { + // no corpus.json — fall back to a directory listing + } + const transcriptsDir = path.join(this.dir, "transcripts"); + let entries: string[] = []; + try { + entries = await readdir(transcriptsDir); + } catch { + return []; + } + const channels: ChannelRef[] = []; + for (const slug of entries.sort()) { + // A channel dir has a manifest.json; skip stray files. + try { + await readFile(path.join(transcriptsDir, slug, "manifest.json"), "utf8"); + channels.push({ key: slug, slug, name: slug }); + } catch { + // not a channel dir + } + } + return channels; + } + + transcriptsManifest(ch: ChannelRef): Promise<ChannelTranscriptsManifest> { + return this.transcriptManifests.take(ch.slug, () => + readLocalJson<ChannelTranscriptsManifest>( + path.join(this.dir, "transcripts", ch.slug, "manifest.json"), + "transcriptsManifest", + ), + ); + } + + transcriptPage(ch: ChannelRef, page: number): Promise<TranscriptDetail[]> { + return this.transcriptPages.take(`${ch.slug}:${page}`, () => + readLocalJsonSized<TranscriptDetail[]>( + path.join(this.dir, "transcripts", ch.slug, pageFileName(page)), + "transcriptPage", + ), + ); + } + + // Local reads are CPU-bound on JSON.parse, so a modest window is all that is + // available to win: it overlaps the next page's read with this page's parse. + readonly pageConcurrency = DEFAULT_PAGE_CONCURRENCY; + + videoIndex(): Promise<VideoIndex> { + this.index ??= buildVideoIndex( + () => + readLocalJson<Manifest>( + path.join(this.dir, "summaries", "manifest.json"), + "summariesManifest", + ), + (page) => + readLocalJson<DisplaySummary[]>( + path.join(this.dir, "summaries", pageFileName(page)), + "summariesPage", + ), + ); + return this.index; + } + + availabilityMap(): Promise<ReadonlyMap<string, VideoAvailability>> { + return this.videoIndex(); + } + + digestsManifest(ch: ChannelRef): Promise<ChannelDigestsManifest | null> { + return this.digestManifests.take(ch.slug, () => + readLocalJson<ChannelDigestsManifest>( + path.join(this.dir, "digests", ch.slug, "manifest.json"), + "digestsManifest", + ).catch(() => null), // channel has no digests + ); + } + + digestPage(ch: ChannelRef, page: number): Promise<VideoDigest[]> { + return readLocalJson<VideoDigest[]>( + path.join(this.dir, "digests", ch.slug, pageFileName(page)), + "digestPage", + ); + } + + duplicateIndex(): Promise<DuplicateIndex> { + this.duplicates ??= readLocalJson<DuplicateReport>( + path.join(this.dir, DUPLICATES_FILENAME), + "duplicates", + ) + .then(buildDuplicateIndex) + .catch(() => buildDuplicateIndex(null)); // no report shipped — no clusters + return this.duplicates; + } + + statsIndex(): Promise<ReadonlyMap<string, VideoStat>> { + this.stats ??= buildStatsIndex( + () => + readLocalJson<StatsManifest>( + path.join(this.dir, "stats", "manifest.json"), + "statsManifest", + ), + (page) => + readLocalJson<VideoStat[]>( + path.join(this.dir, "stats", pageFileName(page)), + "statsPage", + ), + ); + return this.stats; + } +} diff --git a/common/lib/archive/reader-hub.ts b/common/lib/archive/reader-hub.ts @@ -0,0 +1,316 @@ +// ─── Hub: federate over every member site listed in the hub corpus.json ─── +// +// Each member is its own RemoteSource; channels are namespaced by site so keys +// stay unique, and manifest/page calls dispatch to the owning member. Browser- +// safe for the same reason reader.ts is: every read is `fetch`. + +import type { ChannelTranscriptsManifest, ChannelSubsManifest } from "../manifest"; +import type { TranscriptDetail } from "../transcripts"; +import type { SubsDetail } from "../subs"; +import type { ChannelPostsManifest, Post } from "../posts"; +import { coerceAliasConfig, type SearchAlias } from "../searchAliases"; +import type { ChannelDigestsManifest, VideoDigest } from "../digests"; +import type { VideoStat } from "../stats"; +import { corpusUrl, rootFileUrl, type HubSite } from "./contract"; +import { + EMPTY_GROUPS, + RemoteSource, + pageCacheBudgetBytes, + type ArchiveReader, + type ChannelGroups, + type ChannelRef, + type DuplicateIndex, + type HubCorpusJson, + type IndexedVideo, + type VideoAvailability, + type VideoIndex, +} from "./reader"; + +export class HubSource implements ArchiveReader { + readonly label: string; + readonly hubBase: string; + private members = new Map<string, RemoteSource>(); // siteId -> source + private aliases?: SearchAlias[]; + private index?: Promise<Map<string, IndexedVideo>>; + private duplicates?: Promise<DuplicateIndex>; + private stats?: Promise<ReadonlyMap<string, VideoStat>>; + private sites?: Promise<HubSite[]>; + + // A hub is N HTTP origins, so the latency argument for a wide window applies + // even harder than for a single remote. + readonly pageConcurrency = 8; + // Optional subset allowlist of member siteIds. Undefined = federate every + // member; a set restricts listChannels() to those members (site discovery via + // listSites() stays unfiltered so a picker can still see all members). + private allowSiteIds?: Set<string>; + + constructor(hubUrl: string, allowSiteIds?: string[]) { + this.hubBase = hubUrl.replace(/\/+$/, ""); + this.allowSiteIds = + allowSiteIds && allowSiteIds.length > 0 + ? new Set(allowSiteIds) + : undefined; + this.label = this.allowSiteIds + ? `hub:${this.hubBase} (${this.allowSiteIds.size} site(s))` + : `hub:${this.hubBase}`; + } + + // Fetch the hub's corpus.json and return its member sites — UNFILTERED (the + // full membership), even when this source is scoped to a subset, so a picker + // (list_sources / use_source) can show every member. + // + // Memoised on the promise: readChannels() and videoIndex() both need it, so + // an un-memoised version fetched the hub roster twice on a cold hub search. + // A failure is not memoised. + listSites(): Promise<HubSite[]> { + this.sites ??= this.readSites().catch((e: unknown) => { + this.sites = undefined; + throw e; + }); + return this.sites; + } + + private async readSites(): Promise<HubSite[]> { + const url = corpusUrl(this.hubBase); + const res = await fetch(url); + if (!res.ok) { + throw new Error(`GET ${url} -> ${res.status} ${res.statusText}`); + } + const hub = (await res.json()) as HubCorpusJson; + return hub.sites ?? []; + } + + // One page-cache budget for the whole hub, divided across its members — N + // members must not each get the full ceiling. Floored so a large federation + // still caches something per member. + private memberBudget(memberCount: number): number { + const total = pageCacheBudgetBytes(); + const floor = 8 * 1024 * 1024; + return Math.max(floor, Math.floor(total / Math.max(1, memberCount))); + } + + // A hub can ship its own /search-aliases.json (the merged federation-wide + // dictionary); if it doesn't, aliases are simply off for hub-wide search. + async loadAliases(): Promise<SearchAlias[]> { + if (this.aliases) return this.aliases; + try { + const res = await fetch(rootFileUrl("search-aliases.json", this.hubBase)); + this.aliases = res.ok + ? coerceAliasConfig(await res.json()).aliases + : []; + } catch { + this.aliases = []; + } + return this.aliases; + } + + // In hub mode each group is itself a federated member site (a different model + // — accent-per-origin), so hub-wide channel-group tokens are deferred: return + // the empty fallback. Multi-channel scoping still works (member groupIds carry + // through listChannels, they just don't resolve against hub-level groups). + async loadGroups(): Promise<ChannelGroups> { + return EMPTY_GROUPS; + } + + private memberFor(siteId: string): RemoteSource { + const m = this.members.get(siteId); + if (!m) throw new Error(`unknown hub member site: ${siteId}`); + return m; + } + + private channelList?: Promise<ChannelRef[]>; + + listChannels(opts: { refresh?: boolean } = {}): Promise<ChannelRef[]> { + if (opts.refresh) { + this.channelList = undefined; + this.sites = undefined; + this.index = undefined; + this.duplicates = undefined; + this.stats = undefined; + this.aliases = undefined; + // Members hold their own manifest/page caches; drop them wholesale so a + // refresh means the same thing federation-wide as it does locally. + this.members.clear(); + } + this.channelList ??= this.readChannels().catch((e: unknown) => { + this.channelList = undefined; // don't memoise a failure + throw e; + }); + return this.channelList; + } + + // Populating `members` must stay INSIDE the memoised call: memberFor() + // depends on it, so a memo that skipped this would leave every + // transcriptPage/postsPage dispatch throwing "unknown hub member site". + private async readChannels(): Promise<ChannelRef[]> { + const sites = (await this.listSites()).filter( + (s) => !this.allowSiteIds || this.allowSiteIds.has(s.siteId), + ); + const budget = this.memberBudget(sites.length); + const all: ChannelRef[] = []; + // Sequential member fetches keep it simple and polite; the channel count is + // small. A failing member is skipped rather than failing the whole list. + for (const site of sites) { + const remote = + this.members.get(site.siteId) ?? new RemoteSource(site.url, budget); + this.members.set(site.siteId, remote); + try { + const channels = await remote.listChannels(); + for (const c of channels) { + all.push({ + ...c, + key: `${site.siteId}/${c.slug}`, + siteId: site.siteId, + siteTitle: site.title, + siteUrl: site.url, + }); + } + } catch { + // skip an unreachable member + } + } + return all; + } + + transcriptsManifest(ch: ChannelRef): Promise<ChannelTranscriptsManifest> { + if (!ch.siteId) throw new Error("hub channel ref missing siteId"); + return this.memberFor(ch.siteId).transcriptsManifest(ch); + } + + transcriptPage(ch: ChannelRef, page: number): Promise<TranscriptDetail[]> { + if (!ch.siteId) throw new Error("hub channel ref missing siteId"); + return this.memberFor(ch.siteId).transcriptPage(ch, page); + } + + // A hub has no single viewer origin — each video's origin is its member + // site's url (carried on the ChannelRef as `siteUrl`), which momentUrl prefers + // per-video. Return null so we never mint a wrong-origin viewer link. + publicOrigin(): string | null { + return null; + } + + async subsManifest(ch: ChannelRef): Promise<ChannelSubsManifest | null> { + if (!ch.siteId) return null; + try { + return await this.memberFor(ch.siteId).subsManifest(ch); + } catch { + return null; // member not yet registered / unreachable + } + } + + subsPage(ch: ChannelRef, page: number): Promise<SubsDetail[]> { + if (!ch.siteId) throw new Error("hub channel ref missing siteId"); + return this.memberFor(ch.siteId).subsPage(ch, page); + } + + async postsManifest(ch: ChannelRef): Promise<ChannelPostsManifest | null> { + if (!ch.siteId) return null; + try { + return await this.memberFor(ch.siteId).postsManifest(ch); + } catch { + return null; // member not yet registered / unreachable + } + } + + postsPage(ch: ChannelRef, page: number): Promise<Post[]> { + if (!ch.siteId) throw new Error("hub channel ref missing siteId"); + return this.memberFor(ch.siteId).postsPage(ch, page); + } + + async digestsManifest(ch: ChannelRef): Promise<ChannelDigestsManifest | null> { + if (!ch.siteId) return null; + try { + return await this.memberFor(ch.siteId).digestsManifest(ch); + } catch { + return null; // member not yet registered / unreachable + } + } + + digestPage(ch: ChannelRef, page: number): Promise<VideoDigest[]> { + if (!ch.siteId) throw new Error("hub channel ref missing siteId"); + return this.memberFor(ch.siteId).digestPage(ch, page); + } + + // Merge each member's video index. Keys are member-local slugs + // (`<channelSlug>/<id>`) — the same slug a member's transcript page records + // carry — so a per-record lookup joins correctly. Built lazily/cached. + videoIndex(): Promise<VideoIndex> { + this.index ??= this.buildMergedIndex(); + return this.index; + } + + private async buildMergedIndex(): Promise<Map<string, IndexedVideo>> { + const merged = new Map<string, IndexedVideo>(); + let sites: HubSite[]; + try { + sites = (await this.listSites()).filter( + (s) => !this.allowSiteIds || this.allowSiteIds.has(s.siteId), + ); + } catch { + return merged; + } + const budget = this.memberBudget(sites.length); + for (const site of sites) { + const remote = + this.members.get(site.siteId) ?? new RemoteSource(site.url, budget); + this.members.set(site.siteId, remote); + try { + for (const [slug, rec] of await remote.videoIndex()) { + merged.set(slug, rec); + } + } catch { + // skip an unreachable member + } + } + return merged; + } + + availabilityMap(): Promise<ReadonlyMap<string, VideoAvailability>> { + return this.videoIndex(); + } + + // Duplicate clusters are detected WITHIN a site, so federating them is a + // merge of per-member maps and nothing more — this deliberately does not try + // to detect mirrors ACROSS member sites. Two sites holding the same recording + // is a real thing, but nothing has compared their transcripts, and inventing + // a cross-site cluster here would be asserting a duplicate no detector ever + // confirmed. + duplicateIndex(): Promise<DuplicateIndex> { + this.duplicates ??= this.mergeMembers((m) => m.duplicateIndex()); + return this.duplicates; + } + + statsIndex(): Promise<ReadonlyMap<string, VideoStat>> { + this.stats ??= this.mergeMembers((m) => m.statsIndex()); + return this.stats; + } + + // Merge one lazily-read layer across every member site, keyed by the + // member-local slug — the same shape and the same tolerance as + // buildMergedIndex (an unreachable member is skipped, not fatal). + private async mergeMembers<T>( + read: (m: RemoteSource) => Promise<ReadonlyMap<string, T>>, + ): Promise<Map<string, T>> { + const merged = new Map<string, T>(); + let sites: HubSite[]; + try { + sites = (await this.listSites()).filter( + (s) => !this.allowSiteIds || this.allowSiteIds.has(s.siteId), + ); + } catch { + return merged; + } + const budget = this.memberBudget(sites.length); + for (const site of sites) { + const remote = + this.members.get(site.siteId) ?? new RemoteSource(site.url, budget); + this.members.set(site.siteId, remote); + try { + for (const [k, v] of await read(remote)) merged.set(k, v); + } catch { + // skip an unreachable member + } + } + return merged; + } +} diff --git a/common/lib/archive/reader.test.ts b/common/lib/archive/reader.test.ts @@ -0,0 +1,454 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import type { ChannelTranscriptsManifest } from "../manifest"; +import type { TranscriptDetail } from "../transcripts"; +import type { ChannelDigestsManifest, VideoDigest } from "../digests"; +import { manifestUrl, pageUrl } from "./contract"; +import { + PageCache, + PromiseMap, + RemoteSource, + type ArchiveReader, + type ChannelRef, +} from "./reader"; + +// ─── an in-memory archive ─── +// +// A Map from URL to JSON body, handed to RemoteSource as `fetch`. That makes +// the thing under test the REAL reader (its caches, its walk, its tolerance) +// rather than a stub that agrees with it — and it makes every assertion below +// about fetches countable, which is the only way to state "one read, not N". + +const ORIGIN = "https://fixture.example"; + +type Archive = { + fetches: string[]; + body: Map<string, unknown>; + // Resolvers for URLs deliberately held open, so two concurrent callers can be + // observed sharing one in-flight read. + gate?: { url: string; release: () => void }; +}; + +function makeArchive(entries: Record<string, unknown>): Archive { + const a: Archive = { fetches: [], body: new Map(Object.entries(entries)) }; + return a; +} + +function installFetch(a: Archive, hold?: string): () => void { + const real = globalThis.fetch; + let release: (() => void) | undefined; + globalThis.fetch = (async (input: RequestInfo | URL) => { + const url = String(input); + a.fetches.push(url); + if (!a.body.has(url)) { + return new Response("not found", { status: 404, statusText: "Not Found" }); + } + const text = JSON.stringify(a.body.get(url)); + if (hold !== undefined && url === hold) { + await new Promise<void>((r) => { + release = r; + }); + } + return new Response(text, { status: 200 }); + }) as typeof fetch; + a.gate = { url: hold ?? "", release: () => release?.() }; + return () => { + globalThis.fetch = real; + }; +} + +function pageRecords(slug: string, ids: string[]): TranscriptDetail[] { + return ids.map( + (id) => + ({ + id, + slug: `${slug}/${id}`, + title: id, + uploadDate: "20240101", + duration: 60, + channel: slug, + channelSlug: slug, + cues: [{ start: 0, end: 1, text: "hello" }], + }) as unknown as TranscriptDetail, + ); +} + +const CORPUS = { + spec: 3, + kind: "site", + site: { id: "fixture", title: "Fixture", url: ORIGIN }, + channels: [{ slug: "alpha", name: "Alpha", videoCount: 3 }], +}; + +const ALPHA_MANIFEST: ChannelTranscriptsManifest = { + version: 1, + channelSlug: "alpha", + pageCount: 2, + maxPageBytes: 1000, + generatedAt: "2026-09-12T00:00:00.000Z", + slugToPage: { "alpha/a1": 0, "alpha/a2": 0, "alpha/a3": 1 }, +}; + +function fullArchive(): Archive { + return makeArchive({ + [`${ORIGIN}/corpus.json`]: CORPUS, + [manifestUrl("transcripts", "alpha", ORIGIN)]: ALPHA_MANIFEST, + [pageUrl("transcripts", "alpha", 0, ORIGIN)]: pageRecords("alpha", ["a1", "a2"]), + [pageUrl("transcripts", "alpha", 1, ORIGIN)]: pageRecords("alpha", ["a3"]), + }); +} + +const ALPHA: ChannelRef = { key: "alpha", slug: "alpha", name: "Alpha" }; + +// ─── the walk ─── + +test("the documented walk: corpus.json -> manifest -> slugToPage -> page", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader: ArchiveReader = new RemoteSource(ORIGIN); + const channels = await reader.listChannels(); + assert.deepEqual( + channels.map((c) => c.slug), + ["alpha"], + ); + // The channel ref carries its owning origin, which is what lets a hub + // attribute a hit to the member site it came from. + assert.equal(channels[0].siteUrl, ORIGIN); + + const manifest = await reader.transcriptsManifest(channels[0]); + const page = manifest.slugToPage["alpha/a3"]; + assert.equal(page, 1); + + const records = await reader.transcriptPage(channels[0], page); + assert.deepEqual( + records.map((r) => r.id), + ["a3"], + ); + + assert.deepEqual(a.fetches, [ + `${ORIGIN}/corpus.json`, + `${ORIGIN}/transcripts/alpha/manifest.json`, + `${ORIGIN}/transcripts/alpha/page-0001.json`, + ]); + } finally { + restore(); + } +}); + +test("the reader reproduces the contract's URL shapes exactly", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + await reader.transcriptsManifest(ALPHA); + await reader.transcriptPage(ALPHA, 0); + assert.deepEqual(a.fetches, [ + manifestUrl("transcripts", "alpha", ORIGIN), + pageUrl("transcripts", "alpha", 0, ORIGIN), + ]); + // …and the shapes really are the published ones, not just self-consistent. + assert.deepEqual(a.fetches, [ + `${ORIGIN}/transcripts/alpha/manifest.json`, + `${ORIGIN}/transcripts/alpha/page-0000.json`, + ]); + } finally { + restore(); + } +}); + +test("a trailing slash on the origin does not double up", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(`${ORIGIN}/`); + await reader.transcriptsManifest(ALPHA); + assert.deepEqual(a.fetches, [`${ORIGIN}/transcripts/alpha/manifest.json`]); + } finally { + restore(); + } +}); + +// ─── caching ─── + +test("the manifest memo is the promise: two concurrent callers, one fetch", async () => { + const a = fullArchive(); + const held = manifestUrl("transcripts", "alpha", ORIGIN); + const restore = installFetch(a, held); + try { + const reader = new RemoteSource(ORIGIN); + const first = reader.transcriptsManifest(ALPHA); + const second = reader.transcriptsManifest(ALPHA); + // Both callers are waiting on the SAME in-flight read — not two reads that + // happen to land on the same answer. + assert.equal(a.fetches.length, 1); + a.gate?.release(); + const [m1, m2] = await Promise.all([first, second]); + assert.equal(m1, m2); + assert.equal(a.fetches.length, 1); + + // A third, after it resolved, is still free. + await reader.transcriptsManifest(ALPHA); + assert.equal(a.fetches.length, 1); + } finally { + restore(); + } +}); + +test("listChannels({refresh}) is the only way to re-read, and it drops the page caches too", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + await reader.listChannels(); + await reader.transcriptsManifest(ALPHA); + await reader.transcriptPage(ALPHA, 0); + const before = a.fetches.length; + await reader.listChannels(); + await reader.transcriptsManifest(ALPHA); + await reader.transcriptPage(ALPHA, 0); + assert.equal(a.fetches.length, before, "everything served from cache"); + + await reader.listChannels({ refresh: true }); + await reader.transcriptsManifest(ALPHA); + await reader.transcriptPage(ALPHA, 0); + assert.equal(a.fetches.length, before * 2); + } finally { + restore(); + } +}); + +test("a failed read is not memoised", async () => { + const a = makeArchive({}); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + await assert.rejects(() => reader.listChannels(), /404/); + // The corpus appears on the second try; a memoised failure would hide it. + a.body.set(`${ORIGIN}/corpus.json`, CORPUS); + const channels = await reader.listChannels(); + assert.equal(channels.length, 1); + } finally { + restore(); + } +}); + +// ─── the sparse layers: a 404 is data, and the negative is cached ─── + +test("a channel with no digests caches the absence — one probe, not one per query", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + assert.equal(await reader.digestsManifest?.(ALPHA), null); + assert.equal(await reader.digestsManifest?.(ALPHA), null); + assert.equal(await reader.digestsManifest?.(ALPHA), null); + assert.deepEqual(a.fetches, [manifestUrl("digests", "alpha", ORIGIN)]); + } finally { + restore(); + } +}); + +test("a channel WITH digests reads them through the same walk", async () => { + const digests: ChannelDigestsManifest = { + version: 1, + channelSlug: "alpha", + pageCount: 1, + maxPageBytes: 1000, + generatedAt: "2026-09-12T00:00:00.000Z", + slugToPage: { "alpha/a1": 0 }, + } as unknown as ChannelDigestsManifest; + const a = fullArchive(); + a.body.set(manifestUrl("digests", "alpha", ORIGIN), digests); + a.body.set(pageUrl("digests", "alpha", 0, ORIGIN), [ + { slug: "alpha/a1", chapters: [] } as unknown as VideoDigest, + ]); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + const m = await reader.digestsManifest?.(ALPHA); + assert.equal(m?.slugToPage["alpha/a1"], 0); + const page = await reader.digestPage?.(ALPHA, 0); + assert.equal(page?.length, 1); + } finally { + restore(); + } +}); + +test("subs and posts are absent-tolerant in the same way", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + assert.equal(await reader.subsManifest(ALPHA), null); + assert.equal(await reader.postsManifest(ALPHA), null); + assert.equal(await reader.subsManifest(ALPHA), null); + assert.equal(await reader.postsManifest(ALPHA), null); + assert.deepEqual(a.fetches, [ + manifestUrl("subs", "alpha", ORIGIN), + manifestUrl("posts", "alpha", ORIGIN), + ]); + } finally { + restore(); + } +}); + +test("an absent summaries index is an empty map, not an error", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + const index = await reader.videoIndex?.(); + assert.equal(index?.size, 0); + // And the availability map IS that index — one read, not two. + const avail = await reader.availabilityMap(); + assert.equal(avail.size, 0); + } finally { + restore(); + } +}); + +test("an absent duplicates.json is no clusters, not a failure", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + const dupes = await reader.duplicateIndex?.(); + assert.equal(dupes?.size, 0); + assert.deepEqual(a.fetches, [`${ORIGIN}/duplicates.json`]); + } finally { + restore(); + } +}); + +test("an absent search-aliases.json leaves search plain", async () => { + const a = fullArchive(); + const restore = installFetch(a); + try { + const reader = new RemoteSource(ORIGIN); + assert.deepEqual(await reader.loadAliases(), []); + assert.deepEqual(await reader.loadAliases(), []); + assert.deepEqual(a.fetches, [`${ORIGIN}/search-aliases.json`]); + } finally { + restore(); + } +}); + +// ─── PromiseMap ─── + +test("PromiseMap: one load per key, and a rejection evicts itself", async () => { + const map = new PromiseMap<number>(); + let loads = 0; + const load = () => { + loads++; + return Promise.resolve(7); + }; + const [a, b] = await Promise.all([map.take("k", load), map.take("k", load)]); + assert.equal(a, 7); + assert.equal(b, 7); + assert.equal(loads, 1); + + let fails = 0; + const boom = () => { + fails++; + return Promise.reject(new Error("nope")); + }; + await assert.rejects(() => map.take("bad", boom)); + await assert.rejects(() => map.take("bad", boom)); + assert.equal(fails, 2, "a transient failure is never cached as a permanent one"); + + map.clear(); + await map.take("k", load); + assert.equal(loads, 2); +}); + +// ─── PageCache ─── + +test("PageCache evicts by RAW BYTES, least-recently-used first", async () => { + // Budget for exactly three 10-byte pages. + const cache = new PageCache<string>(30); + const loaded: string[] = []; + const put = (k: string) => + cache.take(k, async () => { + loaded.push(k); + return { value: k, bytes: 10 }; + }); + + await put("p1"); + await put("p2"); + await put("p3"); + assert.deepEqual(loaded, ["p1", "p2", "p3"], "three pages fit the budget exactly"); + + // A hit moves its entry to the most-recently-used end. Touch p1 so the + // least-recently-used is p2, not the insertion order's p1 — that is the whole + // difference between an LRU and a queue. + await put("p1"); + assert.deepEqual(loaded, ["p1", "p2", "p3"], "all three still resident"); + + // A fourth page takes the total to 40 > 30, so one goes: p2. + await put("p4"); + assert.deepEqual(loaded, ["p1", "p2", "p3", "p4"]); + await put("p3"); + await put("p1"); + assert.equal(loaded.length, 4, "p3 and p1 survived"); + await put("p2"); + assert.deepEqual(loaded, ["p1", "p2", "p3", "p4", "p2"], "p2 was the victim"); +}); + +test("PageCache keeps MIN_CACHED_PAGES even when one page blows the budget", async () => { + // One page is four times the whole budget: thrashing on it would mean never + // caching anything, which is worse than going over. + const cache = new PageCache<string>(10); + const loaded: string[] = []; + const put = (k: string) => + cache.take(k, async () => { + loaded.push(k); + return { value: k, bytes: 40 }; + }); + await put("big1"); + await put("big2"); + await put("big1"); + await put("big2"); + assert.deepEqual(loaded, ["big1", "big2"]); +}); + +test("PageCache with a zero budget caches nothing but still serves", async () => { + const cache = new PageCache<string>(0); + const loaded: string[] = []; + const put = (k: string) => + cache.take(k, async () => { + loaded.push(k); + return { value: k, bytes: 10 }; + }); + assert.equal(await put("p"), "p"); + assert.equal(await put("p"), "p"); + assert.deepEqual(loaded, ["p", "p"]); +}); + +test("PageCache coalesces concurrent callers and drops a rejection", async () => { + const cache = new PageCache<string>(1000); + let loads = 0; + let release!: () => void; + const gate = new Promise<void>((r) => { + release = r; + }); + const load = async () => { + loads++; + await gate; + return { value: "v", bytes: 10 }; + }; + const a = cache.take("k", load); + const b = cache.take("k", load); + assert.equal(loads, 1); + release(); + assert.deepEqual(await Promise.all([a, b]), ["v", "v"]); + + let fails = 0; + const boom = async () => { + fails++; + throw new Error("nope"); + }; + await assert.rejects(() => cache.take("bad", boom)); + await assert.rejects(() => cache.take("bad", boom)); + assert.equal(fails, 2); +}); diff --git a/common/lib/archive/reader.ts b/common/lib/archive/reader.ts @@ -0,0 +1,779 @@ +// THE ARCHIVE READER — one walk of the published shard scheme. +// +// Moved verbatim from mcp/src/source.ts, which was the best of the five +// hand-rolled readers in the repo (three transports, promise-coalescing caches, +// a byte-budgeted LRU, tolerant of every layer a site legitimately does not +// ship). `mcp/src/source.ts` is now a re-export of this file, so nothing in the +// MCP server changed shape; what changed is that the viewer, the offline cache +// and umtool can reach the same implementation instead of re-deriving it. +// +// ZERO node imports. This file and reader-hub.ts run in a browser; the only +// node-touching transport is reader-fs.ts (LocalSource), which imports FROM +// here and is never value-imported BY here — a `next build` of the export app +// is the test of that, because tsc cannot tell a `node:fs` import from any +// other. +// +// `process` appears once, behind a `typeof process` guard, to read the page +// cache budget knob. Everything else is transport-agnostic. + +import { + type ChannelTranscriptsManifest, + type ChannelSubsManifest, + type Manifest, +} from "../manifest"; +import type { TranscriptDetail, DisplaySummary } from "../transcripts"; +import { summaryState, type VideoState } from "../availability"; +import type { SubsDetail } from "../subs"; +import type { ChannelPostsManifest, Post } from "../posts"; +import { coerceAliasConfig, type SearchAlias } from "../searchAliases"; +import { + resolveCanonicalSlug, + DUPLICATES_FILENAME, + type DuplicateReport, +} from "../duplicates"; +import type { ChannelDigestsManifest, VideoDigest } from "../digests"; +import type { StatsManifest, VideoStat } from "../stats"; +import { + parseChannelGroups, + resolveDefaultGroupId, + DEFAULT_GROUP_FALLBACK_ID, + type ChannelGroup, +} from "../channelGroups"; +import { manifestUrl, pageUrl, rootFileUrl, type HubSite } from "./contract"; +import { recordRead } from "./io-stats"; + +// A channel the source can serve. `siteId`/`siteUrl` are only populated in hub +// mode (so results can be attributed to the owning member site); `key` is the +// stable, source-unique handle a tool passes back to fetch this channel's data. +// `groupId` is the channel's raw group membership as shipped in corpus.json (it +// may be unknown/absent — resolve it against loadGroups() with +// resolveChannelGroupId before using it). +export type ChannelRef = { + key: string; + slug: string; + name: string; + videoCount?: number; + groupId?: string; + siteId?: string; + siteTitle?: string; + siteUrl?: string; +}; + +// The site's channel-group definitions, as read from summaries/manifest.json — +// the same groups the viewer's channel filter renders. `defaultGroupId` is the +// bucket unknown/absent channel groupIds fold onto (resolveChannelGroupId). +export type ChannelGroups = { + groups: ChannelGroup[]; + defaultGroupId: string; +}; + +// What a source returns when it has no group definitions (absent/malformed +// manifest, or hub mode where per-site groups are a different model). +export const EMPTY_GROUPS: ChannelGroups = { + groups: [], + defaultGroupId: DEFAULT_GROUP_FALLBACK_ID, +}; + +// Per-video availability, joined in from the summaries shards (a +// TranscriptDetail record does not carry presence state). Keyed by the +// member-local video slug (`<channelSlug>/<id>`) — the same slug a transcript +// page record carries — so the search engine can apply the `fav` filter. +export type VideoAvailability = { state: VideoState }; + +// One video as the summaries shards describe it. A strict superset of +// VideoAvailability, so the same map serves both the availability join and the +// filter-first page planner — there is one index, not two parallel reads of the +// same files. +// +// Everything here comes from `summaries/`, which is a GLOBAL index of every +// video in the corpus: 1.4 MB and ~0.5 s to parse, against 1.3 GB and ~48 s for +// the transcripts. That ratio is the whole basis of filter-first scanning — the +// exact page set a filtered query needs is computable from this plus each +// channel manifest's `slugToPage`, before a single transcript byte is read. +export type IndexedVideo = { + state: VideoState; + id: string; + channelSlug: string; + title: string; + uploadDate: string; + isLivestream: boolean; + ageRestricted: boolean; +}; + +export type VideoIndex = ReadonlyMap<string, IndexedVideo>; + +// One video's membership in a cross-platform duplicate cluster, as shipped in +// duplicates.json. Keyed by the member-local slug, like the video index. +// +// `aligned` is carried per SIBLING, not per cluster, and is the gate on ever +// translating a timestamp from one copy to another. Absent means NOT MEASURED, +// which must be read as not aligned — a mirror with a longer intro matches on +// text at shifted times, so a plausible-looking citation would land in the +// wrong place in the wrong upload. That is the failure mode that looks like +// success, and the only defence is refusing to guess. +export type ClusterMembership = { + clusterId: string; + // The member that owns derived work for this cluster (resolveCanonicalSlug), + // or null when a human marked the cluster not-a-duplicate. + canonicalSlug: string | null; + isCanonical: boolean; + // Every OTHER member of the cluster. + siblings: { + slug: string; + id: string; + channelSlug: string; + channel: string; + platform: string; + title: string; + duration: number; + uploadDate: string; + hasTranscript: boolean; + aligned: boolean; + offsetSeconds: number | null; + }[]; + // The cluster is a clip-of-a-longer-video relationship: the members overlap + // only partially, so nothing may be mapped across wholesale. + contained: boolean; + // Title+duration only — nothing compared the actual content. An unconfirmed + // suspect, not an established duplicate. + needsReview: boolean; +}; + +export type DuplicateIndex = ReadonlyMap<string, ClusterMembership>; + +// Fold a duplicates.json report into a slug → membership map. Tolerant of +// absence throughout: `compose-site.ts` only writes the file when there is at +// least one publishable cluster, and corpus.json doesn't even declare it, so a +// site legitimately ships none. +export function buildDuplicateIndex( + report: DuplicateReport | null, +): Map<string, ClusterMembership> { + const map = new Map<string, ClusterMembership>(); + if (!report || !Array.isArray(report.clusters)) return map; + for (const cluster of report.clusters) { + const refs = cluster.videoRefs ?? []; + if (refs.length < 2) continue; + const canonicalSlug = resolveCanonicalSlug(cluster); + // A human marked it not-a-duplicate — it is not a cluster any more. + if (canonicalSlug === null) continue; + for (const ref of refs) { + map.set(ref.slug, { + clusterId: cluster.clusterId, + canonicalSlug, + isCanonical: ref.slug === canonicalSlug, + contained: cluster.contained === true, + needsReview: cluster.needsReview === true, + siblings: refs + .filter((o) => o.slug !== ref.slug) + .map((o) => ({ + slug: o.slug, + id: o.id, + channelSlug: o.channelSlug, + channel: o.channel, + platform: o.platform, + title: o.title, + duration: o.duration, + uploadDate: o.uploadDate, + hasTranscript: o.hasTranscript === true, + // Alignment is a property of the PAIR, and the report records it + // against each member relative to the cluster's canonical. Both + // sides must be measured-and-aligned before a timestamp may cross. + aligned: ref.aligned === true && o.aligned === true, + offsetSeconds: o.offsetSeconds ?? null, + })), + }); + } + } + return map; +} + +// Fold the shipped stats shards into a slug → stat map. Same tolerant shape as +// the summaries read: absent or malformed is an empty map, never an error. +// +// Note the shape difference that forces a full read rather than a targeted one: +// StatsManifest carries `channels` and `pageCount` but NO `slugToPage`, so +// there is no way to jump to the page holding one video. Pages are large (up to +// STATS_MAX_PAGE_BYTES = 20 MB), which is exactly why this is lazy — nothing +// reads it until a tool asks for a stat. +export async function buildStatsIndex( + readManifest: () => Promise<StatsManifest | null>, + readPage: (page: number) => Promise<VideoStat[] | null>, +): Promise<Map<string, VideoStat>> { + const map = new Map<string, VideoStat>(); + let manifest: StatsManifest | null; + try { + manifest = await readManifest(); + } catch { + return map; + } + if (!manifest || typeof manifest.pageCount !== "number") return map; + for (let page = 0; page < manifest.pageCount; page++) { + let records: VideoStat[] | null; + try { + records = await readPage(page); + } catch { + continue; + } + if (!records) continue; + for (const r of records) { + if (typeof r.slug === "string") map.set(r.slug, r); + } + } + return map; +} + +// Read a site's global summaries shards (summaries/manifest.json + +// summaries/page-NNNN.json) via `readPage` and fold them into a slug → video +// map. Tolerant: an absent/malformed manifest yields an empty map, and a page +// that fails to read is skipped. `readManifest`/`readPage` throw or return null +// on absence per the source's transport. +// +// Tolerance is load-bearing for the page planner, not just politeness: a +// summaries set that is missing, partial, or older than the transcripts must +// degrade to "I don't know about this video", and the planner's rule for +// don't-know is to scan the page anyway. +export async function buildVideoIndex( + readManifest: () => Promise<Manifest | null>, + readPage: (page: number) => Promise<DisplaySummary[] | null>, +): Promise<Map<string, IndexedVideo>> { + const map = new Map<string, IndexedVideo>(); + let manifest: Manifest | null; + try { + manifest = await readManifest(); + } catch { + return map; + } + if (!manifest || typeof manifest.pageCount !== "number") return map; + for (let page = 0; page < manifest.pageCount; page++) { + let records: DisplaySummary[] | null; + try { + records = await readPage(page); + } catch { + continue; + } + if (!records) continue; + for (const r of records) { + if (typeof r.slug !== "string") continue; + // summaryState falls back to the legacy isDeleted/isUnlisted booleans, + // which matters here more than anywhere: a hub reads summaries pages + // from member origins it does not control, so some of them will have + // been built before `state` existed. + map.set(r.slug, { + state: summaryState(r), + id: r.id, + channelSlug: r.channelSlug, + title: r.title ?? "", + uploadDate: r.uploadDate ?? "", + isLivestream: r.isLivestream === true, + ageRestricted: r.ageRestricted === true, + }); + } + } + return map; +} + +// Parse a summaries/manifest.json blob into channel-group defs, tolerating any +// missing/malformed shape (→ empty fallback). +export function parseGroupsManifest(raw: unknown): ChannelGroups { + const m = (raw ?? {}) as { groups?: unknown; defaultGroupId?: unknown }; + const groups = parseChannelGroups(m.groups); + return { groups, defaultGroupId: resolveDefaultGroupId(m.defaultGroupId, groups) }; +} + +// A read-only view over a transcript corpus's paginated JSON shards. Three +// implementations (local dir / remote origin / federated hub) all speak the +// same three-call contract, which mirrors the documented shard scheme in +// corpus.json: list channels, get a channel's manifest (slug -> page map), get +// a page of full transcript records. +// +// NOTE the method that is deliberately absent: there is no `record(layer, slug, +// id)`. The contract's walk is manifest -> shard -> record and a caller that +// wants one record already holds the shard, so a per-record fetch would be a +// structural regression (N reads where the plan needs one) dressed up as +// convenience. +export interface ArchiveReader { + readonly label: string; + // The channel list, memoised per source instance. searchTranscripts, + // findVideo and findPost all call it, so a 20-id get_transcripts batch used + // to cost 20 corpus.json fetches over HTTP. The memo is the promise, so + // concurrent callers share one fetch. Pass `refresh` to drop it and re-read — + // an explicit staleness escape hatch, deliberately not a TTL. + listChannels(opts?: { refresh?: boolean }): Promise<ChannelRef[]>; + transcriptsManifest(ch: ChannelRef): Promise<ChannelTranscriptsManifest>; + transcriptPage(ch: ChannelRef, page: number): Promise<TranscriptDetail[]>; + // The site's shipped curated search aliases (the same /search-aliases.json the + // viewer reads). Returns [] when the file is absent or malformed. Used to make + // caption search alias-aware, so a query for a term with a curated regex + // (e.g. "k cups" → "(k|cake)[ -]?cup") also matches the mis-transcribed + // spellings. Result is cached per source. + loadAliases(): Promise<SearchAlias[]>; + // The site's channel-group definitions (summaries/manifest.json). Returns the + // empty fallback when absent/malformed, or in hub mode (federated per-site + // groups are a different model — deferred). Cached per source. + loadGroups(): Promise<ChannelGroups>; + // The public origin of the viewer that owns this source's videos, or null + // when there isn't one (a local dir on disk). Used to build archilyzer viewer + // deep links for cited moments (momentUrl). A single-site remote returns its + // base URL; a hub returns null because each video's origin is its member + // site's url (carried on the ChannelRef as `siteUrl`) — prefer that per-video. + publicOrigin(): string | null; + // A channel's live-chat/subs manifest (subs/<slug>/manifest.json), or null + // when the channel ships no subs shards. Same slugToPage/pageCount shape as + // the transcripts manifest. Fetched lazily — only the chat search scope needs + // it. + subsManifest(ch: ChannelRef): Promise<ChannelSubsManifest | null>; + // A page of a channel's subs records (subs/<slug>/page-NNNN.json). Each record + // inlines its per-track cues under `tracks` (e.g. `tracks.live_chat`). + subsPage(ch: ChannelRef, page: number): Promise<SubsDetail[]>; + // A channel's social-posts manifest (posts/<slug>/manifest.json), or null + // when the channel ships no posts shards (i.e. it is a video channel). Same + // slugToPage/pageCount shape as the transcripts manifest. This is the MCP's + // single abstraction boundary, so adding it here yields the posts corpus on + // all three transports (local / remote / hub) at once. + postsManifest(ch: ChannelRef): Promise<ChannelPostsManifest | null>; + // A page of a channel's posts (posts/<slug>/page-NNNN.json). + postsPage(ch: ChannelRef, page: number): Promise<Post[]>; + // A map of every video's availability (deleted/unlisted), keyed by the + // member-local video slug (`<channelSlug>/<id>`), built from the summaries + // shards. Fetched lazily and cached — only the `fav` availability filter needs + // it. Empty when the source ships no summaries. + // + // ReadonlyMap so the richer videoIndex() can BE this map rather than a + // projection of it: ReadonlyMap is covariant in its value type, so one + // Map<string, IndexedVideo> satisfies both and the two can never drift. + availabilityMap(): Promise<ReadonlyMap<string, VideoAvailability>>; + // The full summaries-backed index, when this source ships one. OPTIONAL: the + // in-memory test stubs don't implement it, and a site that ships no + // summaries/ genuinely has no index — callers must degrade to a full scan + // rather than assume an empty index means an empty corpus. + videoIndex?(): Promise<VideoIndex>; + // How many shard pages this source is willing to have in flight at once. + // Optional; callers use `?? DEFAULT_PAGE_CONCURRENCY`. Local is CPU-bound on + // JSON.parse (measured: 42 ms read vs 389 ms parse for an 8 MB page), so + // concurrency there only overlaps read with parse and saturates quickly. + // Remote is latency-bound, where it is the dominant win. + readonly pageConcurrency?: number; + // The shipped cross-platform duplicate report, folded to slug → membership. + // OPTIONAL and empty-when-absent: compose-site only writes duplicates.json + // when there is at least one publishable cluster, corpus.json does not + // declare it, and the in-memory test stubs have no such concept. A site that + // ships none must behave exactly as it does today. + duplicateIndex?(): Promise<DuplicateIndex>; + // The shipped per-video stats index (view/like counts, cueCount, transcript + // coverage), slug-keyed. OPTIONAL for the same reasons. Lazy: stats/ is one + // ~3.4 MB page here and up to 20 MB elsewhere, so it is only read when a tool + // actually asks for it. + statsIndex?(): Promise<ReadonlyMap<string, VideoStat>>; + // A channel's AI-digest manifest (digests/<slug>/manifest.json), or null when + // the channel has none. Mirrors the posts pair, including the cached negative + // — the digest corpus is SPARSE BY DESIGN (a channel with zero digests gets + // no manifest at all), so probing it per query must not cost a read per + // channel per call. OPTIONAL on the interface for the usual reason. + digestsManifest?(ch: ChannelRef): Promise<ChannelDigestsManifest | null>; + digestPage?(ch: ChannelRef, page: number): Promise<VideoDigest[]>; +} + +// Used when a source states no preference. Deliberately modest: each in-flight +// page costs its raw bytes plus ~2.7× that once parsed, and this box is shared. +export const DEFAULT_PAGE_CONCURRENCY = 4; + +// Shape of the channels we read out of a site corpus.json (Layer 1). Kept loose +// — we only need slug/name/count/group. +export type CorpusJsonChannel = { + slug: string; + name?: string; + videoCount?: number; + groupId?: string; +}; +export type SiteCorpusJson = { + channels?: CorpusJsonChannel[]; + // The composed site's own declared public origin — the deployed archilyzer + // viewer these shards were built for. Present in every spec-3 corpus.json. + site?: { id?: string; title?: string; url?: string }; +}; + +export type HubCorpusJson = { + kind?: string; + sites?: HubSite[]; +}; + +// One server-only env read, guarded so this module loads in a browser. The +// dynamic `process.env?.[…]` form is deliberately NOT the inlinable +// `process.env.NAME` member expression — a page cache budget is a server knob, +// and a browser reader takes the default. +function envVar(name: string): string | undefined { + return typeof process !== "undefined" ? process.env?.[name] : undefined; +} + +// ─── Bounded promise caches ─── +// +// Two different caching problems, so two different structures: +// +// manifests — ~551 KB for the whole corpus (29 channels). Small, hot, and +// re-read constantly: an un-hinted 20-id get_transcripts batch +// used to cost ~300 manifest reads because findVideo walks every +// channel per id. Cached OUTRIGHT, no bound. +// +// pages — ~7.4 MB of raw JSON each, several times that once parsed. An +// unbounded map of these is gigabytes, so this is a small LRU. +// Its job is the 20-id batch that lands on ONE shared page (20 +// reads → 1); it is deliberately NOT sized to hold a scan, which +// visits each page exactly once and would only be paying memory +// for evictions. +// +// The bound is a RAW-BYTE budget, not an entry count, because page sizes differ +// by an order of magnitude across corpora (a 13-record VOD page is 8 MB; a +// shorts channel's page is a fraction of that). Default 48 MB, configurable +// with TRANSCRIPT_MCP_PAGE_CACHE_MB (0 disables). +// +// Why 48: a parsed page retains about 2.7× its file bytes (measured — an +// 8.09 MB page holds 21.8 MB of JS heap), so 48 MB of raw budget is roughly +// 130 MB resident. That is the most I am willing to hold on a box that also +// runs a GPU digest sweep and other agents' jobs. It is ~6 pages of this +// corpus, which covers the working set this cache exists for (a 20-id batch +// from an enumerate worklist arrives in page order and lands on 1–3 pages). +// Note what it deliberately does NOT cover: a full-corpus scan is 170 pages ≈ +// 3.7 GB retained, so there is no cache size between "6 pages" and "impossible" +// that changes the full-scan story. Filter-first scanning changes that instead. +const DEFAULT_PAGE_CACHE_MB = 48; +// Always keep at least this many entries, so a corpus whose single page exceeds +// the whole budget still caches that page rather than thrashing on it. +const MIN_CACHED_PAGES = 2; + +export function pageCacheBudgetBytes(): number { + const raw = envVar("TRANSCRIPT_MCP_PAGE_CACHE_MB"); + const mb = + raw === undefined || raw.trim() === "" ? DEFAULT_PAGE_CACHE_MB : Number(raw); + const safe = Number.isFinite(mb) && mb >= 0 ? mb : DEFAULT_PAGE_CACHE_MB; + return Math.floor(safe * 1024 * 1024); +} + +// What a cached loader reports back: the parsed value plus the raw byte size it +// was parsed from, which is what the budget is denominated in. +export type Sized<T> = { value: T; bytes: number }; + +// An LRU keyed by string, holding PROMISES rather than values so that N +// concurrent callers for the same page coalesce onto one read — the pattern +// makeChatFetcher already uses. A rejected promise evicts itself, so a +// transient failure is never cached as a permanent one. +// +// Sizes are only known once a read resolves, so an in-flight entry counts as 0 +// and the budget is enforced on resolve. An entry evicted while still in flight +// resolves normally for whoever already holds its promise; it just isn't +// remembered. +export class PageCache<T> { + private map = new Map<string, { p: Promise<T>; bytes: number }>(); + private total = 0; + constructor(private readonly maxBytes: number) {} + + take(key: string, load: () => Promise<Sized<T>>): Promise<T> { + const hit = this.map.get(key); + if (hit !== undefined) { + this.map.delete(key); + this.map.set(key, hit); // most-recently used goes last + return hit.p; + } + if (this.maxBytes <= 0) return load().then((s) => s.value); + + const entry: { p: Promise<T>; bytes: number } = { p: null as never, bytes: 0 }; + entry.p = load() + .then((s) => { + // Only account for it if we're still the live entry for this key — + // a refresh() between issue and resolve must not resurrect it. + if (this.map.get(key) === entry) { + entry.bytes = s.bytes; + this.total += s.bytes; + this.evict(); + } + return s.value; + }) + .catch((e: unknown) => { + this.drop(key, entry); + throw e; + }); + this.map.set(key, entry); + return entry.p; + } + + private drop(key: string, entry: { bytes: number }): void { + if (this.map.get(key) === entry) { + this.map.delete(key); + this.total -= entry.bytes; + } + } + + private evict(): void { + while (this.total > this.maxBytes && this.map.size > MIN_CACHED_PAGES) { + const oldest = this.map.entries().next().value; + if (oldest === undefined) break; + this.map.delete(oldest[0]); + this.total -= oldest[1].bytes; + } + } + + clear(): void { + this.map.clear(); + this.total = 0; + } +} + +// The unbounded sibling, for the small-and-hot caches (manifests). Same +// don't-memoise-a-failure rule. +export class PromiseMap<T> { + private map = new Map<string, Promise<T>>(); + + take(key: string, load: () => Promise<T>): Promise<T> { + const hit = this.map.get(key); + if (hit !== undefined) return hit; + const p = load().catch((e: unknown) => { + this.map.delete(key); + throw e; + }); + this.map.set(key, p); + return p; + } + + clear(): void { + this.map.clear(); + } +} + +// ─── Remote: fetch shards from a deployed site origin over HTTP ─── +export class RemoteSource implements ArchiveReader { + readonly label: string; + private base: string; + private aliases?: SearchAlias[]; + private groups?: ChannelGroups; + private index?: Promise<Map<string, IndexedVideo>>; + private duplicates?: Promise<DuplicateIndex>; + private stats?: Promise<ReadonlyMap<string, VideoStat>>; + private subsManifests = new PromiseMap<ChannelSubsManifest | null>(); + private digestManifests = new PromiseMap<ChannelDigestsManifest | null>(); + private subsPages: PageCache<SubsDetail[]>; + private postsManifests = new PromiseMap<ChannelPostsManifest | null>(); + private transcriptManifests = new PromiseMap<ChannelTranscriptsManifest>(); + private transcriptPages: PageCache<TranscriptDetail[]>; + + // Over HTTP the cost is latency, not parse, so a wider window is the dominant + // win — this is where bounded concurrency actually pays. + readonly pageConcurrency = 8; + + // `budgetBytes` lets a hub divide one memory ceiling across its members + // instead of granting each member the full budget (N members × 48 MB is not a + // budget, it's N budgets). + constructor(baseUrl: string, budgetBytes = pageCacheBudgetBytes()) { + this.base = baseUrl.replace(/\/+$/, ""); + this.label = `remote:${this.base}`; + this.subsPages = new PageCache<SubsDetail[]>(budgetBytes); + this.transcriptPages = new PageCache<TranscriptDetail[]>(budgetBytes); + } + + // The deployed site origin — the archilyzer viewer that owns these videos. + publicOrigin(): string | null { + return this.base; + } + + private resetCaches(): void { + this.subsManifests.clear(); + this.digestManifests.clear(); + this.subsPages.clear(); + this.postsManifests.clear(); + this.transcriptManifests.clear(); + this.transcriptPages.clear(); + this.aliases = undefined; + this.groups = undefined; + this.index = undefined; + this.duplicates = undefined; + this.stats = undefined; + } + + subsManifest(ch: ChannelRef): Promise<ChannelSubsManifest | null> { + return this.subsManifests.take(ch.slug, async () => { + try { + const res = await fetch(manifestUrl("subs", ch.slug, this.base)); + return res.ok ? ((await res.json()) as ChannelSubsManifest) : null; + } catch { + return null; + } + }); + } + + subsPage(ch: ChannelRef, page: number): Promise<SubsDetail[]> { + return this.subsPages.take(`${ch.slug}:${page}`, () => + this.getJsonSized<SubsDetail[]>(pageUrl("subs", ch.slug, page), "subsPage"), + ); + } + + // Cached per channel including the negative answer — otherwise every + // posts-covering search costs one 404 per video-only channel. + postsManifest(ch: ChannelRef): Promise<ChannelPostsManifest | null> { + return this.postsManifests.take(ch.slug, async () => { + try { + const res = await fetch(manifestUrl("posts", ch.slug, this.base)); + return res.ok ? ((await res.json()) as ChannelPostsManifest) : null; + } catch { + return null; + } + }); + } + + postsPage(ch: ChannelRef, page: number): Promise<Post[]> { + return this.getJson(pageUrl("posts", ch.slug, page), "postsPage"); + } + + videoIndex(): Promise<VideoIndex> { + this.index ??= buildVideoIndex( + async () => { + const res = await fetch(manifestUrl("summaries", undefined, this.base)); + return res.ok ? ((await res.json()) as Manifest) : null; + }, + async (page) => { + const res = await fetch(pageUrl("summaries", undefined, page, this.base)); + return res.ok ? ((await res.json()) as DisplaySummary[]) : null; + }, + ); + return this.index; + } + + availabilityMap(): Promise<ReadonlyMap<string, VideoAvailability>> { + return this.videoIndex(); + } + + digestsManifest(ch: ChannelRef): Promise<ChannelDigestsManifest | null> { + return this.digestManifests.take(ch.slug, async () => { + try { + const res = await fetch(manifestUrl("digests", ch.slug, this.base)); + return res.ok ? ((await res.json()) as ChannelDigestsManifest) : null; + } catch { + return null; + } + }); + } + + digestPage(ch: ChannelRef, page: number): Promise<VideoDigest[]> { + return this.getJson(pageUrl("digests", ch.slug, page), "digestPage"); + } + + duplicateIndex(): Promise<DuplicateIndex> { + this.duplicates ??= (async () => { + try { + const res = await fetch(rootFileUrl(DUPLICATES_FILENAME, this.base)); + return buildDuplicateIndex( + res.ok ? ((await res.json()) as DuplicateReport) : null, + ); + } catch { + return buildDuplicateIndex(null); + } + })(); + return this.duplicates; + } + + statsIndex(): Promise<ReadonlyMap<string, VideoStat>> { + this.stats ??= buildStatsIndex( + async () => { + const res = await fetch(manifestUrl("stats", undefined, this.base)); + return res.ok ? ((await res.json()) as StatsManifest) : null; + }, + async (page) => { + const res = await fetch(pageUrl("stats", undefined, page, this.base)); + return res.ok ? ((await res.json()) as VideoStat[]) : null; + }, + ); + return this.stats; + } + + async loadAliases(): Promise<SearchAlias[]> { + if (this.aliases) return this.aliases; + try { + const res = await fetch(rootFileUrl("search-aliases.json", this.base)); + this.aliases = res.ok + ? coerceAliasConfig(await res.json()).aliases + : []; + } catch { + this.aliases = []; + } + return this.aliases; + } + + async loadGroups(): Promise<ChannelGroups> { + if (this.groups) return this.groups; + try { + const res = await fetch(manifestUrl("summaries", undefined, this.base)); + this.groups = res.ok + ? parseGroupsManifest(await res.json()) + : EMPTY_GROUPS; + } catch { + this.groups = EMPTY_GROUPS; + } + return this.groups; + } + + // Fetches as TEXT so the byte size is knowable — the page cache's budget is + // denominated in raw bytes, and res.json() throws the length away. + // + // `p` is a ROOT-RELATIVE path from the contract builders; the base is joined + // here so the thrown error names the full URL, exactly as before. + private async getJsonSized<T>( + p: string, + kind: string, + ): Promise<{ value: T; bytes: number }> { + const res = await fetch(`${this.base}${p}`); + if (!res.ok) { + throw new Error(`GET ${this.base}${p} -> ${res.status} ${res.statusText}`); + } + const raw = await res.text(); + recordRead(kind, raw.length); + return { value: JSON.parse(raw) as T, bytes: raw.length }; + } + + private async getJson<T>(p: string, kind = "json"): Promise<T> { + return (await this.getJsonSized<T>(p, kind)).value; + } + + private channelList?: Promise<ChannelRef[]>; + + listChannels(opts: { refresh?: boolean } = {}): Promise<ChannelRef[]> { + if (opts.refresh) { + this.channelList = undefined; + this.resetCaches(); + } + this.channelList ??= this.readChannels().catch((e: unknown) => { + this.channelList = undefined; // don't memoise a failure + throw e; + }); + return this.channelList; + } + + private async readChannels(): Promise<ChannelRef[]> { + const corpus = await this.getJson<SiteCorpusJson>( + rootFileUrl("corpus.json"), + "corpus", + ); + return (corpus.channels ?? []).map((c) => ({ + key: c.slug, + slug: c.slug, + name: c.name ?? c.slug, + videoCount: c.videoCount, + groupId: c.groupId, + siteUrl: this.base, + })); + } + + transcriptsManifest(ch: ChannelRef): Promise<ChannelTranscriptsManifest> { + return this.transcriptManifests.take(ch.slug, () => + this.getJson<ChannelTranscriptsManifest>( + manifestUrl("transcripts", ch.slug), + "transcriptsManifest", + ), + ); + } + + transcriptPage(ch: ChannelRef, page: number): Promise<TranscriptDetail[]> { + return this.transcriptPages.take(`${ch.slug}:${page}`, () => + this.getJsonSized<TranscriptDetail[]>( + pageUrl("transcripts", ch.slug, page), + "transcriptPage", + ), + ); + } +}