// Remote-worker SERVER side: the logic behind /api/worker/transcribe on an // instance that is acting as a worker for another instance. Audio is uploaded to // a scratch dir, transcribed through THIS instance's own worker pool (so the // remote picks among its local workers by priority), and the produced // transcript.json is held in scratch until the requester pulls it. The thin // route handlers wrap these functions with token auth. import path from "node:path"; import fs from "fs-extra"; import { getPaths, type Paths } from "../lib/paths"; import { getRegistry, newJobId, type JobStatus } from "../jobs/registry"; import { getWorkerPool } from "../jobs/workerPool"; import { runManagedFunction } from "../jobs/streamCommand"; import { makeTaskTracker } from "../jobs/taskHooks"; import { transcribeWithWorker } from "./transcribeOne"; import { getOperation, type OperationRunOutcome, } from "../lib/operations"; import { CUES_JSON_FILENAME } from "../lib/videoStatus"; import type { AttributionSettings, DiarizationSettings, } from "../lib/settings"; import type { DigestAppConfig } from "../lib/digest"; import type { DigestContext } from "../lib/digestContext-server"; type WorkerJob = { remoteJobId: string; // Dir the transcription runs in. For an upload this is a scratch dir we own and // delete on cleanup; for a shared-fs job it's the requester's actual video dir // on the shared mount, which we must NOT delete. workDir: string; inPlace: boolean; managedJobId: string; // Byte offset already streamed to the requester, so /events can return only // new log output on each poll. logOffset: number; // Set for a UNIT job (a backfill kind run against a scratch corpus). The op // decides which output files the result endpoint reads back. unit?: { op: string; videoDir: string }; }; const ID_RE = /^[A-Za-z0-9_-]+$/; // Shared by both entry points: spawn the managed transcription job over a work // dir and register it. Returns the managed job result. async function startManaged(opts: { remoteJobId: string; workDir: string; inPlace: boolean; videoId: string; audioFilename: string; strictAudio: boolean; channelSlug?: string; }): Promise { const paths = getPaths(); const res = await runManagedFunction({ kind: "worker-transcribe", queueKey: "", // immediate; this instance's worker pool throttles concurrency paths, channelSlug: opts.channelSlug, videoId: opts.videoId, fn: async (onLog, signal, _setProgress, ctx) => { await transcribeWithWorker({ paths, videoDir: opts.workDir, videoId: opts.videoId, audioFilename: opts.audioFilename, strictAudio: opts.strictAudio, tracker: makeTaskTracker(ctx, onLog), taskId: opts.videoId, taskLabel: opts.videoId, onLog, signal, }); }, }); if (!res.ok) { if (!opts.inPlace) await fs.remove(opts.workDir).catch(() => {}); throw new Error(res.error); } res.stream.cancel().catch(() => {}); jobs().set(opts.remoteJobId, { remoteJobId: opts.remoteJobId, workDir: opts.workDir, inPlace: opts.inPlace, managedJobId: res.jobId, logOffset: 0, }); } declare global { // eslint-disable-next-line no-var var __yttWorkerJobs__: Map | undefined; } function jobs(): Map { if (!globalThis.__yttWorkerJobs__) globalThis.__yttWorkerJobs__ = new Map(); return globalThis.__yttWorkerJobs__; } // Only [A-Za-z0-9._-], and force a basename — the upload name is attacker // controlled across the trust boundary. function sanitizeAudioName(name: string): string { const base = path.basename(name || "audio"); const cleaned = base.replace(/[^A-Za-z0-9._-]/g, "_"); return cleaned && cleaned !== "." && cleaned !== ".." ? cleaned : "audio"; } export type StartWorkerInput = { audio: Buffer; audioName: string; // Optional label shown on the remote's own Workers/jobs UI. label?: string; }; // Write the uploaded audio to a fresh scratch dir and kick off a managed // transcription job that runs through this instance's worker pool. Returns the // remoteJobId the requester polls. Does not block on the transcription. export async function startWorkerTranscription( input: StartWorkerInput, ): Promise<{ remoteJobId: string }> { const paths = getPaths(); // Pick up the latest local worker config before transcribing (the request may // arrive long after this instance last touched its pool). getWorkerPool().reconfigure(); const remoteJobId = newJobId(); const scratchDir = path.join(paths.workerScratchDir, remoteJobId); await fs.ensureDir(scratchDir); const audioName = sanitizeAudioName(input.audioName); await fs.writeFile(path.join(scratchDir, audioName), input.audio); const label = input.label ?? remoteJobId; await startManaged({ remoteJobId, workDir: scratchDir, inPlace: false, videoId: label, audioFilename: audioName, strictAudio: true, }); return { remoteJobId }; } export type StartWorkerSharedInput = { channelSlug: string; videoId: string; audioFilename: string; }; // Shared-filesystem fast path: the requester mounts the same transcripts dir, so // instead of uploading audio we transcribe its video dir IN PLACE (writing // transcript.json back onto the shared mount, where the requester already sees // it). The requester never pulls a result. Validates slug/videoId — they cross // the trust boundary and index into the channels tree. export async function startWorkerTranscriptionShared( input: StartWorkerSharedInput, ): Promise<{ remoteJobId: string }> { const paths = getPaths(); if (!ID_RE.test(input.channelSlug) || !ID_RE.test(input.videoId)) { throw new Error("invalid channelSlug or videoId"); } const audioFilename = sanitizeAudioName(input.audioFilename); getWorkerPool().reconfigure(); const remoteJobId = newJobId(); const videoDir = path.join( paths.channelsDir, input.channelSlug, "data", input.videoId, ); if (!(await fs.pathExists(videoDir))) { throw new Error(`video dir not found on shared mount: ${videoDir}`); } await startManaged({ remoteJobId, workDir: videoDir, inPlace: true, videoId: input.videoId, audioFilename, strictAudio: false, channelSlug: input.channelSlug, }); return { remoteJobId }; } // --------------------------------------------------------------------------- // Unit jobs: run ONE backfill kind (attribution, diarization, …) against a // scratch corpus materialized from the uploaded inputs. The generalisation of // the transcription protocol above: same auth, same poll/pull/DELETE shape, // but the payload is a backfill kind's declared inputs and the result its // declared outputs. The primary applies those outputs through the kinds' // guarded writers — the executor never touches a real corpus. // --------------------------------------------------------------------------- const UNIT_OUTCOME_FILENAME = "unit-outcome.json"; export type StartWorkerUnitInput = { // A backfill kind id. `download`/`transcription` are ExternalOperations, not // kinds, so getOperation() refusing them is the door guard that keeps // download politeness (and the transcription pool) single-scheduler. op: string; channelSlug: string; videoId: string; // Input files by (base)name, base64-encoded. Written in arrival order with // transcript.cues.json FORCED LAST: isCuesJsonFresh compares mtimes, and a // cues file older than its metadata/raw transcript makes the run skip with // stale-cues — a unit that silently does nothing. files: Record; // The primary's resolved freshness target, verbatim. target: unknown; force?: boolean; // The primary's IDENTITY config, verbatim — model, promptVersion, numCtx, // context — with any baseUrl stripped so the executor localises its own // endpoint. Mandatory in spirit: a kind that answers "disabled" or // "not-configured" here is treated as an ERROR, because on a bare box that // answer means the injection was dropped and default settings leaked in. config?: { attribution?: AttributionSettings; diarization?: DiarizationSettings; appConfig?: DigestAppConfig; context?: DigestContext; }; }; export async function startWorkerUnit( input: StartWorkerUnitInput, ): Promise<{ remoteJobId: string }> { const kind = getOperation(input.op); if (!kind) { throw new Error( `operation "${input.op}" is not a backfill kind this executor can run`, ); } if (!ID_RE.test(input.channelSlug) || !ID_RE.test(input.videoId)) { throw new Error("invalid channelSlug or videoId"); } const paths = getPaths(); const remoteJobId = newJobId(); const scratchRoot = path.join(paths.workerScratchDir, remoteJobId); const scratchChannels = path.join(scratchRoot, "channels"); const videoDir = path.join( scratchChannels, input.channelSlug, "data", input.videoId, ); await fs.ensureDir(videoDir); // Materialize the inputs, cues LAST (see StartWorkerUnitInput.files). The // names cross the trust boundary — basename + allowlist, like uploads. const names = Object.keys(input.files ?? {}); const ordered = [ ...names.filter((n) => n !== CUES_JSON_FILENAME), ...names.filter((n) => n === CUES_JSON_FILENAME), ]; for (const name of ordered) { const safe = sanitizeAudioName(name); await fs.writeFile( path.join(videoDir, safe), Buffer.from(input.files[name], "base64"), ); } // The kind runs against a SCRATCH corpus root — the real channelsDir is // never touched, and a bare executor has none anyway. const scratchPaths: Paths = { ...paths, channelsDir: scratchChannels }; const res = await runManagedFunction({ kind: "worker-unit", queueKey: "", // immediate; the primary's scheduler already throttled this paths, // Deliberately NO channelSlug: the completion hook would request a channel // snapshot for a channel this box does not have. videoId: input.videoId, fn: async (onLog, signal) => { const outcome: OperationRunOutcome = await kind.run({ paths: scratchPaths, videoDir, videoId: input.videoId, channelSlug: input.channelSlug, target: input.target, force: input.force, attributionSettings: input.config?.attribution, diarizationSettings: input.config?.diarization, appConfig: input.config?.appConfig, digestContext: input.config?.context, onLog, signal, }); if (outcome === "disabled" || outcome === "not-configured") { throw new Error( `unit ${input.op} reported "${outcome}" on the executor — the config ` + `injection was dropped and this box's default settings leaked in`, ); } if (outcome === "failed") { throw new Error(`unit ${input.op} failed for ${input.videoId}`); } await fs.writeFile( path.join(scratchRoot, UNIT_OUTCOME_FILENAME), JSON.stringify({ outcome }), ); }, }); if (!res.ok) { await fs.remove(scratchRoot).catch(() => {}); throw new Error(res.error); } res.stream.cancel().catch(() => {}); jobs().set(remoteJobId, { remoteJobId, workDir: scratchRoot, inPlace: false, managedJobId: res.jobId, logOffset: 0, unit: { op: input.op, videoDir }, }); return { remoteJobId }; } export type WorkerUnitResult = { // The kind's OperationRunOutcome ("done", "already-present", "skipped", …). outcome: string; // The kind's declared output files, as UTF-8 JSON text by name. May be empty // for a no-op outcome. files: Record; }; // The generalisation of readWorkerResult: instead of one hardcoded // transcript.json, the unit's kind declares its outputs. null until the run // has written its outcome marker. export async function readWorkerUnitResult( remoteJobId: string, ): Promise { const job = jobs().get(remoteJobId); if (!job?.unit) return null; const outcomeRaw = await fs .readFile(path.join(job.workDir, UNIT_OUTCOME_FILENAME), "utf8") .catch(() => null); if (outcomeRaw === null) return null; let outcome = "done"; try { const parsed = JSON.parse(outcomeRaw) as { outcome?: string }; if (typeof parsed.outcome === "string") outcome = parsed.outcome; } catch { // keep "done" — the marker only exists after a successful run } const kind = getOperation(job.unit.op); const files: Record = {}; for (const name of kind?.outputs ?? []) { const p = path.join(job.unit.videoDir, name); if (await fs.pathExists(p)) { files[name] = await fs.readFile(p, "utf8"); } } return { outcome, files }; } export type WorkerJobEvents = { status: JobStatus | "unknown"; fraction?: number; detail?: string; // New log output since the last poll. log: string; }; // Status + incremental log for a worker job. Reads the managed job's record // (status, single task's parsed progress) and tails its log file from the // offset last returned to this requester. export async function getWorkerJobEvents( remoteJobId: string, ): Promise { const job = jobs().get(remoteJobId); if (!job) return { status: "unknown", log: "" }; const rec = getRegistry().get(job.managedJobId); let log = ""; try { const buf = await fs.readFile(rec?.logPath ?? ""); if (buf.length > job.logOffset) { log = buf.toString("utf8", job.logOffset); job.logOffset = buf.length; } } catch { // log file not created yet / already cleaned — fine } if (!rec) return { status: "unknown", log }; const task = rec.tasks?.[0]; return { status: rec.status, fraction: task?.fraction, detail: task?.detail, log, }; } // The produced transcript.json bytes, once the job is done. null if missing. export async function readWorkerResult( remoteJobId: string, ): Promise { const job = jobs().get(remoteJobId); if (!job) return null; const p = path.join(job.workDir, "transcript.json"); if (!(await fs.pathExists(p))) return null; return fs.readFile(p); } // Hard-cancel a running worker job (the requester aborted). Idempotent. export function cancelWorkerJob(remoteJobId: string): void { const job = jobs().get(remoteJobId); if (!job) return; getRegistry().cancel(job.managedJobId); } // Remove a worker job's scratch dir and forget it. Called when the requester has // pulled the result or given up. Idempotent. export async function cleanupWorkerJob(remoteJobId: string): Promise { const job = jobs().get(remoteJobId); if (!job) return; jobs().delete(remoteJobId); // Never delete an in-place (shared-fs) job's dir — that's the requester's // actual video dir. Only scratch dirs we created are removed. if (!job.inPlace) await fs.remove(job.workDir).catch(() => {}); }