commit 7cb1d765993386693547b561922cb8c63b614a66
parent 27ea3025161a1c577b69d03202a32b38a640f564
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 12 Sep 2026 02:06:53 -0400
mcp: source.ts is a re-export — 1,386 lines to 47, no importer changed
The reader now lives in the core, so this file states the names this server
already spells and where they come from: `ShardSource` is `ArchiveReader`,
`LocalSource` / `RemoteSource` / `HubSource` are the three transports,
`HubSite` is the contract's projection of a published hub member entry, and
`ioStatsEnabled` / `ioStatsSnapshot` / `recordRead` are the counters bench reads
off stderr.
Every name the seven `./source` importers pull resolves unchanged —
sourceRegistry.ts (HubSource, HubSite, ShardSource), sources.ts (the three
constructors + ShardSource), server.ts (ioStatsEnabled, ioStatsSnapshot,
ChannelRef, HubSite, ShardSource), search.ts (DEFAULT_PAGE_CONCURRENCY,
ChannelRef, ShardSource, VideoAvailability) and the three test files
(ChannelGroups, ChannelRef, DuplicateIndex, HubSite, LocalSource, ShardSource,
VideoAvailability, VideoIndex). `mcp/bench` spawns the server over stdio and
never imports this file, so it is untouched by construction.
Gates: mcp 205/205; tsc --noEmit clean in all six packages.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
| M | mcp/src/source.ts | | | 1429 | +++---------------------------------------------------------------------------- |
1 file changed, 45 insertions(+), 1384 deletions(-)
diff --git a/mcp/src/source.ts b/mcp/src/source.ts
@@ -1,1386 +1,47 @@
-import { readFile, readdir } from "node:fs/promises";
-import path from "node:path";
-import {
- pageFileName,
- type ChannelTranscriptsManifest,
- type ChannelSubsManifest,
- type Manifest,
-} from "yt-dlp-transcript-common/lib/manifest";
-import type {
- TranscriptDetail,
- DisplaySummary,
-} from "yt-dlp-transcript-common/lib/transcripts";
-import {
- summaryState,
- type VideoState,
-} from "yt-dlp-transcript-common/lib/availability";
-import type { SubsDetail } from "yt-dlp-transcript-common/lib/subs";
-import {
- postsPageFileName,
- type ChannelPostsManifest,
- type Post,
-} from "yt-dlp-transcript-common/lib/posts";
-import {
- coerceAliasConfig,
- type SearchAlias,
-} from "yt-dlp-transcript-common/lib/searchAliases";
-import {
- resolveCanonicalSlug,
- DUPLICATES_FILENAME,
- type DuplicateReport,
-} from "yt-dlp-transcript-common/lib/duplicates";
-import {
- digestPageFileName,
- type ChannelDigestsManifest,
- type VideoDigest,
-} from "yt-dlp-transcript-common/lib/digests";
-import {
- statsPageFileName,
- type StatsManifest,
- type VideoStat,
-} from "yt-dlp-transcript-common/lib/stats";
-import {
- parseChannelGroups,
- resolveDefaultGroupId,
- DEFAULT_GROUP_FALLBACK_ID,
- type ChannelGroup,
-} from "yt-dlp-transcript-common/lib/channelGroups";
-
-// 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).
-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.
-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.
-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).
-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.
-export interface ShardSource {
- 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;
-
-// ─── I/O instrumentation (opt-in, for mcp/bench) ───
-//
-// A process-wide counter of shard reads and parsed bytes, so the benchmark can
-// report the STRUCTURAL cost of a query (how many pages, how many bytes) next
-// to its wall time. That matters on this box specifically: wall time is only
-// meaningful when the machine is idle, but read counts and byte counts are
-// properties of the query plan and hold under any load.
-//
-// Off unless MCP_IO_STATS=1, and even then it is two integer adds per read.
-export type IoStats = { reads: number; bytes: number };
-
-const IO_STATS_ON = process.env.MCP_IO_STATS === "1";
-
-const ioTotals: Record<string, IoStats> = {};
-
-export function recordRead(kind: string, bytes: number): void {
- if (!IO_STATS_ON) return;
- const slot = (ioTotals[kind] ??= { reads: 0, bytes: 0 });
- slot.reads++;
- slot.bytes += bytes;
-}
-
-// A snapshot of every counter so far, for diffing across one tool call.
-export function ioStatsSnapshot(): Record<string, IoStats> {
- const out: Record<string, IoStats> = {};
- for (const [k, v] of Object.entries(ioTotals)) out[k] = { ...v };
- return out;
-}
-
-export function ioStatsEnabled(): boolean {
- return IO_STATS_ON;
-}
-
-// 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;
-}
-
-// Shape of the channels we read out of a site corpus.json (Layer 1). Kept loose
-// — we only need slug/name/count/group.
-type CorpusJsonChannel = {
- slug: string;
- name?: string;
- videoCount?: number;
- groupId?: string;
-};
-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 };
-};
-
-// ─── 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;
-
-function pageCacheBudgetBytes(): number {
- const raw = process.env.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.
-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.
-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.
-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();
- }
-}
-
-// 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";
-type HubCorpusJson = {
- kind?: string;
- sites?: { siteId: string; title: string; url: string }[];
-};
+// The archive reader, as this server has always known it.
+//
+// This file WAS the reader: 1,386 lines holding the only good walk of the
+// published shard scheme in the repo, plus three transports and two caches. It
+// now lives in `yt-dlp-transcript-common/lib/archive/`, where the viewer, the
+// offline cache and umtool can reach the same implementation instead of each
+// re-deriving the walk (one-core plan, phase 2). Nothing about it changed in
+// the move — the interface is the same 18 members under a new name, and there
+// is deliberately still no per-record fetch.
+//
+// `ShardSource` stays as the name every importer here already spells, so this
+// server's seven `./source` importers and `mcp/bench` (which spawns the server
+// and never imports this file) are untouched.
+
+export type {
+ ArchiveReader as ShardSource,
+ ChannelGroups,
+ ChannelRef,
+ ClusterMembership,
+ DuplicateIndex,
+ IndexedVideo,
+ VideoAvailability,
+ VideoIndex,
+} from "yt-dlp-transcript-common/lib/archive/reader";
+
+export {
+ DEFAULT_PAGE_CONCURRENCY,
+ RemoteSource,
+ buildDuplicateIndex,
+} from "yt-dlp-transcript-common/lib/archive/reader";
+
+export { LocalSource } from "yt-dlp-transcript-common/lib/archive/reader-fs";
+export { HubSource } from "yt-dlp-transcript-common/lib/archive/reader-hub";
// One member site of a hub, as listed in the hub's corpus.json. Exposed so the
-// source controller can resolve a `site`/`sites` token to a member origin.
-export type HubSite = { siteId: string; title: string; url: string };
-
-// ─── Local: read composed shards from a directory on disk ───
-// `dir` is a composed public dir (or any dir containing transcripts/<slug>/…).
-// 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 ShardSource {
- 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, postsPageFileName(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, digestPageFileName(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", statsPageFileName(page)),
- "statsPage",
- ),
- );
- return this.stats;
- }
-}
-
-// ─── Remote: fetch shards from a deployed site origin over HTTP ───
-export class RemoteSource implements ShardSource {
- 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(`${this.base}/subs/${ch.slug}/manifest.json`);
- 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[]>(
- `/subs/${ch.slug}/${pageFileName(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(`${this.base}/posts/${ch.slug}/manifest.json`);
- return res.ok ? ((await res.json()) as ChannelPostsManifest) : null;
- } catch {
- return null;
- }
- });
- }
-
- postsPage(ch: ChannelRef, page: number): Promise<Post[]> {
- return this.getJson(`/posts/${ch.slug}/${postsPageFileName(page)}`, "postsPage");
- }
-
- videoIndex(): Promise<VideoIndex> {
- this.index ??= buildVideoIndex(
- async () => {
- const res = await fetch(`${this.base}/summaries/manifest.json`);
- return res.ok ? ((await res.json()) as Manifest) : null;
- },
- async (page) => {
- const res = await fetch(`${this.base}/summaries/${pageFileName(page)}`);
- 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(`${this.base}/digests/${ch.slug}/manifest.json`);
- return res.ok ? ((await res.json()) as ChannelDigestsManifest) : null;
- } catch {
- return null;
- }
- });
- }
-
- digestPage(ch: ChannelRef, page: number): Promise<VideoDigest[]> {
- return this.getJson(
- `/digests/${ch.slug}/${digestPageFileName(page)}`,
- "digestPage",
- );
- }
-
- duplicateIndex(): Promise<DuplicateIndex> {
- this.duplicates ??= (async () => {
- try {
- const res = await fetch(`${this.base}/${DUPLICATES_FILENAME}`);
- 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(`${this.base}/stats/manifest.json`);
- return res.ok ? ((await res.json()) as StatsManifest) : null;
- },
- async (page) => {
- const res = await fetch(`${this.base}/stats/${statsPageFileName(page)}`);
- 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(`${this.base}/search-aliases.json`);
- 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(`${this.base}/summaries/manifest.json`);
- 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.
- 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>("/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>(
- `/transcripts/${ch.slug}/manifest.json`,
- "transcriptsManifest",
- ),
- );
- }
-
- transcriptPage(ch: ChannelRef, page: number): Promise<TranscriptDetail[]> {
- return this.transcriptPages.take(`${ch.slug}:${page}`, () =>
- this.getJsonSized<TranscriptDetail[]>(
- `/transcripts/${ch.slug}/${pageFileName(page)}`,
- "transcriptPage",
- ),
- );
- }
-}
-
-// ─── 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.
-export class HubSource implements ShardSource {
- 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 res = await fetch(`${this.hubBase}/corpus.json`);
- if (!res.ok) {
- throw new Error(
- `GET ${this.hubBase}/corpus.json -> ${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(`${this.hubBase}/search-aliases.json`);
- 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;
- }
-}
+// source controller can resolve a `site`/`sites` token to a member origin. Now
+// a projection of the entry the hub actually publishes (`HubCorpusSite`), so
+// the read and write spellings cannot drift.
+export type { HubSite } from "yt-dlp-transcript-common/lib/archive/contract";
+
+// The opt-in read/byte counters mcp/bench reads off stderr.
+export {
+ ioStatsEnabled,
+ ioStatsSnapshot,
+ recordRead,
+ type IoStats,
+} from "yt-dlp-transcript-common/lib/archive/io-stats";