// ATTACH LOCAL MEDIA TO HELD VIDEOS (release 21 D1, the `attach-media` job). // // A channel whose platform copies are gone (a deleted YouTube channel) can // still be clipped when somebody archived its videos: this takes such an // archive — a directory, a .zip read IN PLACE, or a .7z — finds the video file // for each held record (lib/attachMediaPlan.ts says how an id is read), and // puts it in the saved-video store as that video's source container, so a // clip window, a report's evidence cut and report-to-video all cut it // locally. Nothing is fetched. // // THE ONE WRITER. The container lands through `persistSourceVideo`, as every // persisted source does: `source-media.` in the store, a pointer in // `data//saved-video.json`, `keepReason: "pin"` (never pruned), and an // origin `{kind: "local-archive", archive, entry, sha256, attachedAt}` — which // archive, which entry, and the hash of the bytes. The bytes are copied into a // staging file on the STORE's disk first (a stored zip entry is a byte range // of the archive; a deflated one is inflated; a 7z entry is extracted by `7z`) // and then renamed into place, so 36 GB never passes through the corpus SSD. // // NEVER WRITES TO THE ARCHIVE. Every read of it is read-only. // // REFUSES OVER AN EXISTING POINTER unless `replace` — a container already in // the store, from a download or an earlier attach, is somebody's decision. A // replaced container that had another name is removed once the new pointer is // written (a same-named one was replaced by the rename). // // NOT HELD. A video in the archive the channel has no record of is listed; // with `createRecords` a record is written first — `metadata.info.json` from // the folder's own yt-dlp `.info.json`, else built from its description.txt // and name (lib/attachMediaPlan.ts `archiveFolderInfoJson`) — through the // metadata history, and the channel's `archive` file gains its line, as a // download would. Transcribing it is a separate step. // // GUARDED: the job kind declares `needsMedia` (runManagedFunction asks the // channel's media guard), every store write asks // `assertSavedVideosStoreWritable`, the store's root must exist (an unmounted // drive is never materialised by a mkdir), and each copy needs its size plus a // margin free on the store's disk — short of it the run stops, and a re-run // resumes (an attached video is "already attached"). // // DRAIN AND CANCEL are honoured between videos; a cancel also aborts the copy // in flight and removes its staging file. import path from "node:path"; import { appendFile, mkdir, readFile, readdir, rm, stat } from "node:fs/promises"; import { createWriteStream } from "node:fs"; import { pipeline } from "node:stream/promises"; import { Transform, type TransformCallback } from "node:stream"; import { createHash } from "node:crypto"; import { execa } from "execa"; import type { Paths } from "../lib/paths"; import type { ChannelConfig } from "../lib/channelConfig"; import { savedVideoDir, savedVideoPath, savedVideoRoot, type SavedVideoOrigin, } from "../lib/savedVideo"; import { loadSavedVideo, persistSourceVideo } from "../lib/savedVideo-server"; import { assertSavedVideosStoreWritable } from "../lib/savedVideoStore"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import { withMetadataHistory } from "../lib/metadataHistory-server"; import { getFreeBytes } from "../lib/diskSpace"; import { formatBytes } from "../lib/format"; import { probeMediaDurationSec } from "../ytdlp/ffprobeDuration"; import { archiveFolderInfoJson, countAttachRows, groupArchiveFolders, planAttach, type AttachItem, type AttachPlan, type AttachRow, type AttachRowClass, type SourceFile, } from "../lib/attachMediaPlan"; import { copyFileHashed, copyZipEntry, readZipEntries, readZipEntryText, type CopiedEntry, type ZipEntry, } from "../lib/zipReader"; export const ATTACH_MEDIA_REQUESTER = "attach-media"; // Free space the store's disk must keep beyond each copy. export const ATTACH_MEDIA_DISK_MARGIN_BYTES = 2 * 1024 ** 3; // --- sources ---------------------------------------------------------------- export type ArchiveSource = { kind: "zip" | "7z" | "dir"; // The archive's absolute path (what the origin records). archive: string; files: SourceFile[]; readText(file: SourceFile): Promise; copyOut(file: SourceFile, dest: string, signal?: AbortSignal): Promise; }; async function openZipSource(archive: string): Promise { const entries = await readZipEntries(archive); const byName = new Map(); for (const e of entries) if (!e.isDirectory) byName.set(e.name, e); const entry = (f: SourceFile) => { const e = byName.get(f.path); if (!e) throw new Error(`${f.path} is not in ${archive}`); return e; }; return { kind: "zip", archive, files: [...byName.values()].map((e) => ({ path: e.name, size: e.size })), readText: (f) => readZipEntryText(archive, entry(f)), copyOut: (f, dest, signal) => copyZipEntry(archive, entry(f), dest, signal ? { signal } : {}), }; } async function walkDir(root: string, rel = ""): Promise { const out: SourceFile[] = []; const entries = await readdir(path.join(root, rel), { withFileTypes: true }); for (const d of entries.sort((a, b) => a.name.localeCompare(b.name))) { const p = rel ? `${rel}/${d.name}` : d.name; if (d.isDirectory()) out.push(...(await walkDir(root, p))); else if (d.isFile() || d.isSymbolicLink()) { const st = await stat(path.join(root, p)).catch(() => null); if (st?.isFile()) out.push({ path: p, size: st.size }); } } return out; } async function openDirSource(dir: string): Promise { const files = await walkDir(dir); const abs = (f: SourceFile) => path.join(dir, ...f.path.split("/")); return { kind: "dir", archive: dir, files, readText: async (f) => { if (f.size > 4 * 1024 * 1024) throw new Error(`${f.path} is too large to read as text`); return readFile(abs(f), "utf8"); }, copyOut: (f, dest, signal) => copyFileHashed(abs(f), dest, signal ? { signal } : {}), }; } class HashCount extends Transform { bytes = 0; readonly hash = createHash("sha256"); _transform(chunk: Buffer, _e: BufferEncoding, cb: TransformCallback): void { this.bytes += chunk.length; this.hash.update(chunk); cb(null, chunk); } } // `7z l -slt` output → files. Blocks are separated by blank lines; the // archive's own block (before "----------") is skipped. export function parse7zListing(text: string): SourceFile[] { const body = text.split(/^-{10,}\s*$/m).slice(1).join("\n"); const out: SourceFile[] = []; for (const block of body.split(/\r?\n\s*\r?\n/)) { const field = (k: string) => new RegExp(`^${k} = (.*)$`, "m").exec(block)?.[1]; const p = field("Path"); if (!p) continue; const folder = field("Folder"); const attrs = field("Attributes") ?? ""; if (folder === "+" || attrs.startsWith("D")) continue; out.push({ path: p.replace(/\\/g, "/"), size: Number(field("Size") ?? 0) || 0 }); } return out; } async function open7zSource(archive: string, sevenZipBin: string): Promise { const listed = await execa(sevenZipBin, ["l", "-slt", archive], { reject: false }); if (listed.exitCode !== 0) { throw new Error(`${sevenZipBin} could not list ${archive}: ${String(listed.stderr).trim() || `exit ${listed.exitCode}`}`); } const files = parse7zListing(String(listed.stdout)); // `-spd`: the entry name is a name, not a wildcard pattern. const extract = (f: SourceFile, signal?: AbortSignal) => execa(sevenZipBin, ["e", "-so", "-spd", archive, f.path], { buffer: false, stdout: "pipe", ...(signal ? { cancelSignal: signal } : {}), }); return { kind: "7z", archive, files, readText: async (f) => { if (f.size > 4 * 1024 * 1024) throw new Error(`${f.path} is too large to read as text`); const r = await execa(sevenZipBin, ["e", "-so", "-spd", archive, f.path], { encoding: "utf8" }); return String(r.stdout); }, copyOut: async (f, dest, signal) => { const child = extract(f, signal); const hc = new HashCount(); await Promise.all([ pipeline(child.stdout!, hc, createWriteStream(dest, { flags: "wx" })), child, ]); if (hc.bytes !== f.size) { throw new Error(`${f.path}: ${hc.bytes} bytes extracted, the archive says ${f.size}`); } return { bytes: hc.bytes, sha256: hc.hash.digest("hex") }; }, }; } export async function openArchiveSource( source: string, opts: { sevenZipBin?: string } = {}, ): Promise { if (!path.isAbsolute(source)) { throw new Error(`"${source}" is not an absolute path`); } const st = await stat(source).catch(() => null); if (!st) throw new Error(`${source} does not exist (is its drive mounted?)`); if (st.isDirectory()) return openDirSource(source); const lower = source.toLowerCase(); if (lower.endsWith(".zip")) return openZipSource(source); if (lower.endsWith(".7z")) return open7zSource(source, opts.sevenZipBin ?? "7z"); throw new Error(`${source}: a source is a directory, a .zip or a .7z`); } // --- the run ---------------------------------------------------------------- export type AttachMediaOptions = { paths: Paths; slug: string; channelConfig: ChannelConfig; source: string; items?: AttachItem[]; match?: string; createRecords?: boolean; replace?: boolean; dryRun?: boolean; onLog: (line: string) => void; signal?: AbortSignal; drainSignal?: AbortSignal; deps?: { now?: () => Date; probeDuration?: (file: string) => Promise; openSource?: (source: string) => Promise; }; }; export type AttachMediaResult = { counts: Record; attached: string[]; created: string[]; failed: { id: string; error: string }[]; notHeld: string[]; lost: string[]; unmatched: string[]; ambiguous: string[]; alreadyAttached: string[]; heldWithoutMedia: string[] | null; stopped?: string; dryRun: boolean; }; const LINE = (s: string) => (s.endsWith("\n") ? s : `${s}\n`); async function heldIds(dataDir: string): Promise> { try { const entries = await readdir(dataDir, { withFileTypes: true }); return new Set(entries.filter((d) => d.isDirectory() && !d.name.startsWith(".")).map((d) => d.name)); } catch (err) { if ((err as NodeJS.ErrnoException).code === "ENOENT") return new Set(); throw err; } } function extOf(p: string): string { const dot = p.lastIndexOf("."); return dot < 0 ? "" : p.slice(dot + 1).toLowerCase(); } async function storeRootExists(root: string): Promise { const st = await stat(root).catch(() => null); return !!st?.isDirectory(); } // The record's metadata.info.json, for a video the channel does not hold. async function recordInfo( source: ArchiveSource, row: AttachRow, opts: AttachMediaOptions, stagedMedia: string, ): Promise> { const folder = row.source; if (folder?.infoJson) { const parsed = JSON.parse(await source.readText(folder.infoJson)) as unknown; if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { throw new Error(`${folder.infoJson.path} is not a JSON object`); } const info = parsed as Record; if (info.id !== row.id) { throw new Error(`${folder.infoJson.path} is for "${String(info.id)}", not ${row.id}`); } return info; } const descriptionText = folder?.descriptionTxt ? await source.readText(folder.descriptionTxt) : undefined; const probe = opts.deps?.probeDuration ?? ((f: string) => probeMediaDurationSec({ ffprobeBin: opts.paths.ffprobeBin, file: f, ...(opts.signal ? { signal: opts.signal } : {}) })); const durationSec = await probe(stagedMedia); return archiveFolderInfoJson({ id: row.id as string, folderName: folder?.name ?? "", ...(descriptionText !== undefined ? { descriptionText } : {}), ...(opts.channelConfig.platform ? { platform: opts.channelConfig.platform } : {}), durationSec, }); } function logPlan(plan: AttachPlan, onLog: (s: string) => void): void { for (const r of plan.rows) { const what = r.media ? ` ← ${r.media.path} (${formatBytes(r.media.size)})` : ""; const note = r.class === "ambiguous" ? ` — ${r.candidates?.length ?? 0} video files: name one with "items"` : ""; onLog(LINE(`${r.class.padEnd(16)} ${r.id ?? "(no id)"}${what}${r.class === "attach" && r.replacing ? " [replacing]" : ""}${!r.media ? ` ${r.folder}` : ""}${note}`)); } if (plan.heldWithoutMedia) { onLog(LINE(`held, no media in the archive: ${plan.heldWithoutMedia.length}${plan.heldWithoutMedia.length ? ` — ${plan.heldWithoutMedia.join(", ")}` : ""}`)); } } export async function attachMedia(opts: AttachMediaOptions): Promise { const onLog = (s: string) => opts.onLog(LINE(s)); const now = opts.deps?.now ?? (() => new Date()); let match: RegExp | undefined; if (opts.match !== undefined) { try { match = new RegExp(opts.match, "i"); } catch (e) { throw new Error(`"match" is not a valid regex: ${(e as Error).message}`); } } const source = await (opts.deps?.openSource ?? ((s) => openArchiveSource(s)))(opts.source); onLog(`${source.kind} ${source.archive}: ${source.files.length} files`); const channelDir = path.join(opts.paths.channelsDir, opts.slug); const dataDir = path.join(channelDir, "data"); const held = await heldIds(dataDir); const attached = new Set(); const pointers = new Map>>(); for (const id of held) { const p = await loadSavedVideo(path.join(dataDir, id)); if (p) { attached.add(id); pointers.set(id, p); } } const folders = groupArchiveFolders(source.files); const plan = planAttach({ folders, held, attached, ...(opts.createRecords ? { createRecords: true } : {}), ...(opts.replace ? { replace: true } : {}), ...(match ? { match } : {}), ...(opts.items ? { items: opts.items, files: source.files } : {}), }); logPlan(plan, onLog); const ids = (c: AttachRowClass) => plan.rows.filter((r) => r.class === c).map((r) => r.id ?? r.folder); const result: AttachMediaResult = { counts: countAttachRows(plan.rows), attached: [], created: [], failed: [], notHeld: ids("not-held"), lost: ids("lost"), unmatched: ids("unmatched"), ambiguous: ids("ambiguous"), alreadyAttached: ids("already-attached"), heldWithoutMedia: plan.heldWithoutMedia, dryRun: !!opts.dryRun, }; const work = plan.rows.filter((r) => r.class === "attach" || r.class === "create"); const total = work.reduce((s, r) => s + (r.media?.size ?? 0), 0); onLog(`${work.length} to attach (${formatBytes(total)}): ${result.counts.attach} held, ${result.counts.create} new records`); if (opts.dryRun || work.length === 0) { onLog(`summary: ${JSON.stringify(summaryOf(result))}`); return result; } const storeRoot = savedVideoRoot(opts.paths, opts.channelConfig); if (!(await storeRootExists(storeRoot))) { throw new Error(`the saved-video store ${storeRoot} is not there (is its drive mounted?) — nothing attached`); } let n = 0; for (const row of work) { n += 1; if (opts.signal?.aborted) { result.stopped = "cancelled"; break; } if (opts.drainSignal?.aborted) { result.stopped = "drained"; onLog(`Drained: ${work.length - n + 1} left for the next run.`); break; } const id = row.id as string; const media = row.media as SourceFile; const storeDir = savedVideoDir(opts.paths, opts.channelConfig, opts.slug, id); try { await assertSavedVideosStoreWritable(opts.paths, storeDir); } catch (err) { result.stopped = (err as Error).message; onLog(`Stopped: ${result.stopped}`); break; } const free = await getFreeBytes(storeRoot); if (free < media.size + ATTACH_MEDIA_DISK_MARGIN_BYTES) { result.stopped = `the store's disk has ${formatBytes(free)} free; ${id} needs ${formatBytes(media.size)} plus ${formatBytes(ATTACH_MEDIA_DISK_MARGIN_BYTES)}`; onLog(`Stopped: ${result.stopped}. A re-run resumes.`); break; } const slugDir = path.dirname(storeDir); const staging = path.join(slugDir, `.${id}.attach-${process.pid}.part`); const ext = extOf(media.path); const videoDir = path.join(dataDir, id); try { onLog(`[${n}/${work.length}] ${id}: copying ${media.path} (${formatBytes(media.size)})`); await mkdir(slugDir, { recursive: true }); await rm(staging, { force: true }); const copied = await source.copyOut(media, staging, opts.signal); if (row.class === "create") { const info = await recordInfo(source, row, opts, staging); await mkdir(videoDir, { recursive: true }); await withMetadataHistory(videoDir, { by: "local-archive", requestedBy: ATTACH_MEDIA_REQUESTER, onLog }, () => writeJsonAtomic(path.join(videoDir, "metadata.info.json"), info, { indent: 0, newline: false }), ); const extractor = typeof info.extractor_key === "string" ? info.extractor_key.toLowerCase() : null; if (extractor) { await appendFile(path.join(channelDir, "archive"), `${extractor} ${id}\n`); } result.created.push(id); onLog(`${id}: record written (${typeof info.title === "string" ? info.title : id})`); } const previous = pointers.get(id) ?? null; const origin: SavedVideoOrigin = { requestedBy: ATTACH_MEDIA_REQUESTER, kind: "local-archive", archive: source.archive, entry: media.path, sha256: copied.sha256, attachedAt: now().toISOString(), }; const pointer = await persistSourceVideo({ videoDir, sourceFilename: `source-media.${ext}`, sourcePath: staging, storeDir, keepReason: "pin", origin, sha256: copied.sha256, paths: opts.paths, }); if (previous && savedVideoPath(previous) !== savedVideoPath(pointer)) { await rm(savedVideoPath(previous), { force: true }); onLog(`${id}: the replaced container ${previous.file} removed`); } result.attached.push(id); onLog(`${id}: attached ${pointer.file} (${formatBytes(pointer.bytes)}, sha256 ${copied.sha256.slice(0, 12)}…)`); } catch (err) { await rm(staging, { force: true }).catch(() => {}); 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(summaryOf(result))}`); return result; } export function summaryOf(r: AttachMediaResult): Record { return { dryRun: r.dryRun, counts: r.counts, attached: r.attached.length, created: r.created, failed: r.failed, notHeld: r.notHeld, lost: r.lost.length, unmatched: r.unmatched, ambiguous: r.ambiguous, alreadyAttached: r.alreadyAttached.length, heldWithoutMedia: r.heldWithoutMedia, ...(r.stopped ? { stopped: r.stopped } : {}), }; }