import { test } from "node:test"; import assert from "node:assert/strict"; import { chmod, mkdir, mkdtemp, readFile, readdir, rm, symlink, writeFile, } from "node:fs/promises"; import { existsSync, statSync } from "node:fs"; import os from "node:os"; import path from "node:path"; import type { Worker } from "../lib/workers"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/transcribeFile.test.ts // // The one-off file transcription (`pnpm ops transcribe`): the body's shape, the // refusals that come before any job, the window's ffmpeg cut and the cue // offset, which engine and model a worker resolves to — and one whole job, run // through the real registry and worker pool against a FAKE ffmpeg and a FAKE // engine (a parakeet worker whose `bin` is a script), in a temp corpus. Nothing // here spawns a real engine or touches a real transcripts dir. const ROOT = await mkdtemp(path.join(os.tmpdir(), "transcribe-file-")); const CORPUS = path.join(ROOT, "transcripts"); const OUTSIDE = path.join(ROOT, "outside"); const TMP = path.join(ROOT, "tmp"); const BIN = path.join(ROOT, "bin"); await mkdir(path.join(CORPUS, "channels", "chan", "data"), { recursive: true }); await mkdir(OUTSIDE, { recursive: true }); await mkdir(TMP, { recursive: true }); await mkdir(BIN, { recursive: true }); // A fake ffmpeg: records its argv, writes a WAV-sized file to its last arg — // an EMPTY one (a bare 44-byte header) for an input named *empty*. const FFMPEG_ARGS = path.join(ROOT, "ffmpeg-args.json"); const ffmpeg = path.join(BIN, "ffmpeg"); await writeFile( ffmpeg, `#!/usr/bin/env node const fs = require("node:fs"); const args = process.argv.slice(2); fs.writeFileSync(${JSON.stringify(FFMPEG_ARGS)}, JSON.stringify(args)); const input = args[args.indexOf("-i") + 1]; const out = args[args.length - 1]; fs.writeFileSync(out, Buffer.alloc(input.includes("empty") ? 44 : 4096)); `, ); await chmod(ffmpeg, 0o755); // A fake engine standing in for the parakeet wrapper: records its argv and // cwd, and writes chough-native JSON to its --output path (cue times from 0, // as a window's always are). const ENGINE_ARGS = path.join(ROOT, "engine-args.json"); const engine = path.join(BIN, "fake-parakeet"); await writeFile( engine, `#!/usr/bin/env node const fs = require("node:fs"); const args = process.argv.slice(2); fs.writeFileSync(${JSON.stringify(ENGINE_ARGS)}, JSON.stringify({ args, cwd: process.cwd() })); const audio = args[args.length - 1]; if (!fs.existsSync(audio)) { console.error("no audio " + audio); process.exit(3); } const out = args[args.indexOf("--output") + 1]; fs.writeFileSync(out, JSON.stringify({ chunk_data: [ { start_time: 0.5, end_time: 1.25, text: " hello" }, { start_time: 2, end_time: 3.5, text: "world " }, ] })); `, ); await chmod(engine, 0o755); const WORKERS: Worker[] = [ { id: "gpu", name: "GPU parakeet", kind: "local", enabled: true, priority: 0, appId: "parakeet", config: { bin: engine, model: "/models/tdt.gguf" }, }, { id: "cpu", name: "CPU parakeet", kind: "local", enabled: true, priority: 5, appId: "parakeet", config: { bin: engine, device: "cpu" }, }, { id: "off", name: "Switched off", kind: "local", enabled: false, priority: 9, appId: "parakeet", config: { bin: engine }, }, { id: "far", name: "A remote", kind: "remote", enabled: true, priority: 1, remote: { baseUrl: "http://127.0.0.1:9", slots: 1 }, }, ]; const SETTINGS_FILE = path.join(ROOT, "settings.json"); const SETTINGS_TEXT = JSON.stringify({ workers: WORKERS }, null, 2); await writeFile(SETTINGS_FILE, SETTINGS_TEXT); // Set before getPaths (which caches) is first reached through the imports. process.env.TRANSCRIPTS_DIR = CORPUS; process.env.SETTINGS_FILE = SETTINGS_FILE; process.env.FFMPEG_BIN = ffmpeg; process.env.PARAKEET_MODEL = "/models/default.gguf"; process.env.TMPDIR = TMP; const { TRANSCRIBE_RESULT_MARKER, checkTranscribeFileRequest, corpusRootContaining, corpusRoots, describeTranscribeWorker, enqueueTranscribeFile, offsetCues, parseTranscribeFileBody, transcribeFileTier, transcribeWorkerFilter, URGENT_MAX_AUDIO_SEC, wavSeconds, windowOf, windowWavArgs, wordsFromTranscript, } = await import("./transcribeFile"); const { getPaths } = await import("../lib/paths"); const paths = getPaths(); test.after(() => rm(ROOT, { recursive: true, force: true })); const MEDIA = path.join(OUTSIDE, "clip.mp4"); await writeFile(MEDIA, "not really media"); // --- the body --------------------------------------------------------------- test("a body needs an absolute path", () => { assert.match( (parseTranscribeFileBody({}) as { error: string }).error, /"path" is required/, ); assert.match( (parseTranscribeFileBody({ path: "clip.mp4" }) as { error: string }).error, /"path" must be absolute/, ); assert.match( (parseTranscribeFileBody({ path: 7 }) as { error: string }).error, /"path" is required/, ); const ok = parseTranscribeFileBody({ path: "/a/../b/clip.mp4" }); assert.deepEqual(ok, { ok: true, value: { path: "/b/clip.mp4" } }); }); test("a window must end after it starts, in non-negative seconds", () => { const err = (b: Record) => (parseTranscribeFileBody({ path: MEDIA, ...b }) as { error?: string }).error; assert.match(err({ start: 10, end: 10 })!, /"end" \(10\) must be after "start" \(10\)/); assert.match(err({ start: 10, end: 4 })!, /must be after/); assert.match(err({ end: 0 })!, /"end" \(0\) must be after "start" \(0\)/); assert.match(err({ start: -1 })!, /"start" must be a number of seconds/); assert.match(err({ end: "30" })!, /"end" must be a number of seconds/); assert.match(err({ start: Number.NaN })!, /"start" must be/); assert.equal(err({ start: 10, end: 10.5 }), undefined); assert.equal(err({ start: 10 }), undefined); assert.equal(err({ end: 10 }), undefined); }); test("workerId and out are checked for shape", () => { const err = (b: Record) => (parseTranscribeFileBody({ path: MEDIA, ...b }) as { error?: string }).error; assert.match(err({ workerId: "" })!, /"workerId" must be a non-empty string/); assert.match(err({ workerId: 3 })!, /"workerId"/); assert.match(err({ out: "result.json" })!, /"out" must be an absolute path/); assert.match(err({ out: 1 })!, /"out" must be a non-empty string/); }); test("words is a boolean, kept only when true", () => { const parsed = (b: Record) => parseTranscribeFileBody({ path: MEDIA, ...b }); assert.match((parsed({ words: "yes" }) as { error: string }).error, /"words" must be true or false/); const on = parsed({ words: true }); assert.ok(on.ok && on.value.words === true); const off = parsed({ words: false }); assert.ok(off.ok && !("words" in off.value)); }); test("wordsFromTranscript shifts the engine's words onto the file's clock", () => { const raw = JSON.stringify({ chunk_data: [{ start_time: 0, end_time: 1, text: "um so" }], words: [ { w: " um", start: 0.12, end: 0.4, conf: 0.9 }, { w: "so", start: 0.5, end: 0.7 }, { w: " ", start: 0.8, end: 0.9 }, { w: "bad", start: "x", end: 1 }, ], }); assert.deepEqual(wordsFromTranscript(raw, 120), [ { w: "um", start: 120.12, end: 120.4, conf: 0.9 }, { w: "so", start: 120.5, end: 120.7 }, ]); // Another engine's document, or an older wrapper: no words, not a failure. assert.deepEqual(wordsFromTranscript(JSON.stringify({ chunk_data: [] }), 0), []); assert.deepEqual(wordsFromTranscript("not json", 0), []); }); // --- the disk checks -------------------------------------------------------- const ctx = { paths, workers: WORKERS }; async function check(body: Record, extra: object = {}) { const parsed = parseTranscribeFileBody({ path: MEDIA, ...body }); assert.ok(parsed.ok, JSON.stringify(parsed)); const res = await checkTranscribeFileRequest(parsed.value, { ...ctx, ...extra }); return res.ok ? null : res.error; } test("a path that is not a readable file is refused", async () => { assert.match((await check({ path: path.join(OUTSIDE, "nope.mp4") }))!, /does not exist/); assert.match((await check({ path: OUTSIDE }))!, /is not a file/); const locked = path.join(OUTSIDE, "locked.wav"); await writeFile(locked, "x"); await chmod(locked, 0o000); // root reads anything; the refusal can only be seen as a user. if (process.getuid?.() !== 0) { assert.match((await check({ path: locked }))!, /is not readable/); } assert.equal(await check({}), null); }); test("an out inside the corpus is refused, through a symlink too", async () => { assert.match( (await check({ out: path.join(CORPUS, "result.json") }))!, /is inside the corpus/, ); assert.match( (await check({ out: path.join(CORPUS, "channels", "chan", "data", "x.json") }))!, /is inside the corpus/, ); // A link OUTSIDE the corpus that points INTO it. const link = path.join(OUTSIDE, "into-corpus"); await symlink(path.join(CORPUS, "channels"), link); assert.match( (await check({ out: path.join(link, "r.json") }))!, /is inside the corpus/, ); // A channel's media linked off to another drive, reached through the corpus. const drive = path.join(ROOT, "drive", "chan", "media"); await mkdir(drive, { recursive: true }); await symlink(drive, path.join(CORPUS, "channels", "chan", "media")); assert.match( (await check({ out: path.join(CORPUS, "channels", "chan", "media", "r.json") }))!, /is inside the corpus/, ); // ...and that drive written to directly is a storage location root's. assert.equal(await check({ out: path.join(drive, "r.json") }), null); assert.match( (await check( { out: path.join(drive, "r.json") }, { locations: [{ root: path.join(ROOT, "drive") }] }, ))!, /is inside the corpus/, ); assert.equal(await check({ out: path.join(OUTSIDE, "r.json") }), null); }); test("an out that is the input, a directory, or in a missing directory is refused", async () => { assert.match((await check({ out: MEDIA }))!, /is the input file itself/); assert.match((await check({ out: OUTSIDE }))!, /is a directory/); assert.match( (await check({ out: path.join(OUTSIDE, "no", "such", "r.json") }))!, /does not exist/, ); }); test("corpusRoots and corpusRootContaining name the root", async () => { const roots = corpusRoots(paths, [{ root: "/mnt/platter" }]); assert.ok(roots.includes(path.resolve(CORPUS))); assert.ok(roots.includes("/mnt/platter")); assert.equal(await corpusRootContaining("/mnt/platter/x/media/a.json", roots), "/mnt/platter"); assert.equal(await corpusRootContaining("/mnt/platterx/a.json", roots), null); }); test("an unknown workerId is refused, naming the known ones", async () => { const err = await check({ workerId: "nope" }); assert.match(err!, /no worker "nope" — known: gpu, cpu, off, far/); }); test("a remote workerId is refused, naming the local ones", async () => { assert.match( (await check({ workerId: "far" }))!, /"far" is a remote worker; a file is transcribed on a local one \(gpu, cpu, off\)/, ); }); test("a named worker switched off on the Workers page is refused, not waited for", async () => { const workerStates = new Map([ ["off", { state: "disabled", degraded: false }], ["cpu", { state: "enabled", degraded: true }], ["gpu", { state: "enabled", degraded: false }], ]); assert.match((await check({ workerId: "off" }, { workerStates }))!, /"off" is disabled/); assert.match((await check({ workerId: "cpu" }, { workerStates }))!, /"cpu" is degraded/); assert.equal(await check({ workerId: "gpu" }, { workerStates }), null); }); test("with no local worker configured, a default request is refused", async () => { const err = await check({}, { workers: WORKERS.filter((w) => w.kind !== "local") }); assert.match(err!, /no local transcription worker is configured/); }); // --- the pure pieces ---------------------------------------------------------- test("the window is cut to 16 kHz mono WAV: input seek, then a duration", () => { assert.deepEqual(windowWavArgs("/in.mp4", "/t/a.wav", { start: 120, end: 150.5 }), [ "-nostdin", "-hide_banner", "-v", "error", "-y", "-ss", "120", "-i", "/in.mp4", "-t", "30.5", "-vn", "-ac", "1", "-ar", "16000", "-c:a", "pcm_s16le", "-f", "wav", "/t/a.wav", ]); const whole = windowWavArgs("/in.mp4", "/t/a.wav"); assert.ok(!whole.includes("-ss") && !whole.includes("-t")); const toEnd = windowWavArgs("/in.mp4", "/t/a.wav", { end: 40 }); assert.ok(!toEnd.includes("-ss")); assert.deepEqual(toEnd.slice(toEnd.indexOf("-t"), toEnd.indexOf("-t") + 2), ["-t", "40"]); }); test("cue times are shifted back onto the source's clock", () => { assert.deepEqual( offsetCues([{ start: 0.5, end: 1.25, text: "a" }, { start: 2.0004, end: 3, text: "b" }], 120), [{ start: 120.5, end: 121.25, text: "a" }, { start: 122, end: 123, text: "b" }], ); assert.deepEqual(offsetCues([{ start: 1, end: 2, text: "x" }], 0), [{ start: 1, end: 2, text: "x" }]); assert.equal(windowOf({}), null); assert.deepEqual(windowOf({ start: 5 }), { start: 5, end: null }); assert.deepEqual(windowOf({ end: 9 }), { start: 0, end: 9 }); }); test("a worker resolves to its engine and model, the app's default model when unset", () => { assert.deepEqual(describeTranscribeWorker(WORKERS[0]), { id: "gpu", name: "GPU parakeet", appId: "parakeet", model: "/models/tdt.gguf", device: null, }); assert.deepEqual(describeTranscribeWorker(WORKERS[1]), { id: "cpu", name: "CPU parakeet", appId: "parakeet", model: "/models/default.gguf", device: "cpu", }); const chough = describeTranscribeWorker({ id: "c", name: "c", kind: "local", enabled: true, priority: 0, appId: "chough", }); assert.equal(chough.appId, "chough"); assert.equal(chough.model, null); }); test("the worker filter keeps to local workers, or to the one named", () => { const any = transcribeWorkerFilter(); assert.deepEqual(WORKERS.filter(any).map((w) => w.id), ["gpu", "cpu", "off"]); assert.deepEqual(WORKERS.filter(transcribeWorkerFilter("cpu")).map((w) => w.id), ["cpu"]); assert.deepEqual(WORKERS.filter(transcribeWorkerFilter("far")).map((w) => w.id), []); }); // --- one whole job ------------------------------------------------------------ async function runJob(body: Record) { const res = await enqueueTranscribeFile(body, { paths }); if (!res.ok) return { error: res.error }; void res.stream.cancel(); const done = await res.done; const log = await readFile(path.join(paths.jobsDir, `${res.jobId}.log`), "utf8"); const line = log.split("\n").find((l) => l.startsWith(TRANSCRIBE_RESULT_MARKER)); return { status: done.status, log, result: line ? JSON.parse(line.slice(TRANSCRIBE_RESULT_MARKER.length)) : null, }; } test("a window is transcribed on the named worker, its cues on the source clock, out written", async () => { const out = path.join(OUTSIDE, "result.json"); const run = await runJob({ path: MEDIA, start: 120, end: 150, workerId: "cpu", out }); assert.equal(run.status, "done", run.log); const r = run.result; assert.deepEqual(r.window, { start: 120, end: 150 }); assert.deepEqual(r.cues, [ { start: 120.5, end: 121.25, text: "hello" }, { start: 122, end: 123.5, text: "world" }, ]); assert.equal(r.text, "hello world"); assert.deepEqual(r.worker, { id: "cpu", name: "CPU parakeet", appId: "parakeet", model: "/models/default.gguf", device: "cpu", }); assert.equal(r.transcriptFormat, "chough-json"); assert.equal(r.path, MEDIA); // The same document in `out`. assert.deepEqual(JSON.parse(await readFile(out, "utf8")), r); // ffmpeg cut the window; the engine got the registry's command line for // THAT worker (its device, the default model) and ran in the scratch dir. const ff = JSON.parse(await readFile(FFMPEG_ARGS, "utf8")) as string[]; assert.deepEqual(ff.slice(ff.indexOf("-ss"), ff.indexOf("-ss") + 4), ["-ss", "120", "-i", MEDIA]); assert.deepEqual(ff.slice(ff.indexOf("-t"), ff.indexOf("-t") + 2), ["-t", "30"]); const eng = JSON.parse(await readFile(ENGINE_ARGS, "utf8")) as { args: string[]; cwd: string }; assert.deepEqual(eng.args.slice(0, 2), ["--model", "/models/default.gguf"]); assert.ok(eng.args.includes("--device") && eng.args.includes("cpu"), eng.args.join(" ")); assert.equal(eng.args[eng.args.length - 1], "audio.wav"); assert.ok(eng.cwd.startsWith(TMP), `engine ran in ${eng.cwd}`); // The scratch dir is gone, and nothing landed in the corpus but the job log. assert.deepEqual( (await readdir(TMP)).filter((n) => n.startsWith("archilyzer-transcribe-")), [], ); assert.deepEqual( (await readdir(CORPUS)).filter((n) => n !== ".jobs" && n !== "channels"), [], ); assert.deepEqual(await readdir(path.join(CORPUS, "channels", "chan", "data")), []); // settings.json is byte-for-byte what it was. assert.equal(await readFile(SETTINGS_FILE, "utf8"), SETTINGS_TEXT); }); test("with no workerId the pool's highest-priority local worker runs it, whole file, no offset", async () => { const run = await runJob({ path: MEDIA }); assert.equal(run.status, "done", run.log); assert.equal(run.result.worker.id, "gpu"); assert.equal(run.result.worker.model, "/models/tdt.gguf"); assert.equal(run.result.window, null); assert.equal(run.result.cues[0].start, 0.5); const ff = JSON.parse(await readFile(FFMPEG_ARGS, "utf8")) as string[]; assert.ok(!ff.includes("-ss") && !ff.includes("-t")); assert.equal(await readFile(SETTINGS_FILE, "utf8"), SETTINGS_TEXT); }); test("a window holding no audio fails the job with a sentence, and cleans up", async () => { const empty = path.join(OUTSIDE, "empty.wav"); await writeFile(empty, "x"); const run = await runJob({ path: empty, start: 9999 }); assert.equal(run.status, "failed"); assert.match(run.log!, /no audio in .*start past the end/); assert.equal(run.result, null); assert.deepEqual( (await readdir(TMP)).filter((n) => n.startsWith("archilyzer-transcribe-")), [], ); }); test("the guards answer before any job exists", async () => { assert.match((await runJob({ path: "rel.mp4" })).error!, /must be absolute/); assert.match((await runJob({ path: MEDIA, start: 5, end: 5 })).error!, /must be after/); assert.match((await runJob({ path: MEDIA, workerId: "ghost" })).error!, /no worker "ghost"/); assert.match( (await runJob({ path: MEDIA, out: path.join(CORPUS, "x.json") })).error!, /inside the corpus/, ); // The pool seeded "off" from settings as disabled. assert.match((await runJob({ path: MEDIA, workerId: "off" })).error!, /"off" is disabled/); assert.equal(await readFile(SETTINGS_FILE, "utf8"), SETTINGS_TEXT); }); // --- the wait, and the tier ---------------------------------------------------- test("a cut's length decides its tier: up to 15 minutes of audio is urgent", () => { assert.equal(wavSeconds(44), 0); assert.equal(wavSeconds(44 + 32_000 * 90), 90); assert.equal(URGENT_MAX_AUDIO_SEC, 900); assert.equal(transcribeFileTier(0.1), "urgent"); assert.equal(transcribeFileTier(900), "urgent"); assert.equal(transcribeFileTier(901), "foreground"); }); const { getWorkerPool } = await import("../jobs/workerPool"); const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); const onGpu = (w: Worker) => w.id === "gpu"; // Enqueue a job and resolve once it is PARKED on the pool (its log says it is // waiting for a worker) — a fixed sleep raced the fake ffmpeg's start on a // loaded machine. `finish` reads the result as runJob does. async function parkedJob(body: Record) { const res = await enqueueTranscribeFile(body, { paths }); if (!res.ok) throw new Error(`refused: ${res.error}`); void res.stream.cancel(); const logFile = path.join(paths.jobsDir, `${res.jobId}.log`); const t0 = Date.now(); while (!/Waiting for worker/.test(await readFile(logFile, "utf8").catch(() => ""))) { if (Date.now() - t0 > 20_000) throw new Error("the job never parked"); await sleep(20); } const finish = async () => { const done = await res.done; const log = await readFile(logFile, "utf8"); const line = log.split("\n").find((l) => l.startsWith(TRANSCRIBE_RESULT_MARKER)); return { status: done.status, log, result: line ? JSON.parse(line.slice(TRANSCRIBE_RESULT_MARKER.length)) : null, }; }; return { finish }; } test("durationMs is the engine's time; the wait for a busy worker is waitedMs", async () => { const pool = getWorkerPool(); // Something else holds the GPU worker's one slot. const held = await pool.acquire(undefined, { only: onGpu }); const job = await parkedJob({ path: MEDIA, workerId: "gpu" }); await sleep(500); const released = Date.now(); held.release(); const run = await job.finish(); assert.equal(run.status, "done", run.log); const r = run.result; assert.ok(r.waitedMs >= 450, `waitedMs ${r.waitedMs}`); assert.ok(r.durationMs >= 0, `durationMs ${r.durationMs}`); // The engine's clock starts when the worker is taken — after the release. assert.ok(r.durationMs <= Date.now() - released, `durationMs ${r.durationMs} counts the wait`); assert.match(run.log, /ahead of queued transcriptions/); }); test("a short file goes ahead of parked transcriptions — the lane's and a manual batch's", async () => { const pool = getWorkerPool(); const held = await pool.acquire(undefined, { only: onGpu }); const t0 = Date.now(); await sleep(20); // so an engine run after this is visibly newer than t0 // Did the file job's engine run before this slot was granted? const engineRan = () => existsSync(ENGINE_ARGS) && statSync(ENGINE_ARGS).mtimeMs >= t0; const order: string[] = []; // Parked first: an auto-lane unit (background) and a manual batch's next // video (foreground), both waiting for the GPU worker. const lane = pool.acquire(undefined, { background: true, only: onGpu }).then((l) => { order.push(`lane${engineRan() ? " after the file" : ""}`); l.release(); }); const batch = pool.acquire(undefined, { only: onGpu }).then((l) => { order.push(`batch${engineRan() ? " after the file" : ""}`); l.release(); }); const job = await parkedJob({ path: MEDIA, workerId: "gpu" }); held.release(); const run = await job.finish(); await Promise.all([lane, batch]); assert.equal(run.status, "done", run.log); // The file job, parked LAST, took the freed slot first; the two parked // before it got it after, in their own order (manual before the lane). assert.deepEqual(order, ["batch after the file", "lane after the file"]); });