// 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 `..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 type { SchedulerTier } from "../jobs/jobKinds"; 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", "words", ] 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; // The cut is 16 kHz mono s16 (windowWavArgs): 32,000 bytes a second. const WAV_BYTES_PER_SEC = 16_000 * 2; // A file transcription at most this long (seconds of audio) waits for a worker // in the pool's "urgent" tier, ahead of every parked transcription — the // lane's background units and a manual batch's next video alike. It is the // quote check an agent is waiting on; one more long transcription in the // queue is not. Longer files keep the manual ("foreground") tier: still ahead // of the lane, behind batches queued before them. The tier only orders // WAITERS: a transcription already running is never interrupted. export const URGENT_MAX_AUDIO_SEC = 15 * 60; /** Seconds of audio in the cut WAV, from its size. */ export function wavSeconds(bytes: number): number { return Math.max(0, bytes - EMPTY_WAV_BYTES) / WAV_BYTES_PER_SEC; } /** The worker-pool tier a file transcription of `audioSec` seconds asks for. */ export function transcribeFileTier(audioSec: number): SchedulerTier { return audioSec <= URGENT_MAX_AUDIO_SEC ? "urgent" : "foreground"; } 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; // True returns the engine's word timestamps as well as the cues. Only an // engine that keeps them (parakeet) answers with any; the rest give none. words?: boolean; }; // One word as the engine timed it, on the source file's clock (seconds). export type TranscribedWord = { w: string; start: number; end: number; conf?: number }; 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; // The engine's own time: from the moment a worker took the job to the // transcript. The wait for a free worker is `waitedMs`, never in here. durationMs: number; // How long the job waited for a free worker (the pool's queue). waitedMs: number; cues: Cue[]; text: string; // Present only when the request asked for words: [] when the engine has none. words?: TranscribedWord[]; }; type Check = { ok: true; value: T } | { ok: false; error: string }; // --- body ------------------------------------------------------------------ function seconds(raw: unknown, key: string): Check { 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, ): Check { 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}")` }; } } const words = body.words; if (words !== undefined && typeof words !== "boolean") { return { ok: false, error: '"words" must be true or false' }; } 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) } : {}), ...(words === true ? { words: true } : {}), }, }; } // --- 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 { 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, locations: readonly Pick[] = [], ): 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 { 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[]; // 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; }; export async function checkTranscribeFileRequest( req: TranscribeFileRequest, ctx: TranscribeFileContext, ): Promise> { 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. // The top-level `words` an engine asked with `words` wrote (parakeet's // wrapper, --words), shifted by `offset` onto the source file's clock. A // document without them -- another engine, or an older wrapper -- gives []. export function wordsFromTranscript(raw: string, offset: number): TranscribedWord[] { let doc: unknown; try { doc = JSON.parse(raw); } catch { return []; } const list = (doc as { words?: unknown } | null)?.words; if (!Array.isArray(list)) return []; const out: TranscribedWord[] = []; for (const w of list) { if (!w || typeof w !== "object") continue; const { w: text, start, end, conf } = w as Record; if (typeof text !== "string" || !text.trim()) continue; if (typeof start !== "number" || typeof end !== "number") continue; if (!Number.isFinite(start) || !Number.isFinite(end)) continue; out.push({ w: text.trim(), start: ms(start + offset), end: ms(end + offset), ...(typeof conf === "number" && Number.isFinite(conf) ? { conf } : {}), }); } return out; } 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, ): 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; }; export async function runTranscribeFile( opts: RunTranscribeFileOpts, ): Promise { const { request: req, onLog, paths } = opts; 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. // `startedAt` is the moment a worker took it: the engine's clock starts // there (the last attempt's, when a transport failure moved it). const used: { worker?: Worker; startedAt?: number } = {}; const audioSec = wavSeconds(wavStat.size); const tier = transcribeFileTier(audioSec); onLog( `${req.workerId ? `Waiting for worker ${req.workerId}` : "Waiting for a free local worker"}` + (tier === "urgent" ? ` (${Math.round(audioSec)}s of audio: ahead of queued transcriptions)…` : "…"), ); const asked = Date.now(); 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), tier, onWorker: (w) => { used.worker = w; used.startedAt = Date.now(); }, skipInlineDiarization: true, ...(req.words ? { words: true } : {}), }); const finished = Date.now(); 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: finished - (used.startedAt ?? asked), waitedMs: (used.startedAt ?? asked) - asked, cues, text: cues.map((c) => c.text.trim()).filter(Boolean).join(" "), ...(req.words ? { words: wordsFromTranscript(raw, req.start ?? 0) } : {}), }; } 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, opts: { paths: Paths }, ): Promise { 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` + (result.waitedMs >= 1000 ? ` (after ${(result.waitedMs / 1000).toFixed(1)}s waiting for the worker)` : ""), ); if (req.out) { await writeJsonAtomic(req.out, result, { indent: 2 }); onLog(`Written to ${req.out}`); } onLog(`${TRANSCRIBE_RESULT_MARKER}${JSON.stringify(result)}`); }, }); }