// ─── 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 { TAGS_FILENAME, type PublishedTag } from "../curatedTags"; import { coercePublishedTags } from "../publishedTags"; 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(); // siteId -> source private aliases?: SearchAlias[]; private tags?: PublishedTag[]; private index?: Promise>; private duplicates?: Promise; private stats?: Promise>; private sites?: Promise; // 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; 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 { this.sites ??= this.readSites().catch((e: unknown) => { this.sites = undefined; throw e; }); return this.sites; } private async readSites(): Promise { 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 { 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; } // A hub can ship its own /tags.json — the federation-wide vocabulary with // counts summed across members. If it doesn't, tags are simply off for // hub-wide search: merging member documents here would produce counts that // no single site can reproduce, and a chip whose number nobody can check is // worse than no chip. async loadTags(): Promise { if (this.tags) return this.tags; try { const res = await fetch(rootFileUrl(TAGS_FILENAME, this.hubBase)); this.tags = res.ok ? coercePublishedTags(await res.json()) : []; } catch { this.tags = []; } return this.tags; } // 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 { 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; listChannels(opts: { refresh?: boolean } = {}): Promise { 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 { 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 { if (!ch.siteId) throw new Error("hub channel ref missing siteId"); return this.memberFor(ch.siteId).transcriptsManifest(ch); } transcriptPage(ch: ChannelRef, page: number): Promise { 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 { 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 { if (!ch.siteId) throw new Error("hub channel ref missing siteId"); return this.memberFor(ch.siteId).subsPage(ch, page); } async postsManifest(ch: ChannelRef): Promise { 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 { if (!ch.siteId) throw new Error("hub channel ref missing siteId"); return this.memberFor(ch.siteId).postsPage(ch, page); } async digestsManifest(ch: ChannelRef): Promise { 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 { 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 // (`/`) — the same slug a member's transcript page records // carry — so a per-record lookup joins correctly. Built lazily/cached. videoIndex(): Promise { this.index ??= this.buildMergedIndex(); return this.index; } private async buildMergedIndex(): Promise> { const merged = new Map(); 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> { 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 { this.duplicates ??= this.mergeMembers((m) => m.duplicateIndex()); return this.duplicates; } statsIndex(): Promise> { 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( read: (m: RemoteSource) => Promise>, ): Promise> { const merged = new Map(); 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; } }