import path from "node:path"; import { readdir } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { loadRawMetadataFromDir } from "../lib/transcripts-server"; import { mapConcurrent } from "../lib/concurrency"; // Metadata reads in flight while keying a channel's videos by upload date. const KEY_READ_CONCURRENCY = 16; // Rolling "keep-latest" window computation. Given a channel's keepLatest config, // returns the ids of the newest N videos (by upload date). Used by both the // cleanup protection (cleanAudioFromTranscribed / channelSnapshot) and the // per-download persistence rule (downloadOneManaged). Kept pure-ish (just fs // reads) so it can be reused everywhere the kept set is needed. export type ComputeKeptOptions = { paths: Paths; channelSlug: string; keepLatest: number; // Runs each video's metadata read. The snapshot passes `onDrive` for a // channel on another drive (lib/storageHealth.ts), which caps the reads in // flight on that drive and gives up on one that does not answer. Default: // the read itself. through?: (read: () => Promise) => Promise; }; // List a channel's data-dir video ids (directories, skipping dotfiles). Mirrors // editor's readDataDirVideoIds, duplicated here so common doesn't depend on the // editor's server-only module. export async function listChannelVideoIds( paths: Paths, channelSlug: string, ): Promise { const dataDir = path.join(paths.channelsDir, channelSlug, "data"); try { const entries = await readdir(dataDir, { withFileTypes: true }); return entries .filter((e) => e.isDirectory() && !e.name.startsWith(".")) .map((e) => e.name); } catch { return []; } } // Recency sort key for a video given its already-known upload_date: the // YYYYMMDD string, falling back to a YYYYMMDD_ prefix on the id, else "" (sorts // oldest). Shared so the live cutoff test for a not-yet-downloaded video keys it // the same way the on-disk window does. export function uploadKeyFor( uploadDate: string | undefined, id: string, ): string { if (uploadDate) return uploadDate; return id.match(/^(\d{8})(?:_|$)/)?.[1] ?? ""; } // Recency sort key for a video: upload_date (YYYYMMDD) from metadata.info.json, // falling back to a YYYYMMDD_ prefix on the dir name, else "" (sorts oldest). async function uploadKey(videoDir: string, id: string): Promise { const meta = await loadRawMetadataFromDir(videoDir); return uploadKeyFor(meta?.upload_date, id); } // Sort a channel's data-dir videos newest-first by upload key (ties broken by id // descending for determinism). Shared by computeKeptVideoIds (the on-disk window) // and computeKeepWindow (the live cutoff used at download time). async function keyedVideosNewestFirst( paths: Paths, channelSlug: string, through: (read: () => Promise) => Promise = (read) => read(), ): Promise> { const ids = await through(() => listChannelVideoIds(paths, channelSlug)); if (ids.length === 0) return []; const dataDir = path.join(paths.channelsDir, channelSlug, "data"); // BOUNDED, like every other corpus-shaped fan-out (lib/concurrency.ts): one // metadata read per video, and the largest channel has eleven thousand. With // a `through` of `onDrive`, at most four of these are on the drive at once // and the rest wait in its queue, which refuses a waiting read only when // nothing on the drive has returned for the watchdog's budget — never for // the queue's depth alone. const keyed = await mapConcurrent(ids, KEY_READ_CONCURRENCY, async (id) => ({ id, key: await through(() => uploadKey(path.join(dataDir, id), id)), })); keyed.sort((a, b) => a.key === b.key ? b.id.localeCompare(a.id) : b.key.localeCompare(a.key), ); return keyed; } // The set of the newest `keepLatest` video ids for a channel (by upload date, // newest first; ties broken by id descending for determinism). Empty when // keepLatest <= 0 or the channel has no data dir. Used by cleanup/snapshot, // which see every kept video already on disk. export async function computeKeptVideoIds({ paths, channelSlug, keepLatest, through, }: ComputeKeptOptions): Promise> { if (!Number.isFinite(keepLatest) || keepLatest <= 0) return new Set(); const keyed = await keyedVideosNewestFirst(paths, channelSlug, through); return new Set(keyed.slice(0, Math.floor(keepLatest)).map((k) => k.id)); } // The keep-latest window expressed as a CUTOFF rather than a fixed id set. This // is what the per-download rule needs: at download time the newest videos aren't // on disk yet, so membership can't come from listing dirs. Instead we capture // - cutoffKey: the upload key of the Nth-newest video currently on disk // - full: whether the channel already holds >= keepLatest videos // and test a candidate by its own upload key (see isInKeepWindow). A candidate // newer than the current Nth would displace it, so it belongs in the window. export type KeepWindow = { keepLatest: number; cutoffKey: string | null; full: boolean; }; export async function computeKeepWindow({ paths, channelSlug, keepLatest, }: ComputeKeptOptions): Promise { const n = Number.isFinite(keepLatest) ? Math.floor(keepLatest) : 0; if (n <= 0) return { keepLatest: 0, cutoffKey: null, full: false }; const keyed = await keyedVideosNewestFirst(paths, channelSlug); const full = keyed.length >= n; const cutoffKey = full ? keyed[n - 1].key : null; return { keepLatest: n, cutoffKey, full }; } // Whether a candidate video (identified by its own upload key) belongs in the // keep-latest window. With fewer than N videos on disk the window isn't full, so // any candidate is kept; once full, a candidate is kept iff it's at least as new // as the current Nth-newest (a tie or newer displaces the boundary video). export function isInKeepWindow( uploadKey: string, window: KeepWindow | undefined, ): boolean { if (!window || window.keepLatest <= 0) return false; if (!window.full) return true; if (window.cutoffKey == null) return true; return uploadKey >= window.cutoffKey; }