import path from "node:path"; import { readFile, rm, stat } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import type { ChannelConfig } from "../lib/channelConfig"; import type { SiteSettings } from "../lib/settings"; import { getSettings } from "../lib/settings"; import { diskGate, type DiskGateStatus } from "../lib/diskSpace"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; import type { DownloadOutcomeRecord } from "../lib/downloadOutcome"; import { savedVideoPath, type SavedVideoPointer, } from "../lib/savedVideo"; import { loadSavedVideo, restoreSavedVideoProvenance, } from "../lib/savedVideo-server"; import { readChannelConfig } from "./channels"; import { findVideoSourceUrl } from "./undownloadedVideos"; import { downloadOneManaged, sourceFetchFailure, } from "../ytdlp/downloadOneManaged"; import { resolveSourceVideoQuality, VIDEO_720_MAX_HEIGHT, type SourceVideoQuality, } from "../ytdlp/downloadFormat"; import { channelPaceSeconds, channelPlatform, pacingPlatformKey, } from "../ytdlp/channelArgs"; import { platformMinGapSeconds, staticSleepRequestsSeconds, } from "../ytdlp/platformArgs.mjs"; import { downloadGapMs } from "../jobs/platformBackoff"; import { recordDownloadBackoff } from "../jobs/downloadBackoff"; // PERSIST A LIST OF SPECIFIC VIDEOS, across channels, to the saved-video store. // // The general form of persistKept: that pass walks one channel's keep-latest // window, this one walks an explicit `[{slug, id}]` list — a report's cited // videos, a hand-picked set. Each video not yet saved has its source container // re-fetched through downloadOneManaged exactly as "Persist source video" does // (keepSourceVideoOverride + forceMedia), at the requested quality. // // RESUMABLE BY RE-RUNNING. Nothing is remembered between runs: a video already // saved is skipped, so the second run of the same list fetches only what the // first did not finish. That is also why the run stops — not skips — at a disk // floor or a rate limit: the rest of the list is left for the next run. // // PACED AND GATED, ONE VIDEO AT A TIME. Before each download: the gap since the // previous one (`gapMs`, default the batch downloads' own gap), the memory wait // (`minFreeMemMb`), and the disk floor. Cancel and drain are honoured between // videos; the one in flight is a cancel's business, as everywhere else. // // REPLACING NEVER DELETES FIRST. `replace: "above-height"` re-fetches a saved // video whose recorded height is unknown or above the requested quality's // ceiling. The old container stays where it is until the new pointer is // written; only then is it removed (when its path differs — a same-named one // was already replaced by the store's atomic move). The replacement keeps the // old pointer's `keepReason` and `origin`, so a keep-latest container stays // evictable and a requested one keeps its requester. export type PersistVideoItem = { slug: string; id: string }; export const PERSIST_REPLACE_POLICIES = ["never", "above-height"] as const; export type PersistReplacePolicy = (typeof PERSIST_REPLACE_POLICIES)[number]; export function isPersistReplacePolicy(v: unknown): v is PersistReplacePolicy { return ( typeof v === "string" && (PERSIST_REPLACE_POLICIES as readonly string[]).includes(v) ); } export type PersistVideosBucket = { count: number; items: PersistVideoItem[] }; // What a run would do with each item, decided before anything is fetched. // saved already in the store at an acceptable height — left alone // wrongHeight in the store, but its height is unknown or above the ceiling // of the requested quality; re-fetched only under "above-height" // toFetch not saved, with a resolvable source URL // noUrl not saved, and no source URL can be resolved // unknown no such channel, or no data// for that video export type PersistVideosPlan = { saved: PersistVideosBucket; wrongHeight: PersistVideosBucket; toFetch: PersistVideosBucket; noUrl: PersistVideosBucket; unknown: PersistVideosBucket; // How many downloads a run would attempt: toFetch, plus wrongHeight when // the replace policy re-fetches it. willFetch: number; }; export type PersistVideosStop = "low-disk" | "rate-limit" | "cancelled" | "drained"; export type PersistVideosResult = { plan: PersistVideosPlan; dryRun: boolean; // Newly saved this run (replacements included). persisted: PersistVideoItem[]; // Of `persisted`, the ones that replaced an existing container. replaced: PersistVideoItem[]; failed: (PersistVideoItem & { error: string })[]; // Due a download but not attempted because the run stopped first. notAttempted: PersistVideoItem[]; // Why the run ended before its list did; null when it did not. stopped: PersistVideosStop | null; }; // The download one item takes. Injectable so the tests never reach yt-dlp. export type PersistVideoDownload = (opts: { paths: Paths; channelSlug: string; channelConfig: ChannelConfig; settings: SiteSettings; videoUrl: string; quality: SourceVideoQuality; onLog: (line: string) => void; signal: AbortSignal; }) => Promise; export type PersistVideosDeps = { download: PersistVideoDownload; sleep: (ms: number, signal?: AbortSignal) => Promise; // MemAvailable in MiB, or null when it cannot be read. readMemAvailableMb: () => Promise; diskGate: (paths: Paths, settings: SiteSettings) => Promise; getSettings: () => SiteSettings; readChannelConfig: (paths: Paths, slug: string) => Promise; // A rate-limit or network failure, recorded against the platform's pacing. onPlatformBackoff: ( platform: string, paths: Paths, failureClass: "rate_limit" | "network", ) => Promise; removeFile: (file: string) => Promise; }; // How often the memory wait re-reads /proc/meminfo. export const MEM_POLL_MS = 10_000; const defaultDownload: PersistVideoDownload = (o) => downloadOneManaged({ channelSlug: o.channelSlug, channelConfig: o.channelConfig, paths: o.paths, videoUrl: o.videoUrl, onLog: o.onLog, signal: o.signal, cookiePolicy: resolveCookiePolicy(o.settings, o.channelConfig), inlineTranscribeOnFallback: o.settings.inlineTranscribeOnFallback, globalSkipLiveDownloads: o.settings.skipLiveDownloads, appendArchive: true, keepSourceVideoOverride: true, // Persist means persist — see persistKept. forceMedia: true, persistFormatPreset: o.quality, }); function abortableSleep(ms: number, signal?: AbortSignal): Promise { if (ms <= 0 || signal?.aborted) return Promise.resolve(); return new Promise((resolve) => { const onAbort = () => { clearTimeout(t); resolve(); }; const t = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); signal?.addEventListener("abort", onAbort, { once: true }); }); } export async function readMemAvailableMb(): Promise { try { const raw = await readFile("/proc/meminfo", "utf8"); const m = /^MemAvailable:\s+(\d+)\s+kB/m.exec(raw); return m ? Math.floor(Number(m[1]) / 1024) : null; } catch { return null; } } const DEFAULT_DEPS: PersistVideosDeps = { download: defaultDownload, sleep: abortableSleep, readMemAvailableMb, diskGate: (paths, settings) => diskGate(paths, settings), getSettings, readChannelConfig, onPlatformBackoff: (platform, paths, failureClass) => recordDownloadBackoff(platform, paths, failureClass), removeFile: (file) => rm(file, { force: true }), }; // The tallest height a quality accepts, or null for "original" (no ceiling). export function sourceVideoQualityMaxHeight( quality: SourceVideoQuality, ): number | null { return quality === "video_720" ? VIDEO_720_MAX_HEIGHT : null; } // A video id is a directory name under data/: one path segment, never a walk. function isPlainVideoId(id: string): boolean { return ( id !== "" && id !== "." && id !== ".." && !id.includes("/") && !id.includes("\\") && !id.includes("\0") ); } async function isDirectory(p: string): Promise { try { return (await stat(p)).isDirectory(); } catch { return false; } } type PlannedItem = PersistVideoItem & { bucket: Exclude; config?: ChannelConfig; quality?: SourceVideoQuality; url?: string; pointer?: SavedVideoPointer; }; function emptyBucket(): PersistVideosBucket { return { count: 0, items: [] }; } // Classify every item. Duplicates collapse (first occurrence wins the order). async function classify( paths: Paths, items: PersistVideoItem[], format: SourceVideoQuality | undefined, settings: SiteSettings, deps: PersistVideosDeps, ): Promise { const configs = new Map(); const seen = new Set(); const out: PlannedItem[] = []; for (const raw of items) { const item = { slug: raw.slug, id: raw.id }; const key = `${item.slug}\0${item.id}`; if (seen.has(key)) continue; seen.add(key); if (!configs.has(item.slug)) { configs.set(item.slug, await deps.readChannelConfig(paths, item.slug)); } const config = configs.get(item.slug) ?? null; if (!config || !isPlainVideoId(item.id)) { out.push({ ...item, bucket: "unknown" }); continue; } const videoDir = path.join(paths.channelsDir, item.slug, "data", item.id); if (!(await isDirectory(videoDir))) { out.push({ ...item, bucket: "unknown" }); continue; } const quality = resolveSourceVideoQuality({ override: format, channel: config.sourceVideoQuality, global: settings.sourceVideoQuality, }); const pointer = await loadSavedVideo(videoDir); const url = (await findVideoSourceUrl(paths, item.slug, item.id, config)) ?? undefined; if (pointer) { const ceiling = sourceVideoQualityMaxHeight(quality); const height = pointer.format?.height; const tooTall = ceiling !== null && (height === undefined || height > ceiling); out.push({ ...item, bucket: tooTall ? "wrongHeight" : "saved", config, quality, url, pointer, }); continue; } out.push({ ...item, bucket: url ? "toFetch" : "noUrl", config, quality, url }); } return out; } function planOf( planned: PlannedItem[], replace: PersistReplacePolicy, ): PersistVideosPlan { const plan: PersistVideosPlan = { saved: emptyBucket(), wrongHeight: emptyBucket(), toFetch: emptyBucket(), noUrl: emptyBucket(), unknown: emptyBucket(), willFetch: 0, }; for (const p of planned) { plan[p.bucket].items.push({ slug: p.slug, id: p.id }); plan[p.bucket].count += 1; } plan.willFetch = plan.toFetch.count + (replace === "above-height" ? plan.wrongHeight.count : 0); return plan; } function label(item: PersistVideoItem): string { return `${item.slug}/${item.id}`; } export async function persistVideos({ paths, items, format, replace = "never", gapMs, minFreeMemMb = 0, dryRun = false, onLog, signal, drainSignal, deps: depsOverride, }: { paths: Paths; items: PersistVideoItem[]; // Absent = each channel's effective quality (resolveSourceVideoQuality). format?: SourceVideoQuality; replace?: PersistReplacePolicy; // The pause before each download after the first. Absent = the batch // downloads' own gap for the item's channel (downloadGapMs). gapMs?: number; // Before each download, wait until MemAvailable is at least this. 0 = off. minFreeMemMb?: number; dryRun?: boolean; onLog?: (line: string) => void; signal?: AbortSignal; drainSignal?: AbortSignal; deps?: Partial; }): Promise { const deps: PersistVideosDeps = { ...DEFAULT_DEPS, ...depsOverride }; const log = (line: string) => onLog?.(line.endsWith("\n") ? line : `${line}\n`); const downloadLog = (line: string) => onLog?.(line); const downloadSignal = signal ?? new AbortController().signal; const settings = deps.getSettings(); const planned = await classify(paths, items, format, settings, deps); const plan = planOf(planned, replace); const result: PersistVideosResult = { plan, dryRun, persisted: [], replaced: [], failed: [], notAttempted: [], stopped: null, }; log( `Persist videos: ${planned.length} item(s) — ${plan.toFetch.count} to fetch, ` + `${plan.saved.count} already saved, ${plan.wrongHeight.count} saved above the ` + `requested height${replace === "above-height" ? " (to replace)" : " (left alone)"}, ` + `${plan.noUrl.count} with no source URL, ${plan.unknown.count} unknown.`, ); for (const p of planned) { if (p.bucket === "unknown") log(` ${label(p)}: unknown channel or video, skipping.`); else if (p.bucket === "noUrl") log(` ${label(p)}: no resolvable source URL, skipping.`); else if (p.bucket === "wrongHeight" && replace === "above-height" && !p.url) { log(` ${label(p)}: no resolvable source URL to replace it from, skipping.`); } } if (dryRun) return result; const work = planned.filter( (p) => p.bucket === "toFetch" || (p.bucket === "wrongHeight" && replace === "above-height" && p.url), ); let attempted = 0; const stop = (why: PersistVideosStop, from: number) => { result.stopped = why; for (const p of work.slice(from)) { result.notAttempted.push({ slug: p.slug, id: p.id }); } }; const interrupted = (): PersistVideosStop | null => signal?.aborted ? "cancelled" : drainSignal?.aborted ? "drained" : null; for (let i = 0; i < work.length; i++) { const p = work[i]; const config = p.config!; const quality = p.quality!; const why = interrupted(); if (why) { stop(why, i); break; } if (attempted > 0) { const gap = gapMs ?? downloadGapMs( config.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds, channelPaceSeconds(config), staticSleepRequestsSeconds(channelPlatform(config)), { minSeconds: platformMinGapSeconds(channelPlatform(config)) }, ); if (gap > 0) { log(`Sleeping ${gap / 1000}s before the next download...`); await deps.sleep(gap, signal); const after = interrupted(); if (after) { stop(after, i); break; } } } if (minFreeMemMb > 0) { let waited = false; for (;;) { const mem = await deps.readMemAvailableMb(); // Unreadable is not low: a platform with no /proc/meminfo is not held. if (mem === null || mem >= minFreeMemMb) break; if (!waited) { log(` ${label(p)}: ${mem} MiB available, waiting for ${minFreeMemMb} MiB...`); waited = true; } await deps.sleep(MEM_POLL_MS, signal); if (interrupted()) break; } const after = interrupted(); if (after) { stop(after, i); break; } } // PER ITEM, as persistKept: each download writes a full container, so the // tenth must not inherit the first one's headroom. A closed gate ends the // run; the rest are left for the next one. const gate = await deps.diskGate(paths, settings); if (!gate.ok) { log(` ${label(p)}: ${gate.message} — stopping; the rest are left for a later run.`); stop("low-disk", i); break; } const replacing = p.bucket === "wrongHeight"; log( ` ${label(p)}: ${replacing ? "replacing the saved container" : "fetching the source container"} (${quality})…`, ); attempted += 1; const videoDir = path.join(paths.channelsDir, p.slug, "data", p.id); let record: DownloadOutcomeRecord; try { record = await deps.download({ paths, channelSlug: p.slug, channelConfig: config, settings, videoUrl: p.url!, quality, onLog: downloadLog, signal: downloadSignal, }); } catch (e) { result.failed.push({ slug: p.slug, id: p.id, error: (e as Error).message }); log(` ${label(p)}: persist failed — ${(e as Error).message}`); continue; } const failure = sourceFetchFailure(record); if (failure) { result.failed.push({ slug: p.slug, id: p.id, error: failure }); log(` ${label(p)}: persist failed — ${failure}`); const cls = record.failureClass; if (cls === "rate_limit" || cls === "network") { await deps .onPlatformBackoff(pacingPlatformKey(config), paths, cls) .catch(() => {}); log(` ${label(p)}: ${cls.replace("_", " ")} — stopping; the rest are left for a later run.`); stop("rate-limit", i + 1); break; } continue; } // A download that returned is not a saved container: the move into the // store can fail and leave it in the data dir. The pointer is the proof — // and for a replacement, a NEW pointer. const pointer = await loadSavedVideo(videoDir); if (!pointer || (p.pointer && pointer.storedAt === p.pointer.storedAt)) { const error = "the container was not saved to the store"; result.failed.push({ slug: p.slug, id: p.id, error }); log(` ${label(p)}: persist failed — ${error}`); continue; } if (p.pointer) { // The same video's container, so it keeps the old pointer's retention // class and requester: the forced persist writes "override", which // would make a keep-latest container permanent. await restoreSavedVideoProvenance(videoDir, p.pointer); const oldFile = savedVideoPath(p.pointer); if (oldFile !== savedVideoPath(pointer)) { await deps.removeFile(oldFile).catch(() => {}); log(` ${label(p)}: removed the replaced container ${oldFile}.`); } result.replaced.push({ slug: p.slug, id: p.id }); } result.persisted.push({ slug: p.slug, id: p.id }); } log( `Persist videos: ${result.persisted.length} persisted` + (result.replaced.length ? ` (${result.replaced.length} replaced)` : "") + `, ${result.failed.length} failed, ${plan.saved.count} already saved` + (result.notAttempted.length ? `, ${result.notAttempted.length} not attempted (stopped: ${result.stopped})` : "") + ".", ); return result; }