import path from "node:path"; import { open, readdir, readFile } from "node:fs/promises"; import pLimit from "p-limit"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { extractVideoId } from "../lib/videoId"; // The ROSTER: every video id this channel has ever been seen to contain, with // the URL it was seen at. `channels//roster.json`. // // Why it exists. Until now the only record that a video belonged to a channel // was the `playlist` file, which a full sweep overwrites wholesale. A video // that appeared in a listing but was never downloaded has no data// dir and // no metadata.info.json, so `playlist` was the ONLY place its URL lived — when // it dropped out of the listing it was erased with no record it had ever // existed, and no way to attempt a direct-link recovery. Downstream, // channelSnapshot's undownloadedIds walks that same file, so the auto-download // runner's whole work-list went with it. // // The roster is the durable answer: append-only, never pruned by a sweep. The // invariant that makes every writer safe is that only a FULL, ACCEPTED // enumeration may compute the "missing" sets (see ./channelSets). Every other // path — the paged sync, the quick availability check, a one-off import — may // only ADD. mergeRoster below has no removal branch at all; that is deliberate. // // Deliberately a leaf apart from Paths and extractVideoId, the same way // ./maybeMissingStore is, so ytdlp/runYtdlp.ts can write it without a cycle. export const ROSTER_FILENAME = "roster.json"; export const ROSTER_VERSION = 1; // Where an id was FIRST learned. Never revised: paired with firstSeenAt, it // records the original sighting, so `source: "disk"` keeps meaning "this only // ever showed up as a directory" even after a later listing confirms it. export type RosterSource = "listing" | "disk" | "import"; // What an enumeration was judged to be. Only "ok" and "shrink-confirmed" are // accepted listings — see ./acceptListing, which reads lastSweep to implement // the two-observation confirmation. export type SweepVerdict = | "ok" | "empty" | "shrink-suspect" | "shrink-confirmed"; export type RosterSweep = { at: string; listedCount: number; verdict: SweepVerdict; }; export type RosterEntry = { // The URL the video was seen at. The whole point of this file: it survives // the video leaving the listing, which is what makes recovery possible. // May be "" for a disk-seeded entry whose URL wasn't cheaply recoverable; // the next accepted enumeration that lists it fills it in. url: string; firstSeenAt: string; // The last ACCEPTED enumeration that contained this id. Advanced only by a // listing observation — an import or a disk scan is not an enumeration. lastListedAt: string; source: RosterSource; }; export type Roster = { version: number; updatedAt: string; lastSweep: RosterSweep | null; entries: Record; }; export type ObservedVideo = { id: string; url: string }; export function emptyRoster(): Roster { return { version: ROSTER_VERSION, updatedAt: "", lastSweep: null, entries: {}, }; } export function rosterPath(paths: Paths, slug: string): string { return path.join(paths.channelsDir, slug, ROSTER_FILENAME); } // Canonicalize a list of listing URLs into roster observations. Nulls are // dropped: an id we can't derive would never match a data-dir name anyway. export function observationsFromUrls( urls: ReadonlyArray, ): ObservedVideo[] { const out: ObservedVideo[] = []; for (const url of urls) { const id = extractVideoId(url); if (id) out.push({ id, url }); } return out; } // Fold observations into the roster. PURE and strictly ADDITIVE — there is no // path here that drops an entry, which is what lets every non-sweep writer call // it without needing to be trusted. // // Returns the SAME object when nothing changed, so callers can skip the write // (a big channel's roster is a couple of MB and most paged syncs see nothing // new). export function mergeRoster( roster: Roster, observed: ReadonlyArray, now: string, source: RosterSource, ): Roster { const fromListing = source === "listing"; const entries: Record = { ...roster.entries }; let changed = false; for (const { id, url } of observed) { if (!id) continue; const prev = entries[id]; if (!prev) { entries[id] = { url: url || "", firstSeenAt: now, lastListedAt: now, source, }; changed = true; continue; } // A listing URL is the freshest we can get, so it wins; an import or disk // scan only fills a gap. firstSeenAt and source never move. const nextUrl = url && (fromListing || !prev.url) ? url : prev.url; const nextLastListed = fromListing ? now : prev.lastListedAt; if (nextUrl !== prev.url || nextLastListed !== prev.lastListedAt) { entries[id] = { ...prev, url: nextUrl, lastListedAt: nextLastListed }; changed = true; } } if (!changed) return roster; return { ...roster, version: ROSTER_VERSION, updatedAt: now, entries }; } // A RENAMED RECORD keeps its roster entry under its new id. Not a removal: // the entry moves (url, firstSeenAt, source and all), which is what a dir // renamed to its canonical id (reconcileVideoDirs.ts) needs — left under the // old id it would read as a video the channel has and nobody downloaded. // When both ids have an entry the new one's stays and the earlier // firstSeenAt wins. Returns the SAME object when nothing moved. export function renameRosterEntries( roster: Roster, renames: ReadonlyArray<{ from: string; to: string }>, now: string, ): Roster { const entries: Record = { ...roster.entries }; let changed = false; for (const { from, to } of renames) { const prev = entries[from]; if (!prev || from === to) continue; const there = entries[to]; entries[to] = there ? { ...there, firstSeenAt: prev.firstSeenAt && prev.firstSeenAt < there.firstSeenAt ? prev.firstSeenAt : there.firstSeenAt, url: there.url || prev.url } : prev; delete entries[from]; changed = true; } if (!changed) return roster; return { ...roster, version: ROSTER_VERSION, updatedAt: now, entries }; } // Stamp the outcome of an enumeration. Kept separate from mergeRoster because a // REJECTED enumeration still merges (additively, losing nothing) while recording // that its listing was not trusted — that record is what the next enumeration // compares against to confirm a genuine mass deletion. export function recordSweep(roster: Roster, sweep: RosterSweep): Roster { return { ...roster, version: ROSTER_VERSION, updatedAt: sweep.at, lastSweep: sweep, }; } function isSweepVerdict(v: unknown): v is SweepVerdict { return ( v === "ok" || v === "empty" || v === "shrink-suspect" || v === "shrink-confirmed" ); } function isRosterSource(v: unknown): v is RosterSource { return v === "listing" || v === "disk" || v === "import"; } // Coerce a stored value into a clean Roster, dropping anything ill-typed. A // hand-edited or half-written file must never crash a sync; the worst case is // that it reads as empty and the next seed rebuilds it from playlist ∪ data/. export function normalizeRoster(value: unknown): Roster { const empty = emptyRoster(); if (!value || typeof value !== "object") return empty; const r = value as Record; const entries: Record = {}; if (r.entries && typeof r.entries === "object") { for (const [id, raw] of Object.entries(r.entries as Record)) { if (!id || !raw || typeof raw !== "object") continue; const e = raw as Record; const firstSeenAt = typeof e.firstSeenAt === "string" ? e.firstSeenAt : ""; entries[id] = { url: typeof e.url === "string" ? e.url : "", firstSeenAt, lastListedAt: typeof e.lastListedAt === "string" ? e.lastListedAt : firstSeenAt, source: isRosterSource(e.source) ? e.source : "listing", }; } } const rawSweep = r.lastSweep as Record | undefined | null; const lastSweep: RosterSweep | null = rawSweep && typeof rawSweep === "object" && typeof rawSweep.at === "string" && typeof rawSweep.listedCount === "number" && Number.isFinite(rawSweep.listedCount) && isSweepVerdict(rawSweep.verdict) ? { at: rawSweep.at, listedCount: Math.max(0, Math.floor(rawSweep.listedCount)), verdict: rawSweep.verdict, } : null; return { version: typeof r.version === "number" ? r.version : ROSTER_VERSION, updatedAt: typeof r.updatedAt === "string" ? r.updatedAt : "", lastSweep, entries, }; } export async function loadRoster(paths: Paths, slug: string): Promise { try { const raw = await readFile(rosterPath(paths, slug), "utf8"); return normalizeRoster(JSON.parse(raw)); } catch { return emptyRoster(); } } // Atomic tmp+rename, like maybe-missing.json: a crashed write must never leave a // truncated roster behind, because a truncated roster is exactly the data loss // this file exists to prevent. Several actors write it (the sync, the quick // availability check, the pipeline's server action — a different bundle), so // it goes through the one per-path chain on globalThis. export async function writeRoster( paths: Paths, slug: string, roster: Roster, ): Promise { await writeJsonAtomic(rosterPath(paths, slug), roster, { mkdir: true }); } // Load, merge, write in one step — the shape every additive writer wants. Skips // the write when the merge changed nothing. export async function mergeRosterFile( paths: Paths, slug: string, observed: ReadonlyArray, now: string, source: RosterSource, ): Promise { const before = await loadRoster(paths, slug); const after = mergeRoster(before, observed, now, source); if (after !== before) await writeRoster(paths, slug, after); return after; } const DISK_URL_CONCURRENCY = 16; // How far into a metadata.info.json we're willing to read looking for // webpage_url. yt-dlp writes it AFTER the formats array, which on a real // YouTube video puts it around 80 KB into a ~550 KB file. Seeding a channel is // a one-off, but multiplying a full 550 KB read by 10,795 videos is not // something to do inside a sync job; a miss just leaves the URL blank, and the // next accepted enumeration fills it in. const METADATA_SCAN_LIMIT_BYTES = 256 * 1024; // Recover a downloaded video's URL from its own dir, cheapest source first. // availability.json is 137 bytes and is written by the availability backfill // that runs on every sync, so it covers nearly everything; the capped scan of // metadata.info.json is the fallback for dirs that predate it. export async function readVideoUrlFromDisk(videoDir: string): Promise { try { const raw = await readFile(path.join(videoDir, "availability.json"), "utf8"); const url = (JSON.parse(raw) as { webpageUrl?: unknown }).webpageUrl; if (typeof url === "string" && url) return url; } catch { /* fall through to the metadata scan */ } return scanMetadataForUrl(path.join(videoDir, "metadata.info.json")); } async function scanMetadataForUrl(file: string): Promise { let handle: Awaited> | null = null; try { handle = await open(file, "r"); const chunk = Buffer.alloc(64 * 1024); // Carry the tail of the previous chunk so a key split across a boundary is // still found. let carry = ""; let read = 0; while (read < METADATA_SCAN_LIMIT_BYTES) { const { bytesRead } = await handle.read(chunk, 0, chunk.length, read); if (bytesRead === 0) break; read += bytesRead; const text = carry + chunk.subarray(0, bytesRead).toString("utf8"); const match = /"webpage_url"\s*:\s*"((?:[^"\\]|\\.)*)"/.exec(text); if (match) { try { return JSON.parse(`"${match[1]}"`) as string; } catch { return match[1]; } } carry = text.slice(-64); } } catch { /* no metadata, or unreadable — the URL simply isn't recoverable here */ } finally { await handle?.close().catch(() => {}); } return ""; } async function readSeedPlaylist(file: string): Promise { try { const raw = await readFile(file, "utf8"); return observationsFromUrls( raw.split("\n").map((s) => s.trim()).filter(Boolean), ); } catch { return []; } } async function readVideoDirNames(dataDir: string): Promise { try { const entries = await readdir(dataDir, { withFileTypes: true }); return entries.filter((e) => e.isDirectory()).map((e) => e.name); } catch { return []; } } // Build the roster for a channel that doesn't have one yet, from the state that // exists RIGHT NOW: the pre-existing `playlist` (ids and URLs) union the data/ // dir names. This has to happen before anything overwrites `playlist`, and it // is the migration's whole value — it captures today's listed-but-never-fetched // entries before the first sweep can erase them, and gives a channel that has // never had a listing file (there is one such channel in the live corpus, with // 617 videos) its first record of what it contains. // // A corrupt or empty roster re-seeds, which is the healing behaviour we want. export async function seedRosterIfAbsent( paths: Paths, slug: string, opts: { now: string; onLog?: (msg: string) => void }, ): Promise { const existing = await loadRoster(paths, slug); if (Object.keys(existing.entries).length > 0) return existing; const channelDir = path.join(paths.channelsDir, slug); const [listed, dirNames] = await Promise.all([ readSeedPlaylist(path.join(channelDir, "playlist")), readVideoDirNames(path.join(channelDir, "data")), ]); let roster = mergeRoster(existing, listed, opts.now, "listing"); // Only dirs the listing didn't already account for need a disk read; on a // channel in a downloaded steady state that is a handful, not thousands. const known = new Set(Object.keys(roster.entries)); const diskOnly = dirNames.filter((id) => !known.has(id)); const limit = pLimit(DISK_URL_CONCURRENCY); const diskObserved = await Promise.all( diskOnly.map((id) => limit(async () => ({ id, url: await readVideoUrlFromDisk(path.join(channelDir, "data", id)), })), ), ); roster = mergeRoster(roster, diskObserved, opts.now, "disk"); const total = Object.keys(roster.entries).length; if (total > 0) { await writeRoster(paths, slug, roster); opts.onLog?.( `Roster: seeded ${total} video(s) — ${listed.length} from the stored playlist, ${diskOnly.length} from data/ dirs not in it.\n`, ); } return roster; }