import { test } from "node:test"; import assert from "node:assert/strict"; import { chmod, mkdir, mkdtemp, readFile, readdir, rm, stat, symlink, utimes, writeFile, } from "node:fs/promises"; import { rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import path from "node:path"; import { assertMirrorDirection, classifyDrift, copyMirrorVerify, CopyVerificationError, describeDifferences, mirrorTree, MirrorDirectionError, makeProgressSink, type RelocationProgress, } from "./relocateDir"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/relocateDir.test.ts // // The sink, on its own: no rsync, no disk. What it has to get right is the two // things a bar claims — that it is about the TREE, and that it only moves one // way. function frame(bytes: number, percent: number): string { return ` ${bytes.toLocaleString("en-US")} ${percent}% 1.03GB/s 0:00:03`; } function collect(opts: { totalBytes: number; alreadyBytes?: number }) { const seen: RelocationProgress[] = []; const lines: string[] = []; const sink = makeProgressSink({ ...opts, log: (m) => lines.push(m), onProgress: (p) => seen.push(p), }); return { sink, seen, lines }; } test("a fresh copy's fraction is transferred over the measured tree", () => { const { sink, seen } = collect({ totalBytes: 100 }); sink(frame(25, 25)); sink(frame(100, 100)); assert.deepEqual( seen.map((p) => p.fraction), [0.25, 1], ); assert.equal(seen[0].bytes, 25); }); // RSYNC COUNTS WHAT *IT* SENT, NOT WHAT IS THERE. Resume a copy that died at // 60 % and the bar climbs 0 → 40 % and stops, having moved every remaining // byte — which an operator reads as a copy that wedged four fifths of the way // through something. The caller measures the destination before the copy; this // is what that measurement is for. test("a resumed copy is offset by what is already on the far side", () => { const { sink, seen } = collect({ totalBytes: 100, alreadyBytes: 60 }); sink(frame(0, 0)); sink(frame(20, 50)); sink(frame(40, 100)); assert.deepEqual( seen.map((p) => p.fraction), [0.6, 0.8, 1], ); assert.deepEqual( seen.map((p) => p.bytes), [60, 80, 100], ); // The detail line says the same thing, so the log and the bar cannot // disagree about the same instant. assert.match(seen[1].detail, /80 B of 100 B · 80 %/); }); // A resumed copy re-sends the partial file it died inside, so // `already + transferred` can exceed the tree by one file's size. That is a // bar reading 112 %, which is a bar nobody believes again. test("the fraction is clamped when a partial file is re-sent", () => { const { sink, seen } = collect({ totalBytes: 100, alreadyBytes: 90 }); sink(frame(30, 100)); assert.equal(seen[0].fraction, 1); assert.equal(seen[0].bytes, 100); }); // ONE LINE PER DECILE, and a resumed copy does not re-announce the deciles it // did not do — which is why the latch lives in the closure and not in a // parameter. test("the log gets one line per decile, starting from the resume point", () => { const { sink, lines } = collect({ totalBytes: 100, alreadyBytes: 60 }); sink(frame(0, 0)); sink(frame(1, 2)); sink(frame(20, 50)); sink(frame(40, 100)); assert.equal(lines.length, 3); assert.match(lines[0], /^Copying… 60 B of 100 B · 60 %/); assert.match(lines[2], /100 B of 100 B · 100 %/); }); test("a line that is not a frame is ignored entirely", () => { const { sink, seen, lines } = collect({ totalBytes: 100 }); sink("sending incremental file list"); sink(".d..t...... 20240101_test1234567/"); assert.deepEqual(seen, []); assert.deepEqual(lines, []); }); // --------------------------------------------------------------------------- // The mirror pass (release 16 slice RM) // --------------------------------------------------------------------------- // // The REAL rsync, over temp dirs, as relocateChannelMedia.test.ts does. The // copy-phase tests reach "after the first pass" through the log: rsyncTree // logs each command line (`$ rsync …`) synchronously before it spawns the // child, so a hook that changes the source on the mirror pass's line changes // it after the copy pass has finished and before the mirror pass starts — // deterministically, with no fake rsync. const MTIME = new Date("2021-03-04T05:06:07.000Z"); async function withTrees( fn: (t: { src: string; dest: string; dir: string }) => Promise, ): Promise { const dir = await mkdtemp(path.join(tmpdir(), "ttb-mirror-")); const src = path.join(dir, "src"); const dest = path.join(dir, "dest"); await mkdir(src, { recursive: true }); await mkdir(dest, { recursive: true }); try { await fn({ src, dest, dir }); } finally { await rm(dir, { recursive: true, force: true }); } } async function put(root: string, rel: string, content: string): Promise { const file = path.join(root, rel); await mkdir(path.dirname(file), { recursive: true }); await writeFile(file, content); await utimes(file, MTIME, MTIME); } async function tree(root: string): Promise { const out: string[] = []; const walk = async (rel: string) => { for (const e of await readdir(path.join(root, rel), { withFileTypes: true })) { const r = rel ? `${rel}/${e.name}` : e.name; if (e.isDirectory()) await walk(r); else out.push(r); } }; await walk(""); return out.sort(); } function cancelled(): Error { return new Error("cancelled"); } test("a file added to the source after the copy pass is carried by the mirror pass", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v1/audio.mp3", "one".repeat(100)); let mirrored = false; const res = await copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, log: (m) => { // The transcription that finished after rsync had passed v1/. if (m.startsWith("$ ") && m.includes("--delete") && !m.includes("--dry-run") && !mirrored) { mirrored = true; writeFileSync(path.join(src, "v1", "transcript.json"), "{}"); } }, }); assert.equal(mirrored, true, "the hook ran before the mirror pass"); assert.equal(res.retried, false, "the mirror pass, not a retry, carried it"); assert.equal(res.files, 2); assert.equal(await readFile(path.join(dest, "v1", "transcript.json"), "utf8"), "{}"); }); }); test("a file removed from the source after the copy pass is removed from the destination", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v1/audio.mp3", "one".repeat(100)); await put(src, "v1/.audio.mp3.parakeet/meta.json", "{}"); await put(src, "v1/.audio.mp3.parakeet/win-0000.json", "[]"); let removed = false; await copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, log: (m) => { // The transcriber cleaning up its scratch dir once it finished. if (m.startsWith("$ ") && m.includes("--delete") && !m.includes("--dry-run") && !removed) { removed = true; rmSync(path.join(src, "v1", ".audio.mp3.parakeet"), { recursive: true }); } }, }); assert.equal(removed, true); assert.deepEqual(await tree(dest), ["v1/audio.mp3"]); // The source is what it was left as — never the target of --delete. assert.deepEqual(await tree(src), ["v1/audio.mp3"]); }); }); test("a stale extra dir on the destination (the realcandaceo case) is removed, the source untouched", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v50t5yt/audio.mp3", "a".repeat(64)); await put(src, "v50t5yt/transcript.json", "{}"); // What the interrupted move left: the full copy, plus a transcriber's // scratch dir copied mid-transcription and since deleted from the source. await put(dest, "v50t5yt/audio.mp3", "a".repeat(64)); await put(dest, "v50t5yt/transcript.json", "{}"); for (const f of ["meta", "win-0000", "win-0001", "win-0002", "win-0003"]) { await put(dest, `v50t5yt/.audio.mp3.parakeet/${f}.json`, "{}"); } const res = await copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, log: () => {}, }); assert.equal(res.files, 2); assert.deepEqual(await tree(dest), ["v50t5yt/audio.mp3", "v50t5yt/transcript.json"]); assert.deepEqual(await tree(src), ["v50t5yt/audio.mp3", "v50t5yt/transcript.json"]); }); }); test("one change between the mirror and the check gets one more pass", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v1/audio.mp3", "one"); let wrote = false; const lines: string[] = []; const res = await copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, log: (m) => { lines.push(m); if (m.startsWith("$ ") && m.includes("--dry-run") && !wrote) { wrote = true; writeFileSync(path.join(src, "v1", "diarization.json"), "{}"); } }, }); assert.equal(res.retried, true); assert.ok(lines.some((l) => l.includes("running one more mirror pass"))); assert.ok(lines.some((l) => l.includes("1 missing on the destination (v1/diarization.json)"))); assert.deepEqual(await tree(dest), ["v1/audio.mp3", "v1/diarization.json"]); }); }); test("a source that changes again after the second pass refuses, naming the paths by kind", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v1/audio.mp3", "one"); let n = 0; await assert.rejects( () => copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, log: (m) => { // Something that keeps writing: a new file before every check. if (m.startsWith("$ ") && m.includes("--dry-run")) { n++; writeFileSync(path.join(src, "v1", `chunk-${n}.json`), "{}"); } }, }), (err: unknown) => { assert.ok(err instanceof CopyVerificationError); assert.deepEqual(err.differences.missing, ["v1/chunk-2.json"]); assert.match(err.message, /1 missing on the destination \(v1\/chunk-2\.json\)/); assert.match(err.message, /0 extra on the destination/); assert.match(err.message, /Reconcile and resume/); assert.match(err.message, /The source has NOT been touched/); return true; }, ); // Every source file is still there. assert.deepEqual(await tree(src), ["v1/audio.mp3", "v1/chunk-1.json", "v1/chunk-2.json"]); }); }); test("a reconcile lists the differences by kind first, then mirrors without a copy pass", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v1/audio.mp3", "one".repeat(10)); await put(src, "v1/transcript.json", '{"new":true}'); await put(dest, "v1/audio.mp3", "one".repeat(10)); await put(dest, "v1/transcript.json", "{}"); await put(dest, "v1/stale.json", "{}"); const lines: string[] = []; await copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, reconcile: true, log: (m) => lines.push(m), }); const said = lines.find((l) => l.startsWith("Reconciling:")) ?? ""; assert.match(said, /1 extra on the destination \(v1\/stale\.json\)/); assert.match(said, /0 missing on the destination/); assert.match(said, /1 changed \(v1\/transcript\.json\)/); // No copy pass: the only non-dry-run rsync is the mirror. assert.equal( lines.filter((l) => l.startsWith("$ ") && l.includes("--partial")).length, 0, ); assert.deepEqual(await tree(dest), ["v1/audio.mp3", "v1/transcript.json"]); assert.equal(await readFile(path.join(dest, "v1", "transcript.json"), "utf8"), '{"new":true}'); }); }); // THE ONE DIRECTION --delete MAY RUN IN, asserted before any rsync is spawned. // `live` comes from the caller (the channel's `data/`, the store), never from // `src` or `dest`, so the guard is not satisfied by construction. The fake // rsync records that it ran; it must never be reached. async function fakeRsync(dir: string): Promise<{ bin: string; ran: string }> { const ran = path.join(dir, "rsync-ran"); const bin = path.join(dir, "fake-rsync.sh"); await writeFile(bin, `#!/bin/sh\ntouch "${ran}"\nexit 0\n`); await chmod(bin, 0o755); return { bin, ran }; } test("source and destination swapped across two roots: refused, nothing spawned", async () => { // Two separate roots, as a real move has: the corpus and the platter. const corpus = await mkdtemp(path.join(tmpdir(), "ttb-live-")); const platter = await mkdtemp(path.join(tmpdir(), "ttb-platter-")); try { const live = path.join(corpus, "channels", "alpha", "data"); const copy = path.join(platter, "alpha", "data"); await put(live, "v1/audio.mp3", "one"); await put(copy, "v1/other.mp3", "two"); const { bin, ran } = await fakeRsync(platter); // OUT, swapped: the platter's copy named as the source, `data/` as the // destination. await assert.rejects( () => copyMirrorVerify({ rsyncBin: bin, src: copy, dest: live, live, cancelled, log: () => {}, }), (err: unknown) => err instanceof MirrorDirectionError && /is the live media/.test(err.message), ); // BACK, swapped: `data/` is a link to the platter copy by then, and the // destination named is the copy it resolves to. const linked = path.join(corpus, "channels", "beta", "data"); await mkdir(path.dirname(linked), { recursive: true }); await symlink(copy, linked); const incoming = path.join(corpus, "channels", "beta", "data.incoming"); await mkdir(incoming, { recursive: true }); await assert.rejects( () => copyMirrorVerify({ rsyncBin: bin, src: incoming, dest: copy, live: linked, cancelled, log: () => {}, }), MirrorDirectionError, ); await assert.rejects(() => stat(ran), "the fake rsync was never spawned"); assert.deepEqual(await tree(live), ["v1/audio.mp3"]); assert.deepEqual(await tree(copy), ["v1/other.mp3"]); // The right way round, both directions, is allowed. await assertMirrorDirection({ src: live, dest: copy, live }); await assertMirrorDirection({ src: copy, dest: incoming, live: linked }); } finally { await rm(corpus, { recursive: true, force: true }); await rm(platter, { recursive: true, force: true }); } }); test("a destination inside the live media, or containing it, is refused", async () => { await withTrees(async ({ src, dest, dir }) => { await put(src, "v1/audio.mp3", "one"); const { bin, ran } = await fakeRsync(dir); // Inside: a copy built under the channel's own data/. await assert.rejects( () => mirrorTree({ rsyncBin: bin, src: dest, dest: path.join(src, "v1"), live: src, log: () => {}, }), /is inside the live media/, ); // Containing: a destination one level above the media. await assert.rejects( () => assertMirrorDirection({ src: dest, dest: dir, live: src }), /contains the live media/, ); await assert.rejects(() => stat(ran), "the fake rsync was never spawned"); assert.deepEqual(await tree(src), ["v1/audio.mp3"]); }); }); test("--delete is refused when source and destination contain one another", async () => { await withTrees(async ({ src, dest, dir }) => { await put(src, "v1/audio.mp3", "one"); await put(dest, "v1/other.mp3", "two"); const { bin, ran } = await fakeRsync(dir); const elsewhere = path.join(dir, "elsewhere"); await assert.rejects( () => mirrorTree({ rsyncBin: bin, src, dest: path.join(src, "v1"), live: elsewhere, log: () => {}, }), /one is inside the other/, ); await assert.rejects( () => assertMirrorDirection({ src: path.join(dest, "x"), dest, live: elsewhere, }), /one is inside the other/, ); await assert.rejects(() => stat(ran), "the fake rsync was never spawned"); assert.deepEqual(await tree(src), ["v1/audio.mp3"]); assert.deepEqual(await tree(dest), ["v1/other.mp3"]); await assertMirrorDirection({ src, dest, live: src }); }); }); test("every removal by the mirror pass is in the log", async () => { await withTrees(async ({ src, dest }) => { await put(src, "v1/audio.mp3", "one"); await put(dest, "v1/audio.mp3", "one"); await put(dest, "v1/stale.json", "{}"); const lines: string[] = []; await copyMirrorVerify({ rsyncBin: "rsync", src, dest, live: src, cancelled, log: (m) => lines.push(m), }); assert.ok( lines.some((l) => /deleting v1\/stale\.json/.test(l)), `the mirror's deletion is logged: ${JSON.stringify(lines)}`, ); }); }); test("rsync's itemized lines sort into extra, missing and changed", () => { const d = classifyDrift([ "*deleting v50t5yt/.audio.mp3.parakeet/meta.json", "*deleting v50t5yt/.audio.mp3.parakeet/", ">f+++++++++ v51/transcript.json", "cd+++++++++ v52/", ">f.st...... v1/transcript.json", ".d..t...... v1/", "rsync: something it wanted to say", ]); assert.deepEqual(d.extra, [ "v50t5yt/.audio.mp3.parakeet/meta.json", "v50t5yt/.audio.mp3.parakeet/", ]); assert.deepEqual(d.missing, ["v51/transcript.json", "v52/"]); assert.deepEqual(d.changed, [ "v1/transcript.json", "v1/", "rsync: something it wanted to say", ]); assert.equal( describeDifferences({ extra: [], missing: ["a", "b", "c", "d", "e", "f", "g"], changed: [] }), "0 extra on the destination, 7 missing on the destination (a, b, c, d, e, … 2 more), 0 changed", ); });