commit 97ed489a1cec5f21920e3505ae550aa04be330b1
parent cabe2faa9518d772f06c1aa7f12a7ea5ec5d9082
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 8 Oct 2026 22:51:30 -0400
Merge branch 'main' into editor/no-lead-paragraphs
# Conflicts:
# editor/CHANGELOG.md
Diffstat:
13 files changed, 1353 insertions(+), 20 deletions(-)
diff --git a/common/controller/transcribeFile.test.ts b/common/controller/transcribeFile.test.ts
@@ -0,0 +1,434 @@
+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 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,
+ transcribeWorkerFilter,
+ windowOf,
+ windowWavArgs,
+} = 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<string, unknown>) =>
+ (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<string, unknown>) =>
+ (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/);
+});
+
+// --- the disk checks --------------------------------------------------------
+
+const ctx = { paths, workers: WORKERS };
+async function check(body: Record<string, unknown>, 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<string, unknown>) {
+ 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);
+});
diff --git a/common/controller/transcribeFile.ts b/common/controller/transcribeFile.ts
@@ -0,0 +1,525 @@
+// ONE-OFF TRANSCRIPTION OF AN ARBITRARY FILE, AS AN EDITOR JOB — what
+// `pnpm ops transcribe` (`POST /api/ops/transcribe`) enqueues.
+//
+// A quote check or "what is audible in this clip" used to mean a hand-run
+// whisper-cli / parakeet-cli with a model path typed from memory. This runs the
+// SAME engine the corpus does, through the same path: the worker pool hands out
+// a lease (so a one-off never oversubscribes the GPU slot auto-transcribe is
+// using, and jumps ahead of queued background work as any manual transcribe
+// does), and `transcribeWithWorker` builds the command line from the worker's
+// config through the transcription-app registry. Nothing here knows an engine's
+// argv.
+//
+// WHAT IT TOUCHES. It reads `path` (anywhere, the corpus included) and writes
+// only to a scratch dir under the OS temp dir, removed afterwards, and to `out`
+// when given — which is refused inside the corpus (the transcripts dir, the
+// saved-video store, the sites dir, every storage location root; compared both
+// as written and resolved through symlinks, since a channel's `media` may be a
+// link to another drive). It never writes settings.json, a sidecar, or anything
+// under a channel. Its job log lands where every job's does.
+//
+// THE AUDIO IS ALWAYS A 16 kHz MONO WAV CUT BY ffmpeg, whole file or window.
+// Two reasons: the engines disagree about containers (whisper-cli wants audio,
+// parakeet's wrapper reads anything ffmpeg does), and parakeet's wrapper keeps
+// its resumable work dir BESIDE its input — given the source file directly, it
+// would create `.<name>.parakeet/` next to it, possibly inside the corpus.
+//
+// THE RESULT is the cue format the index uses ({start, end, text}, seconds),
+// shifted back into the SOURCE file's clock when a window was cut, plus which
+// worker, engine and model produced it. It is written to `out` when given and
+// always logged as ONE line starting with TRANSCRIBE_RESULT_MARKER, which is
+// how `pnpm ops transcribe --wait` prints it without reading any file on the
+// editor's disk.
+//
+// LOCAL WORKERS ONLY. A remote worker delegates by channel and video id (or by
+// an upload into ITS pool); a one-off file has neither identity, and the point
+// of the command is "this machine's engine". The default is the worker
+// auto-transcribe would get — the pool's highest-priority free one — among the
+// local workers.
+
+import os from "node:os";
+import path from "node:path";
+import { access, constants, mkdtemp, readFile, realpath, rm, stat } from "node:fs/promises";
+import { execa } from "execa";
+import type { Paths } from "../lib/paths";
+import type { Worker } from "../lib/workers";
+import type { Cue } from "../lib/vtt";
+import type { StorageLocation } from "../lib/storageLocations";
+import {
+ getTranscriptionApp,
+ type TranscriptOutputFormat,
+} from "../lib/transcriptionApps";
+import { detectTranscriptFormat, parseTranscriptJson } from "../lib/whisper";
+import { writeJsonAtomic } from "../lib/jsonFile-server";
+import { getSettings } from "../lib/settings";
+import { getWorkerPool, type WorkerFilter } from "../jobs/workerPool";
+import { makeTaskTracker } from "../jobs/taskHooks";
+import {
+ runManagedFunction,
+ type JobRunContext,
+ type StreamActionResult,
+} from "../jobs/streamCommand";
+import { transcribeWithWorker } from "./transcribeOne";
+
+export const TRANSCRIBE_FILE_JOB_KIND = "transcribe-file";
+
+// The log line carrying the result. One line, compact JSON after the marker.
+// scripts/archilyzer-ops.mjs matches the same string.
+export const TRANSCRIBE_RESULT_MARKER = "@@transcribe-result ";
+
+// The body `/api/ops/transcribe` accepts; anything else is a 400.
+export const TRANSCRIBE_FILE_BODY_KEYS = [
+ "path",
+ "start",
+ "end",
+ "workerId",
+ "out",
+] as const;
+
+const AUDIO_NAME = "audio.wav";
+// A WAV header with no samples after it: the window held no audio.
+const EMPTY_WAV_BYTES = 44;
+
+export type TranscribeFileRequest = {
+ path: string;
+ // Seconds into the source. Absent start = 0; absent end = to the end.
+ start?: number;
+ end?: number;
+ workerId?: string;
+ out?: string;
+};
+
+export type TranscribeWorkerInfo = {
+ id: string;
+ name: string;
+ appId: string;
+ model: string | null;
+ device: string | null;
+};
+
+export type TranscribeFileResult = {
+ version: 1;
+ path: string;
+ // Null when the whole file was transcribed. `end: null` = to the end.
+ window: { start: number; end: number | null } | null;
+ worker: TranscribeWorkerInfo;
+ transcriptFormat: TranscriptOutputFormat;
+ transcribedAt: string;
+ durationMs: number;
+ cues: Cue[];
+ text: string;
+};
+
+type Check<T> = { ok: true; value: T } | { ok: false; error: string };
+
+// --- body ------------------------------------------------------------------
+
+function seconds(raw: unknown, key: string): Check<number | undefined> {
+ if (raw === undefined || raw === null) return { ok: true, value: undefined };
+ if (typeof raw !== "number" || !Number.isFinite(raw) || raw < 0) {
+ return { ok: false, error: `"${key}" must be a number of seconds, zero or more` };
+ }
+ return { ok: true, value: raw };
+}
+
+// The body's shape, judged without touching the disk. Every sentence here is
+// the refusal an ops caller reads.
+export function parseTranscribeFileBody(
+ body: Record<string, unknown>,
+): Check<TranscribeFileRequest> {
+ const p = body.path;
+ if (typeof p !== "string" || !p.trim()) {
+ return { ok: false, error: '"path" is required: an absolute path to an audio or video file' };
+ }
+ if (!path.isAbsolute(p)) {
+ return { ok: false, error: `"path" must be absolute (got "${p}")` };
+ }
+ const start = seconds(body.start, "start");
+ if (!start.ok) return start;
+ const end = seconds(body.end, "end");
+ if (!end.ok) return end;
+ if (end.value !== undefined && end.value <= (start.value ?? 0)) {
+ return {
+ ok: false,
+ error: `"end" (${end.value}) must be after "start" (${start.value ?? 0})`,
+ };
+ }
+ const workerId = body.workerId;
+ if (workerId !== undefined && (typeof workerId !== "string" || !workerId.trim())) {
+ return { ok: false, error: '"workerId" must be a non-empty string' };
+ }
+ const out = body.out;
+ if (out !== undefined) {
+ if (typeof out !== "string" || !out.trim()) {
+ return { ok: false, error: '"out" must be a non-empty string' };
+ }
+ if (!path.isAbsolute(out)) {
+ return { ok: false, error: `"out" must be an absolute path (got "${out}")` };
+ }
+ }
+ return {
+ ok: true,
+ value: {
+ path: path.resolve(p),
+ ...(start.value !== undefined ? { start: start.value } : {}),
+ ...(end.value !== undefined ? { end: end.value } : {}),
+ ...(typeof workerId === "string" ? { workerId: workerId.trim() } : {}),
+ ...(typeof out === "string" ? { out: path.resolve(out) } : {}),
+ },
+ };
+}
+
+// --- the corpus fence --------------------------------------------------------
+
+// `p` with every symlink in its EXISTING prefix resolved; the part that does
+// not exist yet is appended as written. An `out` is usually a file that does
+// not exist, in a directory that does.
+export async function realpathDeep(p: string): Promise<string> {
+ const abs = path.resolve(p);
+ const rest: string[] = [];
+ let cur = abs;
+ for (;;) {
+ try {
+ const real = await realpath(cur);
+ return rest.length ? path.join(real, ...rest.reverse()) : real;
+ } catch {
+ const parent = path.dirname(cur);
+ if (parent === cur) return abs;
+ rest.push(path.basename(cur));
+ cur = parent;
+ }
+ }
+}
+
+function within(child: string, parent: string): boolean {
+ const rel = path.relative(parent, child);
+ return rel === "" || (!rel.startsWith("..") && !path.isAbsolute(rel));
+}
+
+// Every directory a write of ours must stay out of.
+export function corpusRoots(
+ paths: Pick<Paths, "transcriptsDir" | "channelsDir" | "savedVideosDir" | "sitesDir">,
+ locations: readonly Pick<StorageLocation, "root">[] = [],
+): string[] {
+ const roots = [
+ paths.transcriptsDir,
+ paths.channelsDir,
+ paths.savedVideosDir,
+ paths.sitesDir,
+ ...locations.map((l) => l.root),
+ ].filter((r): r is string => typeof r === "string" && r.trim() !== "");
+ return [...new Set(roots.map((r) => path.resolve(r)))];
+}
+
+// The root `target` lies under, or null. Both sides are compared as written
+// AND resolved: `channels/x/media` may be a link to another drive, and a
+// location root may itself be reached through a link.
+export async function corpusRootContaining(
+ target: string,
+ roots: readonly string[],
+): Promise<string | null> {
+ const forms = [path.resolve(target), await realpathDeep(target)];
+ for (const root of roots) {
+ const rootForms = [path.resolve(root), await realpathDeep(root)];
+ for (const t of forms) {
+ for (const r of rootForms) {
+ if (within(t, r)) return root;
+ }
+ }
+ }
+ return null;
+}
+
+// --- workers -----------------------------------------------------------------
+
+// Which workers may take the transcription: local ones, or the one named.
+export function transcribeWorkerFilter(workerId?: string): WorkerFilter {
+ return (w) => w.kind === "local" && (workerId === undefined || w.id === workerId);
+}
+
+// What a result records about the worker that produced it: the engine (app id)
+// and the model AFTER the app's own default, resolved by the same registry that
+// built the command line.
+export function describeTranscribeWorker(worker: Worker): TranscribeWorkerInfo {
+ const app = getTranscriptionApp(worker.appId);
+ const config = worker.config ?? {};
+ return {
+ id: worker.id,
+ name: worker.name,
+ appId: app.id,
+ model: app.resolveModel(config) ?? null,
+ device: config.device?.trim() || null,
+ };
+}
+
+// --- the disk checks ---------------------------------------------------------
+
+export type TranscribeFileContext = {
+ paths: Paths;
+ workers: readonly Worker[];
+ locations?: readonly Pick<StorageLocation, "root">[];
+ // The pool's runtime view (id → state). A named worker switched off on the
+ // Workers page is refused rather than waited for: the wait would be forever.
+ workerStates?: ReadonlyMap<string, { state: string; degraded: boolean }>;
+};
+
+export async function checkTranscribeFileRequest(
+ req: TranscribeFileRequest,
+ ctx: TranscribeFileContext,
+): Promise<Check<TranscribeFileRequest>> {
+ let st;
+ try {
+ st = await stat(req.path);
+ } catch {
+ return { ok: false, error: `"path" ${req.path} does not exist` };
+ }
+ if (!st.isFile()) {
+ return { ok: false, error: `"path" ${req.path} is not a file` };
+ }
+ try {
+ await access(req.path, constants.R_OK);
+ } catch {
+ return { ok: false, error: `"path" ${req.path} is not readable` };
+ }
+
+ if (req.out !== undefined) {
+ const root = await corpusRootContaining(
+ req.out,
+ corpusRoots(ctx.paths, ctx.locations ?? []),
+ );
+ if (root) {
+ return {
+ ok: false,
+ error: `"out" ${req.out} is inside the corpus (${root}) — this command never writes there; pick a path outside it`,
+ };
+ }
+ if ((await realpathDeep(req.out)) === (await realpathDeep(req.path))) {
+ return { ok: false, error: '"out" is the input file itself' };
+ }
+ const outStat = await stat(req.out).catch(() => null);
+ if (outStat?.isDirectory()) {
+ return { ok: false, error: `"out" ${req.out} is a directory; name a file` };
+ }
+ const dirStat = await stat(path.dirname(req.out)).catch(() => null);
+ if (!dirStat?.isDirectory()) {
+ return {
+ ok: false,
+ error: `"out": the directory ${path.dirname(req.out)} does not exist`,
+ };
+ }
+ }
+
+ const local = ctx.workers.filter((w) => w.kind === "local");
+ if (req.workerId !== undefined) {
+ const named = ctx.workers.find((w) => w.id === req.workerId);
+ if (!named) {
+ return {
+ ok: false,
+ error: `no worker "${req.workerId}" — known: ${
+ ctx.workers.map((w) => w.id).join(", ") || "none"
+ }`,
+ };
+ }
+ if (named.kind !== "local") {
+ return {
+ ok: false,
+ error: `worker "${req.workerId}" is a ${named.kind} worker; a file is transcribed on a local one (${
+ local.map((w) => w.id).join(", ") || "none configured"
+ })`,
+ };
+ }
+ const runtime = ctx.workerStates?.get(req.workerId);
+ if (runtime && (runtime.state !== "enabled" || runtime.degraded)) {
+ return {
+ ok: false,
+ error: `worker "${req.workerId}" is ${
+ runtime.degraded ? "degraded" : runtime.state
+ } on the Workers page — enable it there, or leave out "workerId"`,
+ };
+ }
+ } else if (local.length === 0) {
+ return { ok: false, error: "no local transcription worker is configured" };
+ }
+ return { ok: true, value: req };
+}
+
+// --- the work ----------------------------------------------------------------
+
+// ffmpeg's arguments for the 16 kHz mono WAV: input seeking for the start, a
+// duration for the end, no video.
+export function windowWavArgs(
+ src: string,
+ dst: string,
+ win: { start?: number; end?: number } = {},
+): string[] {
+ const start = win.start ?? 0;
+ const args = ["-nostdin", "-hide_banner", "-v", "error", "-y"];
+ if (start > 0) args.push("-ss", String(start));
+ args.push("-i", src);
+ if (win.end !== undefined) args.push("-t", String(win.end - start));
+ args.push("-vn", "-ac", "1", "-ar", "16000", "-c:a", "pcm_s16le", "-f", "wav", dst);
+ return args;
+}
+
+const ms = (n: number) => Math.round(n * 1000) / 1000;
+
+// Cues from a window start at zero; shift them back onto the source's clock.
+export function offsetCues(cues: readonly Cue[], offset: number): Cue[] {
+ return cues.map((c) => ({
+ start: ms(c.start + offset),
+ end: ms(c.end + offset),
+ text: c.text,
+ }));
+}
+
+export function windowOf(
+ req: Pick<TranscribeFileRequest, "start" | "end">,
+): TranscribeFileResult["window"] {
+ if (req.start === undefined && req.end === undefined) return null;
+ return { start: req.start ?? 0, end: req.end ?? null };
+}
+
+function describeRequest(req: TranscribeFileRequest): string {
+ const win = windowOf(req);
+ return win
+ ? `${req.path} [${win.start}s – ${win.end === null ? "end" : `${win.end}s`}]`
+ : req.path;
+}
+
+export type RunTranscribeFileOpts = {
+ paths: Paths;
+ request: TranscribeFileRequest;
+ onLog: (line: string) => void;
+ signal?: AbortSignal;
+ ctx?: Pick<JobRunContext, "jobId" | "addTask" | "updateTask" | "removeTask" | "recordTaskDone">;
+};
+
+export async function runTranscribeFile(
+ opts: RunTranscribeFileOpts,
+): Promise<TranscribeFileResult> {
+ const { request: req, onLog, paths } = opts;
+ const started = Date.now();
+ const scratch = await mkdtemp(path.join(os.tmpdir(), "archilyzer-transcribe-"));
+ try {
+ const wav = path.join(scratch, AUDIO_NAME);
+ onLog(`Cutting ${describeRequest(req)} to 16 kHz mono WAV…`);
+ try {
+ await execa(paths.ffmpegBin, windowWavArgs(req.path, wav, req), {
+ cancelSignal: opts.signal,
+ });
+ } catch (err) {
+ if (opts.signal?.aborted) throw err;
+ const e = err as { stderr?: string; shortMessage?: string; message: string };
+ throw new Error(
+ `ffmpeg could not read ${req.path}: ${(e.stderr || e.shortMessage || e.message).trim()}`,
+ );
+ }
+ const wavStat = await stat(wav).catch(() => null);
+ if (!wavStat || wavStat.size <= EMPTY_WAV_BYTES) {
+ throw new Error(
+ `no audio in ${describeRequest(req)}${req.start ? " — does the window start past the end of the file?" : ""}`,
+ );
+ }
+
+ // Set by onWorker; a holder, so the closure's write is seen after the await.
+ const used: { worker?: Worker } = {};
+ onLog(
+ req.workerId
+ ? `Waiting for worker ${req.workerId}…`
+ : "Waiting for a free local worker…",
+ );
+ const label = `${path.basename(req.path)}${windowOf(req) ? " (window)" : ""}`;
+ const outcome = await transcribeWithWorker({
+ paths,
+ videoDir: scratch,
+ videoId: path.basename(req.path),
+ audioFilename: AUDIO_NAME,
+ strictAudio: true,
+ tracker: opts.ctx ? makeTaskTracker(opts.ctx, onLog) : undefined,
+ taskId: opts.ctx ? `file:${opts.ctx.jobId}` : undefined,
+ taskLabel: label,
+ onLog,
+ signal: opts.signal,
+ only: transcribeWorkerFilter(req.workerId),
+ onWorker: (w) => {
+ used.worker = w;
+ },
+ skipInlineDiarization: true,
+ });
+ const worker = used.worker;
+ if (outcome !== "transcribed" || !worker) {
+ throw new Error(
+ outcome === "paused"
+ ? "the transcription was stopped before it finished"
+ : `the transcription did not run (${outcome})`,
+ );
+ }
+ const raw = await readFile(path.join(scratch, "transcript.json"), "utf8");
+ const transcriptFormat = detectTranscriptFormat(raw) ?? "whisper-json";
+ const cues = offsetCues(parseTranscriptJson(raw, transcriptFormat), req.start ?? 0);
+ return {
+ version: 1,
+ path: req.path,
+ window: windowOf(req),
+ worker: describeTranscribeWorker(worker),
+ transcriptFormat,
+ transcribedAt: new Date().toISOString(),
+ durationMs: Date.now() - started,
+ cues,
+ text: cues.map((c) => c.text.trim()).filter(Boolean).join(" "),
+ };
+ } finally {
+ await rm(scratch, { recursive: true, force: true }).catch(() => {});
+ }
+}
+
+// --- the job -----------------------------------------------------------------
+
+// Validate, then enqueue. Every refusal comes back as `{ ok: false, error }`
+// BEFORE a job exists, so the ops route answers it as a 400.
+export async function enqueueTranscribeFile(
+ body: Record<string, unknown>,
+ opts: { paths: Paths },
+): Promise<StreamActionResult> {
+ const parsed = parseTranscribeFileBody(body);
+ if (!parsed.ok) return parsed;
+ const settings = getSettings();
+ const pool = getWorkerPool();
+ const workerStates = new Map(
+ pool.summary().map((w) => [w.id, { state: w.state, degraded: w.degraded }]),
+ );
+ const checked = await checkTranscribeFileRequest(parsed.value, {
+ paths: opts.paths,
+ workers: settings.workers,
+ locations: settings.storage?.locations ?? [],
+ workerStates,
+ });
+ if (!checked.ok) return checked;
+ const req = checked.value;
+ return runManagedFunction({
+ kind: TRANSCRIBE_FILE_JOB_KIND,
+ // Parallel: the worker pool is what serialises transcriptions.
+ queueKey: "",
+ paths: opts.paths,
+ fn: async (onLog, signal, _setProgress, ctx) => {
+ onLog(`Transcribe ${describeRequest(req)}`);
+ const result = await runTranscribeFile({
+ paths: opts.paths,
+ request: req,
+ onLog,
+ signal,
+ ctx,
+ });
+ onLog(
+ `Transcribed with ${result.worker.appId} [${result.worker.id}]` +
+ `${result.worker.model ? ` model ${result.worker.model}` : ""}: ` +
+ `${result.cues.length} cue(s) in ${(result.durationMs / 1000).toFixed(1)}s`,
+ );
+ if (req.out) {
+ await writeJsonAtomic(req.out, result, { indent: 2 });
+ onLog(`Written to ${req.out}`);
+ }
+ onLog(`${TRANSCRIBE_RESULT_MARKER}${JSON.stringify(result)}`);
+ },
+ });
+}
diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts
@@ -5,7 +5,11 @@ import type { Paths } from "../lib/paths";
import type { Worker } from "../lib/workers";
import { getTranscriptionApp } from "../lib/transcriptionApps";
import { detectTranscriptFormat } from "../lib/whisper";
-import { WORKER_DEGRADE_THRESHOLD, getWorkerPool } from "../jobs/workerPool";
+import {
+ WORKER_DEGRADE_THRESHOLD,
+ getWorkerPool,
+ type WorkerFilter,
+} from "../jobs/workerPool";
import type { TaskTracker } from "../jobs/taskHooks";
import { normalizeTranscript } from "./normalizeTranscript";
import { diarizeOneVideo } from "./diarizeOne";
@@ -91,6 +95,10 @@ export type TranscribeOneOptions = {
// resume. transcribeOneVideo returns "paused" so the batch leaves the video
// untranscribed (next run resumes). Ignored by engines without partial support.
partialSignal?: AbortSignal;
+ // True skips the inline diarization pass even when settings turn it on. A
+ // one-off file transcription (controller/transcribeFile.ts) runs in a scratch
+ // dir that is deleted afterwards, so a diarization there is work thrown away.
+ skipInlineDiarization?: boolean;
};
export type TranscribeOneOutcome = "transcribed" | "already-exists" | "paused";
@@ -302,7 +310,11 @@ export async function transcribeOneVideo(
// The "paused" early returns above correctly bypass this — there is no
// transcript yet, so there is nothing to diarize alongside.
const diarization = getSettings().diarization;
- if (diarization.enabled && diarization.inlineAfterTranscribe) {
+ if (
+ !opts.skipInlineDiarization &&
+ diarization.enabled &&
+ diarization.inlineAfterTranscribe
+ ) {
try {
const outcome = await diarizeOneVideo({
paths: opts.paths,
@@ -360,6 +372,14 @@ export type TranscribeWithWorkerOptions = {
// Auto-runner units pass true so they park BEHIND any manual (foreground)
// acquire in the worker pool — a manual transcribe preempts queued auto work.
background?: boolean;
+ // Narrows WHICH workers may take this video (the pool's `only`): a one-off
+ // file transcription keeps to local workers, or to the one it was told to use.
+ only?: WorkerFilter;
+ // Called with the worker each attempt runs on, before it starts — how a
+ // caller learns which engine and model produced the transcript.
+ onWorker?: (worker: Worker) => void;
+ // See TranscribeOneOptions.skipInlineDiarization.
+ skipInlineDiarization?: boolean;
};
// Acquire a worker from the global pool and transcribe one video through it,
@@ -383,11 +403,18 @@ export async function transcribeWithWorker(
if (opts.signal?.aborted || opts.drainSignal?.aborted) throw abortError();
let lease;
try {
- lease = await pool.acquire(acquireSignal, { background: opts.background });
- } catch {
+ lease = await pool.acquire(acquireSignal, {
+ background: opts.background,
+ ...(opts.only ? { only: opts.only } : {}),
+ });
+ } catch (err) {
+ // An `only` no configured worker passes is refused at once — that is
+ // not a cancel, and must not read as one.
+ if (opts.only && !acquireSignal?.aborted) throw err;
throw abortError(); // cancelled or drained while parked
}
const worker = lease.worker;
+ opts.onWorker?.(worker);
const task = opts.tracker?.start({
id: opts.taskId ?? opts.videoId,
label: opts.taskLabel ?? opts.videoId,
@@ -416,6 +443,7 @@ export async function transcribeWithWorker(
onProgress: task ? task.update : undefined,
signal: opts.signal,
partialSignal: partialController?.signal,
+ skipInlineDiarization: opts.skipInlineDiarization,
});
pool.markSuccess(worker.id);
return outcome;
diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts
@@ -92,6 +92,8 @@ const ADDED_KINDS: Record<string, { label: string; drainable: boolean }> = {
"fetch-windows": { label: "Fetch windows", drainable: true },
// One video's metadata re-read: one spawn, cancelled rather than drained.
"refresh-metadata": { label: "Refresh metadata", drainable: false },
+ // One file through a local worker (pnpm ops transcribe): one engine run.
+ "transcribe-file": { label: "Transcribe file", drainable: false },
};
test("added kinds carry their pinned label and drainability", () => {
diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts
@@ -763,6 +763,20 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
queueKeyStrategy: "custom",
needsMedia: true,
},
+ // ONE ARBITRARY FILE (or a window of it) through a local worker, for a quote
+ // check — `pnpm ops transcribe` (controller/transcribeFile.ts). No channel:
+ // it reads the file it was given and writes only to a scratch dir and the
+ // caller's `out`, never into the corpus, so neither guard applies. Parallel
+ // (queueKey ""): the worker pool serialises it against every other
+ // transcription. One engine run — cancelled, not drained; not replayable,
+ // since the file it names may be gone by the time anyone retries.
+ "transcribe-file": {
+ kind: "transcribe-file",
+ label: "Transcribe file",
+ drainable: false,
+ replayable: false,
+ queueKeyStrategy: "parallel",
+ },
"download-one-pipeline": {
kind: "download-one-pipeline",
drainable: false,
diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts
@@ -460,3 +460,42 @@ test("legacy { background: true } maps to the background tier", async () => {
await b;
assert.deepEqual(order, ["F", "B"]);
});
+
+// opts.only (WorkerFilter): the one-off file transcription names a worker, or
+// keeps to local ones. The filter narrows the free-list and the parked grant;
+// a filter no configured worker passes is refused at once rather than parked.
+test("an acquire narrowed by `only` takes the named worker, even when a better one is free", async () => {
+ const pool = new WorkerPool();
+ pool.reconfigure([worker("gpu", 0), worker("cpu", 5)], { applyEnabled: true });
+ const lease = await pool.acquire(undefined, { only: (w) => w.id === "cpu" });
+ assert.equal(lease.worker.id, "cpu");
+ lease.release();
+});
+
+test("an acquire narrowed by `only` parks for its worker and skips others that free", async () => {
+ const pool = new WorkerPool();
+ pool.reconfigure([worker("gpu", 0), worker("cpu", 5)], { applyEnabled: true });
+ const gpu = await pool.acquire();
+ const cpu = await pool.acquire();
+ assert.equal(gpu.worker.id, "gpu");
+ let got: string | null = null;
+ const p = pool
+ .acquire(undefined, { only: (w) => w.id === "cpu" })
+ .then((l) => ((got = l.worker.id), l));
+ gpu.release(); // the wrong worker frees: the narrowed waiter stays parked
+ await new Promise((r) => setImmediate(r));
+ assert.equal(got, null);
+ cpu.release();
+ const l = await p;
+ assert.equal(l.worker.id, "cpu");
+ l.release();
+});
+
+test("an `only` filter no configured worker passes is refused, not parked", async () => {
+ const pool = new WorkerPool();
+ pool.reconfigure([worker("w1")], { applyEnabled: true });
+ await assert.rejects(
+ pool.acquire(undefined, { only: (w) => w.id === "nope" }),
+ /no configured worker can take this work/,
+ );
+});
diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts
@@ -101,8 +101,16 @@ type Waiter = {
// FIRST waiter it matches, not blindly to the head — pump() scans past a
// waiter whose requirement this slot cannot satisfy.
requires?: readonly string[];
+ // A further restriction on WHICH worker may serve this waiter (acquire's
+ // opts.only). Matched alongside `requires`, never instead of it.
+ only?: WorkerFilter;
};
+// Narrows an acquire to particular workers — by id, by kind — on top of the
+// capability rule. The one-off file transcription uses it to stay on LOCAL
+// workers, or on the one worker the caller named.
+export type WorkerFilter = (worker: Worker) => boolean;
+
export class WorkerPool {
// Insertion order is the tiebreak for equal priority, so use a Map (ordered).
private entries = new Map<string, PoolEntry>();
@@ -316,21 +324,29 @@ export class WorkerPool {
return this.pausedSnapshot !== null;
}
- private eligible(entry: PoolEntry, requires?: readonly string[]): boolean {
+ private eligible(
+ entry: PoolEntry,
+ requires?: readonly string[],
+ only?: WorkerFilter,
+ ): boolean {
return (
entry.state === "enabled" &&
!entry.degraded &&
!entry.busy &&
- workerMatches(entry.config, requires)
+ workerMatches(entry.config, requires) &&
+ (!only || only(entry.config))
);
}
// Highest-priority free entry matching the requirement:
// (priority asc, insertion order asc).
- private pickFree(requires?: readonly string[]): PoolEntry | null {
+ private pickFree(
+ requires?: readonly string[],
+ only?: WorkerFilter,
+ ): PoolEntry | null {
let best: PoolEntry | null = null;
for (const entry of this.entries.values()) {
- if (!this.eligible(entry, requires)) continue;
+ if (!this.eligible(entry, requires, only)) continue;
if (!best || entry.config.priority < best.config.priority) best = entry;
}
return best;
@@ -365,7 +381,7 @@ export class WorkerPool {
let granted = false;
for (let i = 0; i < this.waiters.length; i++) {
const waiter = this.waiters[i];
- const entry = this.pickFree(waiter.requires);
+ const entry = this.pickFree(waiter.requires, waiter.only);
if (!entry) continue;
this.waiters.splice(i, 1);
waiter.onAbort?.();
@@ -405,12 +421,15 @@ export class WorkerPool {
// granted only on a worker that matches it. A requirement NO configured worker
// matches rejects immediately — parking on it would never resolve, and the
// caller should treat the work as unrunnable rather than wedge.
+ // opts.only narrows the lease further, to the workers the filter accepts (see
+ // WorkerFilter); a filter no configured worker passes rejects the same way.
acquire(
signal?: AbortSignal,
opts?: {
tier?: SchedulerTier;
background?: boolean;
requires?: readonly string[];
+ only?: WorkerFilter;
},
): Promise<Lease> {
this.ensureInit();
@@ -418,13 +437,14 @@ export class WorkerPool {
return Promise.reject(abortError());
}
const requires = opts?.requires;
- if (requires && requires.length > 0) {
+ const only = opts?.only;
+ if ((requires && requires.length > 0) || only) {
let satisfiable = false;
for (const e of this.entries.values()) {
// Config-level, ignoring busy/disabled/degraded on purpose: a matching
// worker that is merely busy or switched off is a reason to park, not
// to give up.
- if (workerMatches(e.config, requires)) {
+ if (workerMatches(e.config, requires) && (!only || only(e.config))) {
satisfiable = true;
break;
}
@@ -432,19 +452,21 @@ export class WorkerPool {
if (!satisfiable) {
return Promise.reject(
new Error(
- `no configured worker matches requirement [${requires.join(", ")}]`,
+ requires && requires.length > 0
+ ? `no configured worker matches requirement [${requires.join(", ")}]`
+ : "no configured worker can take this work",
),
);
}
}
- const entry = this.pickFree(requires);
+ const entry = this.pickFree(requires, only);
if (entry && this.waiters.length === 0) {
return Promise.resolve(this.grant(entry));
}
const tier: SchedulerTier =
opts?.tier ?? (opts?.background ? "background" : "foreground");
return new Promise<Lease>((resolve, reject) => {
- const waiter: Waiter = { resolve, reject, tier, requires };
+ const waiter: Waiter = { resolve, reject, tier, requires, only };
if (signal) {
const onAbort = () => {
const idx = this.waiters.indexOf(waiter);
@@ -493,7 +515,13 @@ export class WorkerPool {
}
if (!best) return null;
const claimed = best;
- if (this.waiters.some((w) => workerMatches(claimed.config, w.requires))) {
+ if (
+ this.waiters.some(
+ (w) =>
+ workerMatches(claimed.config, w.requires) &&
+ (!w.only || w.only(claimed.config)),
+ )
+ ) {
return null;
}
return this.grant(claimed);
diff --git a/common/lib/transcriptionApps.ts b/common/lib/transcriptionApps.ts
@@ -88,6 +88,10 @@ export type TranscriptionApp = {
};
// Default binary when the per-app `bin` override is empty.
defaultBin: () => string;
+ // The model this config runs with, after the app's own default — what a
+ // record of "which model produced this" should say. Undefined when the app
+ // picks one itself (chough with no model set).
+ resolveModel: (config: AppInstanceConfig) => string | undefined;
build: (input: TranscribeBuildInput) => TranscribeBuild;
makeProgressParser: () => { feed: (line: string) => ProgressUpdate | null };
// True when the engine can be stopped mid-run and still produce a usable
@@ -175,12 +179,13 @@ const whisperCpp: TranscriptionApp = {
label: "whisper.cpp (whisper-cli)",
fields: { model: true, customArgs: true },
defaultBin: () => getPaths().whisperBin,
+ resolveModel: (config) => config.model ?? getPaths().whisperModel,
build({ audioFile, outputBase, config }) {
const template =
config.customArgs && config.customArgs.length > 0
? config.customArgs
: DEFAULT_TRANSCRIBE_ARGS;
- const model = config.model ?? getPaths().whisperModel;
+ const model = whisperCpp.resolveModel(config) as string;
const argv = substitutePlaceholders(template, {
audioFile,
outputBase,
@@ -197,6 +202,7 @@ const chough: TranscriptionApp = {
label: "chough",
fields: { model: true, remoteUrl: true, chunkSize: true },
defaultBin: () => process.env.CHOUGH_BIN ?? "chough",
+ resolveModel: (config) => config.model?.trim() || undefined,
build({ audioFile, outputBase, config }) {
const argv = ["-f", "json", "-o", outputBase];
if (typeof config.chunkSize === "number" && config.chunkSize > 0) {
@@ -208,7 +214,7 @@ const chough: TranscriptionApp = {
argv.push("-r");
env.CHOUGH_URL = remoteUrl;
}
- const model = config.model?.trim();
+ const model = chough.resolveModel(config);
if (model) env.CHOUGH_MODEL = model;
argv.push(audioFile);
// chough writes EXACTLY the `-o` path — no ".json" is appended.
@@ -237,9 +243,10 @@ const parakeet: TranscriptionApp = {
// a partial transcript on SIGTERM — so a long run is interruptible.
supportsPartialStop: true,
defaultBin: () => getPaths().parakeetBin,
+ resolveModel: (config) => config.model?.trim() || getPaths().parakeetModel,
build({ audioFile, outputBase, config }) {
const paths = getPaths();
- const model = config.model?.trim() || paths.parakeetModel;
+ const model = parakeet.resolveModel(config) as string;
const argv = ["--model", model, "--output", outputBase];
if (typeof config.chunkSize === "number" && config.chunkSize > 0) {
argv.push("--segment", String(Math.floor(config.chunkSize)));
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **A single file can be transcribed through the editor.** `pnpm ops transcribe` (`POST /api/ops/transcribe`) takes `{"path"}` — an absolute path to a local audio or video file — with `"start"`/`"end"` (seconds) for a window, `"workerId"` (a configured local worker; default: the one auto-transcribe would get) and `"out"` (an absolute path for the result). It runs as a "Transcribe file" job on `/jobs`, through the worker pool and the same engine, model and command line the corpus is transcribed with. A window is cut by ffmpeg to a temporary 16 kHz mono WAV, removed afterwards. The result is `{path, window, worker: {id, name, appId, model, device}, transcriptFormat, cues: [{start, end, text}], text, …}`, cue times on the source file's clock; it ends the job's log, is written to `out` when given, and `--wait` prints it on stdout. Refused before any job: a relative or unreadable path, an `end` not after `start`, an unknown or remote worker, a named worker switched off, an `out` inside the corpus (resolved through symlinks, storage locations included). Nothing is written to the corpus or to settings.
- **Page titles no longer carry a paragraph of explanation.** /channels used to open with three lines of prose about where its numbers come from, and at a narrow window with the sidebar open its four buttons took the row and squeezed that prose to one word a line. It is now one status line — how many channels, how old the oldest report is, how many have none — and when the buttons do not fit beside the title they drop below it. Workers, Sites, Review and Operations lose their lead paragraph; Tags, Storage and Saved videos keep one line each, the instruction (rule edits apply at the next index build; media moves from a channel's Storage panel; the keep-latest window is in a channel's settings); the Monitor widget builder loses its own, which repeated the board's. Needs a restart of the editor.
- **Clip windows can be fetched as a list, paced, in one job per platform.** `pnpm ops fetch-windows` (`POST /api/ops/fetch-windows`) takes `{"siteId"}` — every window a site's published reports cite and the disk does not hold — or `{"items": [{"slug", "id", "from", "to", "clipId"?, "reason"?}], "requestedBy", "manifest"?}`, with `"maxHeight"` and `"dryRun"`. A window already on disk is answered at once and joins no job; the rest are grouped by platform queue and fetched as one `fetch-windows` job per platform, so YouTube and Rumble run side by side, each window through the same managed fetch as a single one (cookie policy, auth and HLS retries, provenance). Between two fetches on a platform the job waits that platform's batch gap, at least 30 s and up to half again at random; before each it checks the platform's cooldown and hold and stops when either is set. A 429 backs the platform off and stops the job; one 403 is that window's failure, two in a row back the platform off and stop it. A platform cooling down or held is refused for its group before anything starts. The job is drainable, shows its progress as clip windows, and Retry or running the same body again fetches only what is still missing. A dry run lists the windows per platform, the ones on disk, and the spans no window can fill (a video the source says is deleted, private or members-only, a channel off the site, an unmounted drive). A site's **Reports** tab has **Fetch missing evidence** and **Preview missing evidence** beside **Prepare evidence media**, and umtool's `fetch-via-editor.mjs --all` sends a manifest's whole timeline as one request. A window fetch whose cookie retry runs into a 429 now records the cooldown too. Needs a restart of the editor.
- **A social channel can be renamed and deleted from its page.** A posts channel's page (X, Bluesky, a forum thread) had no Danger zone, so it could not be renamed or deleted in the editor at all. It now has the one a video channel has, collapsed under the posts panel and opened by the same `?stage=danger` link: **Rename channel** and **Delete channel**, typing the slug to confirm, refused while the channel is busy, in the same sentences. Needs a restart of the editor.
diff --git a/editor/app/api/ops/transcribe/route.test.ts b/editor/app/api/ops/transcribe/route.test.ts
@@ -0,0 +1,86 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/transcribe/route.test.ts"
+//
+// The adapter's surface: the token, the body's keys, and the controller's
+// refusals answered as 400s before any job exists. Every request here is
+// refused — no job is queued and no engine runs (the job itself is covered by
+// common/controller/transcribeFile.test.ts).
+
+const ROOT = await mkdtemp(path.join(os.tmpdir(), "transcribe-route-"));
+const CORPUS = path.join(ROOT, "transcripts");
+await mkdir(CORPUS, { recursive: true });
+const SETTINGS_FILE = path.join(ROOT, "settings.json");
+const SETTINGS_TEXT = JSON.stringify({
+ workers: [
+ {
+ id: "gpu",
+ name: "GPU",
+ kind: "local",
+ enabled: true,
+ priority: 0,
+ appId: "parakeet",
+ config: { model: "/models/x.gguf" },
+ },
+ ],
+});
+await writeFile(SETTINGS_FILE, SETTINGS_TEXT);
+const MEDIA = path.join(ROOT, "clip.wav");
+await writeFile(MEDIA, "x");
+// Set before the route (and getPaths, which caches) is first imported.
+process.env.WORKER_TOKEN = "test-token";
+process.env.TRANSCRIPTS_DIR = CORPUS;
+process.env.SETTINGS_FILE = SETTINGS_FILE;
+const { POST } = await import("./route");
+test.after(() => rm(ROOT, { recursive: true, force: true }));
+
+async function post(
+ body: Record<string, unknown>,
+ token = "test-token",
+): Promise<{ status: number; json: { ok?: boolean; error?: string; jobId?: string } }> {
+ const res = await POST(
+ new Request("http://localhost/api/ops/transcribe", {
+ method: "POST",
+ headers: {
+ authorization: `Bearer ${token}`,
+ "content-type": "application/json",
+ },
+ body: JSON.stringify(body),
+ }),
+ );
+ return { status: res.status, json: await res.json() };
+}
+
+test("a wrong token is a 401", async () => {
+ const { status } = await post({ path: MEDIA }, "wrong");
+ assert.equal(status, 401);
+});
+
+test("an unknown key is a 400 naming the accepted ones", async () => {
+ const { status, json } = await post({ path: MEDIA, model: "/x.gguf" });
+ assert.equal(status, 400);
+ assert.match(json.error!, /unknown key\(s\): model — this route accepts path, start, end, workerId, out/);
+});
+
+test("the controller's refusals come back as 400s", async () => {
+ const cases: [Record<string, unknown>, RegExp][] = [
+ [{ path: "clip.wav" }, /"path" must be absolute/],
+ [{ path: path.join(ROOT, "missing.wav") }, /does not exist/],
+ [{ path: MEDIA, start: 30, end: 10 }, /"end" \(10\) must be after "start" \(30\)/],
+ [{ path: MEDIA, workerId: "nope" }, /no worker "nope" — known: gpu/],
+ [{ path: MEDIA, out: path.join(CORPUS, "r.json") }, /inside the corpus/],
+ ];
+ for (const [body, re] of cases) {
+ const { status, json } = await post(body);
+ assert.equal(status, 400, JSON.stringify(body));
+ assert.equal(json.ok, false);
+ assert.match(json.error!, re);
+ assert.equal(json.jobId, undefined);
+ }
+ assert.equal(await readFile(SETTINGS_FILE, "utf8"), SETTINGS_TEXT);
+});
diff --git a/editor/app/api/ops/transcribe/route.ts b/editor/app/api/ops/transcribe/route.ts
@@ -0,0 +1,32 @@
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import {
+ TRANSCRIBE_FILE_BODY_KEYS,
+ enqueueTranscribeFile,
+} from "yt-dlp-transcript-common/controller/transcribeFile";
+import { jobResponse, ops } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { path, start?, end?, workerId?, out? } -> { ok: true, jobId }
+//
+// Transcribe ONE local file — or the window [start, end) seconds of it — on a
+// local transcription worker, as a job on /jobs: the engine, model and
+// command line the corpus uses, never a hand-run whisper or parakeet. No
+// `workerId` = the worker auto-transcribe would get (the pool's highest-priority
+// free local one).
+//
+// The job's log ends with the result as ONE line, `@@transcribe-result
+// {json}`: `{version, path, window, worker: {id, name, appId, model, device},
+// transcriptFormat, transcribedAt, durationMs, cues: [{start, end, text}],
+// text}`, cue times on the SOURCE file's clock. `pnpm ops transcribe --wait`
+// prints that JSON on stdout. `out` writes it to that file as well.
+//
+// AN ADAPTER: every refusal is the controller's (controller/transcribeFile.ts)
+// — a relative or unreadable `path`, `end` not after `start`, an unknown or
+// non-local `workerId`, a named worker switched off, an `out` inside the corpus
+// (resolved through symlinks) or in a directory that does not exist.
+export async function POST(request: Request) {
+ return ops(request, TRANSCRIBE_FILE_BODY_KEYS, async (body) =>
+ jobResponse(await enqueueTranscribeFile(body, { paths: getPaths() })),
+ );
+}
diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs
@@ -60,6 +60,8 @@
// pnpm ops fetch-windows --file windows.json --wait
// pnpm ops get tags eva-collab
// pnpm ops cut-release --json '{"workspace":"all","version":"next","commit":true}'
+// pnpm ops transcribe --json '{"path":"/abs/clip.mp4","start":120,"end":150}' --wait
+// pnpm ops transcribe --json '{"path":"/abs/a.wav","workerId":"parakeet-cpu","out":"/tmp/a.json"}' --wait
//
// --file reads the BODY from a JSON file, which is how a big one gets sent: a
// four-thousand-id tag-videos body is written by a script, not typed by a model
@@ -182,8 +184,17 @@ const ACTIONS = [
// needs no editor at all — this route exists only on an editor built from
// release 10 or later.
"cut-release",
+ // ONE local file, or a window of it, through a local transcription worker
+ // ({path, start?, end?, workerId?, out?}) — the corpus's own engine and
+ // model, as a job. With --wait the result JSON is what stdout carries.
+ "transcribe",
];
+// The log line a transcribe job ends with: this marker, then the result as
+// compact JSON. The same string is TRANSCRIBE_RESULT_MARKER in
+// common/controller/transcribeFile.ts.
+export const TRANSCRIBE_RESULT_MARKER = "@@transcribe-result ";
+
// The provenance a tag write from this CLI carries. Everything else ignores it.
function agentSource() {
return `agent:${process.env.ARCHILYZER_AGENT || "cli"}`;
@@ -307,6 +318,8 @@ export function parseArgs(argv) {
// a name are told apart; the route's own default (`agent:ops`) covers a
// caller that is neither.
...(action === "tag-videos" ? { defaultSource: agentSource() } : {}),
+ // A job whose log carries a RESULT, which --wait prints on stdout.
+ ...(action === "transcribe" ? { resultMarker: TRANSCRIBE_RESULT_MARKER } : {}),
wait,
quiet,
waitTimeout,
@@ -339,6 +352,39 @@ function printPreviewUrls(payload) {
for (const url of previewUrlsIn(payload)) console.error(`preview: ${url}`);
}
+// Pull a job's RESULT line out of its log as the log streams by. `feed` takes
+// each chunk and returns the text to echo — every complete line except the
+// result line; `finish` returns what is left and the parsed result (null when
+// the log never carried one). Line-buffered, because a poll may end mid-line.
+export function makeResultCapture(marker) {
+ let pending = "";
+ let result = null;
+ const take = (line) => {
+ if (line.startsWith(marker)) {
+ try {
+ result = JSON.parse(line.slice(marker.length));
+ return "";
+ } catch {
+ return `${line}\n`;
+ }
+ }
+ return `${line}\n`;
+ };
+ return {
+ feed(chunk) {
+ pending += chunk;
+ const lines = pending.split("\n");
+ pending = lines.pop() ?? "";
+ return lines.map(take).join("");
+ },
+ finish() {
+ const echo = pending ? take(pending).replace(/\n$/, "") : "";
+ pending = "";
+ return { echo, result };
+ },
+ };
+}
+
export function usage() {
return [
"Usage: pnpm ops <action> [--json '<body>' | --file <path>] [--wait]",
@@ -477,6 +523,15 @@ export function usage() {
" (an older one answers 404); with no editor running, `archilyzer release",
" cut` does the same locally.",
"",
+ 'transcribe runs ONE local file through a local transcription worker, as a',
+ ' job: {"path": "/abs/file"} (audio or video), "start"/"end" (seconds) for a',
+ ' window, "workerId" (a settings worker id; default: the one auto-transcribe',
+ ' would get), "out" (an absolute path for the result JSON, never inside the',
+ " corpus). The result is {path, window, worker: {id, appId, model, device},",
+ " cues: [{start, end, text}], text, ...}, cue times on the file's own clock.",
+ " With --wait it is printed on stdout (the response and the log go to",
+ " stderr), so `pnpm ops transcribe ... --wait | jq -r .text` works.",
+ "",
"Env: ARCHILYZER_EDITOR_URL (default http://localhost:3001), WORKER_TOKEN,",
" ARCHILYZER_AGENT (provenance of a tag write; default \"cli\")",
].join("\n");
@@ -540,7 +595,13 @@ export async function followJob(jobId, quiet, opts = {}) {
const res = await doFetch(`${base}/api/jobs/${id}/log?from=${from}`);
if (!res.ok) return null;
const payload = await res.json();
- if (payload.content && !quiet) process.stderr.write(payload.content);
+ if (payload.content) {
+ // onContent sees every chunk, quiet or not, and says what to echo.
+ const echo = opts.onContent
+ ? opts.onContent(payload.content)
+ : payload.content;
+ if (echo && !quiet) process.stderr.write(echo);
+ }
// Only advance once the chunk is in hand: a poll that failed halfway
// must re-ask for the same offset.
from = payload.nextOffset ?? from;
@@ -683,7 +744,10 @@ async function main() {
console.error(`HTTP ${res.status}: ${text.slice(0, 500)}`);
return 1;
}
- console.log(JSON.stringify(payload, null, 2));
+ // A result-carrying action under --wait keeps stdout for the RESULT: the
+ // response goes to stderr with the log.
+ const resultMode = Boolean(parsed.wait && parsed.resultMarker);
+ (resultMode ? console.error : console.log)(JSON.stringify(payload, null, 2));
if (!res.ok || payload.ok === false) return 1;
if (!parsed.wait) {
printPreviewUrls(payload);
@@ -708,9 +772,21 @@ async function main() {
}
let worst = 0;
for (const jobId of jobIds) {
+ const capture = resultMode ? makeResultCapture(parsed.resultMarker) : null;
const status = await followJob(jobId, parsed.quiet, {
timeoutSeconds: parsed.waitTimeout ?? 0,
+ ...(capture ? { onContent: capture.feed } : {}),
});
+ if (capture) {
+ const { echo, result } = capture.finish();
+ if (echo && !parsed.quiet) process.stderr.write(echo);
+ if (result !== null) {
+ console.log(JSON.stringify(result, null, 2));
+ } else if (status === "done") {
+ console.error(`[${jobId}] finished but its log carries no result`);
+ worst = 1;
+ }
+ }
console.error(`[${jobId}] ${status}`);
if (status !== "done") worst = 1;
}
diff --git a/scripts/archilyzer-ops.test.mjs b/scripts/archilyzer-ops.test.mjs
@@ -6,7 +6,9 @@
import assert from "node:assert/strict";
import test from "node:test";
import {
+ TRANSCRIBE_RESULT_MARKER,
followJob,
+ makeResultCapture,
parseArgs,
previewUrlsIn,
usage,
@@ -477,3 +479,62 @@ test("refresh-metadata is a POST to its route, named in the usage", () => {
assert.match(usage(), /Actions:.*metadata-scan, refresh-metadata/);
assert.match(usage(), /refresh-metadata re-reads ONE video/);
});
+
+// --- transcribe: a job whose log carries a RESULT ---------------------------
+
+test("transcribe posts its body and carries the result marker", () => {
+ const p = parseArgs([
+ "transcribe",
+ "--json",
+ '{"path":"/abs/clip.mp4","start":120,"end":150}',
+ "--wait",
+ ]);
+ assert.equal(p.method, "POST");
+ assert.equal(p.path, "/api/ops/transcribe");
+ assert.deepEqual(p.body, { path: "/abs/clip.mp4", start: 120, end: 150 });
+ assert.equal(p.wait, true);
+ assert.equal(p.resultMarker, TRANSCRIBE_RESULT_MARKER);
+ // No other action carries one.
+ assert.equal(parseArgs(["sync", "--json", '{"slug":"x"}']).resultMarker, undefined);
+ assert.match(usage(), /transcribe runs ONE local file through a local transcription worker/);
+ assert.match(usage(), /it is printed on stdout/);
+});
+
+test("the result line is taken out of the log, across chunk boundaries", () => {
+ const result = { cues: [{ start: 120.5, end: 121, text: "hi" }], text: "hi" };
+ const line = `${TRANSCRIBE_RESULT_MARKER}${JSON.stringify(result)}`;
+ const cap = makeResultCapture(TRANSCRIBE_RESULT_MARKER);
+ let echoed = "";
+ // The log arrives in polls that split lines anywhere — the result line too.
+ const log = `Transcribe /abs/clip.mp4\nprogress 50%\n${line}\nafter\n`;
+ for (const chunk of [log.slice(0, 7), log.slice(7, 40), log.slice(40, 60), log.slice(60)]) {
+ echoed += cap.feed(chunk);
+ }
+ const { echo, result: got } = cap.finish();
+ echoed += echo;
+ assert.deepEqual(got, result);
+ assert.equal(echoed, "Transcribe /abs/clip.mp4\nprogress 50%\nafter\n");
+});
+
+test("a log with no result line captures null and echoes everything", () => {
+ const cap = makeResultCapture(TRANSCRIBE_RESULT_MARKER);
+ const echoed = cap.feed("one\ntwo");
+ const { echo, result } = cap.finish();
+ assert.equal(result, null);
+ assert.equal(echoed + echo, "one\ntwo");
+});
+
+test("followJob hands each chunk to onContent", async () => {
+ const seen = [];
+ const editor = fakeEditor({
+ log: [
+ { content: "a\n", nextOffset: 2, status: "running" },
+ { content: "b\n", nextOffset: 4, status: "done" },
+ ],
+ });
+ const status = await follow("j1", editor, {
+ onContent: (c) => (seen.push(c), ""),
+ });
+ assert.equal(status, "done");
+ assert.deepEqual(seen, ["a\n", "b\n"]);
+});