Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 53896f7b6e47f77c1244722dd410617ef94e69c0
parent 0c1a529882eb6bdb7037de9f8d1b346a2f019169
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun,  4 Oct 2026 16:25:56 -0400

common: persistVideos — persist an explicit list of videos to the store

persistVideos takes [{slug, id}] across channels, a quality (default each
channel's, through resolveSourceVideoQuality) and a replace policy ("never",
or "above-height": re-fetch a saved container whose recorded height is unknown
or above the quality's ceiling). Before each download it pays the gap (gapMs,
default the batch downloads' own), waits for MemAvailable when minFreeMemMb is
set, and checks the disk floor; a closed floor or a rate limit stops the run,
cancel and drain are honoured between items. A dry run returns the buckets:
saved, wrongHeight, toFetch, noUrl, unknown. A replacement removes the old
container only after the new pointer is written, and keeps the old pointer's
keepReason and origin (restoreSavedVideoProvenance).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
Acommon/controller/persistVideos.test.ts | 501+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/persistVideos.ts | 523+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/savedVideo-server.ts | 19+++++++++++++++++++
3 files changed, 1043 insertions(+), 0 deletions(-)

diff --git a/common/controller/persistVideos.test.ts b/common/controller/persistVideos.test.ts @@ -0,0 +1,501 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, readFile, rm, stat, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import type { ChannelConfig } from "../lib/channelConfig"; +import type { SiteSettings } from "../lib/settings"; +import type { DiskGateStatus } from "../lib/diskSpace"; +import type { DownloadOutcomeRecord } from "../lib/downloadOutcome"; +import { loadSavedVideo, persistSourceVideo } from "../lib/savedVideo-server"; +import { + persistVideos, + type PersistVideoDownload, + type PersistVideosDeps, +} from "./persistVideos"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/persistVideos.test.ts +// +// Every dependency that would reach the network, the clock or the machine's +// memory is injected: the download writes a fake container and persists it the +// way downloadOneManaged does, sleeps are recorded, MemAvailable is scripted. + +const CONFIG: ChannelConfig = { + handling: "transcribe", + platform: "youtube", + url: "https://www.youtube.com/@demo", +}; + +const SETTINGS = { + sourceVideoQuality: "original", + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + inlineTranscribeOnFallback: false, + skipLiveDownloads: false, +} as unknown as SiteSettings; + +const OPEN_GATE: DiskGateStatus = { + ok: true, + enabled: false, + freeBytes: Number.POSITIVE_INFINITY, + thresholdBytes: 0, + resumeBytes: 0, + reason: "ok", + message: "", +}; + +const CLOSED_GATE: DiskGateStatus = { + ...OPEN_GATE, + ok: false, + enabled: true, + freeBytes: 1, + thresholdBytes: 10, + reason: "below-floor", + message: "Low disk space", +} as DiskGateStatus; + +async function withPaths(fn: (paths: Paths) => Promise<void>): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-persistvideos-")); + const paths = { + transcriptsDir: dir, + channelsDir: path.join(dir, "channels"), + savedVideosDir: path.join(dir, "saved"), + } as Paths; + try { + await fn(paths); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +function videoDir(paths: Paths, slug: string, id: string): string { + return path.join(paths.channelsDir, slug, "data", id); +} + +// A downloaded video: its dir and the metadata that names its source URL. +async function seedVideo(paths: Paths, slug: string, id: string): Promise<void> { + const dir = videoDir(paths, slug, id); + await mkdir(dir, { recursive: true }); + await writeFile( + path.join(dir, "metadata.info.json"), + JSON.stringify({ id, webpage_url: `https://www.youtube.com/watch?v=${id}` }), + ); +} + +async function seedSaved( + paths: Paths, + slug: string, + id: string, + opts: { height?: number; file?: string; keepReason?: "keep-latest" | "pin" } = {}, +): Promise<void> { + await seedVideo(paths, slug, id); + const dir = videoDir(paths, slug, id); + const file = opts.file ?? "source-media.mp4"; + await writeFile(path.join(dir, file), `old ${id}`); + await persistSourceVideo({ + videoDir: dir, + sourceFilename: file, + storeDir: path.join(paths.savedVideosDir, slug, id), + keepReason: opts.keepReason ?? "keep-latest", + format: { + preset: "original", + ...(opts.height ? { height: opts.height } : {}), + }, + }); +} + +type Harness = { + deps: Partial<PersistVideosDeps>; + downloads: { slug: string; id: string; quality: string }[]; + sleeps: number[]; + events: string[]; +}; + +// The default harness: every channel exists, the gate is open, memory is +// plentiful, and a download persists a 720p container named `file`. +function harness( + paths: Paths, + over: Partial<PersistVideosDeps> & { file?: string; height?: number } = {}, +): Harness { + const downloads: Harness["downloads"] = []; + const sleeps: number[] = []; + const events: string[] = []; + const download: PersistVideoDownload = async (o) => { + const id = o.videoUrl.split("v=")[1]; + downloads.push({ slug: o.channelSlug, id, quality: o.quality }); + events.push(`download ${id}`); + const dir = videoDir(paths, o.channelSlug, id); + const file = over.file ?? "source-media.mp4"; + await writeFile(path.join(dir, file), `new ${id}`); + await persistSourceVideo({ + videoDir: dir, + sourceFilename: file, + storeDir: path.join(paths.savedVideosDir, o.channelSlug, id), + keepReason: "override", + format: { preset: o.quality, height: over.height ?? 720 }, + }); + return { + videoId: id, + status: "ok", + startedAt: "", + finishedAt: "", + attempts: [], + } satisfies DownloadOutcomeRecord; + }; + return { + downloads, + sleeps, + events, + deps: { + download, + sleep: async (ms) => { + sleeps.push(ms); + events.push(`sleep ${ms}`); + }, + readMemAvailableMb: async () => 99_999, + diskGate: async () => OPEN_GATE, + getSettings: () => SETTINGS, + readChannelConfig: async (_p, slug) => (slug === "nope" ? null : CONFIG), + onPlatformBackoff: async () => {}, + ...over, + }, + }; +} + +test("already-saved items are skipped; the rest are fetched in order", async () => { + await withPaths(async (paths) => { + await seedSaved(paths, "demo-channel", "aaa111"); + await seedVideo(paths, "demo-channel", "bbb222"); + await seedVideo(paths, "other-channel", "ccc333"); + const h = harness(paths); + const res = await persistVideos({ + paths, + items: [ + { slug: "demo-channel", id: "aaa111" }, + { slug: "demo-channel", id: "bbb222" }, + { slug: "other-channel", id: "ccc333" }, + ], + deps: h.deps, + }); + assert.deepEqual( + h.downloads.map((d) => d.id), + ["bbb222", "ccc333"], + ); + assert.equal(res.plan.saved.count, 1); + assert.equal(res.persisted.length, 2); + assert.equal(res.failed.length, 0); + assert.equal(res.stopped, null); + // Re-running the same list is the resume: everything is saved now. + const again = await persistVideos({ + paths, + items: [ + { slug: "demo-channel", id: "aaa111" }, + { slug: "demo-channel", id: "bbb222" }, + { slug: "other-channel", id: "ccc333" }, + ], + deps: harness(paths).deps, + }); + assert.equal(again.plan.saved.count, 3); + assert.equal(again.plan.willFetch, 0); + }); +}); + +test("format defaults to each channel's quality, and an explicit one wins", async () => { + await withPaths(async (paths) => { + await seedVideo(paths, "demo-channel", "aaa111"); + await seedVideo(paths, "hd-channel", "bbb222"); + const configs: Record<string, ChannelConfig> = { + "demo-channel": CONFIG, + "hd-channel": { ...CONFIG, sourceVideoQuality: "video_720" }, + }; + const h = harness(paths, { readChannelConfig: async (_p, s) => configs[s] ?? null }); + await persistVideos({ + paths, + items: [ + { slug: "demo-channel", id: "aaa111" }, + { slug: "hd-channel", id: "bbb222" }, + ], + dryRun: false, + deps: h.deps, + }); + assert.deepEqual( + h.downloads.map((d) => d.quality), + ["original", "video_720"], + ); + + await seedVideo(paths, "demo-channel", "ccc333"); + const h2 = harness(paths, { readChannelConfig: async (_p, s) => configs[s] ?? null }); + await persistVideos({ + paths, + items: [{ slug: "demo-channel", id: "ccc333" }], + format: "video_720", + deps: h2.deps, + }); + assert.deepEqual(h2.downloads.map((d) => d.quality), ["video_720"]); + }); +}); + +test("above-height replaces a too-tall or unmeasured container, keeping its provenance", async () => { + await withPaths(async (paths) => { + await seedSaved(paths, "demo-channel", "tall111", { height: 1080, file: "source-media.webm" }); + await seedSaved(paths, "demo-channel", "nohgt22"); + await seedSaved(paths, "demo-channel", "ok33333", { height: 720 }); + const items = [ + { slug: "demo-channel", id: "tall111" }, + { slug: "demo-channel", id: "nohgt22" }, + { slug: "demo-channel", id: "ok33333" }, + ]; + + // "never" (the default) leaves every saved item alone. + const never = harness(paths); + const left = await persistVideos({ paths, items, format: "video_720", deps: never.deps }); + assert.equal(never.downloads.length, 0); + assert.equal(left.plan.wrongHeight.count, 2); + assert.equal(left.plan.saved.count, 1); + assert.equal(left.plan.willFetch, 0); + + const oldTall = await loadSavedVideo(videoDir(paths, "demo-channel", "tall111")); + const h = harness(paths); + const res = await persistVideos({ + paths, + items, + format: "video_720", + replace: "above-height", + deps: h.deps, + }); + assert.deepEqual( + h.downloads.map((d) => d.id), + ["tall111", "nohgt22"], + ); + assert.equal(res.replaced.length, 2); + const tall = await loadSavedVideo(videoDir(paths, "demo-channel", "tall111")); + assert.equal(tall?.file, "source-media.mp4"); + assert.equal(tall?.format?.height, 720); + // The replacement is still the keep-latest container it replaced. + assert.equal(tall?.keepReason, "keep-latest"); + // The old, differently-named container went only after the new one landed. + await assert.rejects(stat(path.join(oldTall!.dir, oldTall!.file))); + assert.equal( + await readFile(path.join(tall!.dir, tall!.file), "utf8"), + "new tall111", + ); + }); +}); + +test("a replacement that does not land leaves the old container alone", async () => { + await withPaths(async (paths) => { + await seedSaved(paths, "demo-channel", "tall111", { height: 1080, file: "source-media.webm" }); + const h = harness(paths, { + // A download that "succeeds" without moving anything into the store. + download: async () => ({ + videoId: "tall111", + status: "ok", + startedAt: "", + finishedAt: "", + attempts: [], + }), + }); + const res = await persistVideos({ + paths, + items: [{ slug: "demo-channel", id: "tall111" }], + format: "video_720", + replace: "above-height", + deps: h.deps, + }); + assert.equal(res.failed.length, 1); + assert.match(res.failed[0].error, /not saved to the store/); + const ptr = await loadSavedVideo(videoDir(paths, "demo-channel", "tall111")); + assert.equal(await readFile(path.join(ptr!.dir, ptr!.file), "utf8"), "old tall111"); + }); +}); + +test("a closed disk gate stops the run; the rest are not attempted", async () => { + await withPaths(async (paths) => { + for (const id of ["aaa111", "bbb222", "ccc333"]) { + await seedVideo(paths, "demo-channel", id); + } + let calls = 0; + const h = harness(paths, { + diskGate: async () => (++calls >= 2 ? CLOSED_GATE : OPEN_GATE), + }); + const res = await persistVideos({ + paths, + items: ["aaa111", "bbb222", "ccc333"].map((id) => ({ slug: "demo-channel", id })), + deps: h.deps, + }); + assert.deepEqual(h.downloads.map((d) => d.id), ["aaa111"]); + assert.equal(res.stopped, "low-disk"); + assert.deepEqual( + res.notAttempted.map((i) => i.id), + ["bbb222", "ccc333"], + ); + }); +}); + +test("a rate limit stops the run and records the platform's backoff", async () => { + await withPaths(async (paths) => { + for (const id of ["aaa111", "bbb222"]) await seedVideo(paths, "demo-channel", id); + const backoffs: string[] = []; + const h = harness(paths, { + download: async () => ({ + videoId: "aaa111", + status: "failed", + startedAt: "", + finishedAt: "", + attempts: [], + failureClass: "rate_limit", + }), + onPlatformBackoff: async (platform, _p, cls) => { + backoffs.push(`${platform}:${cls}`); + }, + }); + const res = await persistVideos({ + paths, + items: ["aaa111", "bbb222"].map((id) => ({ slug: "demo-channel", id })), + deps: h.deps, + }); + assert.equal(res.stopped, "rate-limit"); + assert.equal(res.failed.length, 1); + assert.deepEqual(res.notAttempted.map((i) => i.id), ["bbb222"]); + assert.deepEqual(backoffs, ["youtube:rate_limit"]); + }); +}); + +test("the gap is paid between downloads, not before the first", async () => { + await withPaths(async (paths) => { + for (const id of ["aaa111", "bbb222", "ccc333"]) await seedVideo(paths, "demo-channel", id); + await seedSaved(paths, "demo-channel", "saved11"); + const h = harness(paths); + await persistVideos({ + paths, + items: ["aaa111", "saved11", "bbb222", "ccc333"].map((id) => ({ + slug: "demo-channel", + id, + })), + gapMs: 5000, + deps: h.deps, + }); + assert.deepEqual(h.events, [ + "download aaa111", + "sleep 5000", + "download bbb222", + "sleep 5000", + "download ccc333", + ]); + }); +}); + +test("the default gap is the batch downloads' own (sleepBetweenDownloadsSeconds)", async () => { + await withPaths(async (paths) => { + for (const id of ["aaa111", "bbb222"]) await seedVideo(paths, "demo-channel", id); + const h = harness(paths, { + readChannelConfig: async () => ({ ...CONFIG, sleepBetweenDownloadsSeconds: 7 }), + }); + await persistVideos({ + paths, + items: ["aaa111", "bbb222"].map((id) => ({ slug: "demo-channel", id })), + deps: h.deps, + }); + assert.deepEqual(h.sleeps, [7000]); + }); +}); + +test("the memory wait holds each download until MemAvailable reaches the floor", async () => { + await withPaths(async (paths) => { + await seedVideo(paths, "demo-channel", "aaa111"); + const readings = [1000, 2000, 5000]; + const h = harness(paths, { + readMemAvailableMb: async () => readings.shift() ?? 5000, + }); + const res = await persistVideos({ + paths, + items: [{ slug: "demo-channel", id: "aaa111" }], + minFreeMemMb: 4096, + deps: h.deps, + }); + assert.deepEqual(h.events, ["sleep 10000", "sleep 10000", "download aaa111"]); + assert.equal(res.persisted.length, 1); + }); +}); + +test("cancel and drain stop between items", async () => { + await withPaths(async (paths) => { + for (const id of ["aaa111", "bbb222", "ccc333"]) await seedVideo(paths, "demo-channel", id); + const items = ["aaa111", "bbb222", "ccc333"].map((id) => ({ slug: "demo-channel", id })); + + const abort = new AbortController(); + const h = harness(paths); + const inner = h.deps.download!; + h.deps.download = async (o) => { + const r = await inner(o); + abort.abort(); + return r; + }; + const cancelled = await persistVideos({ paths, items, signal: abort.signal, deps: h.deps }); + assert.equal(cancelled.stopped, "cancelled"); + assert.equal(cancelled.persisted.length, 1); + assert.deepEqual(cancelled.notAttempted.map((i) => i.id), ["bbb222", "ccc333"]); + + const drain = new AbortController(); + const h2 = harness(paths); + const inner2 = h2.deps.download!; + h2.deps.download = async (o) => { + const r = await inner2(o); + drain.abort(); + return r; + }; + const drained = await persistVideos({ paths, items, drainSignal: drain.signal, deps: h2.deps }); + assert.equal(drained.stopped, "drained"); + assert.deepEqual(h2.downloads.map((d) => d.id), ["bbb222"]); + assert.deepEqual(drained.notAttempted.map((i) => i.id), ["ccc333"]); + }); +}); + +test("a dry run buckets every item and fetches nothing", async () => { + await withPaths(async (paths) => { + await seedSaved(paths, "demo-channel", "saved11", { height: 720 }); + await seedSaved(paths, "demo-channel", "tall111", { height: 1080 }); + await seedVideo(paths, "demo-channel", "fetch11"); + // A video dir with no metadata, on a channel with no platform: no URL. + await mkdir(videoDir(paths, "bare-channel", "nourl11"), { recursive: true }); + const h = harness(paths, { + readChannelConfig: async (_p, slug) => + slug === "nope" + ? null + : slug === "bare-channel" + ? { handling: "transcribe" } + : CONFIG, + }); + const res = await persistVideos({ + paths, + items: [ + { slug: "demo-channel", id: "saved11" }, + { slug: "demo-channel", id: "tall111" }, + { slug: "demo-channel", id: "fetch11" }, + { slug: "bare-channel", id: "nourl11" }, + { slug: "nope", id: "xyz" }, + { slug: "demo-channel", id: "missing" }, + { slug: "demo-channel", id: "../escape" }, + // A duplicate collapses. + { slug: "demo-channel", id: "fetch11" }, + ], + format: "video_720", + replace: "above-height", + dryRun: true, + deps: h.deps, + }); + assert.equal(h.downloads.length, 0); + assert.equal(res.dryRun, true); + assert.deepEqual(res.plan.saved, { count: 1, items: [{ slug: "demo-channel", id: "saved11" }] }); + assert.deepEqual(res.plan.wrongHeight.items, [{ slug: "demo-channel", id: "tall111" }]); + assert.deepEqual(res.plan.toFetch.items, [{ slug: "demo-channel", id: "fetch11" }]); + assert.deepEqual(res.plan.noUrl.items, [{ slug: "bare-channel", id: "nourl11" }]); + assert.deepEqual( + res.plan.unknown.items.map((i) => `${i.slug}/${i.id}`), + ["nope/xyz", "demo-channel/missing", "demo-channel/../escape"], + ); + assert.equal(res.plan.willFetch, 2); + }); +}); diff --git a/common/controller/persistVideos.ts b/common/controller/persistVideos.ts @@ -0,0 +1,523 @@ +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 { 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/<id>/ 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<DownloadOutcomeRecord>; + +export type PersistVideosDeps = { + download: PersistVideoDownload; + sleep: (ms: number, signal?: AbortSignal) => Promise<void>; + // MemAvailable in MiB, or null when it cannot be read. + readMemAvailableMb: () => Promise<number | null>; + diskGate: (paths: Paths, settings: SiteSettings) => Promise<DiskGateStatus>; + getSettings: () => SiteSettings; + readChannelConfig: (paths: Paths, slug: string) => Promise<ChannelConfig | null>; + // A rate-limit or network failure, recorded against the platform's pacing. + onPlatformBackoff: ( + platform: string, + paths: Paths, + failureClass: "rate_limit" | "network", + ) => Promise<void>; + removeFile: (file: string) => Promise<void>; +}; + +// 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<void> { + 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<number | null> { + 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<boolean> { + try { + return (await stat(p)).isDirectory(); + } catch { + return false; + } +} + +type PlannedItem = PersistVideoItem & { + bucket: Exclude<keyof PersistVideosPlan, "willFetch">; + 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<PlannedItem[]> { + const configs = new Map<string, ChannelConfig | null>(); + const seen = new Set<string>(); + 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<PersistVideosDeps>; +}): Promise<PersistVideosResult> { + 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)), + ); + 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; +} diff --git a/common/lib/savedVideo-server.ts b/common/lib/savedVideo-server.ts @@ -151,6 +151,25 @@ export async function updateSavedVideoChecksum( await writePointer(videoDir, { ...pointer, sha256 }); } +// Put a replaced pointer's provenance back on its replacement: the old +// `keepReason` (absent stays absent — a legacy pointer the prune already treats +// as non-evictable) and `origin`. A persist that re-fetched a container is the +// same video's container, not a new decision about why it is kept. No-op when +// there is no pointer. +export async function restoreSavedVideoProvenance( + videoDir: string, + from: Pick<SavedVideoPointer, "keepReason" | "origin">, +): Promise<void> { + const pointer = await loadSavedVideo(videoDir); + if (!pointer) return; + const { keepReason: _k, origin: _o, ...rest } = pointer; + await writePointer(videoDir, { + ...rest, + ...(from.keepReason ? { keepReason: from.keepReason } : {}), + ...(from.origin ?? pointer.origin ? { origin: from.origin ?? pointer.origin } : {}), + }); +} + // Best-effort removal of a now-empty store dir (and its empty <slug> parent). async function pruneEmptyStoreDirs(storeDir: string): Promise<void> { await rm(storeDir, { recursive: false, force: true }).catch(() => {});