// PLAYABLE COPIES AND THEIR TORRENTS (release 21 D3, the `prepare-playable` // job). // // For each of a channel's videos with a saved source container: a browser- // playable copy made by a LOSSLESS remux (lib/playableMedia.ts decides the // container and the ffmpeg arguments), and one single-file torrent of that // copy. Written under the playable root — // // ///. the copy // ///.torrent its torrent // //playable.json the manifest (lib/playableMedia.ts) // // — where is, unless named, `playable/` BESIDE the saved-video store's // real directory (the store is on the archive drive; so are the copies). // // KEYED BY THE SOURCE'S SHA-256, so a re-run is a no-op: an entry whose // `sourceSha256` matches the pointer's and whose files are on disk at their // recorded size is left alone. A pointer with no hash (a download's, before // the backup step reached it) is hashed once and the hash recorded on the // pointer, as the backup step would. // // THE TORRENT. `create-torrent`, from the WebTorrent ecosystem: one file, // named `.`; piece length scaled to the size (256 KiB–1 MiB); no web // seed (`urlList`), no comment, no `created by`, a fixed creation date — the // .torrent names nobody and no private URL. The announce list is the // `trackers` parameter (none by default): it is outside the info dict, so the // infohash depends only on the bytes and the piece length, and a site can // announce the same infohash to other trackers later without a re-hash. // // Each copy is written to a `.part` and renamed; the manifest is rewritten // after each video, so a cancelled or drained run resumes where it stopped. /// import path from "node:path"; import { createReadStream } from "node:fs"; import { mkdir, readdir, readFile, realpath, rename, rm, stat, writeFile } from "node:fs/promises"; import { createHash } from "node:crypto"; import { execa } from "execa"; import createTorrent from "create-torrent"; import parseTorrent from "parse-torrent"; import type { Paths } from "../lib/paths"; import type { ChannelConfig } from "../lib/channelConfig"; import { savedVideoPath, savedVideoRoot } from "../lib/savedVideo"; import { loadSavedVideo, updateSavedVideoChecksum } from "../lib/savedVideo-server"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import { getFreeBytes } from "../lib/diskSpace"; import { formatBytes } from "../lib/format"; import { PLAYABLE_MANIFEST_FILENAME, decidePlayable, parsePlayableManifest, parseProbe, pieceLengthFor, type MediaProbe, type PlayableEntry, type PlayableManifest, } from "../lib/playableMedia"; // The creation date every torrent carries: fixed, so the same copy makes the // same .torrent file byte for byte (the infohash never depended on it). export const PLAYABLE_TORRENT_CREATION_DATE = new Date("2026-01-01T00:00:00Z"); export const PLAYABLE_DISK_MARGIN_BYTES = 2 * 1024 ** 3; export async function defaultPlayableRoot( paths: Paths, config: Pick | null, ): Promise { const store = savedVideoRoot(paths, config); const real = await realpath(store).catch(() => store); return path.join(path.dirname(real), "playable"); } export function playableManifestPath(root: string, slug: string): string { return path.join(root, slug, PLAYABLE_MANIFEST_FILENAME); } export async function loadPlayableManifest(root: string, slug: string): Promise { try { return parsePlayableManifest(JSON.parse(await readFile(playableManifestPath(root, slug), "utf8")), slug); } catch { return parsePlayableManifest(null, slug); } } // --- the torrent ------------------------------------------------------------ export type TorrentMade = { torrent: Buffer; infoHash: string; pieceLength: number }; // One single-file torrent of `file`, named `name`. Deterministic: the same // bytes and piece length give the same infohash, whatever `trackers` is. export async function makePlayableTorrent(opts: { file: string; name: string; bytes: number; trackers?: readonly string[]; pieceLength?: number; }): Promise { const pieceLength = opts.pieceLength ?? pieceLengthFor(opts.bytes); // What is NOT passed matters as much as what is. `createdBy`, `comment` // and `urlList` absent: create-torrent then writes no such key. `private` // absent: even `false` would put `private: 0` INSIDE the info dict and so // change the infohash — the one thing that must not move when the torrent // goes public later. const torrent = await new Promise((resolve, reject) => { createTorrent( opts.file, { name: opts.name, pieceLength, creationDate: PLAYABLE_TORRENT_CREATION_DATE, // A list, even an empty one, not create-torrent's default public // trackers: the pilot announces only where it is told to. A fresh // array each call (create-torrent appends to the one it is given). announceList: opts.trackers && opts.trackers.length ? [[...opts.trackers]] : [], }, (err, buf) => (err ? reject(err) : resolve(Buffer.from(buf))), ); }); const parsed = await parseTorrent(torrent); if (!parsed.infoHash) throw new Error(`no infohash in the torrent made for ${opts.name}`); return { torrent, infoHash: parsed.infoHash, pieceLength }; } // --- the run ---------------------------------------------------------------- export type PreparePlayableOptions = { paths: Paths; slug: string; channelConfig: ChannelConfig; root?: string; ids?: string[]; trackers?: string[]; dryRun?: boolean; onLog: (line: string) => void; signal?: AbortSignal; drainSignal?: AbortSignal; deps?: { now?: () => Date; probe?: (file: string) => Promise; remux?: (src: string, dest: string, args: string[], signal?: AbortSignal) => Promise; }; }; export type PreparePlayableResult = { root: string; prepared: string[]; unchanged: string[]; unplayable: { id: string; reason: string }[]; noContainer: string[]; failed: { id: string; error: string }[]; stopped?: string; dryRun: boolean; }; async function sha256File(file: string, signal?: AbortSignal): Promise { const h = createHash("sha256"); for await (const chunk of createReadStream(file, { highWaterMark: 1 << 20, ...(signal ? { signal } : {}) })) { h.update(chunk as Buffer); } return h.digest("hex"); } async function ffprobe(paths: Paths, file: string, signal?: AbortSignal): Promise { const r = await execa( paths.ffprobeBin, ["-v", "error", "-show_streams", "-show_format", "-of", "json", file], { ...(signal ? { cancelSignal: signal } : {}) }, ); return parseProbe(JSON.parse(String(r.stdout)), path.extname(file)); } // EACH REMUX TAKES THE HEAVY SLOT (release 19 B1, scripts/queue-lock.mjs // --heavy — the gate `pnpm heavy` and a publish stage's `next build` use): // one heavy job machine-wide, above the memory floor. Per video, not per run, // so a 36 GB channel interleaves with builds and e2e instead of holding the // slot for hours. The gate forwards a cancel to ffmpeg. A checkout without // the gate (a test's temp root) runs ffmpeg directly; HEAVY=0 bypasses it. async function ffmpegRemux(paths: Paths, src: string, dest: string, args: string[], signal?: AbortSignal): Promise { const ffmpegArgs = ["-nostdin", "-hide_banner", "-loglevel", "error", "-y", "-i", src, ...args, dest]; const gate = paths.monorepoRoot ? path.join(paths.monorepoRoot, "scripts", "queue-lock.mjs") : null; const gated = gate !== null && (await stat(gate).then((s) => s.isFile()).catch(() => false)); await execa( gated ? process.execPath : paths.ffmpegBin, gated ? [gate, "--heavy", "--", paths.ffmpegBin, ...ffmpegArgs] : ffmpegArgs, { ...(signal ? { cancelSignal: signal } : {}) }, ); } async function sizeOf(file: string): Promise { const st = await stat(file).catch(() => null); return st?.isFile() ? st.size : null; } async function heldWithPointers(dataDir: string): Promise { const entries = await readdir(dataDir, { withFileTypes: true }).catch((err: NodeJS.ErrnoException) => { if (err.code === "ENOENT") return []; throw err; }); const out: string[] = []; for (const d of entries) { if (!d.isDirectory() || d.name.startsWith(".")) continue; if (await loadSavedVideo(path.join(dataDir, d.name))) out.push(d.name); } return out.sort(); } export async function preparePlayable(opts: PreparePlayableOptions): Promise { const onLog = (s: string) => opts.onLog(s.endsWith("\n") ? s : `${s}\n`); const now = opts.deps?.now ?? (() => new Date()); const probe = opts.deps?.probe ?? ((f: string) => ffprobe(opts.paths, f, opts.signal)); const remux = opts.deps?.remux ?? ((src: string, dest: string, args: string[], signal?: AbortSignal) => ffmpegRemux(opts.paths, src, dest, args, signal)); const root = opts.root ?? (await defaultPlayableRoot(opts.paths, opts.channelConfig)); const result: PreparePlayableResult = { root, prepared: [], unchanged: [], unplayable: [], noContainer: [], failed: [], dryRun: !!opts.dryRun, }; // The root's parent must exist: a copy aimed at an unmounted drive is never // materialised on the root filesystem by a mkdir. const parent = await stat(path.dirname(root)).catch(() => null); if (!parent?.isDirectory()) { throw new Error(`${path.dirname(root)} is not there (is its drive mounted?) — the playable root ${root} cannot be made`); } const dataDir = path.join(opts.paths.channelsDir, opts.slug, "data"); const ids = opts.ids ?? (await heldWithPointers(dataDir)); const slugDir = path.join(root, opts.slug); const manifest = await loadPlayableManifest(root, opts.slug); onLog(`${ids.length} video(s) with a saved container; playable root ${root}`); let n = 0; for (const id of ids) { n += 1; if (opts.signal?.aborted) { result.stopped = "cancelled"; break; } if (opts.drainSignal?.aborted) { result.stopped = "drained"; onLog(`Drained: ${ids.length - n + 1} left for the next run.`); break; } const videoDir = path.join(dataDir, id); const pointer = await loadSavedVideo(videoDir); const src = pointer ? savedVideoPath(pointer) : null; const srcBytes = src ? await sizeOf(src) : null; if (!pointer || !src || srcBytes === null) { result.noContainer.push(id); onLog(`${id}: no saved container${pointer ? ` (${src} is missing)` : ""}`); continue; } try { let sourceSha256 = pointer.sha256; if (!sourceSha256) { onLog(`${id}: hashing ${pointer.file} (${formatBytes(srcBytes)})`); sourceSha256 = await sha256File(src, opts.signal); if (!opts.dryRun) await updateSavedVideoChecksum(videoDir, sourceSha256); } const prior = manifest.entries[id]; if (prior && prior.sourceSha256 === sourceSha256) { const [fb, tb] = await Promise.all([ sizeOf(path.join(slugDir, prior.file)), sizeOf(path.join(slugDir, prior.torrent)), ]); if (fb === prior.bytes && tb !== null) { result.unchanged.push(id); continue; } } const decision = decidePlayable(await probe(src)); if (!decision.playable) { result.unplayable.push({ id, reason: decision.reason }); onLog(`${id}: not playable — ${decision.reason}`); continue; } if (opts.dryRun) { onLog(`[${n}/${ids.length}] ${id}: would remux — ${decision.reason}`); result.prepared.push(id); continue; } const free = await getFreeBytes(path.dirname(root)); if (free < srcBytes + PLAYABLE_DISK_MARGIN_BYTES) { result.stopped = `the playable root's disk has ${formatBytes(free)} free; ${id} needs about ${formatBytes(srcBytes)} plus ${formatBytes(PLAYABLE_DISK_MARGIN_BYTES)}`; onLog(`Stopped: ${result.stopped}. A re-run resumes.`); break; } const outDir = path.join(slugDir, id); await mkdir(outDir, { recursive: true }); const name = `${id}.${decision.ext}`; const out = path.join(outDir, name); const part = path.join(outDir, `.${name}.part`); onLog(`[${n}/${ids.length}] ${id}: ${decision.reason}`); await rm(part, { force: true }); try { await remux(src, part, decision.args, opts.signal); await rename(part, out); } catch (err) { await rm(part, { force: true }).catch(() => {}); throw err; } // A copy of another extension from an earlier run is not this video's any more. for (const f of await readdir(outDir)) { if (f !== name && f !== `${id}.torrent` && f.startsWith(`${id}.`)) await rm(path.join(outDir, f), { force: true }); } const bytes = (await sizeOf(out)) ?? 0; const made = await makePlayableTorrent({ file: out, name, bytes, ...(opts.trackers ? { trackers: opts.trackers } : {}), }); const torrentPath = path.join(outDir, `${id}.torrent`); await writeFile(`${torrentPath}.part`, made.torrent); await rename(`${torrentPath}.part`, torrentPath); const probed = await probe(out).catch(() => null); const entry: PlayableEntry = { sourceSha256, file: `${id}/${name}`, torrent: `${id}/${id}.torrent`, ext: decision.ext, mime: decision.mime, bytes, infoHash: made.infoHash, pieceLength: made.pieceLength, vcodec: probed?.vcodec ?? "", acodec: probed?.acodec ?? null, preparedAt: now().toISOString(), }; manifest.entries[id] = entry; await writeJsonAtomic(playableManifestPath(root, opts.slug), manifest); result.prepared.push(id); onLog(`${id}: ${name} ${formatBytes(bytes)}, infohash ${made.infoHash}, ${made.pieceLength / 1024} KiB pieces`); } catch (err) { if (opts.signal?.aborted) { result.stopped = "cancelled"; break; } result.failed.push({ id, error: (err as Error).message }); onLog(`${id}: FAILED — ${(err as Error).message}`); } } onLog( `summary: ${JSON.stringify({ dryRun: result.dryRun, root, prepared: result.prepared.length, unchanged: result.unchanged.length, unplayable: result.unplayable, noContainer: result.noContainer, failed: result.failed, ...(result.stopped ? { stopped: result.stopped } : {}), })}`, ); return result; }