import path from "node:path"; import { readdir, readFile, rm, stat } from "node:fs/promises"; import type { Dirent } from "node:fs"; import type { ChannelConfig } from "../lib/channelConfig"; import { channelConfigSchema } from "../lib/channelConfigSchema"; import { readJsonFile, withJsonFileLock, writeJsonAtomic, } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { mapConcurrent } from "../lib/concurrency"; import { isVideoDownloaded, isVideoTranscribed, readVideoFiles, } from "../lib/videoStatus"; import { loadDigest } from "../lib/digest-server"; import { channelTextStall, readRelocationMarker } from "../lib/channelMedia"; import { isDriveNotAnswering, onDrive } from "../lib/storageHealth"; // TYPE-ONLY, and it must stay that way: ./channelSnapshot imports // readChannelConfig from this module, and it drags in the snapshot generator's // whole dependency graph (lmdb, the archive reader, the digest layer). A value // import here would both close an import cycle and make every consumer of // channels.ts pay for the generator. `import type` is erased, so it does // neither — which is why the snapshot *reader* lives down here (below) rather // than beside the generator. import type { ChannelSnapshot } from "./channelSnapshot"; export type ChannelStat = { slug: string; config: ChannelConfig; playlistCount: number | null; videoCount: number; transcriptCount: number; downloadCount: number; // Videos carrying a non-empty ai-digest.json. The Active Jobs progress bar // re-counts `current` from disk rather than trusting the runner, so a digest // job needs this counter to have a bar at all. Optional so a caller reading an // older serialized stat still type-checks. digestCount?: number; }; // A channel slug is also its directory name under transcripts/channels/, so it // must be a safe single path segment. Channel slugs are mixed-case (e.g. // "HasanAbiVODs3") — unlike site ids — so this is broader than SITE_ID_RE, but // it still forbids slashes, dots-only names, and leading punctuation, which // blocks path traversal (no "/", no ".."). export const CHANNEL_SLUG_RE = /^[A-Za-z0-9][A-Za-z0-9._-]*$/; export function isValidChannelSlug(slug: unknown): slug is string { return ( typeof slug === "string" && slug !== ".." && CHANNEL_SLUG_RE.test(slug) ); } // Fan-out widths. Channels number in the dozens, videos in the tens of // thousands, so both lists get the same cap for the same reason (see // ../lib/concurrency) — the channel one just never reaches it. const CHANNEL_READ_CONCURRENCY = 64; const VIDEO_READ_CONCURRENCY = 64; async function exists(p: string): Promise { try { await stat(p); return true; } catch { return false; } } // Cheap existence/coverage check for the channel digestCount below, which runs // over every video dir in a channel. It is identity-BLIND on purpose — it // answers "is there a digest at all?", never "is it current"; the work list is // the digest operation's `state()`. Reads the file (a few KB) rather than // statting, because a digest whose sections are all empty is not coverage. // // PRIVATE, and it lives here rather than in lib/digest-server.ts because this // is its only caller: exported, it read as a general "does this video have a // digest" helper, which is exactly the second definition of "digested" the // registry's state() exists to be the only one of. async function hasDigestWithItems(videoDir: string): Promise { const record = await loadDigest(videoDir); if (!record) return false; return ( (record.sections.chapters?.items.length ?? 0) > 0 || (record.sections.tags?.items.length ?? 0) > 0 ); } // `drive` is the retired `config.dataDir` of a `legacy` channel, the one layout // whose TEXT is on another drive (release 17: a relocated channel's text stays // on the corpus disk, and this walk reads only names and text sidecars, so its // `mediaDir` is never a reason to wrap it — onDrive is keyed by file kind). // For a legacy channel every read goes through `onDrive`: none while that // location is stalled, at most `inFlightPerLocation` (4) in flight on it, and // one that has not answered within the budget (`storage.health.budgetMs`, 3 s // by default) marks it stalled — and the walk answers null ("the drive did // not answer") instead of counts. The rest of the walk is refused without a // call. An in-place channel's walk is on the corpus disk and is not wrapped. async function countDataFiles( dataDir: string, drive?: string, ): Promise<{ videos: number; transcripts: number; downloads: number; digests: number; } | null> { const through = (call: () => Promise): Promise => drive ? onDrive(drive, call) : call(); let dirs: Dirent[]; try { dirs = await through(() => readdir(dataDir, { withFileTypes: true })); } catch (err) { if (isDriveNotAnswering(err)) return null; return { videos: 0, transcripts: 0, downloads: 0, digests: 0 }; } const videoDirs = dirs.filter((d) => d.isDirectory()); // Bounded: the largest channel has 11,224 video dirs and this used to open // them all at once. let flags: Array<{ transcript: boolean; download: boolean; digest: boolean }>; try { flags = await mapConcurrent(videoDirs, VIDEO_READ_CONCURRENCY, (d) => // One video directory's few reads are one call through the watchdog. through(async () => { const dir = path.join(dataDir, d.name); const files = await readVideoFiles(dir); return { transcript: isVideoTranscribed(files), download: isVideoDownloaded(files), // Only transcribed videos can carry a digest, so the sidecar read is // skipped for the rest — the same conditional per-video sidecar-read // pattern channelSnapshot.ts uses for coverage and VTT provenance. digest: isVideoTranscribed(files) ? await hasDigestWithItems(dir) : false, }; }), ); } catch (err) { if (isDriveNotAnswering(err)) return null; throw err; } let transcripts = 0; let downloads = 0; let digests = 0; for (const f of flags) { if (f.transcript) transcripts++; if (f.download) downloads++; if (f.digest) digests++; } return { videos: videoDirs.length, transcripts, downloads, digests }; } // How many videos the channel's `playlist` file names, or null when there is no // such file. ONE small readFile — cheap enough for a render path (the channel // line's first denominator), which is why it is exported rather than folded into // the corpus walk that noCorpusWalkInRenderPaths.test.ts bans. export async function countPlaylist(p: string): Promise { let raw: string; try { raw = await readFile(p, "utf8"); } catch { return null; } let count = 0; for (const line of raw.split("\n")) { if (line.trim().length > 0) count++; } return count; } // THE CHANNEL config.json READER, WRITER AND PATCHER (one-core phase 3 slice // 4b). The shape is lib/channelConfigSchema.ts; everything that reads or // writes a config.json on the server goes through these three, except the two // places that touch the RAW file on purpose: lib/channelMedia.ts's `dataDir` // guard reads it; the legacy `group` / `excludeFromSync` migrations need keys // the schema no longer names, and the latter also writes it back raw. export function channelConfigPath(paths: Paths, slug: string): string { return path.join(paths.channelsDir, slug, "config.json"); } // One config.json by path: the parsed config, or null when the file is absent, // unreadable, not JSON, or not a channel. Never throws. export async function readChannelConfigFile( file: string, ): Promise { const read = await readJsonFile(file); return read.ok ? channelConfigSchema.parse(read.value) : null; } export async function readChannelConfig( paths: Paths, slug: string, ): Promise { return readChannelConfigFile(channelConfigPath(paths, slug)); } export async function channelExists( paths: Paths, slug: string, ): Promise { return exists(path.join(paths.channelsDir, slug, "config.json")); } export async function readChannelStat( paths: Paths, slug: string, ): Promise { const config = await readChannelConfig(paths, slug); if (!config) return null; // A WALK OF `data/` ON A DRIVE THAT IS NOT ANSWERING IS NOT STARTED. The // one-second job-list poll asks this for every channel with a job listed, and // on a stalled drive each readdir and stat in the walk would hold an I/O // thread until the drive came back. No counts is what a caller already // handles (the row draws no progress bar). Only a legacy channel's text is // on a drive that can stall; a stalled MEDIA drive does not stop the count. if (channelTextStall(config)) return null; const channelDir = path.join(paths.channelsDir, slug); const counts = await countDataFiles( path.join(channelDir, "data"), config.dataDir?.trim() || undefined, ); if (!counts) return null; return { slug, config, playlistCount: await countPlaylist(path.join(channelDir, "playlist")), videoCount: counts.videos, transcriptCount: counts.transcripts, downloadCount: counts.downloads, digestCount: counts.digests, }; } async function countMissing( paths: Paths, slug: string, ids: ReadonlyArray, predicate: (files: Awaited>) => boolean, ): Promise { if (ids.length === 0) return 0; const dataDir = path.join(paths.channelsDir, slug, "data"); let missing = 0; for (const id of ids) { const files = await readVideoFiles(path.join(dataDir, id)); if (predicate(files)) missing++; } return missing; } export function countNotYetDownloaded( paths: Paths, slug: string, ids: ReadonlyArray, ): Promise { return countMissing(paths, slug, ids, (files) => !isVideoDownloaded(files)); } export function countNotYetTranscribed( paths: Paths, slug: string, ids: ReadonlyArray, ): Promise { return countMissing(paths, slug, ids, (files) => !isVideoTranscribed(files)); } // ⚠️ WALKS THE ENTIRE CORPUS. One readdir per video directory (78,350 of them // on the live instance) plus a digest sidecar read per transcribed video — // ~474,559 files touched, ~4.4 SECONDS per call. // // This is ground truth, so the batch jobs that must not trust a stale snapshot // (archiveTranscripts, archiveLiveChat, normalizeAll*) still call it. It must // NEVER appear in a page, layout, or API route: this single function, reached // from the root layout's reclaimable-disk badge, was a 4.4 s floor on every // document load in the editor and re-ran every 5 s on the auto-refresh timer. // // For render paths use listChannelBriefs / listChannelConfigs / // listChannelStatsFromSnapshots below, which read 65 small files instead. // A guard test (editor/e2e or common) asserts this name stays out of the app // directory — if you are here to re-add it to a page, that is the sign to stop. export async function listChannelStatsFromDisk( paths: Paths, ): Promise { let entries: Dirent[]; try { entries = await readdir(paths.channelsDir, { withFileTypes: true }); } catch { return []; } const out: ChannelStat[] = []; for (const e of entries) { if (!e.isDirectory()) continue; const slug = e.name; const config = await readChannelConfig(paths, slug); if (!config) continue; const channelDir = path.join(paths.channelsDir, slug); const dataDir = path.join(channelDir, "data"); // A batch job's ground truth: no drive passed, so nothing is raced and the // answer is never null. const counts = (await countDataFiles(dataDir)) ?? { videos: 0, transcripts: 0, downloads: 0, digests: 0, }; out.push({ slug, config, playlistCount: await countPlaylist(path.join(channelDir, "playlist")), videoCount: counts.videos, transcriptCount: counts.transcripts, downloadCount: counts.downloads, // `countDataFiles` has always computed this and the corpus walk has // always discarded it — unlike readChannelStat, which emits it. One line, // and every cross-channel surface (the dashboard, /channels, the // widget) gets a coverage counter it was already paying the I/O for. digestCount: counts.digests, }); } return out.sort((a, b) => a.slug.localeCompare(b.slug)); } // --- Cheap reads: 65 small files instead of half a million ----------------- export const SNAPSHOT_FILENAME = "snapshot.json"; export function snapshotPath(paths: Paths, slug: string): string { return path.join(paths.channelsDir, slug, SNAPSHOT_FILENAME); } // Read one channel's precomputed report. Lives here rather than beside the // generator in ./channelSnapshot so that reading a snapshot doesn't require // loading the machinery that writes one (see the type-only import at the top). // ./channelSnapshot re-exports it, so every existing import still resolves. export async function readChannelSnapshot( paths: Paths, slug: string, ): Promise { try { const raw = await readFile(snapshotPath(paths, slug), "utf8"); return JSON.parse(raw) as ChannelSnapshot; } catch { return null; } } // THE SAME READ, MEMOIZED ON (mtime, size) — for the ONE path that does it on a // timer. // // The auto-queue runner projects every channel's snapshot on each scheduling // tick, and `/api/auto-queue/status` folds FOUR lanes doing the same on a // three-second poll. Since slice 1.2 the digest and backfill lanes draw their // operation work lists from `snapshot.backfill[op].ids` — up to 11,329 strings // on one channel — so that poll went from parsing nothing for two lanes to // parsing all 68 snapshots (~6.5 MB) twice a tick, to produce an object // identical to the one it produced three seconds earlier. // // Keyed on mtimeMs AND size, not on time: a snapshot regenerated by // snapshotScheduler must be seen on the next tick, and it always is — the // generator writes through a rename, so the mtime moves. A stat is ~5 µs // against a parse of ~100 KB. // // DELIBERATELY NOT `readChannelSnapshot` ITSELF. That one is called by the // snapshot GENERATOR and by half a dozen editor render paths, and a shared // parsed object is only safe for callers that treat it as read-only. This // function's callers do (they copy id arrays out and never write back); making // the general reader share objects would be a much larger claim about code that // has never had to honour it. const snapshotMemo = new Map< string, { mtimeMs: number; size: number; snapshot: ChannelSnapshot } >(); // One entry per channel; the corpus has 68. The cap only bounds a pathological // corpus, and dropping the whole map is correct because every entry is a cache. const SNAPSHOT_MEMO_CAP = 512; export async function readChannelSnapshotShared( paths: Paths, slug: string, ): Promise { const file = snapshotPath(paths, slug); try { const st = await stat(file); const hit = snapshotMemo.get(file); if (hit && hit.mtimeMs === st.mtimeMs && hit.size === st.size) { return hit.snapshot; } const snapshot = JSON.parse(await readFile(file, "utf8")) as ChannelSnapshot; if (snapshotMemo.size >= SNAPSHOT_MEMO_CAP) snapshotMemo.clear(); snapshotMemo.set(file, { mtimeMs: st.mtimeMs, size: st.size, snapshot }); return snapshot; } catch { return null; } } // Test seam. Nothing in production drops the memo — an mtime change already // invalidates an entry — but a test that writes a snapshot twice inside one // filesystem timestamp tick needs to. export function resetChannelSnapshotMemo(): void { snapshotMemo.clear(); } // Slug + config + the precomputed snapshot for every channel. One readdir plus // two small file reads per channel, all concurrent: ~60 ms against the corpus // the walk above needs 4.4 s for. export type ChannelBrief = { slug: string; config: ChannelConfig; snapshot: ChannelSnapshot | null; }; async function listChannelSlugs(paths: Paths): Promise { try { const entries = await readdir(paths.channelsDir, { withFileTypes: true }); return entries.filter((e) => e.isDirectory()).map((e) => e.name); } catch { return []; } } export async function listChannelBriefs(paths: Paths): Promise { const slugs = await listChannelSlugs(paths); const briefs = await mapConcurrent( slugs, CHANNEL_READ_CONCURRENCY, async (slug): Promise => { const [config, snapshot] = await Promise.all([ readChannelConfig(paths, slug), readChannelSnapshot(paths, slug), ]); return config ? { slug, config, snapshot } : null; }, ); return briefs .filter((b): b is ChannelBrief => b !== null) .sort((a, b) => a.slug.localeCompare(b.slug)); } // Slug + config only — for the many callers that just need names, URLs and // sync/handling flags. One readdir plus one config read per channel. export async function listChannelConfigs( paths: Paths, ): Promise> { const slugs = await listChannelSlugs(paths); const rows = await mapConcurrent( slugs, CHANNEL_READ_CONCURRENCY, async (slug) => { const config = await readChannelConfig(paths, slug); return config ? { slug, config } : null; }, ); return rows .filter((r): r is { slug: string; config: ChannelConfig } => r !== null) .sort((a, b) => a.slug.localeCompare(b.slug)); } // Digest coverage for one channel, from its snapshot. // // ⚠️ The one field that is NOT a clean substitution for the walk. `totals` has // videos/transcribed/downloaded but no `digests` key, so the count has to be // summed out of `digestEngines` — which 11 of the 65 live snapshots don't have // at all (it postdates them). Those default to 0 and will read as "no digests" // until the channel's next report refresh, rather than reporting a wrong number. export function digestCountOf(snapshot: ChannelSnapshot | null): number { const engines = snapshot?.digestEngines; if (!engines) return 0; let n = 0; for (const count of Object.values(engines)) n += count; return n; } // ChannelStat projected from the precomputed snapshots instead of a corpus // walk. Same shape, ~40× cheaper, and one honest trade: the counts are what the // LAST SNAPSHOT saw, not what is on disk this millisecond. The snapshot // scheduler regenerates on a ~1 s debounce after any report-changing action, so // the lag is that debounce — surfaced in the UI via `snapshot.generatedAt` // rather than left for the reader to assume. export async function listChannelStatsFromSnapshots( paths: Paths, // Pass briefs you have already read to avoid a second pass — a caller that // also wants `generatedAt` for a freshness readout needs the same rows. preread?: ReadonlyArray, ): Promise { const briefs = preread ?? (await listChannelBriefs(paths)); // playlistCount has no snapshot equivalent, but it is one small file per // channel — cheap enough to keep reading directly. const playlistCounts = await mapConcurrent( briefs, CHANNEL_READ_CONCURRENCY, (b) => countPlaylist(path.join(paths.channelsDir, b.slug, "playlist")), ); return briefs.map((b, i) => ({ slug: b.slug, config: b.config, playlistCount: playlistCounts[i], videoCount: b.snapshot?.totals.videos ?? 0, transcriptCount: b.snapshot?.totals.transcribed ?? 0, downloadCount: b.snapshot?.totals.downloaded ?? 0, digestCount: digestCountOf(b.snapshot), })); } // WRITE STAYS STRICT, like writeSettings: the config is parsed on the way out // — which is what drops an unknown or invalid key — and a value that is not a // channel at all (no valid `handling`) THROWS rather than writing a file every // reader would then treat as absent. Every caller passes a parsed config, so // this only ever fires on a bug. The e2e fixtures seed their configs raw, on // purpose, and are unaffected. Returns what it wrote. export async function writeChannelConfig( paths: Paths, slug: string, config: ChannelConfig, ): Promise { const parsed = channelConfigSchema.parse(config); if (!parsed) { throw new Error( `Refusing to write channels/${slug}/config.json: not a channel config ` + `(handling ${JSON.stringify((config as { handling?: unknown })?.handling)})`, ); } await writeJsonAtomic(channelConfigPath(paths, slug), parsed, { mkdir: true }); return parsed; } export type PatchChannelConfigOptions = { // Keys to delete before the patch is applied — how a caller says "clear // this override" (a key absent from `patch` is left alone). unset?: ReadonlyArray; }; // READ-MODIFY-WRITE ONE CHANNEL'S CONFIG, and the only way to change some keys // of it: re-read the file NOW (not a copy the caller read earlier — that was // the stale-spread bug, where a fetch that took minutes wrote back the config // it had read at its start over every edit made meanwhile), delete `unset`, // assign `patch`, write strictly. A channel with no readable config is not // created: the answer is null and nothing is written (the old sync-stamp rule). // Two patches to one channel in this process serialise, so both land. export async function patchChannelConfig( paths: Paths, slug: string, patch: Partial, opts: PatchChannelConfigOptions = {}, ): Promise { const file = channelConfigPath(paths, slug); return withJsonFileLock(file, async () => { const current = await readChannelConfigFile(file); if (!current) return null; const next: ChannelConfig = { ...current }; for (const key of opts.unset ?? []) delete next[key]; Object.assign(next, patch); return writeChannelConfig(paths, slug, next); }); } export async function createChannel( paths: Paths, slug: string, config: ChannelConfig, ): Promise { if (await channelExists(paths, slug)) { throw new Error(`Channel ${slug} already exists`); } await writeChannelConfig(paths, slug, config); } // DELETING A RELOCATED CHANNEL HAS TO REACH THE OTHER DRIVE. `rm -r` on the // channel dir removes the symlink, not what it points at (node's rm does not // follow links), so without this the media would survive the channel as an // orphan nothing in the app can see or reclaim. The target is removed FIRST, so // a target that cannot be removed for a reason OTHER than being absent takes // the delete down with it and leaves the channel intact and retryable. An // unmounted drive is NOT that case — `force: true` makes a missing path a // no-op, and the channel is deleted with its media left on the platter. That is // deliberate: refusing would make a channel whose drive is gone for good // undeletable, and the marker check below is the refusal that matters. export async function deleteChannel( paths: Paths, slug: string, ): Promise { const marker = await readRelocationMarker(paths, slug); if (marker) { throw new Error( `Channel "${slug}" has a media relocation in progress (phase ` + `"${marker.phase}", target ${marker.target}). Finish or cancel it ` + `before deleting the channel.`, ); } const dir = path.join(paths.channelsDir, slug); const config = await readChannelConfig(paths, slug); // The media tier's target (`mediaDir`, release 17) and, on a channel never // migrated off the retired whole-directory layout, its `dataDir` — each on // another drive, each removed before the channel dir. An in-place `media/` // is a real directory inside `dir` and goes with it. for (const target of [config?.mediaDir?.trim(), config?.dataDir?.trim()]) { if (!target) continue; await rm(target, { recursive: true, force: true }); // / is ours by construction (relocatedMediaDir fixes the // suffix), so take it too — but only when nothing else landed in it. const slugRoot = path.dirname(target); if (path.basename(slugRoot) === slug) { const left = await readdir(slugRoot).catch(() => ["keep"]); if (left.length === 0) await rm(slugRoot, { recursive: true, force: true }); } } await rm(dir, { recursive: true, force: true }); }