commit 6d9cb23d1cee35c03e999dc33224207f4570ae86
parent 93e59abdd569b113a1728de0c31b50f4681da298
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 9 Oct 2026 13:08:08 -0400
common: attach-media — a channel's held videos get their media from a local archive (directory, zip read in place, 7z) into the saved-video store, with local-archive provenance
D1 of release 21. The id from the folder's "(<id>)" or the file's yt-dlp
suffix, or explicit items; written through persistSourceVideo (a staging
file on the store's disk, then a rename), keepReason pin, origin {kind:
"local-archive", archive, entry, sha256, attachedAt}; refuses over a pointer
unless replace; createRecords writes metadata.info.json through the metadata
history (writer "local-archive"); [LOST] folders listed, never attached.
Job kind attach-media (drainable, needsMedia, ingest), ops route + action,
pnpm ops row, archilyzer media attach.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
19 files changed, 2356 insertions(+), 1 deletion(-)
diff --git a/common/bin/archilyzer.ts b/common/bin/archilyzer.ts
@@ -468,6 +468,36 @@ export const COMMANDS: Command[] = [
});
},
},
+ {
+ path: ["media", "attach"],
+ usage:
+ "<slug> <source> [--items <json|@file>] [--match <regex>] [--create-records] [--replace] [--dry-run] give a channel's held videos their media from a LOCAL archive (a directory, a .zip read in place, a .7z) into the saved-video store, provenance recorded (archive, entry, sha256); the id from the folder's \"(<id>)\" or the file's yt-dlp suffix; --create-records writes a record for a video not held; --dry-run lists each folder's class and writes nothing; offline, nothing fetched",
+ flags: {
+ items: "string",
+ match: "string",
+ "create-records": "boolean",
+ replace: "boolean",
+ "dry-run": "boolean",
+ },
+ maxPositionals: 2,
+ run: async ({ positionals, flags }) => {
+ const [slug, source] = positionals;
+ if (!slug || !source) {
+ console.error("media attach: which channel, from which archive? Pass <slug> <source>.");
+ return 2;
+ }
+ return (await import("./media-attach")).main({
+ slug,
+ source,
+ ...(typeof flags.items === "string" ? { items: flags.items } : {}),
+ ...(typeof flags.match === "string" ? { match: flags.match } : {}),
+ createRecords: flags["create-records"] === true,
+ replace: flags.replace === true,
+ dryRun: flags["dry-run"] === true,
+ signal: interrupted(),
+ });
+ },
+ },
// The bins that parse their own flags, run as children with their argv
// verbatim (_spawnBin.ts says why). Each row is the whole integration.
script(["duplicates"], "duplicate-shorts.ts",
diff --git a/common/bin/media-attach.ts b/common/bin/media-attach.ts
@@ -0,0 +1,76 @@
+// `archilyzer media attach <slug> <source> [--items <json>] [--match <re>]
+// [--create-records] [--replace] [--dry-run]` — a channel's held videos given
+// their media from a local archive, offline. The work and its rules are
+// controller/attachMedia.ts; the editor runs the same function as the
+// `attach-media` job (`POST /api/ops/attach-media`).
+//
+// Nothing on the network. It copies into the saved-video store and writes
+// pointers (and, with --create-records, records) atomically; it does not see
+// the editor's queues, so not beside a job on the same channel. The channel's
+// media guard is asked here, as the job kind's `needsMedia` asks it there.
+
+import { readFile } from "node:fs/promises";
+import { getPaths } from "../lib/paths";
+import { assertChannelMediaReachable } from "../lib/channelMedia";
+import { isValidChannelSlug, readChannelConfig } from "../controller/channels";
+import { attachMedia } from "../controller/attachMedia";
+import type { AttachItem } from "../lib/attachMediaPlan";
+
+function parseItems(text: string): AttachItem[] {
+ const raw = JSON.parse(text) as unknown;
+ if (!Array.isArray(raw) || raw.length === 0) throw new Error("--items: a non-empty JSON array of {id, path}");
+ return raw.map((e, i) => {
+ const r = e as Record<string, unknown>;
+ if (typeof r?.id !== "string" || !r.id || /[/\\\0]/.test(r.id) || typeof r.path !== "string" || !r.path) {
+ throw new Error(`--items[${i}]: an object {"id": "<video id>", "path": "<entry in the source>"}`);
+ }
+ return { id: r.id, path: r.path };
+ });
+}
+
+export async function main(opts: {
+ slug: string;
+ source: string;
+ items?: string;
+ match?: string;
+ createRecords: boolean;
+ replace: boolean;
+ dryRun: boolean;
+ signal?: AbortSignal;
+}): Promise<number> {
+ if (!isValidChannelSlug(opts.slug)) {
+ console.error(`media attach: "${opts.slug}" is not a channel slug`);
+ return 2;
+ }
+ try {
+ const paths = getPaths();
+ const channelConfig = await readChannelConfig(paths, opts.slug);
+ if (!channelConfig) {
+ console.error(`media attach: no channel "${opts.slug}"`);
+ return 2;
+ }
+ if (!opts.dryRun) await assertChannelMediaReachable(paths, opts.slug, channelConfig);
+ // --items is a JSON array, inline or "@file".
+ const items =
+ opts.items === undefined
+ ? undefined
+ : parseItems(opts.items.startsWith("@") ? await readFile(opts.items.slice(1), "utf8") : opts.items);
+ const result = await attachMedia({
+ paths,
+ slug: opts.slug,
+ channelConfig,
+ source: opts.source,
+ ...(items ? { items } : {}),
+ ...(opts.match !== undefined ? { match: opts.match } : {}),
+ ...(opts.createRecords ? { createRecords: true } : {}),
+ ...(opts.replace ? { replace: true } : {}),
+ ...(opts.dryRun ? { dryRun: true } : {}),
+ onLog: (line) => console.log(line.replace(/\n$/, "")),
+ ...(opts.signal ? { signal: opts.signal } : {}),
+ });
+ return result.failed.length > 0 || (result.stopped && result.stopped !== "drained") ? 1 : 0;
+ } catch (err) {
+ console.error(`media attach: ${(err as Error).message}`);
+ return 1;
+ }
+}
diff --git a/common/controller/attachMedia.test.ts b/common/controller/attachMedia.test.ts
@@ -0,0 +1,297 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { createHash } from "node:crypto";
+import { mkdir, mkdtemp, readFile, readdir, rm, stat, writeFile } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import path from "node:path";
+import { execa } from "execa";
+import type { Paths } from "../lib/paths";
+import type { ChannelConfig } from "../lib/channelConfig";
+import { makeZip, type ZipFixtureEntry } from "../testing/makeZip";
+import { loadSavedVideo } from "../lib/savedVideo-server";
+import { savedVideoPath } from "../lib/savedVideo";
+import { SAVED_VIDEOS_MARKER_FILENAME } from "../lib/savedVideoStore";
+import { attachMedia, parse7zListing, type AttachMediaOptions } from "./attachMedia";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test controller/attachMedia.test.ts
+//
+// A synthetic archive in the pilot's shape (release 21 D1): yt-dlp folders,
+// bare-id folders with description.txt, a [LOST] folder. Every byte is made
+// here; nothing reads a real archive or corpus.
+
+const SLUG = "demo";
+const sha = (b: Buffer) => createHash("sha256").update(b).digest("hex");
+const bytes = (n: number, seed: number) => Buffer.from(Array.from({ length: n }, (_, i) => (i * seed + 3) & 0xff));
+
+const A = bytes(50_000, 7); // held, yt-dlp folder, stored
+const C = bytes(40_000, 11); // not held, bare id + description.txt, deflated
+const D = bytes(30_000, 13); // not held, yt-dlp folder with info.json
+const INFO_D = JSON.stringify({ id: "ddddddddddd", title: "From yt-dlp", extractor_key: "Youtube", upload_date: "20190101" });
+
+const ENTRIES: ZipFixtureEntry[] = [
+ { name: "YouTube/" },
+ { name: "YouTube/001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.mkv", data: A },
+ { name: "YouTube/001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.info.json", data: '{"id":"aaaaaaaaaaa"}' },
+ { name: "YouTube/002 [LOST] - Gone - (bbbbbbbbbbb)/description.txt", data: "lost" },
+ { name: "YouTube/003 - Episode 2 - (ccccccccccc)/ccccccccccc.mp4", data: C, deflate: true },
+ { name: "YouTube/003 - Episode 2 - (ccccccccccc)/description.txt", data: "Published on July 13th 2018\n\nTips" },
+ { name: "YouTube/003 - Episode 2 - (ccccccccccc)/source.txt", data: "Video source: somewhere" },
+ { name: "YouTube/004 - Other - (ddddddddddd)/Other-ddddddddddd.webm", data: D },
+ { name: "YouTube/004 - Other - (ddddddddddd)/Other-ddddddddddd.info.json", data: INFO_D },
+];
+
+type Ctx = { dir: string; paths: Paths; zip: string; lines: string[]; opts: (o?: Partial<AttachMediaOptions>) => AttachMediaOptions };
+
+async function withFixture(fn: (ctx: Ctx) => Promise<void>, entries = ENTRIES): Promise<void> {
+ const dir = await mkdtemp(path.join(tmpdir(), "ttb-attach-"));
+ const transcriptsDir = path.join(dir, "corpus");
+ const paths = {
+ transcriptsDir,
+ channelsDir: path.join(transcriptsDir, "channels"),
+ savedVideosDir: path.join(dir, "store"),
+ ffprobeBin: "ffprobe-not-used",
+ } as Paths;
+ await mkdir(paths.savedVideosDir, { recursive: true });
+ // Held: aaaaaaaaaaa (with media in the archive), bbbbbbbbbbb (LOST), hhhhhhhhhhh (no folder).
+ for (const id of ["aaaaaaaaaaa", "bbbbbbbbbbb", "hhhhhhhhhhh"]) {
+ const vd = path.join(paths.channelsDir, SLUG, "data", id);
+ await mkdir(vd, { recursive: true });
+ await writeFile(path.join(vd, "metadata.info.json"), JSON.stringify({ id, title: id }));
+ }
+ const zip = path.join(dir, "archive", "YouTube.zip");
+ await mkdir(path.dirname(zip), { recursive: true });
+ await writeFile(zip, makeZip(entries));
+ const lines: string[] = [];
+ const opts = (o: Partial<AttachMediaOptions> = {}): AttachMediaOptions => ({
+ paths,
+ slug: SLUG,
+ channelConfig: { name: "Demo", handling: "youtube", platform: "youtube" } as unknown as ChannelConfig,
+ source: zip,
+ onLog: (l) => lines.push(l),
+ deps: { probeDuration: async () => 61.5, now: () => new Date("2026-10-09T12:00:00Z") },
+ ...o,
+ });
+ try {
+ await fn({ dir, paths, zip, lines, opts });
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+}
+
+const videoDir = (paths: Paths, id: string) => path.join(paths.channelsDir, SLUG, "data", id);
+
+test("dryRun lists every class and writes nothing", async () => {
+ await withFixture(async ({ paths, zip, opts, lines }) => {
+ const zipBefore = sha(await readFile(zip));
+ const r = await attachMedia(opts({ dryRun: true }));
+ assert.equal(r.counts.attach, 1);
+ assert.equal(r.counts["not-held"], 2);
+ assert.equal(r.counts.lost, 1);
+ assert.deepEqual(r.notHeld.sort(), ["ccccccccccc", "ddddddddddd"]);
+ assert.deepEqual(r.heldWithoutMedia, ["bbbbbbbbbbb", "hhhhhhhhhhh"]);
+ assert.deepEqual(r.attached, []);
+ assert.equal(await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa")), null);
+ assert.deepEqual(await readdir(paths.savedVideosDir), []);
+ assert.equal(sha(await readFile(zip)), zipBefore, "the archive is never written");
+ assert.ok(lines.some((l) => l.startsWith("summary: ")));
+ });
+});
+
+test("attaches a held video through the store, with local-archive provenance", async () => {
+ await withFixture(async ({ paths, zip, opts }) => {
+ const zipBefore = sha(await readFile(zip));
+ const r = await attachMedia(opts());
+ assert.deepEqual(r.attached, ["aaaaaaaaaaa"]);
+ assert.deepEqual(r.failed, []);
+ const pointer = await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa"));
+ assert.ok(pointer);
+ assert.equal(pointer.file, "source-media.mkv");
+ assert.equal(pointer.dir, path.join(paths.savedVideosDir, SLUG, "aaaaaaaaaaa"));
+ assert.equal(pointer.keepReason, "pin");
+ assert.equal(pointer.bytes, A.length);
+ assert.equal(pointer.sha256, sha(A));
+ assert.deepEqual(pointer.origin, {
+ requestedBy: "attach-media",
+ kind: "local-archive",
+ archive: zip,
+ entry: "YouTube/001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.mkv",
+ sha256: sha(A),
+ attachedAt: "2026-10-09T12:00:00.000Z",
+ });
+ assert.deepEqual(await readFile(savedVideoPath(pointer)), A);
+ // No staging file left beside the store dir; nothing copied into data/.
+ assert.deepEqual(await readdir(path.join(paths.savedVideosDir, SLUG)), ["aaaaaaaaaaa"]);
+ assert.deepEqual((await readdir(videoDir(paths, "aaaaaaaaaaa"))).sort(), ["metadata.info.json", "saved-video.json"]);
+ assert.equal(sha(await readFile(zip)), zipBefore);
+ // Not held, no createRecords: listed, nothing written for them.
+ await assert.rejects(stat(videoDir(paths, "ccccccccccc")));
+ });
+});
+
+test("a re-run is a no-op (already attached), and refuses over a pointer unless replace", async () => {
+ await withFixture(async ({ paths, opts }) => {
+ await attachMedia(opts());
+ const first = await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa"));
+ const again = await attachMedia(opts());
+ assert.deepEqual(again.attached, []);
+ assert.deepEqual(again.alreadyAttached, ["aaaaaaaaaaa"]);
+ assert.deepEqual(await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa")), first);
+
+ const replaced = await attachMedia(opts({ replace: true, deps: { now: () => new Date("2026-10-10T00:00:00Z") } }));
+ assert.deepEqual(replaced.attached, ["aaaaaaaaaaa"]);
+ const after = await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa"));
+ assert.equal(after?.origin?.attachedAt, "2026-10-10T00:00:00.000Z");
+ assert.deepEqual(await readFile(savedVideoPath(after!)), A);
+ });
+});
+
+test("replace removes a previous container of another name", async () => {
+ await withFixture(async ({ paths, opts }) => {
+ // A container persisted earlier as .mp4 (a download, say).
+ const vd = videoDir(paths, "aaaaaaaaaaa");
+ const storeDir = path.join(paths.savedVideosDir, SLUG, "aaaaaaaaaaa");
+ await mkdir(storeDir, { recursive: true });
+ await writeFile(path.join(storeDir, "source-media.mp4"), "old");
+ await writeFile(
+ path.join(vd, "saved-video.json"),
+ JSON.stringify({ storedAt: "x", dir: storeDir, file: "source-media.mp4", bytes: 3, keepReason: "override" }),
+ );
+ const refused = await attachMedia(opts());
+ assert.deepEqual(refused.alreadyAttached, ["aaaaaaaaaaa"]);
+ await attachMedia(opts({ replace: true }));
+ assert.deepEqual((await readdir(storeDir)).sort(), ["source-media.mkv"]);
+ });
+});
+
+test("createRecords: a record from the folder's .info.json, or from description.txt and its name", async () => {
+ await withFixture(async ({ paths, opts }) => {
+ const r = await attachMedia(opts({ createRecords: true }));
+ assert.deepEqual(r.created.sort(), ["ccccccccccc", "ddddddddddd"]);
+ assert.deepEqual(r.attached.sort(), ["aaaaaaaaaaa", "ccccccccccc", "ddddddddddd"]);
+
+ const d = await readFile(path.join(videoDir(paths, "ddddddddddd"), "metadata.info.json"), "utf8");
+ assert.deepEqual(JSON.parse(d), JSON.parse(INFO_D));
+
+ const c = JSON.parse(await readFile(path.join(videoDir(paths, "ccccccccccc"), "metadata.info.json"), "utf8"));
+ assert.equal(c.title, "Episode 2");
+ assert.equal(c.upload_date, "20180713");
+ assert.equal(c.duration, 61.5);
+ assert.equal(c.webpage_url, "https://www.youtube.com/watch?v=ccccccccccc");
+
+ // The deflated entry came out whole.
+ const pc = await loadSavedVideo(videoDir(paths, "ccccccccccc"));
+ assert.deepEqual(await readFile(savedVideoPath(pc!)), C);
+ assert.equal(pc?.file, "source-media.mp4");
+
+ // The channel's archive file gained a line per new record, as a download's would.
+ const archive = await readFile(path.join(paths.channelsDir, SLUG, "archive"), "utf8");
+ assert.deepEqual(archive.trim().split("\n").sort(), ["youtube ccccccccccc", "youtube ddddddddddd"]);
+ });
+});
+
+test("a .info.json for another id is refused, and nothing is left behind", async () => {
+ const entries = ENTRIES.map((e) =>
+ e.name.endsWith("Other-ddddddddddd.info.json") ? { ...e, data: '{"id":"zzzzzzzzzzz"}' } : e,
+ );
+ await withFixture(async ({ paths, opts }) => {
+ const r = await attachMedia(opts({ createRecords: true }));
+ assert.deepEqual(r.failed.map((x) => x.id), ["ddddddddddd"]);
+ assert.match(r.failed[0].error, /not ddddddddddd/);
+ await assert.rejects(stat(videoDir(paths, "ddddddddddd")));
+ assert.deepEqual((await readdir(path.join(paths.savedVideosDir, SLUG))).sort(), ["aaaaaaaaaaa", "ccccccccccc"]);
+ }, entries);
+});
+
+test("an explicit item attaches a named file to a named id", async () => {
+ await withFixture(async ({ paths, opts }) => {
+ const r = await attachMedia(
+ opts({ items: [{ id: "hhhhhhhhhhh", path: "YouTube/004 - Other - (ddddddddddd)/Other-ddddddddddd.webm" }] }),
+ );
+ assert.deepEqual(r.attached, ["hhhhhhhhhhh"]);
+ const p = await loadSavedVideo(videoDir(paths, "hhhhhhhhhhh"));
+ assert.deepEqual(await readFile(savedVideoPath(p!)), D);
+ assert.equal(p?.origin?.entry, "YouTube/004 - Other - (ddddddddddd)/Other-ddddddddddd.webm");
+ });
+});
+
+test("a directory source attaches the same way", async () => {
+ await withFixture(async ({ dir, paths, opts }) => {
+ const src = path.join(dir, "unpacked");
+ const folder = path.join(src, "001 - Intro - (aaaaaaaaaaa)");
+ await mkdir(folder, { recursive: true });
+ await writeFile(path.join(folder, "Intro-aaaaaaaaaaa.mp4"), A);
+ const r = await attachMedia(opts({ source: src }));
+ assert.deepEqual(r.attached, ["aaaaaaaaaaa"]);
+ const p = await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa"));
+ assert.equal(p?.origin?.archive, src);
+ assert.equal(p?.origin?.entry, "001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.mp4");
+ // The directory source is read, never moved.
+ assert.deepEqual(await readFile(path.join(folder, "Intro-aaaaaaaaaaa.mp4")), A);
+ });
+});
+
+test("guards: a store being moved stops the run; a store root that is not there refuses", async () => {
+ await withFixture(async ({ paths, opts }) => {
+ // The in-corpus store path is what the marker guards; point the store there.
+ const inCorpus = { ...paths, savedVideosDir: path.join(paths.transcriptsDir, "saved-videos") } as Paths;
+ await mkdir(inCorpus.savedVideosDir, { recursive: true });
+ await writeFile(
+ path.join(paths.transcriptsDir, SAVED_VIDEOS_MARKER_FILENAME),
+ JSON.stringify({ target: "/elsewhere", direction: "out", startedAt: "2026-10-09T00:00:00Z", phase: "copy" }),
+ );
+ const held = await attachMedia(opts({ paths: inCorpus }));
+ assert.deepEqual(held.attached, []);
+ assert.ok(held.stopped);
+
+ const gone = { ...paths, savedVideosDir: path.join(paths.transcriptsDir, "no-such-store") } as Paths;
+ await assert.rejects(attachMedia(opts({ paths: gone })), /is not there/);
+ await assert.rejects(stat(gone.savedVideosDir), "an unmounted store is never materialised");
+ });
+});
+
+test("drain stops before the next video", async () => {
+ await withFixture(async ({ opts }) => {
+ const drain = new AbortController();
+ drain.abort();
+ const r = await attachMedia(opts({ createRecords: true, drainSignal: drain.signal }));
+ assert.equal(r.stopped, "drained");
+ assert.deepEqual(r.attached, []);
+ });
+});
+
+test("7z listing parse, and a .7z source when 7z is installed", async (t) => {
+ const listing = [
+ "Path = /x/a.7z",
+ "Type = 7z",
+ "",
+ "----------",
+ "Path = dir",
+ "Folder = +",
+ "Size = 0",
+ "",
+ "Path = dir/001 - X - (aaaaaaaaaaa)/x.mp4",
+ "Folder = -",
+ "Size = 50000",
+ "Attributes = A",
+ "",
+ ].join("\n");
+ assert.deepEqual(parse7zListing(listing), [{ path: "dir/001 - X - (aaaaaaaaaaa)/x.mp4", size: 50000 }]);
+
+ const probe = await execa("7z", ["i"], { reject: false }).catch(() => null);
+ if (!probe || probe.exitCode !== 0) {
+ t.skip("7z is not installed");
+ return;
+ }
+ await withFixture(async ({ dir, paths, opts }) => {
+ const src = path.join(dir, "pack");
+ await mkdir(path.join(src, "001 - Intro [x] - (aaaaaaaaaaa)"), { recursive: true });
+ await writeFile(path.join(src, "001 - Intro [x] - (aaaaaaaaaaa)", "Intro-aaaaaaaaaaa.mp4"), A);
+ const archive = path.join(dir, "pack.7z");
+ await execa("7z", ["a", "-bd", archive, "."], { cwd: src });
+ const r = await attachMedia(opts({ source: archive }));
+ assert.deepEqual(r.attached, ["aaaaaaaaaaa"]);
+ const p = await loadSavedVideo(videoDir(paths, "aaaaaaaaaaa"));
+ assert.deepEqual(await readFile(savedVideoPath(p!)), A);
+ assert.equal(p?.sha256, sha(A));
+ });
+});
diff --git a/common/controller/attachMedia.ts b/common/controller/attachMedia.ts
@@ -0,0 +1,503 @@
+// 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.<ext>` in the store, a pointer in
+// `data/<id>/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<string>;
+ copyOut(file: SourceFile, dest: string, signal?: AbortSignal): Promise<CopiedEntry>;
+};
+
+async function openZipSource(archive: string): Promise<ArchiveSource> {
+ const entries = await readZipEntries(archive);
+ const byName = new Map<string, ZipEntry>();
+ 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<SourceFile[]> {
+ 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<ArchiveSource> {
+ 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<ArchiveSource> {
+ 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<ArchiveSource> {
+ 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<number | null>;
+ openSource?: (source: string) => Promise<ArchiveSource>;
+ };
+};
+
+export type AttachMediaResult = {
+ counts: Record<AttachRowClass, number>;
+ 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<Set<string>> {
+ 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<boolean> {
+ 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<Record<string, unknown>> {
+ 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<string, unknown>;
+ 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<AttachMediaResult> {
+ 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<string>();
+ const pointers = new Map<string, Awaited<ReturnType<typeof loadSavedVideo>>>();
+ 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<string, unknown> {
+ 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 } : {}),
+ };
+}
diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts
@@ -293,3 +293,17 @@ test("release 18: every drainable kind but the lane runners and the clip-window
}
assert.equal(isIngestKind("no-such-kind"), false);
});
+
+// RELEASE 21 D1: local media attached to held videos. It writes the big file
+// (the store), stops between videos, and its pointers and created records are
+// what the index reads — so media, drainable, and ingest.
+test("release 21: attach-media is a drainable media ingest kind, not replayable", () => {
+ const meta = getJobKind("attach-media");
+ assert.ok(meta, "registered");
+ assert.equal(jobKindLabel("attach-media"), "Attach local media");
+ assert.equal(isDrainableKind("attach-media"), true);
+ assert.equal(meta.replayable, false);
+ assert.equal(kindNeedsMedia("attach-media"), true);
+ assert.equal(kindNeedsText("attach-media"), false);
+ assert.equal(isIngestKind("attach-media"), true);
+});
diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts
@@ -960,6 +960,23 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
queueKeyStrategy: "custom",
needsMedia: true,
},
+ // LOCAL MEDIA ATTACHED TO HELD VIDEOS (release 21 D1,
+ // controller/attachMedia.ts): each video's file copied out of a local
+ // archive (a directory, a zip read in place, a 7z) into the saved-video
+ // store. Media: it writes the big file. Drainable: it stops between videos,
+ // and a re-run resumes. An INGEST kind: it writes `saved-video.json` beside
+ // each record and, with `createRecords`, whole new records the index reads.
+ // On the channel's own queue (`channel:<slug>`): nothing is fetched, so no
+ // platform queue is owed a turn. Not replayable: the archive it names is a
+ // path on some drive, and a retry is a re-run.
+ "attach-media": {
+ kind: "attach-media",
+ label: "Attach local media",
+ drainable: true,
+ replayable: false,
+ queueKeyStrategy: "custom",
+ needsMedia: true,
+ },
};
export function getJobKind(kind: string): JobKindMeta | undefined {
@@ -1024,6 +1041,7 @@ const INGEST_KINDS: ReadonlySet<string> = new Set([
"check-availability",
"quick-availability-check",
"check-maybe-missing",
+ "attach-media",
]);
export function isIngestKind(kind: string | undefined): boolean {
diff --git a/common/lib/attachMediaPlan.test.ts b/common/lib/attachMediaPlan.test.ts
@@ -0,0 +1,178 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import {
+ archiveFolderInfoJson,
+ countAttachRows,
+ descriptionFromText,
+ groupArchiveFolders,
+ idFromFolderName,
+ idFromMediaName,
+ planAttach,
+ publishedDateFromText,
+ titleFromFolderName,
+ type SourceFile,
+} from "./attachMediaPlan";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test lib/attachMediaPlan.test.ts
+//
+// The folder shapes are the pilot archive's (release 21), with synthetic ids
+// and titles: `NNN - <title> - (<id>)` holding either yt-dlp's files
+// (`<title>-<id>.<ext>` + `.info.json`) or a bare `<id>.<ext>` beside
+// description.txt / source.txt, and `[LOST]` folders with no media.
+
+const f = (p: string, size = 10): SourceFile => ({ path: p, size });
+
+const ARCHIVE: SourceFile[] = [
+ f("YouTube/README.txt"),
+ f("YouTube/001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.mkv", 1000),
+ f("YouTube/001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.info.json"),
+ f("YouTube/001 - Intro - (aaaaaaaaaaa)/Intro-aaaaaaaaaaa.jpg"),
+ f("YouTube/002 [LOST] - Gone - (bbbbbbbbbbb)/description.txt"),
+ f("YouTube/002 [LOST] - Gone - (bbbbbbbbbbb)/source.txt"),
+ f("YouTube/003 - Episode 2 - (ccccccccccc)/ccccccccccc.mp4", 2000),
+ f("YouTube/003 - Episode 2 - (ccccccccccc)/description.txt"),
+ f("YouTube/003 - Episode 2 - (ccccccccccc)/Sources.txt"),
+ f("YouTube/004 - Text only - (ddddddddddd)/description.txt"),
+ f("YouTube/005 - Two files - (eeeeeeeeeee)/a.mp4"),
+ f("YouTube/005 - Two files - (eeeeeeeeeee)/b.webm"),
+ f("YouTube/006 - Odd name - (fffffffffff)/Odd name [watch-fffffffffff-1].mp4", 3000),
+ f("loose/Some title [ggggggggggg].webm", 4000),
+ f("loose/no id here.mp4"),
+];
+
+test("ids: the folder's trailing (id), else the file's yt-dlp suffix", () => {
+ assert.equal(idFromFolderName("001 - Intro - (aaaaaaaaaaa)"), "aaaaaaaaaaa");
+ assert.equal(idFromFolderName("Intro (too-short)"), null);
+ assert.equal(idFromMediaName("x/Intro-aaaaaaaaaaa.mkv"), "aaaaaaaaaaa");
+ assert.equal(idFromMediaName("ccccccccccc.mp4"), "ccccccccccc");
+ assert.equal(idFromMediaName("Some title [ggggggggggg].webm"), "ggggggggggg");
+ assert.equal(idFromMediaName("no id here.mp4"), null);
+});
+
+test("folders: grouped by parent, sidecars found, [LOST] marked", () => {
+ const folders = groupArchiveFolders(ARCHIVE);
+ const by = (id: string) => folders.find((x) => x.id === id)!;
+ assert.equal(by("aaaaaaaaaaa").idFrom, "folder");
+ assert.equal(by("aaaaaaaaaaa").media.length, 1);
+ assert.ok(by("aaaaaaaaaaa").infoJson);
+ assert.equal(by("bbbbbbbbbbb").lost, true);
+ assert.ok(by("ccccccccccc").descriptionTxt);
+ assert.ok(by("ccccccccccc").sourceTxt, "Sources.txt, any case");
+ assert.equal(by("ggggggggggg").idFrom, "file");
+ // The archive's root holds no media and yields a folder with no id.
+ assert.ok(folders.some((x) => x.folder === "YouTube" && x.id === null && x.media.length === 0));
+ // A folder with no id splits per media file.
+ assert.ok(folders.some((x) => x.folder === "loose" && x.id === null && x.media.length === 1));
+});
+
+test("plan: every class, and held videos the archive has no media for", () => {
+ const folders = groupArchiveFolders(ARCHIVE);
+ const plan = planAttach({
+ folders,
+ held: new Set(["aaaaaaaaaaa", "bbbbbbbbbbb", "fffffffffff", "hhhhhhhhhhh", "ddddddddddd"]),
+ attached: new Set(["fffffffffff"]),
+ });
+ const cls = Object.fromEntries(plan.rows.map((r) => [r.id ?? r.folder, r.class]));
+ assert.equal(cls.aaaaaaaaaaa, "attach");
+ assert.equal(cls.bbbbbbbbbbb, "lost");
+ assert.equal(cls.ccccccccccc, "not-held");
+ assert.equal(cls.ddddddddddd, "no-media");
+ assert.equal(cls.eeeeeeeeeee, "ambiguous");
+ assert.equal(cls.fffffffffff, "already-attached");
+ assert.equal(cls.ggggggggggg, "not-held");
+ assert.equal(cls.loose, "unmatched");
+ assert.deepEqual(plan.heldWithoutMedia, ["bbbbbbbbbbb", "ddddddddddd", "hhhhhhhhhhh"]);
+ const counts = countAttachRows(plan.rows);
+ assert.equal(counts.attach, 1);
+ assert.equal(counts["not-held"], 2);
+});
+
+test("plan: createRecords turns not-held into create; replace re-attaches", () => {
+ const folders = groupArchiveFolders(ARCHIVE);
+ const plan = planAttach({
+ folders,
+ held: new Set(["fffffffffff"]),
+ attached: new Set(["fffffffffff"]),
+ createRecords: true,
+ replace: true,
+ });
+ const row = (id: string) => plan.rows.find((r) => r.id === id)!;
+ assert.equal(row("ccccccccccc").class, "create");
+ assert.ok(row("ccccccccccc").source?.descriptionTxt, "the folder rides along for the record");
+ assert.equal(row("fffffffffff").class, "attach");
+ assert.equal(row("fffffffffff").replacing, true);
+});
+
+test("plan: two folders claiming one id are both ambiguous", () => {
+ const folders = groupArchiveFolders([
+ f("a/001 - X - (aaaaaaaaaaa)/x.mp4"),
+ f("b/001 - X again - (aaaaaaaaaaa)/y.mp4"),
+ ]);
+ const plan = planAttach({ folders, held: new Set(["aaaaaaaaaaa"]), attached: new Set() });
+ assert.deepEqual(plan.rows.map((r) => r.class), ["ambiguous", "ambiguous"]);
+ assert.deepEqual(plan.heldWithoutMedia, [], "an ambiguous id is not 'no media'");
+});
+
+test("plan: match narrows the folders and drops the held-without-media list", () => {
+ const plan = planAttach({
+ folders: groupArchiveFolders(ARCHIVE),
+ held: new Set(["aaaaaaaaaaa", "ccccccccccc"]),
+ attached: new Set(),
+ match: /episode 2/i,
+ });
+ assert.deepEqual(plan.rows.map((r) => [r.id, r.class]), [["ccccccccccc", "attach"]]);
+ assert.equal(plan.heldWithoutMedia, null);
+});
+
+test("plan: explicit items are the whole plan, by path", () => {
+ const plan = planAttach({
+ folders: groupArchiveFolders(ARCHIVE),
+ held: new Set(["zzzzzzzzzzz"]),
+ attached: new Set(),
+ files: ARCHIVE,
+ items: [
+ { id: "zzzzzzzzzzz", path: "YouTube/005 - Two files - (eeeeeeeeeee)/b.webm" },
+ { id: "yyyyyyyyyyy", path: "nope.mp4" },
+ ],
+ });
+ assert.deepEqual(plan.rows.map((r) => [r.id, r.class, r.media?.path ?? null]), [
+ ["zzzzzzzzzzz", "attach", "YouTube/005 - Two files - (eeeeeeeeeee)/b.webm"],
+ ["yyyyyyyyyyy", "no-media", null],
+ ]);
+});
+
+test("records: title from the folder name, date and description from description.txt", () => {
+ assert.equal(titleFromFolderName("006 - The Show Episode 002 - (6ijPdSuZdws)"), "The Show Episode 002");
+ assert.equal(titleFromFolderName("002 [LOST] - Gone - (bbbbbbbbbbb)"), "Gone");
+ assert.equal(titleFromFolderName("Deadspin, Halloween & More! - (VcFA_LOnQoc)"), "Deadspin, Halloween & More!");
+ assert.equal(publishedDateFromText("Published on July 13th 2018\n\nTips"), "20180713");
+ assert.equal(publishedDateFromText("Published on March 1, 2019"), "20190301");
+ assert.equal(publishedDateFromText("no date"), null);
+ assert.equal(descriptionFromText("Published on July 13th 2018\n\nTips & comments\nCategory: X"), "Tips & comments\nCategory: X");
+
+ const info = archiveFolderInfoJson({
+ id: "ccccccccccc",
+ folderName: "003 - Episode 2 - (ccccccccccc)",
+ descriptionText: "Published on July 13th 2018\n\nTips",
+ platform: "youtube",
+ durationSec: 61.234,
+ });
+ assert.deepEqual(info, {
+ id: "ccccccccccc",
+ title: "Episode 2",
+ description: "Tips",
+ duration: 61.23,
+ webpage_url: "https://www.youtube.com/watch?v=ccccccccccc",
+ original_url: "https://www.youtube.com/watch?v=ccccccccccc",
+ extractor: "youtube",
+ extractor_key: "Youtube",
+ upload_date: "20180713",
+ });
+ const keys = Object.keys(info);
+ assert.equal(keys[1], "title", "title near the head");
+ assert.equal(keys.at(-1), "upload_date", "upload_date last, in the tail");
+ // Another platform: no invented YouTube URL.
+ const other = archiveFolderInfoJson({ id: "x", folderName: "x", platform: "rumble" });
+ assert.equal(other.webpage_url, undefined);
+ assert.equal(other.extractor, undefined);
+});
diff --git a/common/lib/attachMediaPlan.ts b/common/lib/attachMediaPlan.ts
@@ -0,0 +1,343 @@
+// WHICH FILE OF A LOCAL ARCHIVE IS WHICH VIDEO (release 21 D1).
+//
+// The pure half of `attach-media` (controller/attachMedia.ts): given the file
+// list of an archive (a zip's entries, a directory's files), group it into
+// per-video folders, find each folder's video id and media file, and classify
+// every folder against what the channel holds. No fs: the controller lists,
+// this decides, and the tests run over plain arrays.
+//
+// THE ID, in order: the folder's trailing `(<11-char id>)` — the archive's own
+// naming, `NNN - <title> - (<id>)` — else the media file's yt-dlp suffix
+// (`<title>-<id>.<ext>`, `<title> [<id>].<ext>`, or a bare `<id>.<ext>`), else
+// an explicit item. The 11-character shape is YouTube's; any other platform's
+// files are attached by explicit `items`.
+//
+// `[LOST]` folders (the archive's mark for a video whose media nobody has)
+// are listed and never attached, whatever they hold.
+
+export const ATTACHABLE_MEDIA_EXTS: readonly string[] = ["mp4", "mkv", "webm", "mov", "m4v"];
+
+const YT_ID = "[A-Za-z0-9_-]{11}";
+const FOLDER_ID_RE = new RegExp(`\\((${YT_ID})\\)\\s*$`);
+const FILE_ID_RES = [
+ new RegExp(`^(${YT_ID})\\.[A-Za-z0-9]+$`),
+ new RegExp(`-(${YT_ID})\\.[A-Za-z0-9]+$`),
+ new RegExp(`\\[(${YT_ID})\\]\\.[A-Za-z0-9]+$`),
+];
+
+export type SourceFile = {
+ // Path inside the archive, "/"-separated, no leading "/".
+ path: string;
+ size: number;
+};
+
+export type ArchiveFolder = {
+ // The folder's path inside the archive ("" for the archive's root).
+ folder: string;
+ // Its last segment (the root's is "").
+ name: string;
+ id: string | null;
+ idFrom: "folder" | "file" | null;
+ lost: boolean;
+ // Every video container directly in the folder.
+ media: SourceFile[];
+ // The sidecars createRecords reads, when present.
+ infoJson?: SourceFile;
+ descriptionTxt?: SourceFile;
+ sourceTxt?: SourceFile;
+};
+
+function extOf(name: string): string {
+ const dot = name.lastIndexOf(".");
+ return dot < 0 ? "" : name.slice(dot + 1).toLowerCase();
+}
+
+function baseName(p: string): string {
+ const i = p.lastIndexOf("/");
+ return i < 0 ? p : p.slice(i + 1);
+}
+
+function dirName(p: string): string {
+ const i = p.lastIndexOf("/");
+ return i < 0 ? "" : p.slice(0, i);
+}
+
+export function isAttachableMedia(name: string): boolean {
+ return ATTACHABLE_MEDIA_EXTS.includes(extOf(name));
+}
+
+export function idFromFolderName(name: string): string | null {
+ return FOLDER_ID_RE.exec(name)?.[1] ?? null;
+}
+
+export function idFromMediaName(name: string): string | null {
+ const base = baseName(name);
+ for (const re of FILE_ID_RES) {
+ const m = re.exec(base);
+ if (m) return m[1];
+ }
+ return null;
+}
+
+export function isLostFolder(name: string): boolean {
+ return /\[LOST\]/i.test(name);
+}
+
+// One folder per directory that holds a file. A media file's folder is its
+// immediate parent; the id comes from that folder's name, else the file's.
+export function groupArchiveFolders(files: readonly SourceFile[]): ArchiveFolder[] {
+ const byFolder = new Map<string, SourceFile[]>();
+ for (const f of files) {
+ if (!f.path || f.path.endsWith("/")) continue;
+ const folder = dirName(f.path);
+ const list = byFolder.get(folder) ?? [];
+ list.push(f);
+ byFolder.set(folder, list);
+ }
+ const out: ArchiveFolder[] = [];
+ for (const [folder, list] of [...byFolder].sort(([a], [b]) => a.localeCompare(b))) {
+ const name = baseName(folder);
+ const media = list.filter((f) => isAttachableMedia(f.path));
+ const pick = (test: (base: string) => boolean) => list.find((f) => test(baseName(f.path)));
+ const infoJson = pick((b) => b.toLowerCase().endsWith(".info.json"));
+ const descriptionTxt = pick((b) => b.toLowerCase() === "description.txt");
+ const sourceTxt = pick((b) => /^sources?\.txt$/i.test(b));
+ const lost = isLostFolder(name);
+ const base = {
+ lost,
+ ...(infoJson ? { infoJson } : {}),
+ ...(descriptionTxt ? { descriptionTxt } : {}),
+ ...(sourceTxt ? { sourceTxt } : {}),
+ };
+ const folderId = folder ? idFromFolderName(name) : null;
+ if (folderId) {
+ out.push({ folder, name, id: folderId, idFrom: "folder", media, ...base });
+ continue;
+ }
+ // No id on the folder (or the archive's root): every media file is its own
+ // video, by its file name. Sidecars are shared only when there is one.
+ if (media.length === 0) {
+ out.push({ folder, name, id: null, idFrom: null, media, ...base });
+ continue;
+ }
+ for (const m of media) {
+ const id = idFromMediaName(m.path);
+ out.push({
+ folder,
+ name,
+ id,
+ idFrom: id ? "file" : null,
+ media: [m],
+ ...(media.length === 1 ? base : { lost }),
+ });
+ }
+ }
+ return out;
+}
+
+// The archive's title for a folder named `NNN - <title> - (<id>)`, with any
+// `[LOST]` mark dropped.
+export function titleFromFolderName(name: string): string {
+ return name
+ .replace(FOLDER_ID_RE, "")
+ .replace(/\s*-\s*$/, "")
+ .replace(/^\s*\d+\s*(?:\[LOST\]\s*)?-\s*/i, "")
+ .replace(/\[LOST\]/gi, "")
+ .trim();
+}
+
+const MONTHS = [
+ "january", "february", "march", "april", "may", "june",
+ "july", "august", "september", "october", "november", "december",
+];
+
+// "Published on July 13th 2018" (the archive's description.txt) → "20180713".
+export function publishedDateFromText(text: string): string | null {
+ const m = /published\s+on\s+([A-Za-z]+)\s+(\d{1,2})(?:st|nd|rd|th)?,?\s+(\d{4})/i.exec(text);
+ if (!m) return null;
+ const month = MONTHS.indexOf(m[1].toLowerCase());
+ const day = Number(m[2]);
+ if (month < 0 || day < 1 || day > 31) return null;
+ return `${m[3]}${String(month + 1).padStart(2, "0")}${String(day).padStart(2, "0")}`;
+}
+
+// The description a record gets from description.txt: the text without the
+// "Published on …" line (that is the record's upload_date).
+export function descriptionFromText(text: string): string {
+ return text
+ .split(/\r?\n/)
+ .filter((l) => !/^\s*published\s+on\s+/i.test(l))
+ .join("\n")
+ .trim();
+}
+
+// A metadata.info.json for a video the channel does not hold, built from the
+// folder's description.txt and name when it has no yt-dlp .info.json. Title
+// first (a reader scans the head for it), upload_date last (another scans the
+// tail).
+export function archiveFolderInfoJson(opts: {
+ id: string;
+ folderName: string;
+ descriptionText?: string;
+ platform?: string;
+ durationSec?: number | null;
+}): Record<string, unknown> {
+ const youtube = opts.platform === undefined || opts.platform === "youtube";
+ const date = opts.descriptionText ? publishedDateFromText(opts.descriptionText) : null;
+ const description = opts.descriptionText ? descriptionFromText(opts.descriptionText) : "";
+ const url = youtube ? `https://www.youtube.com/watch?v=${opts.id}` : undefined;
+ return {
+ id: opts.id,
+ title: titleFromFolderName(opts.folderName) || opts.id,
+ ...(description ? { description } : {}),
+ ...(typeof opts.durationSec === "number" && Number.isFinite(opts.durationSec) && opts.durationSec > 0
+ ? { duration: Math.round(opts.durationSec * 100) / 100 }
+ : {}),
+ ...(url ? { webpage_url: url, original_url: url } : {}),
+ ...(youtube ? { extractor: "youtube", extractor_key: "Youtube" } : {}),
+ ...(date ? { upload_date: date } : {}),
+ };
+}
+
+// --- classification ---------------------------------------------------------
+
+export const ATTACH_ROW_CLASSES = [
+ // Held, media found, no pointer (or `replace`): attached.
+ "attach",
+ // Not held, media found, `createRecords`: a record is written, then attached.
+ "create",
+ // Held and already has a saved container: left alone without `replace`.
+ "already-attached",
+ // Media found for a video the channel does not hold (no `createRecords`).
+ "not-held",
+ // A `[LOST]` folder: listed, never attached.
+ "lost",
+ // A folder with an id but no video file.
+ "no-media",
+ // Media with no id the rules can read.
+ "unmatched",
+ // More than one video file for one id: name the file with `items`.
+ "ambiguous",
+] as const;
+export type AttachRowClass = (typeof ATTACH_ROW_CLASSES)[number];
+
+export type AttachRow = {
+ class: AttachRowClass;
+ id: string | null;
+ folder: string;
+ // The file to copy ("attach"/"create"), or the candidates ("ambiguous").
+ media?: SourceFile;
+ candidates?: SourceFile[];
+ // The folder, for createRecords.
+ source?: ArchiveFolder;
+ // `attach` over an existing pointer (replace: true).
+ replacing?: boolean;
+};
+
+export type AttachPlan = {
+ rows: AttachRow[];
+ // Held videos no folder supplied media for (only when the whole archive was
+ // considered: no `items`, no `match`).
+ heldWithoutMedia: string[] | null;
+};
+
+export type AttachItem = { id: string; path: string };
+
+export function planAttach(opts: {
+ folders: readonly ArchiveFolder[];
+ held: ReadonlySet<string>;
+ attached: ReadonlySet<string>;
+ createRecords?: boolean;
+ replace?: boolean;
+ match?: RegExp;
+ // Explicit id → file pairs. When given, they are the whole plan.
+ items?: readonly AttachItem[];
+ // Every file in the archive, to resolve `items` paths.
+ files?: readonly SourceFile[];
+}): AttachPlan {
+ const classify = (id: string, media: SourceFile, folder: string, source?: ArchiveFolder): AttachRow => {
+ if (opts.held.has(id)) {
+ if (opts.attached.has(id) && !opts.replace) {
+ return { class: "already-attached", id, folder, media };
+ }
+ return {
+ class: "attach",
+ id,
+ folder,
+ media,
+ ...(source ? { source } : {}),
+ ...(opts.attached.has(id) ? { replacing: true } : {}),
+ };
+ }
+ return {
+ class: opts.createRecords ? "create" : "not-held",
+ id,
+ folder,
+ media,
+ ...(source ? { source } : {}),
+ };
+ };
+
+ if (opts.items) {
+ const byPath = new Map((opts.files ?? []).map((f) => [f.path, f]));
+ const folders = new Map(opts.folders.map((f) => [f.folder, f]));
+ const rows = opts.items.map((item): AttachRow => {
+ const media = byPath.get(item.path);
+ const folder = dirName(item.path);
+ if (!media) return { class: "no-media", id: item.id, folder };
+ return classify(item.id, media, folder, folders.get(folder));
+ });
+ return { rows, heldWithoutMedia: null };
+ }
+
+ const considered = opts.match
+ ? opts.folders.filter(
+ (f) => opts.match!.test(f.folder) || f.media.some((m) => opts.match!.test(m.path)),
+ )
+ : opts.folders;
+
+ // Two folders claiming one id: neither is attached on a guess.
+ const idCount = new Map<string, number>();
+ for (const f of considered) {
+ if (f.id && f.media.length > 0 && !f.lost) idCount.set(f.id, (idCount.get(f.id) ?? 0) + 1);
+ }
+
+ const rows: AttachRow[] = [];
+ for (const f of considered) {
+ if (f.lost) {
+ rows.push({ class: "lost", id: f.id, folder: f.folder });
+ continue;
+ }
+ if (!f.id) {
+ if (f.media.length > 0) {
+ rows.push({ class: "unmatched", id: null, folder: f.folder, candidates: f.media });
+ }
+ continue;
+ }
+ if (f.media.length === 0) {
+ rows.push({ class: "no-media", id: f.id, folder: f.folder });
+ continue;
+ }
+ if (f.media.length > 1 || (idCount.get(f.id) ?? 0) > 1) {
+ rows.push({ class: "ambiguous", id: f.id, folder: f.folder, candidates: f.media });
+ continue;
+ }
+ rows.push(classify(f.id, f.media[0], f.folder, f));
+ }
+
+ let heldWithoutMedia: string[] | null = null;
+ if (!opts.match) {
+ const supplied = new Set(
+ rows.filter((r) => r.media && r.id).map((r) => r.id as string),
+ );
+ for (const r of rows) if (r.class === "ambiguous" && r.id) supplied.add(r.id);
+ heldWithoutMedia = [...opts.held].filter((id) => !supplied.has(id)).sort();
+ }
+ return { rows, heldWithoutMedia };
+}
+
+export function countAttachRows(rows: readonly AttachRow[]): Record<AttachRowClass, number> {
+ const out = Object.fromEntries(ATTACH_ROW_CLASSES.map((c) => [c, 0])) as Record<AttachRowClass, number>;
+ for (const r of rows) out[r.class] += 1;
+ return out;
+}
diff --git a/common/lib/metadataHistory.ts b/common/lib/metadataHistory.ts
@@ -78,6 +78,10 @@ export const METADATA_HISTORY_WRITERS = [
// changes after the first fetch — a livestream VOD that gains its processed
// formats and auto-captions hours after it ends.
"refresh",
+ // A record created from a local archive's folder (release 21 D1,
+ // controller/attachMedia.ts `createRecords`): the folder's own yt-dlp
+ // .info.json, or one built from its description.txt / source.txt and name.
+ "local-archive",
] as const;
export type MetadataHistoryWriter = (typeof METADATA_HISTORY_WRITERS)[number];
diff --git a/common/lib/savedVideo-server.ts b/common/lib/savedVideo-server.ts
@@ -118,11 +118,21 @@ export async function persistSourceVideo(opts: {
// `saved-videos` as a real directory in the instant between the mover's
// rename and its symlink.
paths?: Pick<Paths, "transcriptsDir" | "savedVideosDir">;
+ // THE FILE TO MOVE, when it is not `<videoDir>/<sourceFilename>`. A container
+ // attached from a local archive (controller/attachMedia.ts) is copied out of
+ // the archive into a staging file BESIDE `storeDir`, on the store's disk, so
+ // the move is a same-filesystem rename — copying 36 GB through the corpus
+ // SSD first would be the wrong disk twice. The stored name is still
+ // `sourceFilename`.
+ sourcePath?: string;
+ // The bytes' sha256, when the caller already hashed them (an attach hashes
+ // while it copies); the backup step otherwise fills it in later.
+ sha256?: string;
}): Promise<SavedVideoPointer> {
if (opts.paths) {
await assertSavedVideosStoreWritable(opts.paths, opts.storeDir);
}
- const src = path.join(opts.videoDir, opts.sourceFilename);
+ const src = opts.sourcePath ?? path.join(opts.videoDir, opts.sourceFilename);
const dest = path.join(opts.storeDir, opts.sourceFilename);
const st = await stat(src);
await moveFileCrossDevice(src, dest);
@@ -132,6 +142,7 @@ export async function persistSourceVideo(opts: {
file: opts.sourceFilename,
bytes: st.size,
...(opts.keepReason ? { keepReason: opts.keepReason } : {}),
+ ...(opts.sha256 ? { sha256: opts.sha256 } : {}),
...(opts.origin ? { origin: opts.origin } : {}),
...(opts.format ? { format: opts.format } : {}),
};
diff --git a/common/lib/savedVideo.ts b/common/lib/savedVideo.ts
@@ -102,14 +102,32 @@ export function savedVideoFormatFromLine(
// The requester of a manually-sourced container: the tool, and the report clip
// it was wanted for. Deliberately the same vocabulary as ClipProvenance
// (lib/clipWindow.ts) so the video page renders both the same way.
+//
+// A CONTAINER ATTACHED FROM A LOCAL ARCHIVE (release 21 D1, `attach-media`)
+// carries `kind: "local-archive"` and says which archive, which entry in it,
+// the bytes' sha256 and when — the provenance of a file nobody downloaded.
+// `requestedBy` is still set (the tool, "attach-media"), so every reader that
+// renders "requested by" keeps working on it unchanged; a pointer without
+// `kind` is the requested-container shape it always was.
export type SavedVideoOrigin = {
requestedBy: string;
manifest?: string;
clipId?: string;
reason?: string;
requestedAt?: string;
+ kind?: SavedVideoOriginKind;
+ // local-archive: the archive (a .zip or .7z file, or a directory) as an
+ // absolute path, the entry's path inside it, the copied bytes' sha256 (hex)
+ // and when the attach wrote the pointer.
+ archive?: string;
+ entry?: string;
+ sha256?: string;
+ attachedAt?: string;
};
+export const SAVED_VIDEO_ORIGIN_KINDS = ["local-archive"] as const;
+export type SavedVideoOriginKind = (typeof SAVED_VIDEO_ORIGIN_KINDS)[number];
+
function parseOrigin(raw: unknown): SavedVideoOrigin | null {
if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null;
const r = raw as Record<string, unknown>;
@@ -119,6 +137,18 @@ function parseOrigin(raw: unknown): SavedVideoOrigin | null {
const v = r[k];
if (typeof v === "string" && v !== "") out[k] = v;
}
+ // An unknown kind is dropped with its fields rather than failing the
+ // pointer: the container is still there, only its provenance goes unread.
+ if (
+ typeof r.kind === "string" &&
+ (SAVED_VIDEO_ORIGIN_KINDS as readonly string[]).includes(r.kind)
+ ) {
+ out.kind = r.kind as SavedVideoOriginKind;
+ for (const k of ["archive", "entry", "sha256", "attachedAt"] as const) {
+ const v = r[k];
+ if (typeof v === "string" && v !== "") out[k] = v;
+ }
+ }
return out;
}
diff --git a/common/lib/zipReader.test.ts b/common/lib/zipReader.test.ts
@@ -0,0 +1,176 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { createHash } from "node:crypto";
+import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import path from "node:path";
+import { execa } from "execa";
+import { makeZip } from "../testing/makeZip";
+import {
+ ZipFormatError,
+ copyFileHashed,
+ copyZipEntry,
+ readZipEntries,
+ readZipEntryText,
+ zipEntryProblem,
+} from "./zipReader";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test lib/zipReader.test.ts
+
+const sha = (b: Buffer) => createHash("sha256").update(b).digest("hex");
+const VIDEO = Buffer.from(Array.from({ length: 200_000 }, (_, i) => (i * 7) & 0xff));
+
+async function withTmp(fn: (dir: string) => Promise<void>): Promise<void> {
+ const dir = await mkdtemp(path.join(tmpdir(), "ttb-zip-"));
+ try {
+ await fn(dir);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+}
+
+for (const zip64 of [false, true]) {
+ test(`reads the central directory${zip64 ? " (ZIP64)" : ""}: names, sizes, methods, directories`, async () => {
+ await withTmp(async (dir) => {
+ const file = path.join(dir, "a.zip");
+ await writeFile(
+ file,
+ makeZip(
+ [
+ { name: "YouTube/" },
+ { name: "YouTube/001 - Ünïcode | name - (abcdefghijk)/v-abcdefghijk.mp4", data: VIDEO },
+ { name: "YouTube/001 - Ünïcode | name - (abcdefghijk)/description.txt", data: "hello", deflate: true },
+ ],
+ { zip64 },
+ ),
+ );
+ const entries = await readZipEntries(file);
+ assert.equal(entries.length, 3);
+ assert.equal(entries[0].isDirectory, true);
+ assert.equal(entries[1].name, "YouTube/001 - Ünïcode | name - (abcdefghijk)/v-abcdefghijk.mp4");
+ assert.equal(entries[1].method, 0);
+ assert.equal(entries[1].size, VIDEO.length);
+ assert.equal(entries[2].method, 8);
+ assert.equal(entries[2].size, 5);
+ });
+ });
+
+ test(`copies a stored and a deflated entry out, hashed and CRC-checked${zip64 ? " (ZIP64)" : ""}`, async () => {
+ await withTmp(async (dir) => {
+ const file = path.join(dir, "a.zip");
+ await writeFile(
+ file,
+ makeZip(
+ [
+ { name: "a/stored.mp4", data: VIDEO },
+ { name: "a/deflated.mkv", data: VIDEO, deflate: true },
+ { name: "a/empty.txt", data: "" },
+ ],
+ { zip64 },
+ ),
+ );
+ const entries = await readZipEntries(file);
+ for (const [i, name] of [[0, "s"], [1, "d"]] as const) {
+ const dest = path.join(dir, name);
+ const got = await copyZipEntry(file, entries[i], dest);
+ assert.equal(got.bytes, VIDEO.length);
+ assert.equal(got.sha256, sha(VIDEO));
+ assert.deepEqual(await readFile(dest), VIDEO);
+ }
+ const empty = await copyZipEntry(file, entries[2], path.join(dir, "e"));
+ assert.equal(empty.bytes, 0);
+ });
+ });
+}
+
+test("a damaged entry is refused by its CRC, and the archive is never written", async () => {
+ await withTmp(async (dir) => {
+ const file = path.join(dir, "a.zip");
+ const zip = makeZip([{ name: "x.mp4", data: VIDEO }]);
+ // Flip one byte of the stored data (the local header is 30 + name bytes).
+ zip[30 + "x.mp4".length + 100] ^= 0xff;
+ await writeFile(file, zip);
+ const before = sha(await readFile(file));
+ const [entry] = await readZipEntries(file);
+ await assert.rejects(copyZipEntry(file, entry, path.join(dir, "out")), /CRC32 mismatch/);
+ assert.equal(sha(await readFile(file)), before);
+ });
+});
+
+test("an existing destination is never overwritten (staging files are exclusive)", async () => {
+ await withTmp(async (dir) => {
+ const file = path.join(dir, "a.zip");
+ await writeFile(file, makeZip([{ name: "x.mp4", data: VIDEO }]));
+ const [entry] = await readZipEntries(file);
+ await writeFile(path.join(dir, "out"), "keep");
+ await assert.rejects(copyZipEntry(file, entry, path.join(dir, "out")));
+ assert.equal(await readFile(path.join(dir, "out"), "utf8"), "keep");
+ });
+});
+
+test("text entries read, stored or deflated; a big one is refused", async () => {
+ await withTmp(async (dir) => {
+ const file = path.join(dir, "a.zip");
+ await writeFile(
+ file,
+ makeZip([
+ { name: "a.info.json", data: '{"id":"x"}' },
+ { name: "description.txt", data: "Published on July 13th 2018", deflate: true },
+ { name: "big.mp4", data: VIDEO },
+ ]),
+ );
+ const [info, desc, big] = await readZipEntries(file);
+ assert.equal(await readZipEntryText(file, info), '{"id":"x"}');
+ assert.equal(await readZipEntryText(file, desc), "Published on July 13th 2018");
+ await assert.rejects(readZipEntryText(file, big, 1000), /too large/);
+ });
+});
+
+test("unsupported methods and encryption are named, not half-read", () => {
+ const base = { name: "x", isDirectory: false, compressedSize: 1, size: 1, crc32: 0, localHeaderOffset: 0, encrypted: false };
+ assert.equal(zipEntryProblem({ ...base, method: 0 }), null);
+ assert.equal(zipEntryProblem({ ...base, method: 8 }), null);
+ assert.match(zipEntryProblem({ ...base, method: 14 }) ?? "", /method 14/);
+ assert.equal(zipEntryProblem({ ...base, method: 0, encrypted: true }), "encrypted");
+});
+
+test("not a zip: a sentence, not a crash", async () => {
+ await withTmp(async (dir) => {
+ const file = path.join(dir, "not.zip");
+ await writeFile(file, "plain text, no end record");
+ await assert.rejects(readZipEntries(file), ZipFormatError);
+ });
+});
+
+test("copyFileHashed hashes a plain file the same way", async () => {
+ await withTmp(async (dir) => {
+ await writeFile(path.join(dir, "src"), VIDEO);
+ const got = await copyFileHashed(path.join(dir, "src"), path.join(dir, "dst"));
+ assert.equal(got.sha256, sha(VIDEO));
+ assert.deepEqual(await readFile(path.join(dir, "dst")), VIDEO);
+ });
+});
+
+// Interop with a real writer: Info-ZIP's `zip`, stored (-0) and deflated, and
+// forced ZIP64 (-fz). Skipped where `zip` is not installed.
+test("reads what Info-ZIP's zip writes, ZIP64 included", async (t) => {
+ const probe = await execa("zip", ["-v"], { reject: false }).catch(() => null);
+ if (!probe || probe.exitCode !== 0) {
+ t.skip("zip is not installed");
+ return;
+ }
+ await withTmp(async (dir) => {
+ await writeFile(path.join(dir, "v.mp4"), VIDEO);
+ await writeFile(path.join(dir, "d.txt"), "some text ".repeat(100));
+ for (const args of [["-0"], ["-9"], ["-0", "-fz"]]) {
+ const file = path.join(dir, `out${args.join("")}.zip`);
+ await execa("zip", ["-q", "-X", ...args, file, "v.mp4", "d.txt"], { cwd: dir });
+ const entries = await readZipEntries(file);
+ const v = entries.find((e) => e.name === "v.mp4")!;
+ const got = await copyZipEntry(file, v, path.join(dir, `copy${args.join("")}`));
+ assert.equal(got.sha256, sha(VIDEO), args.join(" "));
+ const d = entries.find((e) => e.name === "d.txt")!;
+ assert.equal(await readZipEntryText(file, d), "some text ".repeat(100));
+ }
+ });
+});
diff --git a/common/lib/zipReader.ts b/common/lib/zipReader.ts
@@ -0,0 +1,308 @@
+// A ZIP READ IN PLACE (release 21 D1, controller/attachMedia.ts).
+//
+// The archives media is attached from are tens of GB and mostly STORED (no
+// compression: video does not deflate), so the useful operation is "copy this
+// entry's bytes out", and for a stored entry that is a byte range of the zip
+// file itself — no temp dir, no extraction tool, no second copy. A deflated
+// entry is inflated on the way out. Anything else (bzip2, LZMA, encryption) is
+// refused by name rather than half-read.
+//
+// Only the CENTRAL DIRECTORY is trusted for sizes and offsets (the local
+// header's sizes are zero when bit 3 is set), with ZIP64 throughout: a 36 GB
+// archive's later offsets do not fit 32 bits. The local header is read only
+// for its own name/extra lengths, which is where the data starts.
+//
+// READ-ONLY. Every open here is `r`; nothing writes to the archive.
+//
+// SERVER-ONLY (node:fs).
+
+import { createReadStream, createWriteStream } from "node:fs";
+import { open, type FileHandle } from "node:fs/promises";
+import { createHash } from "node:crypto";
+import { Transform, type TransformCallback } from "node:stream";
+import { pipeline } from "node:stream/promises";
+import { createInflateRaw, crc32, inflateRawSync } from "node:zlib";
+import { Readable } from "node:stream";
+
+export type ZipEntry = {
+ // The entry's path inside the archive, "/"-separated, as stored.
+ name: string;
+ isDirectory: boolean;
+ // 0 = stored, 8 = deflated; anything else is refused at extraction.
+ method: number;
+ compressedSize: number;
+ size: number;
+ crc32: number;
+ localHeaderOffset: number;
+ encrypted: boolean;
+};
+
+const EOCD_SIG = 0x06054b50;
+const ZIP64_LOCATOR_SIG = 0x07064b50;
+const ZIP64_EOCD_SIG = 0x06064b50;
+const CDH_SIG = 0x02014b50;
+const LFH_SIG = 0x04034b50;
+const U32_MAX = 0xffffffff;
+const U16_MAX = 0xffff;
+
+export class ZipFormatError extends Error {}
+
+async function readAt(fh: FileHandle, offset: number, length: number): Promise<Buffer> {
+ const buf = Buffer.alloc(length);
+ let got = 0;
+ while (got < length) {
+ const { bytesRead } = await fh.read(buf, got, length - got, offset + got);
+ if (bytesRead === 0) break;
+ got += bytesRead;
+ }
+ return got === length ? buf : buf.subarray(0, got);
+}
+
+// A 64-bit little-endian value as a JS number. Every offset and size in a
+// real archive is far below 2^53; one that is not is a corrupt file.
+function u64(buf: Buffer, at: number): number {
+ const v = buf.readBigUInt64LE(at);
+ if (v > BigInt(Number.MAX_SAFE_INTEGER)) {
+ throw new ZipFormatError("a ZIP64 value exceeds 2^53");
+ }
+ return Number(v);
+}
+
+const utf8Strict = new TextDecoder("utf-8", { fatal: true });
+
+// Bit 11 says UTF-8. Without it the spec says CP437, but in practice archives
+// written on Linux and by most tools are UTF-8 without the flag; valid UTF-8
+// is read as UTF-8, anything else as latin1 (identical to CP437 for ASCII).
+function decodeName(raw: Buffer, utf8Flag: boolean): string {
+ if (utf8Flag) return raw.toString("utf8");
+ try {
+ return utf8Strict.decode(raw);
+ } catch {
+ return raw.toString("latin1");
+ }
+}
+
+type Extra = { id: number; data: Buffer };
+
+function parseExtras(buf: Buffer): Extra[] {
+ const out: Extra[] = [];
+ let at = 0;
+ while (at + 4 <= buf.length) {
+ const id = buf.readUInt16LE(at);
+ const len = buf.readUInt16LE(at + 2);
+ const data = buf.subarray(at + 4, at + 4 + len);
+ out.push({ id, data });
+ at += 4 + len;
+ }
+ return out;
+}
+
+async function findEocd(fh: FileHandle, fileSize: number): Promise<{ buf: Buffer; at: number }> {
+ // EOCD is 22 bytes plus a comment of up to 65535.
+ const tail = Math.min(fileSize, 22 + 0xffff);
+ const buf = await readAt(fh, fileSize - tail, tail);
+ for (let i = buf.length - 22; i >= 0; i--) {
+ if (buf.readUInt32LE(i) === EOCD_SIG) {
+ return { buf: buf.subarray(i), at: fileSize - tail + i };
+ }
+ }
+ throw new ZipFormatError("not a zip file (no end-of-central-directory record)");
+}
+
+// Every entry of the archive, from its central directory.
+export async function readZipEntries(file: string): Promise<ZipEntry[]> {
+ const fh = await open(file, "r");
+ try {
+ const { size: fileSize } = await fh.stat();
+ const eocd = await findEocd(fh, fileSize);
+ let count = eocd.buf.readUInt16LE(10);
+ let cdSize = eocd.buf.readUInt32LE(12);
+ let cdOffset = eocd.buf.readUInt32LE(16);
+ if (count === U16_MAX || cdSize === U32_MAX || cdOffset === U32_MAX) {
+ if (eocd.at < 20) throw new ZipFormatError("ZIP64 locator missing");
+ const loc = await readAt(fh, eocd.at - 20, 20);
+ if (loc.readUInt32LE(0) !== ZIP64_LOCATOR_SIG) {
+ throw new ZipFormatError("ZIP64 locator missing");
+ }
+ const z64At = u64(loc, 8);
+ const z64 = await readAt(fh, z64At, 56);
+ if (z64.readUInt32LE(0) !== ZIP64_EOCD_SIG) {
+ throw new ZipFormatError("ZIP64 end-of-central-directory record missing");
+ }
+ count = u64(z64, 32);
+ cdSize = u64(z64, 40);
+ cdOffset = u64(z64, 48);
+ }
+ const cd = await readAt(fh, cdOffset, cdSize);
+ const entries: ZipEntry[] = [];
+ let at = 0;
+ for (let i = 0; i < count; i++) {
+ if (at + 46 > cd.length || cd.readUInt32LE(at) !== CDH_SIG) {
+ throw new ZipFormatError(`central directory entry ${i} is malformed`);
+ }
+ const flags = cd.readUInt16LE(at + 8);
+ const method = cd.readUInt16LE(at + 10);
+ const crc = cd.readUInt32LE(at + 16);
+ let compressedSize = cd.readUInt32LE(at + 20);
+ let size = cd.readUInt32LE(at + 24);
+ const nameLen = cd.readUInt16LE(at + 28);
+ const extraLen = cd.readUInt16LE(at + 30);
+ const commentLen = cd.readUInt16LE(at + 32);
+ let localHeaderOffset = cd.readUInt32LE(at + 42);
+ const rawName = cd.subarray(at + 46, at + 46 + nameLen);
+ const extras = parseExtras(cd.subarray(at + 46 + nameLen, at + 46 + nameLen + extraLen));
+ let name = decodeName(rawName, (flags & 0x800) !== 0);
+ for (const ex of extras) {
+ if (ex.id === 0x0001) {
+ // ZIP64: only the fields whose 32-bit slot is saturated, in order.
+ let p = 0;
+ if (size === U32_MAX) { size = u64(ex.data, p); p += 8; }
+ if (compressedSize === U32_MAX) { compressedSize = u64(ex.data, p); p += 8; }
+ if (localHeaderOffset === U32_MAX) { localHeaderOffset = u64(ex.data, p); p += 8; }
+ } else if (ex.id === 0x7075 && ex.data.length > 5 && ex.data[0] === 1) {
+ // Info-ZIP Unicode Path: version 1, CRC32 of the raw name, UTF-8 name.
+ if (ex.data.readUInt32LE(1) === crc32(rawName)) {
+ name = ex.data.subarray(5).toString("utf8");
+ }
+ }
+ }
+ name = name.replace(/\\/g, "/");
+ entries.push({
+ name,
+ isDirectory: name.endsWith("/"),
+ method,
+ compressedSize,
+ size,
+ crc32: crc,
+ localHeaderOffset,
+ encrypted: (flags & 0x1) !== 0,
+ });
+ at += 46 + nameLen + extraLen + commentLen;
+ }
+ return entries;
+ } finally {
+ await fh.close();
+ }
+}
+
+// Where an entry's data begins: past its local header, whose name and extra
+// lengths may differ from the central directory's.
+export async function zipEntryDataOffset(file: string, entry: ZipEntry): Promise<number> {
+ const fh = await open(file, "r");
+ try {
+ const lfh = await readAt(fh, entry.localHeaderOffset, 30);
+ if (lfh.length < 30 || lfh.readUInt32LE(0) !== LFH_SIG) {
+ throw new ZipFormatError(`no local header for ${entry.name}`);
+ }
+ return entry.localHeaderOffset + 30 + lfh.readUInt16LE(26) + lfh.readUInt16LE(28);
+ } finally {
+ await fh.close();
+ }
+}
+
+// Why an entry cannot be copied out, or null when it can.
+export function zipEntryProblem(entry: ZipEntry): string | null {
+ if (entry.isDirectory) return "a directory";
+ if (entry.encrypted) return "encrypted";
+ if (entry.method !== 0 && entry.method !== 8) {
+ return `compression method ${entry.method} (only stored and deflated are read)`;
+ }
+ return null;
+}
+
+// Counts bytes, CRC32 and sha256 of what passes through it.
+class Digest extends Transform {
+ bytes = 0;
+ crc = 0;
+ readonly hash = createHash("sha256");
+ _transform(chunk: Buffer, _enc: BufferEncoding, cb: TransformCallback): void {
+ this.bytes += chunk.length;
+ this.crc = crc32(chunk, this.crc);
+ this.hash.update(chunk);
+ cb(null, chunk);
+ }
+}
+
+export type CopiedEntry = { bytes: number; sha256: string };
+
+// Copy one entry's (uncompressed) bytes to `dest`, verifying its size and
+// CRC32 against the central directory. A stored entry is a byte range of the
+// archive; a deflated one is inflated. `dest` is written and left in place on
+// success; on any failure it is removed by the caller (it is a staging path).
+export async function copyZipEntry(
+ file: string,
+ entry: ZipEntry,
+ dest: string,
+ opts: { signal?: AbortSignal } = {},
+): Promise<CopiedEntry> {
+ const problem = zipEntryProblem(entry);
+ if (problem) throw new ZipFormatError(`${entry.name}: ${problem}`);
+ const start = await zipEntryDataOffset(file, entry);
+ const digest = new Digest();
+ const out = createWriteStream(dest, { flags: "wx" });
+ if (entry.compressedSize === 0) {
+ await pipeline(Readable.from([]), digest, out, { signal: opts.signal });
+ } else {
+ const src = createReadStream(file, {
+ start,
+ end: start + entry.compressedSize - 1,
+ highWaterMark: 1 << 20,
+ });
+ if (entry.method === 0) {
+ await pipeline(src, digest, out, { signal: opts.signal });
+ } else {
+ await pipeline(src, createInflateRaw(), digest, out, { signal: opts.signal });
+ }
+ }
+ if (digest.bytes !== entry.size) {
+ throw new ZipFormatError(
+ `${entry.name}: ${digest.bytes} bytes copied, the archive says ${entry.size}`,
+ );
+ }
+ if (digest.crc >>> 0 !== entry.crc32 >>> 0) {
+ throw new ZipFormatError(`${entry.name}: CRC32 mismatch — the archive is damaged here`);
+ }
+ return { bytes: digest.bytes, sha256: digest.hash.digest("hex") };
+}
+
+// A small entry's text (an .info.json, a description.txt). Refuses anything
+// over `maxBytes` so a mislabelled video is never read into memory.
+export async function readZipEntryText(
+ file: string,
+ entry: ZipEntry,
+ maxBytes = 4 * 1024 * 1024,
+): Promise<string> {
+ const problem = zipEntryProblem(entry);
+ if (problem) throw new ZipFormatError(`${entry.name}: ${problem}`);
+ if (entry.size > maxBytes) {
+ throw new ZipFormatError(`${entry.name}: ${entry.size} bytes is too large to read as text`);
+ }
+ const start = await zipEntryDataOffset(file, entry);
+ const fh = await open(file, "r");
+ let raw: Buffer;
+ try {
+ raw = await readAt(fh, start, entry.compressedSize);
+ } finally {
+ await fh.close();
+ }
+ if (entry.method === 8) {
+ raw = inflateRawSync(raw);
+ }
+ return raw.toString("utf8");
+}
+
+// Hash a plain file the same way (a directory source's entry).
+export async function copyFileHashed(
+ src: string,
+ dest: string,
+ opts: { signal?: AbortSignal } = {},
+): Promise<CopiedEntry> {
+ const digest = new Digest();
+ await pipeline(
+ createReadStream(src, { highWaterMark: 1 << 20 }),
+ digest,
+ createWriteStream(dest, { flags: "wx" }),
+ { signal: opts.signal },
+ );
+ return { bytes: digest.bytes, sha256: digest.hash.digest("hex") };
+}
diff --git a/common/testing/makeZip.ts b/common/testing/makeZip.ts
@@ -0,0 +1,95 @@
+// A TINY ZIP WRITER, for the unit tests of lib/zipReader.ts and
+// controller/attachMedia.ts: stored and deflated entries, UTF-8 names, and a
+// forced-ZIP64 variant (every size and offset in the ZIP64 extra field, the
+// ZIP64 end records present), so the reader's 64-bit paths run without a
+// 4 GB fixture. Synthetic bytes only — never a real archive.
+
+import { crc32, deflateRawSync } from "node:zlib";
+
+export type ZipFixtureEntry = {
+ name: string;
+ data?: Buffer | string;
+ // Default: stored. A directory entry has a name ending in "/" and no data.
+ deflate?: boolean;
+};
+
+export function makeZip(entries: ZipFixtureEntry[], opts: { zip64?: boolean } = {}): Buffer {
+ const zip64 = !!opts.zip64;
+ const locals: Buffer[] = [];
+ const centrals: Buffer[] = [];
+ let offset = 0;
+ for (const e of entries) {
+ const name = Buffer.from(e.name, "utf8");
+ const raw = e.data === undefined ? Buffer.alloc(0) : Buffer.isBuffer(e.data) ? e.data : Buffer.from(e.data, "utf8");
+ const method = e.deflate ? 8 : 0;
+ const body = e.deflate ? deflateRawSync(raw) : raw;
+ const crc = crc32(raw) >>> 0;
+ const flags = 0x800;
+
+ const localExtra = zip64 ? z64Extra([raw.length, body.length]) : Buffer.alloc(0);
+ const lfh = Buffer.alloc(30);
+ lfh.writeUInt32LE(0x04034b50, 0);
+ lfh.writeUInt16LE(zip64 ? 45 : 20, 4);
+ lfh.writeUInt16LE(flags, 6);
+ lfh.writeUInt16LE(method, 8);
+ lfh.writeUInt32LE(crc, 14);
+ lfh.writeUInt32LE(zip64 ? 0xffffffff : body.length, 18);
+ lfh.writeUInt32LE(zip64 ? 0xffffffff : raw.length, 22);
+ lfh.writeUInt16LE(name.length, 26);
+ lfh.writeUInt16LE(localExtra.length, 28);
+ const local = Buffer.concat([lfh, name, localExtra, body]);
+
+ const cdExtra = zip64 ? z64Extra([raw.length, body.length, offset]) : Buffer.alloc(0);
+ const cdh = Buffer.alloc(46);
+ cdh.writeUInt32LE(0x02014b50, 0);
+ cdh.writeUInt16LE(zip64 ? 45 : 20, 4);
+ cdh.writeUInt16LE(zip64 ? 45 : 20, 6);
+ cdh.writeUInt16LE(flags, 8);
+ cdh.writeUInt16LE(method, 10);
+ cdh.writeUInt32LE(crc, 16);
+ cdh.writeUInt32LE(zip64 ? 0xffffffff : body.length, 20);
+ cdh.writeUInt32LE(zip64 ? 0xffffffff : raw.length, 24);
+ cdh.writeUInt16LE(name.length, 28);
+ cdh.writeUInt16LE(cdExtra.length, 30);
+ cdh.writeUInt32LE(zip64 ? 0xffffffff : offset, 42);
+ centrals.push(Buffer.concat([cdh, name, cdExtra]));
+
+ locals.push(local);
+ offset += local.length;
+ }
+ const cd = Buffer.concat(centrals);
+ const cdOffset = offset;
+ const tail: Buffer[] = [];
+ if (zip64) {
+ const rec = Buffer.alloc(56);
+ rec.writeUInt32LE(0x06064b50, 0);
+ rec.writeBigUInt64LE(44n, 4);
+ rec.writeUInt16LE(45, 12);
+ rec.writeUInt16LE(45, 14);
+ rec.writeBigUInt64LE(BigInt(entries.length), 24);
+ rec.writeBigUInt64LE(BigInt(entries.length), 32);
+ rec.writeBigUInt64LE(BigInt(cd.length), 40);
+ rec.writeBigUInt64LE(BigInt(cdOffset), 48);
+ const loc = Buffer.alloc(20);
+ loc.writeUInt32LE(0x07064b50, 0);
+ loc.writeBigUInt64LE(BigInt(cdOffset + cd.length), 8);
+ loc.writeUInt32LE(1, 16);
+ tail.push(rec, loc);
+ }
+ const eocd = Buffer.alloc(22);
+ eocd.writeUInt32LE(0x06054b50, 0);
+ eocd.writeUInt16LE(zip64 ? 0xffff : entries.length, 8);
+ eocd.writeUInt16LE(zip64 ? 0xffff : entries.length, 10);
+ eocd.writeUInt32LE(zip64 ? 0xffffffff : cd.length, 12);
+ eocd.writeUInt32LE(zip64 ? 0xffffffff : cdOffset, 16);
+ tail.push(eocd);
+ return Buffer.concat([...locals, cd, ...tail]);
+}
+
+function z64Extra(values: number[]): Buffer {
+ const b = Buffer.alloc(4 + 8 * values.length);
+ b.writeUInt16LE(0x0001, 0);
+ b.writeUInt16LE(8 * values.length, 2);
+ values.forEach((v, i) => b.writeBigUInt64LE(BigInt(v), 4 + 8 * i));
+ return b;
+}
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **A channel's videos can get their media from a local archive.** `pnpm ops attach-media` (`POST /api/ops/attach-media`) and `archilyzer media attach <slug> <source>` take `{"slug", "source"}` — an absolute path to a directory, a `.zip` (read in place: a stored entry is copied straight out of it, with no temp dir) or a `.7z` — and put each held video's file into the saved-video store as its source container, so clip windows and report clips can be cut from it with nothing fetched. The id is the folder's trailing `(<id>)`, else the file's yt-dlp suffix; `"items": [{"id", "path"}]` names exact files and `"match"` narrows the folders. The pointer records where the file came from: `origin: {kind: "local-archive", archive, entry, sha256, attachedAt}`. A video that already has a saved container is left alone unless `"replace": true`; `"createRecords": true` writes a record for a video the channel does not hold (from the folder's yt-dlp `.info.json`, else its `description.txt` and name); `[LOST]` folders are listed and never attached. `"dryRun": true` lists what each folder is, and the held videos the archive has no media for, and writes nothing. The archive is never written to.
- **A curated tag can exist on some sites only.** A tag's new **Sites** field on /tags (`sites` in `transcripts/tags.json`; `pnpm ops tags` takes it in a define) names the sites it exists on. Its rules then fire, and its pins apply, only to videos on those sites' channels, and every other site drops it from its records, its counts and its `/tags.json` — where **Hidden** only hid the chip. Empty is every site, as before. Setting it, or changing the channels of those sites, re-derives the corpus's tags once at the next index update. The Eva tags are what this is for: they belong on Anilyzer alone.
- **The publish lane.** Publishing can run itself: turn it on at **/operations/publish** (the runner's Start, Drain and Stop, the hold, and the lane's settings; or `publish.enabled` in settings) and the lane checks every `checkEveryMinutes` (10) whether the index is stale; when it is — and its last update is at least `refreshEveryMinutes` (360) old — it updates it, then builds every site whose channels changed or whose data the new index moved, one stage at a time on the `publish` queue. What it may do with a site is the site's own — the **Publish policy** on the site's settings form, `site.json` `publish.auto` —: `off` (the default: left alone), `build`, `preview` (built and deployed to the preview branch `publish.previewBranch`) or `production`; the hub and the homepage have `publish.hub` and `publish.homepage`. A private site is only ever built, and a site needs its Cloudflare Pages project before it may deploy. Hold the lane and the stage running finishes and no next one starts; quiet hours (`publish.quietHours`) do the same; Drain finishes the stage and ends the runner. The lane never forces a stage: a stage that finds its target current does nothing. On /jobs every stage of one run reads `run <id> · <target>`, and a stage still queued when the editor restarts is cancelled, never re-queued — the lane works out again what is stale from what is on disk. `archilyzer publish now` runs the same plan from the command line, one stage after another in its own process.
- **One index for every site.** The index is updated once and every site, the hub and the homepage are built from it; `archilyzer publish status` says, per site, whether its build is current — "stale: 3 channels changed (a, b, c)" as soon as a download, transcription or digest on one of its channels finishes, before any index runs; "stale: data changed" once the index has run and the site's data moved; "stale: config changed" after its site.json, tags or aliases changed — and whether what is deployed is that build, with a build made by older code marked "code newer" but not stale.
diff --git a/editor/app/api/ops/attach-media/route.test.ts b/editor/app/api/ops/attach-media/route.test.ts
@@ -0,0 +1,91 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { mkdir, mkdtemp, readFile, readdir, rm, writeFile } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+import { makeZip } from "../../../../../common/testing/makeZip";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/attach-media/route.test.ts"
+//
+// The body's shape, the refusals that come before any job, and one dry run
+// over a synthetic zip, answered from a temp corpus. Nothing is attached.
+
+const ROOT = await mkdtemp(path.join(os.tmpdir(), "attach-media-route-"));
+// Set before the route (and getPaths, which caches) is first imported.
+process.env.WORKER_TOKEN = "test-token";
+process.env.TRANSCRIPTS_DIR = ROOT;
+process.env.SETTINGS_FILE = path.join(ROOT, "settings.json");
+const { POST } = await import("./route");
+const { getRegistry } = await import("yt-dlp-transcript-common/jobs/registry");
+const { getPaths } = await import("yt-dlp-transcript-common/lib/paths");
+test.after(() => rm(ROOT, { recursive: true, force: true }));
+
+const CH = path.join(ROOT, "channels");
+await mkdir(path.join(CH, "tism", "data", "aaaaaaaaaaa"), { recursive: true });
+await writeFile(path.join(CH, "tism", "config.json"), JSON.stringify({ handling: "youtube", name: "tism", platform: "youtube" }));
+await writeFile(path.join(CH, "tism", "data", "aaaaaaaaaaa", "metadata.info.json"), JSON.stringify({ id: "aaaaaaaaaaa" }));
+const ZIP = path.join(ROOT, "archive.zip");
+await writeFile(
+ ZIP,
+ makeZip([
+ { name: "YouTube/001 - A - (aaaaaaaaaaa)/A-aaaaaaaaaaa.mp4", data: Buffer.alloc(1000, 1) },
+ { name: "YouTube/002 [LOST] - B - (bbbbbbbbbbb)/source.txt", data: "x" },
+ { name: "YouTube/003 - C - (ccccccccccc)/ccccccccccc.mp4", data: Buffer.alloc(500, 2) },
+ ]),
+);
+
+async function post(
+ body: Record<string, unknown>,
+): Promise<{ status: number; json: { ok?: boolean; error?: string; jobId?: string } }> {
+ const res = await POST(
+ new Request("http://localhost/api/ops/attach-media", {
+ method: "POST",
+ headers: { authorization: "Bearer test-token", "content-type": "application/json" },
+ body: JSON.stringify(body),
+ }),
+ );
+ return { status: res.status, json: (await res.json()) as { ok?: boolean; error?: string; jobId?: string } };
+}
+
+test("refusals before any job: keys, slug, source, items, match, channel", async () => {
+ const cases: [Record<string, unknown>, RegExp][] = [
+ [{ slug: "tism" }, /"source" is required/],
+ [{ source: ZIP }, /"slug" is required/],
+ [{ slug: "../x", source: ZIP }, /not a valid channel slug/],
+ [{ slug: "tism", source: ZIP, force: true }, /unknown key\(s\): force/],
+ [{ slug: "tism", source: "relative/archive.zip" }, /must be an absolute path/],
+ [{ slug: "tism", source: ZIP, items: [] }, /non-empty array/],
+ [{ slug: "tism", source: ZIP, items: [{ id: "..", path: "x" }] }, /not a video id/],
+ [{ slug: "tism", source: ZIP, items: [{ id: "a", path: "x", size: 1 }] }, /unknown key\(s\): size/],
+ [{ slug: "tism", source: ZIP, items: [{ id: "a", path: "x" }], match: "x" }, /"items" or "match", not both/],
+ [{ slug: "tism", source: ZIP, match: "(" }, /not a valid regex/],
+ [{ slug: "tism", source: ZIP, dryRun: "yes" }, /"dryRun" must be a boolean/],
+ [{ slug: "no-such", source: ZIP }, /Channel "no-such" not found/],
+ ];
+ for (const [body, re] of cases) {
+ const r = await post(body);
+ assert.equal(r.status, 400, JSON.stringify(body));
+ assert.match(r.json.error!, re, JSON.stringify(body));
+ assert.equal(r.json.jobId, undefined);
+ }
+});
+
+test("a dry run is a job whose log classifies the archive and writes nothing", async () => {
+ const r = await post({ slug: "tism", source: ZIP, dryRun: true });
+ assert.equal(r.status, 200, JSON.stringify(r.json));
+ assert.ok(r.json.jobId);
+ for (let i = 0; i < 100; i++) {
+ const status = getRegistry().get(r.json.jobId!)?.status;
+ if (status === "done" || status === "failed") break;
+ await new Promise((res) => setTimeout(res, 50));
+ }
+ assert.equal(getRegistry().get(r.json.jobId!)?.status, "done");
+ const log = await readFile(path.join(getPaths().jobsDir, `${r.json.jobId}.log`), "utf8");
+ const summary = JSON.parse(log.split("\n").find((l) => l.startsWith("summary: "))!.slice("summary: ".length));
+ assert.equal(summary.dryRun, true);
+ assert.equal(summary.counts.attach, 1);
+ assert.deepEqual(summary.notHeld, ["ccccccccccc"]);
+ assert.equal(summary.lost, 1);
+ assert.deepEqual(await readdir(path.join(CH, "tism", "data", "aaaaaaaaaaa")), ["metadata.info.json"]);
+});
diff --git a/editor/app/api/ops/attach-media/route.ts b/editor/app/api/ops/attach-media/route.ts
@@ -0,0 +1,70 @@
+import { attachMediaAction } from "../../../channels/[slug]/attachMediaActions";
+import {
+ OpsInputError,
+ jobResponse,
+ ops,
+ optBool,
+ optString,
+ reqSlug,
+ reqString,
+ reqVideoId,
+ type OpsBody,
+} from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { slug, source, items?: [{id, path}], match?, createRecords?, replace?, dryRun? }
+// -> { ok: true, jobId }
+//
+// Local media attached to a channel's held videos (release 21 D1): each
+// video's file copied out of `source` — an absolute path to a directory, a
+// .zip (read in place) or a .7z — into the saved-video store, as one
+// `attach-media` job on the channel's queue (controller/attachMedia.ts). The
+// id comes from the folder's trailing "(<id>)", else the file's yt-dlp
+// suffix; `items` names exact files instead (each `path` is the entry's path
+// inside the source), `match` narrows the folders by a case-insensitive regex.
+// `createRecords` writes a record for a video the channel does not hold;
+// `replace` re-attaches over an existing saved container. `dryRun` logs what
+// each folder is (attach, not held, already attached, [LOST], no media,
+// unmatched, ambiguous) and the held videos the archive has no media for, and
+// writes nothing. The job's log ends with a `summary: {…}` line.
+function optItems(body: OpsBody): { id: string; path: string }[] | undefined {
+ const v = body.items;
+ if (v === undefined) return undefined;
+ if (!Array.isArray(v) || v.length === 0) {
+ throw new OpsInputError('"items" must be a non-empty array of { "id", "path" }');
+ }
+ return v.map((entry, i) => {
+ if (typeof entry !== "object" || entry === null || Array.isArray(entry)) {
+ throw new OpsInputError(`"items[${i}]" must be an object { "id", "path" }`);
+ }
+ const extra = Object.keys(entry).filter((k) => k !== "id" && k !== "path");
+ if (extra.length) {
+ throw new OpsInputError(
+ `"items[${i}]" has unknown key(s): ${extra.join(", ")} — an item is { "id", "path" }`,
+ );
+ }
+ return {
+ id: reqVideoId(entry as OpsBody, "id"),
+ path: reqString(entry as OpsBody, "path"),
+ };
+ });
+}
+
+export async function POST(request: Request) {
+ return ops(
+ request,
+ ["slug", "source", "items", "match", "createRecords", "replace", "dryRun"],
+ async (body) =>
+ jobResponse(
+ await attachMediaAction(reqSlug(body, "slug"), {
+ source: reqString(body, "source"),
+ items: optItems(body),
+ match: optString(body, "match"),
+ createRecords: optBool(body, "createRecords"),
+ replace: optBool(body, "replace"),
+ dryRun: optBool(body, "dryRun"),
+ }),
+ ),
+ );
+}
diff --git a/editor/app/channels/[slug]/attachMediaActions.ts b/editor/app/channels/[slug]/attachMediaActions.ts
@@ -0,0 +1,94 @@
+"use server";
+
+import path from "node:path";
+import { safeRevalidate } from "../../lib/safeRevalidate";
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import { channelQueueKey } from "yt-dlp-transcript-common/lib/queueKeys";
+import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels";
+import {
+ attachMedia,
+ summaryOf,
+} from "yt-dlp-transcript-common/controller/attachMedia";
+import {
+ runManagedFunction,
+ type StreamActionResult,
+} from "yt-dlp-transcript-common/jobs/streamCommand";
+
+// LOCAL MEDIA ATTACHED TO A CHANNEL'S HELD VIDEOS (release 21 D1): the
+// `attach-media` job, as `pnpm ops attach-media` posts it
+// (`/api/ops/attach-media`). The work and its rules are
+// common/controller/attachMedia.ts; this is the enqueue and the refusals a
+// caller gets before any job exists.
+//
+// On the channel's own queue: nothing is fetched, so no platform queue is owed
+// a turn, and a second attach on the same channel waits behind the first
+// rather than racing it for the same pointers. The kind's `needsMedia` is what
+// keeps it off a channel whose media is moving or unreachable.
+export async function attachMediaAction(
+ slug: string,
+ opts: {
+ source: string;
+ items?: { id: string; path: string }[];
+ match?: string;
+ createRecords?: boolean;
+ replace?: boolean;
+ dryRun?: boolean;
+ },
+): Promise<StreamActionResult> {
+ const paths = getPaths();
+ const channelConfig = await readChannelConfig(paths, slug);
+ if (!channelConfig) {
+ return { ok: false, error: `Channel "${slug}" not found` };
+ }
+ const source = opts.source.trim();
+ if (!path.isAbsolute(source)) {
+ return { ok: false, error: `"source" must be an absolute path (a directory, a .zip or a .7z)` };
+ }
+ if (opts.items !== undefined && opts.match !== undefined) {
+ return { ok: false, error: 'Send "items" or "match", not both: items name the files exactly' };
+ }
+ if (opts.match !== undefined) {
+ try {
+ new RegExp(opts.match, "i");
+ } catch (e) {
+ return { ok: false, error: `"match" is not a valid regex: ${(e as Error).message}` };
+ }
+ }
+ return runManagedFunction({
+ kind: "attach-media",
+ queueKey: channelQueueKey(slug),
+ paths,
+ channelSlug: slug,
+ fn: async (onLog, signal, _setProgress, ctx) => {
+ const result = await attachMedia({
+ paths,
+ slug,
+ channelConfig,
+ source,
+ ...(opts.items ? { items: opts.items } : {}),
+ ...(opts.match !== undefined ? { match: opts.match } : {}),
+ ...(opts.createRecords ? { createRecords: true } : {}),
+ ...(opts.replace ? { replace: true } : {}),
+ ...(opts.dryRun ? { dryRun: true } : {}),
+ onLog,
+ signal,
+ drainSignal: ctx.drainSignal,
+ });
+ if (!opts.dryRun) {
+ safeRevalidate([
+ `/channels/${slug}`,
+ "/channels",
+ ...result.attached.map((id) => `/channels/${slug}/videos/${id}`),
+ ]);
+ }
+ if (result.failed.length > 0 && result.attached.length === 0) {
+ throw new Error(
+ `${result.failed.length} failed, none attached: ${JSON.stringify(summaryOf(result).failed)}`,
+ );
+ }
+ if (result.stopped && result.stopped !== "drained" && result.stopped !== "cancelled") {
+ throw new Error(`stopped: ${result.stopped}`);
+ }
+ },
+ });
+}
diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs
@@ -130,6 +130,10 @@ const ACTIONS = [
// Chosen media files of ONE archive.org item ({slug, item, files: [...] |
// match: "<regex>", dryRun?}): one job, one file at a time, paced.
"import-archive-org",
+ // A channel's held videos given their media from a LOCAL archive — a
+ // directory, a zip read in place, a 7z ({slug, source, items?, match?,
+ // createRecords?, replace?, dryRun?}): one job, nothing fetched.
+ "attach-media",
// A podcast channel's records completed from its RSS feed ({slug, dryRun?}):
// one fetch of the feed, no media.
"feed-metadata",
@@ -511,6 +515,18 @@ export function usage() {
" queue; a low disk or a rate limit stops it, and running the same body",
" again resumes — saved videos are skipped.",
"",
+ 'attach-media copies each held video\'s file out of a LOCAL archive into the',
+ ' saved-video store, as its source container (nothing fetched): {"slug",',
+ ' "source"} — an absolute path to a directory, a .zip (read in place) or a',
+ ' .7z. The id is the folder\'s trailing "(<id>)", else the file\'s yt-dlp',
+ ' suffix; "items": [{"id", "path"}] names exact files (path inside the',
+ ' source), "match" narrows the folders by regex. "createRecords": true',
+ ' writes a record for a video the channel does not hold; "replace": true',
+ ' re-attaches over a saved container. "dryRun": true logs each folder\'s',
+ " class (attach, not-held, already-attached, lost, no-media, unmatched,",
+ " ambiguous) and the held videos with no media in it, and writes nothing.",
+ " The log ends with a summary: line; a re-run resumes.",
+ "",
'fetch-windows fetches clip windows, one paced job per platform queue',
' (YouTube and Rumble side by side): {"siteId"} fetches every window the',
' site\'s published reports cite and the disk does not hold; {"items":',