Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 33afe87ac54ab81361db57ec41158202087e2915
parent 147f0527761b773c422ba47174cf99464df3c49c
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 14 Jun 2026 01:26:56 -0400

Transcription workers: one slot per worker + remote protocol (Phase 4)

Two changes landed together since both reshape the worker model/UX.

One slot per worker (per user request):
- Drop the `slots` field — each worker IS one slot (one transcription at a
  time). Run N in parallel by defining N workers; the settings editor gains a
  Copy button to duplicate one. Each slot is independently togglable on the
  Workers page, so you can free *some* of a CPU/GPU and keep the rest going.
- Pool tracks a boolean `busy` per worker instead of inUse/slots; summary
  exposes `busy`. Migration emits one enabled worker per old
  parallelTranscriptions slot (preserving parallelism), plus a disabled worker
  per other configured engine.
- Settings editor: removed the Slots input, added Copy, and surfaced chough's
  server-URL field (point two copies at one chough --server for two togglable
  slots). Workers page shows idle/busy rather than slot counts.

Remote worker protocol (upload + shared-fs):
- common/controller/workerServer.ts: accept a job (uploaded audio → scratch, or
  shared-fs {slug,videoId} transcribed in place), run it through THIS instance's
  pool, expose status/log/result, cancel + cleanup.
- common/controller/remoteTranscribe.ts: client that uploads, polls events
  (forwarding log + parsed progress to the task bar via the new TaskHandle.update),
  pulls the result, and classifies transport vs transcription failures.
- editor/app/api/worker/* routes (transcribe, [id], [id]/events, [id]/result,
  health), guarded by a WORKER_TOKEN bearer (constant-time compare; endpoints
  disabled unless the token is set). transcribeOne's kind:"remote" branch wired.
- TranscribeError moved to its own module to avoid an import cycle.
- paths: workerScratchDir. test server runs with WORKER_TOKEN for e2e.

e2e: worker-remote.spec (auth, acceptor round-trip, full self-pointing dispatch);
settings/workers specs updated for one-slot + Copy + chough-server. All green;
both packages typecheck clean.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

Diffstat:
Acommon/controller/remoteTranscribe.ts | 210+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/transcribeError.ts | 16++++++++++++++++
Mcommon/controller/transcribeOne.ts | 54++++++++++++++++++++++++++++++++++--------------------
Acommon/controller/workerServer.ts | 237+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/taskHooks.ts | 11+++++++++--
Mcommon/jobs/workerPool.ts | 38++++++++++++++++----------------------
Mcommon/lib/paths.ts | 5+++++
Acommon/lib/workerToken.ts | 43+++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/workers.ts | 51+++++++++++++++++++++------------------------------
Meditor/CHANGELOG.md | 5+++--
Aeditor/app/api/worker/health/route.ts | 20++++++++++++++++++++
Aeditor/app/api/worker/transcribe/[id]/events/route.ts | 25+++++++++++++++++++++++++
Aeditor/app/api/worker/transcribe/[id]/result/route.ts | 33+++++++++++++++++++++++++++++++++
Aeditor/app/api/worker/transcribe/[id]/route.ts | 29+++++++++++++++++++++++++++++
Aeditor/app/api/worker/transcribe/route.ts | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/settings/components/SettingsForm.tsx | 13+++++++------
Meditor/app/settings/components/WorkersField.tsx | 57++++++++++++++++++++++++++++++++++++++++++---------------
Meditor/app/workers/components/WorkersView.tsx | 10+++-------
Meditor/app/workers/page.tsx | 10+++++-----
Meditor/e2e/settings.spec.ts | 58++++++++++++++++++++++++++++++++++++++++++----------------
Aeditor/e2e/worker-remote.spec.ts | 149+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/e2e/workers.spec.ts | 15++++++++-------
Meditor/package.json | 4++--
23 files changed, 1017 insertions(+), 134 deletions(-)

diff --git a/common/controller/remoteTranscribe.ts b/common/controller/remoteTranscribe.ts @@ -0,0 +1,210 @@ +// Remote-worker CLIENT side: delegate one transcription to another instance of +// this app over HTTP. Uploads the audio, polls progress (forwarding the remote's +// log + parsed {fraction, detail}), pulls back the transcript.json bytes, and +// cleans up the remote scratch. transcribeOne calls this for kind:"remote". +// +// Failure classes (see TranscribeError): a network/5xx/poll problem is +// "transport" (the batch retries on another worker); a remote job that finishes +// "failed" is "transcription" (the audio itself failed — no point retrying). + +import fs from "fs-extra"; +import type { Worker } from "../lib/workers"; +import { TranscribeError } from "./transcribeError"; + +const POLL_MS = 1000; +// Stop polling if the remote goes silent this long (no successful poll), to +// avoid a wedged remote parking a batch slot forever. +const POLL_TIMEOUT_MS = 10 * 60 * 1000; + +export type RemoteTranscribeOptions = { + worker: Worker; // kind === "remote" + audioPath: string; + audioName: string; + videoId: string; + // Required for the shared-fs fast path (worker.remote.sharedFs): the remote + // transcribes channelsDir/<slug>/data/<videoId> in place on the shared mount. + channelSlug?: string; + onLog?: (line: string) => void; + onProgress?: (patch: { fraction?: number; detail?: string }) => void; + signal?: AbortSignal; +}; + +type EventsResponse = { + status: "queued" | "running" | "done" | "failed" | "cancelled" | "unknown"; + fraction?: number; + detail?: string; + log?: string; +}; + +function authHeaders(worker: Worker): Record<string, string> { + const token = worker.remote?.token?.trim(); + return token ? { Authorization: `Bearer ${token}` } : {}; +} + +function base(worker: Worker): string { + return (worker.remote?.baseUrl ?? "").replace(/\/+$/, ""); +} + +// Delegate a transcription to the remote and poll to completion. Returns the +// transcript.json bytes for the upload path, or null for the shared-fs path +// (the remote wrote the transcript directly onto the shared mount, so the caller +// already has it in place and only needs to normalize). +export async function transcribeViaRemote( + opts: RemoteTranscribeOptions, +): Promise<Buffer | null> { + const { worker } = opts; + const baseUrl = base(worker); + if (!baseUrl) { + throw new TranscribeError( + `remote worker "${worker.name}" has no base URL`, + "transport", + ); + } + const log = opts.onLog ?? (() => {}); + const sharedFs = worker.remote?.sharedFs === true && !!opts.channelSlug; + + // 1. Submit the job → remoteJobId. Shared-fs sends a JSON path reference; + // otherwise the raw audio bytes are uploaded. + let remoteJobId: string; + try { + const res = sharedFs + ? await fetch(`${baseUrl}/api/worker/transcribe`, { + method: "POST", + headers: { ...authHeaders(worker), "Content-Type": "application/json" }, + body: JSON.stringify({ + channelSlug: opts.channelSlug, + videoId: opts.videoId, + audioFilename: opts.audioName, + }), + signal: opts.signal, + }) + : await fetch(`${baseUrl}/api/worker/transcribe`, { + method: "POST", + headers: { + ...authHeaders(worker), + "Content-Type": "application/octet-stream", + "X-Audio-Name": opts.audioName, + "X-Job-Label": opts.videoId, + }, + body: await fs.readFile(opts.audioPath), + signal: opts.signal, + }); + if (!res.ok) { + throw new TranscribeError( + `remote worker rejected the job (${res.status} ${res.statusText})`, + res.status >= 400 && res.status < 500 && res.status !== 429 + ? "transcription" + : "transport", + ); + } + const body = (await res.json()) as { remoteJobId?: string }; + if (!body.remoteJobId) { + throw new TranscribeError("remote worker returned no job id", "transport"); + } + remoteJobId = body.remoteJobId; + } catch (err) { + if (err instanceof TranscribeError) throw err; + if (opts.signal?.aborted) throw err; + throw new TranscribeError( + `could not reach remote worker "${worker.name}": ${(err as Error).message}`, + "transport", + ); + } + + const jobUrl = `${baseUrl}/api/worker/transcribe/${remoteJobId}`; + + // 2. Poll events until terminal, forwarding log + progress. + try { + let lastOk = Date.now(); + for (;;) { + if (opts.signal?.aborted) throw abortError(); + await delay(POLL_MS, opts.signal); + let ev: EventsResponse; + try { + const res = await fetch(`${jobUrl}/events`, { + headers: authHeaders(worker), + signal: opts.signal, + }); + if (!res.ok) { + throw new Error(`${res.status} ${res.statusText}`); + } + ev = (await res.json()) as EventsResponse; + lastOk = Date.now(); + } catch (err) { + if (opts.signal?.aborted) throw abortError(); + if (Date.now() - lastOk > POLL_TIMEOUT_MS) { + throw new TranscribeError( + `remote worker "${worker.name}" stopped responding`, + "transport", + ); + } + continue; // transient — keep polling until the timeout + } + if (ev.log) for (const line of ev.log.split(/\r?\n/)) if (line) log(line); + if (ev.fraction !== undefined || ev.detail !== undefined) { + opts.onProgress?.({ fraction: ev.fraction, detail: ev.detail }); + } + if (ev.status === "done") break; + if (ev.status === "failed") { + throw new TranscribeError( + `remote transcription failed on "${worker.name}"`, + "transcription", + ); + } + if (ev.status === "cancelled") { + throw new TranscribeError( + `remote transcription was cancelled on "${worker.name}"`, + "transport", + ); + } + if (ev.status === "unknown") { + throw new TranscribeError( + `remote worker "${worker.name}" forgot the job`, + "transport", + ); + } + } + + // 3. Shared-fs: the transcript is already on the shared mount, nothing to + // pull. Otherwise fetch the produced bytes. + if (sharedFs) return null; + const res = await fetch(`${jobUrl}/result`, { + headers: authHeaders(worker), + signal: opts.signal, + }); + if (!res.ok) { + throw new TranscribeError( + `remote worker "${worker.name}" had no result (${res.status})`, + "transport", + ); + } + return Buffer.from(await res.arrayBuffer()); + } finally { + // 4. Best-effort cleanup of the remote scratch (also cancels if still running + // — e.g. when we abort). Fire-and-forget; never block on it. + void fetch(jobUrl, { method: "DELETE", headers: authHeaders(worker) }).catch( + () => {}, + ); + } +} + +function abortError(): Error { + const err = new Error("remote transcription aborted"); + err.name = "AbortError"; + return err; +} + +function delay(ms: number, signal?: AbortSignal): Promise<void> { + return new Promise((resolve, reject) => { + if (signal?.aborted) return reject(abortError()); + const t = setTimeout(resolve, ms); + signal?.addEventListener( + "abort", + () => { + clearTimeout(t); + reject(abortError()); + }, + { once: true }, + ); + }); +} diff --git a/common/controller/transcribeError.ts b/common/controller/transcribeError.ts @@ -0,0 +1,16 @@ +// Why a transcription failed, so a batch can decide whether to retry it on a +// different worker. "transport" = the worker/transport was at fault (remote +// unreachable, 5xx) and the audio may still transcribe elsewhere; "transcription" +// = the audio itself failed and retrying on another worker is pointless. +// +// In its own module (no deps) so both transcribeOne.ts and the remote client +// (remoteTranscribe.ts) can throw it without an import cycle. +export class TranscribeError extends Error { + constructor( + message: string, + readonly failureClass: "transport" | "transcription", + ) { + super(message); + this.name = "TranscribeError"; + } +} diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -4,12 +4,15 @@ import { execa } from "execa"; import type { Paths } from "../lib/paths"; import type { Worker } from "../lib/workers"; import { getTranscriptionApp } from "../lib/transcriptionApps"; +import { detectTranscriptFormat } from "../lib/whisper"; import { getWorkerPool } from "../jobs/workerPool"; import type { TaskTracker } from "../jobs/taskHooks"; import { normalizeTranscript } from "./normalizeTranscript"; +import { transcribeViaRemote } from "./remoteTranscribe"; +import { TranscribeError } from "./transcribeError"; import { isRealAudioFile } from "../lib/videoStatus"; -const { pathExists, readdir, rename } = fs; +const { pathExists, readdir, rename, writeFile } = fs; const AUDIO_PREFERENCE = ["audio.mp3", "audio.m4a", "audio.opus"]; @@ -40,24 +43,18 @@ export type TranscribeOneOptions = { // via execa; a remote worker delegates to another app instance. worker: Worker; onLog?: (msg: string) => void; + // Parsed progress callback, used by the remote branch to forward the remote's + // {fraction, detail} to the per-task bar (the local branch reports progress via + // onLog → the engine's parser instead). + onProgress?: (patch: { fraction?: number; detail?: string }) => void; signal?: AbortSignal; }; export type TranscribeOneOutcome = "transcribed" | "already-exists"; -// Why a transcription failed, so a batch can decide whether to retry it on a -// different worker. "transport" = the worker/transport was at fault (remote -// unreachable, 5xx) and the audio may still transcribe elsewhere; "transcription" -// = the audio itself failed and retrying on another worker is pointless. -export class TranscribeError extends Error { - constructor( - message: string, - readonly failureClass: "transport" | "transcription", - ) { - super(message); - this.name = "TranscribeError"; - } -} +// TranscribeError lives in its own module so the remote client can throw it +// without an import cycle; re-exported here for existing import sites. +export { TranscribeError } from "./transcribeError"; export async function transcribeOneVideo( opts: TranscribeOneOptions, @@ -97,13 +94,29 @@ export async function transcribeOneVideo( let outputFormat: "whisper-json" | "chough-json" | "vtt"; if (worker.kind === "remote") { - // Phase 4 wires this to common/controller/remoteTranscribe.ts. Until then a - // remote worker can't run; surface it as a transport failure so a batch - // re-routes to a local worker instead of marking the audio failed. - throw new TranscribeError( - `remote worker "${worker.name}" is not yet supported`, - "transport", + // Delegate to another instance of this app over HTTP: upload the audio, poll + // progress, pull back the transcript.json bytes, and write them locally. The + // remote ran its OWN worker pool to produce them; we normalize here. + log( + `Transcribe ${opts.videoId} start (${resolvedAudio}) via remote [${worker.id}] ${worker.remote?.baseUrl ?? ""}`, ); + const bytes = await transcribeViaRemote({ + worker, + audioPath: path.join(opts.videoDir, resolvedAudio), + audioName: resolvedAudio, + videoId: opts.videoId, + channelSlug: deriveChannelSlug(opts.paths, opts.videoDir) ?? undefined, + onLog: (line: string) => log(line), + onProgress: opts.onProgress, + signal: opts.signal, + }); + // Upload path: write the pulled bytes. Shared-fs path (bytes === null): the + // remote already wrote transcript.json onto the shared mount at this path. + if (bytes) await writeFile(transcriptPath, bytes); + // The remote didn't tell us the engine format; sniff it from the file (the + // same content sniff the indexer uses for mixed corpora). + const raw = await fs.readFile(transcriptPath, "utf8").catch(() => ""); + outputFormat = detectTranscriptFormat(raw) ?? "whisper-json"; } else { const app = getTranscriptionApp(worker.appId); const appConfig = worker.config ?? {}; @@ -232,6 +245,7 @@ export async function transcribeWithWorker( strictAudio: opts.strictAudio, worker, onLog: task ? task.onLog : log, + onProgress: task ? task.update : undefined, signal: opts.signal, }); pool.markSuccess(worker.id); diff --git a/common/controller/workerServer.ts b/common/controller/workerServer.ts @@ -0,0 +1,237 @@ +// 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 } 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"; + +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; +}; + +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<void> { + 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<string, WorkerJob> | undefined; +} + +function jobs(): Map<string, WorkerJob> { + 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 }; +} + +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<WorkerJobEvents> { + 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<Buffer | null> { + 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<void> { + 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(() => {}); +} diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts @@ -10,6 +10,10 @@ import { getTranscriptionApp } from "../lib/transcriptionApps"; // failures don't leak a task into record.tasks). export type TaskHandle = { onLog: (line: string) => void; + // Push a parsed progress update directly (bypassing the text parser). Used by + // remote workers, which stream already-parsed {fraction, detail} back rather + // than raw engine output. No-op when the caller didn't opt into tracking. + update: (patch: { fraction?: number; detail?: string }) => void; end: () => void; }; @@ -42,7 +46,7 @@ export function makeTaskTracker( return { start({ id, label, kind, appId, workerId }) { if (!ctx) { - return { onLog: forwardLog, end: () => {} }; + return { onLog: forwardLog, update: () => {}, end: () => {} }; } const startedAt = Date.now(); ctx.addTask({ id, label, kind, startedAt, workerId }); @@ -68,13 +72,16 @@ export function makeTaskTracker( if (update) ctx.updateTask(id, update); } }; + const update = (patch: { fraction?: number; detail?: string }) => { + ctx.updateTask(id, patch); + }; const end = () => { if (ended) return; ended = true; ctx.recordTaskDone(Date.now() - startedAt); ctx.removeTask(id); }; - return { onLog, end }; + return { onLog, update, end }; }, }; } diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts @@ -37,8 +37,8 @@ export type WorkerSummary = { kind: Worker["kind"]; appId?: string; priority: number; - slots: number; - inUse: number; + // A worker is one slot: busy is true while its single transcription runs. + busy: boolean; state: WorkerRuntimeState; degraded: boolean; enabled: boolean; // persisted intent (config.enabled) @@ -46,7 +46,8 @@ export type WorkerSummary = { type PoolEntry = { config: Worker; - inUse: number; + // One slot per worker: busy while its single transcription runs. + busy: boolean; state: WorkerRuntimeState; degraded: boolean; consecutiveFailures: number; @@ -109,7 +110,7 @@ class WorkerPool { existing.state = "enabled"; existing.degraded = false; existing.consecutiveFailures = 0; - } else if (existing.inUse > 0) { + } else if (existing.busy) { existing.state = "draining"; } else { existing.state = "disabled"; @@ -119,7 +120,7 @@ class WorkerPool { } else { this.entries.set(w.id, { config: w, - inUse: 0, + busy: false, state: w.enabled ? "enabled" : "disabled", degraded: false, consecutiveFailures: 0, @@ -130,7 +131,7 @@ class WorkerPool { // Workers no longer in settings: drop if idle, else retire when idle. for (const [id, entry] of this.entries) { if (seen.has(id)) continue; - if (entry.inUse > 0) { + if (entry.busy) { entry.retireWhenIdle = true; entry.state = "draining"; } else { @@ -182,11 +183,7 @@ class WorkerPool { } private eligible(entry: PoolEntry): boolean { - return ( - entry.state === "enabled" && - !entry.degraded && - entry.inUse < entry.config.slots - ); + return entry.state === "enabled" && !entry.degraded && !entry.busy; } // Highest-priority free entry: (priority asc, insertion order asc). @@ -200,18 +197,16 @@ class WorkerPool { } private grant(entry: PoolEntry): Lease { - entry.inUse++; + entry.busy = true; let released = false; const release = () => { if (released) return; released = true; - entry.inUse = Math.max(0, entry.inUse - 1); - if (entry.inUse === 0) { - if (entry.retireWhenIdle) { - this.entries.delete(entry.config.id); - } else if (entry.state === "draining") { - entry.state = "disabled"; - } + entry.busy = false; + if (entry.retireWhenIdle) { + this.entries.delete(entry.config.id); + } else if (entry.state === "draining") { + entry.state = "disabled"; } this.pump(); }; @@ -300,7 +295,7 @@ class WorkerPool { drainWorker(id: string): boolean { const entry = this.entries.get(id); if (!entry) return false; - entry.state = entry.inUse > 0 ? "draining" : "disabled"; + entry.state = entry.busy ? "draining" : "disabled"; this.syncSnapshot(id, "disabled"); return true; } @@ -348,8 +343,7 @@ class WorkerPool { kind: e.config.kind, appId: e.config.appId, priority: e.config.priority, - slots: e.config.slots, - inUse: e.inUse, + busy: e.busy, state: e.state, degraded: e.degraded, enabled: e.config.enabled, diff --git a/common/lib/paths.ts b/common/lib/paths.ts @@ -11,6 +11,10 @@ export type Paths = { // single global channel pool; channel downloads are never duplicated per site. sitesDir: string; jobsDir: string; + // Scratch area for remote-worker transcription requests: uploaded audio and + // the produced transcript.json live under workerScratchDir/<remoteJobId>/ and + // are cleaned up once the requesting instance pulls the result (or cancels). + workerScratchDir: string; lmdbPath: string; exportDir: string; // The dir Next.js serves at "/" (also holds checked-in static assets). For @@ -68,6 +72,7 @@ export function getPaths(): Paths { channelsDir: path.join(transcriptsDir, "channels"), sitesDir: process.env.SITES_DIR ?? path.join(transcriptsDir, "sites"), jobsDir: path.join(transcriptsDir, ".jobs"), + workerScratchDir: path.join(transcriptsDir, ".worker-scratch"), lmdbPath: path.join(transcriptsDir, "index.mdb"), exportDir, exportPublicDir, diff --git a/common/lib/workerToken.ts b/common/lib/workerToken.ts @@ -0,0 +1,43 @@ +import crypto from "node:crypto"; + +// Shared-secret auth for the /api/worker/* endpoints (the LAN remote-worker +// protocol). The ACCEPTING instance validates the incoming bearer token against +// its own WORKER_TOKEN env var — never against settings — so a leaked settings +// file can't reveal the acceptor's secret. The endpoints are DISABLED unless +// WORKER_TOKEN is set, so an instance is never an open transcription server by +// accident; you opt in by setting the env var. + +export function getWorkerToken(): string { + return process.env.WORKER_TOKEN ?? ""; +} + +export function workerEndpointEnabled(): boolean { + return getWorkerToken().length > 0; +} + +export type WorkerAuth = { ok: true } | { ok: false; status: number; error: string }; + +// Validate an Authorization header against WORKER_TOKEN with a constant-time +// compare. 503 when the endpoint is disabled (no token configured), 401 on a +// missing/malformed/mismatched token. +export function authorizeWorkerRequest(authHeader: string | null): WorkerAuth { + const token = getWorkerToken(); + if (!token) { + return { + ok: false, + status: 503, + error: "worker endpoint disabled (set WORKER_TOKEN to enable)", + }; + } + const match = /^Bearer\s+(.+)$/i.exec(authHeader ?? ""); + if (!match) return { ok: false, status: 401, error: "missing bearer token" }; + const provided = Buffer.from(match[1]); + const expected = Buffer.from(token); + if ( + provided.length !== expected.length || + !crypto.timingSafeEqual(provided, expected) + ) { + return { ok: false, status: 401, error: "invalid worker token" }; + } + return { ok: true }; +} diff --git a/common/lib/workers.ts b/common/lib/workers.ts @@ -35,6 +35,10 @@ export type RemoteWorkerConfig = { sharedFs?: boolean; }; +// A worker is ONE processing slot — one transcription at a time. To run N in +// parallel on the same engine, define N workers (the settings editor's "Copy" +// button duplicates one). This makes each slot independently togglable on the +// Workers page (free up one core/GPU while the rest keep going). export type Worker = { // Stable slug; used in settings, task ids, and logs. id: string; @@ -44,8 +48,6 @@ export type Worker = { enabled: boolean; // Lower = preferred. Ties broken by array order in the scheduler. priority: number; - // Concurrent transcriptions this worker handles (1..WORKER_SLOTS_MAX). - slots: number; // Reserved for future capability routing (e.g. "only long audio to the GPU"). // Stored but NOT consulted by the scheduler in v1. tags?: string[]; @@ -58,19 +60,6 @@ export type Worker = { remote?: RemoteWorkerConfig; }; -export const WORKER_SLOTS_MAX = 16; -export const WORKER_SLOTS_DEFAULT = 1; - -export function clampWorkerSlots(value: unknown): number { - const n = - typeof value === "number" && Number.isFinite(value) - ? Math.floor(value) - : WORKER_SLOTS_DEFAULT; - if (n < 1) return 1; - if (n > WORKER_SLOTS_MAX) return WORKER_SLOTS_MAX; - return n; -} - // Coerce a raw per-worker engine config into a clean AppInstanceConfig, dropping // unknown/ill-typed fields. Shared with settings.ts's sanitizeTranscriptionApps. export function sanitizeWorkerConfig(value: unknown): AppInstanceConfig { @@ -136,7 +125,6 @@ export function sanitizeWorkers(value: unknown): Worker[] { typeof r.priority === "number" && Number.isFinite(r.priority) ? Math.floor(r.priority) : index, - slots: clampWorkerSlots(r.slots), }; if (Array.isArray(r.tags)) { const tags = r.tags.filter((t): t is string => typeof t === "string"); @@ -195,33 +183,37 @@ export function validateWorkers(workers: Worker[]): string | null { } // Synthesize the default worker list for a settings.json that predates the worker -// model. The active app becomes one enabled worker (priority 0) whose slot count -// preserves the old global parallelTranscriptions; every other configured app -// becomes a disabled worker, ready to enable. Keeps existing installs behaving -// identically while surfacing prior configuration. +// model. Since a worker is now one slot, the old global parallelTranscriptions is +// preserved by emitting that many copies of the active app (priority 0), each a +// separate enabled slot. Every other configured app becomes one disabled worker, +// ready to enable. Keeps existing installs' parallelism while surfacing prior +// configuration. export function defaultWorkersFromApps( activeAppId: string, appConfigs: Record<string, AppInstanceConfig>, parallelSlots: number, ): Worker[] { const activeApp = getTranscriptionApp(activeAppId); - const workers: Worker[] = [ - { - id: "default", - name: activeApp.label, + const slots = Math.max(1, Math.floor(parallelSlots) || 1); + const workers: Worker[] = []; + const seen = new Set<string>(); + for (let i = 0; i < slots; i++) { + const id = uniqueId(slots > 1 ? `${activeApp.id}-${i + 1}` : activeApp.id, seen); + seen.add(id); + workers.push({ + id, + name: slots > 1 ? `${activeApp.label} #${i + 1}` : activeApp.label, kind: "local", enabled: true, priority: 0, - slots: clampWorkerSlots(parallelSlots), appId: activeApp.id, config: appConfigs[activeApp.id] ?? {}, - }, - ]; + }); + } let priority = 1; - const seen = new Set<string>(["default", activeApp.id]); for (const [appId, config] of Object.entries(appConfigs)) { if (appId === activeApp.id || !TRANSCRIPTION_APPS[appId]) continue; - const id = seen.has(appId) ? `${appId}-worker` : appId; + const id = uniqueId(appId, seen); seen.add(id); workers.push({ id, @@ -229,7 +221,6 @@ export function defaultWorkersFromApps( kind: "local", enabled: false, priority: priority++, - slots: 1, appId, config, }); diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -3,8 +3,9 @@ ## [Unreleased] - **Sites can link to each other.** A site's form gained a **Public URL** field (the absolute URL it's served at, e.g. `https://jeralyzer.com`) and a **Related sites** section. The export footer automatically links to every *other* site that has a Public URL, so filling these in is all that's needed for cross-site links; a site left without a URL is simply omitted from the lists. The **Related sites** editor lets a site pull closely-related siblings to the front under named groups (e.g. Jeralyzer featuring Rekietalyzer under "MTG drama") — add a group, give it an optional heading, and check which sibling sites belong; everything you don't feature falls into a trailing "Other sites" group on its own. Groups reorder with ↑/↓. The picker only lists sites that actually exist, and featured ids for sites that were since deleted are dropped on save (with a heads-up note). It's a subtle, secondary feature — see the matching note in the export changelog for how it renders. - **The editor refreshes itself on a timer so its data stays live without a manual reload.** Every page now passively re-fetches its own server-rendered data on a configurable interval — so the sidebar badges (active/running job counts, changelog dot), channel reports, and any other on-screen figures keep up to date on their own. It uses Next's `router.refresh()` (the same mechanism the jobs list already used) mounted once globally in the root layout, so it covers every page and the shared sidebar with no per-page wiring. To avoid wasting work when you're not looking, it **pauses entirely while the browser tab is hidden** and does **one immediate refresh the moment you return** to the tab (rather than waiting out the interval); it also skips a tick while a previous refresh is still settling, so refreshes can't pile up. The cadence is set in **Settings → Auto-refresh interval (seconds)**: default **5s** (clamped 1–600), or **0 to disable** passive refresh completely. This replaces the jobs page's old bespoke 2.5s auto-refresh (the `/jobs/active` page keeps its faster 1s progress-bar polling, which animates per-task bars without a full re-render). -- **Transcription is now driven by configurable workers instead of one global engine.** The old single **App** dropdown in **Settings → Transcription** is replaced by a **Transcription workers** list: each worker is a named processing slot with its own engine (whisper.cpp / chough / parakeet) and config, a **slots** count (how many videos it transcribes at once), and a priority given by its position in the list (top = preferred). A batch ("Transcribe missing", bucket, bulk, single-video) hands each video — per task — to the highest-priority worker with a free slot, so a fast GPU worker and a slower CPU worker (e.g. parakeet on the GPU + chough on the CPU) run side by side instead of one engine doing everything. Total parallelism is the sum of every enabled worker's slots; the old per-run **Concurrency** control and the global **Parallel transcriptions** setting are gone (worker slots replace them — edit a worker's slots, or disable/drain one, to change load). A pre-worker `settings.json` migrates automatically to one enabled worker built from the previously-selected app (its slots seeded from the old parallel-transcriptions value), plus a disabled worker for any other engine you had configured, so existing installs behave identically. Scheduling is a single process-wide pool, so two batches can't oversubscribe the same GPU. (Remote workers — delegating to another instance of this app on the LAN — are wired into the data model and UI; the network protocol lands in a follow-up.) -- **New Workers page (`/workers`) with live status and runtime controls.** Lists every worker with its live slot usage (`1/2 slots`), state (idle / busy / draining / disabled / degraded), and the videos it's currently transcribing with per-task progress. Each worker can be **disabled** (stop taking new work immediately; in-flight transcriptions keep running), **drained** (stop taking new work but let the current video finish — the graceful "free up the GPU when it's done" path), or **enabled** again — without editing settings, so you can hand a CPU/GPU back to other programs and reclaim it later. A **Pause all** button disables every worker at once and remembers each one's state; **Resume all** restores them exactly. These runtime controls are transient (a restart returns workers to their configured enabled state); the Settings list is where the persisted defaults live. +- **Transcription is now driven by configurable workers instead of one global engine.** The old single **App** dropdown in **Settings → Transcription** is replaced by a **Transcription workers** list. Each worker is **one processing slot** — one transcription at a time — with its own engine (whisper.cpp / chough / parakeet) and config, and a priority given by its position in the list (top = preferred). To run several in parallel, add more workers; a **Copy** button duplicates one (e.g. point two copies at the same chough `--server` for two togglable server slots). A batch ("Transcribe missing", bucket, bulk, single-video) hands each video — per task — to the highest-priority free worker, so a fast GPU worker and a slower CPU worker (e.g. parakeet on the GPU + chough on the CPU) run side by side instead of one engine doing everything. Total parallelism is the number of enabled workers; the old per-run **Concurrency** control and the global **Parallel transcriptions** setting are gone (add/remove workers, or disable/drain one, to change load). A pre-worker `settings.json` migrates automatically to one worker per slot of the previously-selected app (the old parallel-transcriptions count becomes that many enabled copies), plus a disabled worker for any other engine you had configured, so existing installs keep their parallelism. Scheduling is a single process-wide pool, so two batches can't oversubscribe the same GPU. One-slot-per-worker also means you can disable a single slot to free *some* of a CPU/GPU while the rest keep transcribing. +- **Remote workers: offload transcription to another instance of this app on your LAN.** Add a **remote** worker in the Settings list with the base URL of another instance (e.g. `http://gpu-box.lan:3001`) and a shared token. When a video is dispatched to it, this instance uploads the audio over HTTP, the remote transcribes it through *its own* worker pool (picking among its local engines), streams progress and log back, and this instance pulls the finished `transcript.json` and normalizes it locally — so the remote needs no knowledge of your channels, just CPU/GPU. The protocol lives under `/api/worker/*` and is **disabled unless `WORKER_TOKEN` is set** in the environment, so an instance is never an open transcription server by accident; every request carries `Authorization: Bearer <token>`, validated with a constant-time compare against the accepting instance's own `WORKER_TOKEN` (never against settings). Uploaded audio and the produced transcript live in a scratch dir that's cleaned up once the result is pulled (or the job is cancelled). If a remote is unreachable or returns a transport error mid-job, the video is automatically retried on another worker; a genuine transcription failure on the remote is not retried. A `GET /api/worker/health` endpoint returns the remote's worker summary for heartbeating. +- **New Workers page (`/workers`) with live status and runtime controls.** Lists every worker with its state (idle / busy / draining / disabled / degraded) and the video it's currently transcribing with per-task progress. Each worker can be **disabled** (stop taking new work immediately; in-flight transcriptions keep running), **drained** (stop taking new work but let the current video finish — the graceful "free up the GPU when it's done" path), or **enabled** again — without editing settings, so you can hand a CPU/GPU back to other programs and reclaim it later. A **Pause all** button disables every worker at once and remembers each one's state; **Resume all** restores them exactly. These runtime controls are transient (a restart returns workers to their configured enabled state); the Settings list is where the persisted defaults live. - **Batches pause instead of failing when no worker is available.** If every worker is disabled (or you hit **Pause all**) while a transcription batch is running, the batch parks — it keeps its in-flight video to completion, starts no new ones, and stays **running** on `/jobs/active` rather than failing the remaining videos. Re-enabling any worker (or **Resume all**) immediately resumes it where it left off. A video whose worker fails for a transport reason (e.g. a remote worker that went away) is automatically retried on another worker before being recorded as failed. - **Git worktrees can run in parallel on non-colliding ports (dev tooling).** Two checkouts of the repo (via `git worktree`) can now run their dev servers and e2e suites at the same time without port clashes. A new helper, `scripts/worktree.mjs` (exposed as `pnpm wt`), assigns each worktree a port block offset by `index * 100` based on its position in `git worktree list` — the main worktree keeps the original defaults (editor 3001, test 3011, export 3010/3000/3020), worktree #1 gets 31xx, and so on. `pnpm dev:editor`, `pnpm dev:export`, `pnpm start:export`, and `pnpm e2e` route through `wt run`, which injects the assigned ports, so they "just work" per worktree; the editor/export `package.json` port flags and the Playwright configs now honor these env vars (previously `pnpm dev:test` hardcoded 3011, so a custom `PORT` only moved the URL Playwright waited on, not the server). E2E specs that hit the editor's test API now derive the base URL from `PLAYWRIGHT_BASE_URL` (centralized in `editor/e2e/baseUrl.ts`) instead of hardcoding `localhost:3011`. `pnpm wt add <branch>` creates a sibling worktree pre-seeded with `settings.json` and prints its ports; `--share-data` links it to the main worktree's downloaded `transcripts/` for read-mostly reuse (with an LMDB concurrent-write caveat). See `WORKTREES.md`. - **Channel reports refresh themselves automatically, on a global debounce.** Any action that changes what a channel report (the per-channel snapshot powering the Channels list, `/actionable`, and the channel page) would say now regenerates that report on its own when it finishes — no more manually clicking **Refresh report** after transcribing, transcoding, cleaning audio, checking availability, editing channel config, or the per-video file operations (delete file, set primary transcript, delete dir, mark untranscribable, archive/do-not-clean). Every report-changing action marks its channel "dirty" and re-arms one shared debounce timer; when activity settles, the scheduler regenerates each dirty channel's snapshot in parallel (reusing the existing `refresh-report` job, deduped against any refresh already running) and revalidates the affected pages. This happens **after each completed sub-operation within a batch**, not only when the whole batch finishes — so a long download or transcription run updates its report incrementally as each video lands, rather than staying stale until the end. The debounce coalesces bursts — videos that finish within the same window collapse into a single regen pass instead of one rewrite per operation. The window is configurable in **Settings → Report refresh debounce**: **Fast** (~1s after activity settles, no cap — the default), **Balanced** (~3s, 30s max), or **Lazy** (~10s, 60s max). The download/sync pipeline's previous behavior of regenerating its report inline is removed in favor of this one uniform mechanism (so its report now lags by the debounce window — ~1s by default — rather than being written synchronously). The manual **Refresh report** / **Update all reports** buttons are unchanged. diff --git a/editor/app/api/worker/health/route.ts b/editor/app/api/worker/health/route.ts @@ -0,0 +1,20 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; + +export const dynamic = "force-dynamic"; + +// Heartbeat for a remote worker: returns this instance's worker summary so a +// primary can skip a dead/empty remote before dispatching to it. Token-guarded. +export async function GET(request: Request) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const pool = getWorkerPool(); + return NextResponse.json({ + ok: true, + paused: pool.isPaused(), + workers: pool.summary(), + }); +} diff --git a/editor/app/api/worker/transcribe/[id]/events/route.ts b/editor/app/api/worker/transcribe/[id]/events/route.ts @@ -0,0 +1,25 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { getWorkerJobEvents } from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +const ID_RE = /^[A-Za-z0-9_-]+$/; + +// Status + incremental log for a worker job. Polled ~1s by the requesting +// instance's remote-transcribe client. +export async function GET( + request: Request, + { params }: { params: Promise<{ id: string }> }, +) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const { id } = await params; + if (!ID_RE.test(id)) { + return NextResponse.json({ error: "invalid id" }, { status: 400 }); + } + const events = await getWorkerJobEvents(id); + return NextResponse.json(events); +} diff --git a/editor/app/api/worker/transcribe/[id]/result/route.ts b/editor/app/api/worker/transcribe/[id]/result/route.ts @@ -0,0 +1,33 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { readWorkerResult } from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +const ID_RE = /^[A-Za-z0-9_-]+$/; + +// Stream the produced transcript.json once the job is done. 404 if not ready. +export async function GET( + request: Request, + { params }: { params: Promise<{ id: string }> }, +) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const { id } = await params; + if (!ID_RE.test(id)) { + return NextResponse.json({ error: "invalid id" }, { status: 400 }); + } + const bytes = await readWorkerResult(id); + if (!bytes) { + return NextResponse.json({ error: "no result yet" }, { status: 404 }); + } + return new NextResponse(new Uint8Array(bytes), { + status: 200, + headers: { + "Content-Type": "application/json", + "Cache-Control": "no-store", + }, + }); +} diff --git a/editor/app/api/worker/transcribe/[id]/route.ts b/editor/app/api/worker/transcribe/[id]/route.ts @@ -0,0 +1,29 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { + cancelWorkerJob, + cleanupWorkerJob, +} from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +const ID_RE = /^[A-Za-z0-9_-]+$/; + +// Cancel (if running) and clean up the job's scratch dir. Called when the +// requester has the result or aborted. Idempotent. +export async function DELETE( + request: Request, + { params }: { params: Promise<{ id: string }> }, +) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const { id } = await params; + if (!ID_RE.test(id)) { + return NextResponse.json({ error: "invalid id" }, { status: 400 }); + } + cancelWorkerJob(id); + await cleanupWorkerJob(id); + return NextResponse.json({ ok: true }); +} diff --git a/editor/app/api/worker/transcribe/route.ts b/editor/app/api/worker/transcribe/route.ts @@ -0,0 +1,58 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { + startWorkerTranscription, + startWorkerTranscriptionShared, +} from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +// Accept a transcription job from another instance. Two body shapes: +// - application/json {channelSlug, videoId, audioFilename}: shared-fs fast path, +// transcribe the video dir in place on the shared mount. +// - else raw audio bytes (X-Audio-Name names the file, X-Job-Label is a label): +// write to scratch and transcribe, requester pulls the result. +// Either way returns { remoteJobId } (202); the requester polls +// /api/worker/transcribe/<id>/events and (upload path) pulls /result. +export async function POST(request: Request) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const contentType = request.headers.get("content-type") ?? ""; + try { + if (contentType.includes("application/json")) { + const body = (await request.json()) as { + channelSlug?: string; + videoId?: string; + audioFilename?: string; + }; + if (!body.channelSlug || !body.videoId) { + return NextResponse.json( + { error: "channelSlug and videoId are required" }, + { status: 400 }, + ); + } + const { remoteJobId } = await startWorkerTranscriptionShared({ + channelSlug: body.channelSlug, + videoId: body.videoId, + audioFilename: body.audioFilename ?? "audio.mp3", + }); + return NextResponse.json({ remoteJobId, inPlace: true }, { status: 202 }); + } + const audioName = request.headers.get("x-audio-name") ?? "audio"; + const label = request.headers.get("x-job-label") ?? undefined; + const audio = Buffer.from(await request.arrayBuffer()); + if (audio.length === 0) { + return NextResponse.json({ error: "empty audio body" }, { status: 400 }); + } + const { remoteJobId } = await startWorkerTranscription({ + audio, + audioName, + label, + }); + return NextResponse.json({ remoteJobId }, { status: 202 }); + } catch (e) { + return NextResponse.json({ error: (e as Error).message }, { status: 500 }); + } +} diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx @@ -49,15 +49,16 @@ export function SettingsForm({ initial, apps }: Props) { Transcription workers </legend> <p className="text-xs text-zinc-500"> - Named processing slots for &ldquo;Transcribe missing&rdquo; and friends. - List order is priority (top = preferred). Each video is handed to the - highest-priority worker with a free slot, so a fast GPU worker and a - CPU worker run side by side. <code>Slots</code> is how many videos that - worker transcribes at once. Enable/disable and drain workers live on the{" "} + Each worker is one processing slot — one transcription at a time. To run + several in parallel on the same engine, add more workers (use{" "} + <strong>Copy</strong> to duplicate one). List order is priority (top = + preferred); each video goes to the highest-priority free worker, so a + fast GPU worker and a CPU worker run side by side. Enable, disable, and + drain individual workers live on the{" "} <a href="/workers" className="underline"> Workers </a>{" "} - page. + page — handy for freeing one core/GPU while the rest keep going. </p> <WorkersField initial={initial.workers} apps={apps} name="workersJson" /> </fieldset> diff --git a/editor/app/settings/components/WorkersField.tsx b/editor/app/settings/components/WorkersField.tsx @@ -25,7 +25,6 @@ function blankLocal(appId: string): Row { kind: "local", enabled: true, priority: 0, - slots: 1, appId, config: {}, }; @@ -39,11 +38,25 @@ function blankRemote(): Row { kind: "remote", enabled: true, priority: 0, - slots: 1, remote: { baseUrl: "" }, }; } +// Duplicate a worker into a new slot: same config, fresh id, "(copy)" name. +function copyOf(row: Row): Row { + const { _key, id, name, ...rest } = row; + void _key; + void id; + return { + ...rest, + _key: keyCounter++, + id: "", + name: name ? `${name} (copy)` : "", + config: row.config ? { ...row.config } : undefined, + remote: row.remote ? { ...row.remote } : undefined, + }; +} + // Strip the React-only _key and set priority from list order before serializing. function serialize(rows: Row[]): string { return JSON.stringify( @@ -89,6 +102,11 @@ export function WorkersField({ initial, apps, name }: Props) { function remove(key: number) { update(rows.filter((r) => r._key !== key)); } + function copy(index: number) { + const next = [...rows]; + next.splice(index + 1, 0, copyOf(rows[index])); + update(next); + } function move(index: number, dir: -1 | 1) { const j = index + dir; if (j < 0 || j >= rows.length) return; @@ -116,6 +134,7 @@ export function WorkersField({ initial, apps, name }: Props) { onPatchConfig={(c) => patchConfig(row._key, c)} onPatchRemote={(c) => patchRemote(row._key, c)} onRemove={() => remove(row._key)} + onCopy={() => copy(i)} onMove={(dir) => move(i, dir)} /> ))} @@ -148,6 +167,7 @@ function WorkerCard({ onPatchConfig, onPatchRemote, onRemove, + onCopy, onMove, }: { row: Row; @@ -158,6 +178,7 @@ function WorkerCard({ onPatchConfig: (c: Partial<NonNullable<Worker["config"]>>) => void; onPatchRemote: (c: Partial<NonNullable<Worker["remote"]>>) => void; onRemove: () => void; + onCopy: () => void; onMove: (dir: -1 | 1) => void; }) { const uid = useId(); @@ -211,6 +232,15 @@ function WorkerCard({ </button> <button type="button" + onClick={onCopy} + aria-label={`copy worker ${index + 1}`} + title="Duplicate this worker into another slot" + className="px-1.5 py-0.5 rounded border border-zinc-300 dark:border-zinc-700 text-xs hover:bg-zinc-100 dark:hover:bg-zinc-800" + > + Copy + </button> + <button + type="button" onClick={onRemove} aria-label={`remove worker ${index + 1}`} className="px-1.5 py-0.5 rounded border border-red-300 dark:border-red-800 text-xs text-red-700 dark:text-red-300 hover:bg-red-50 dark:hover:bg-red-950" @@ -238,20 +268,8 @@ function WorkerCard({ </select> </label> )} - <label className="flex flex-col gap-1 text-xs"> - <span className="font-medium">Slots</span> - <input - type="number" - min={1} - max={16} - value={row.slots} - onChange={(e) => onPatch({ slots: Number(e.target.value) || 1 })} - aria-label={`worker ${index + 1} slots`} - className="w-20 rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm" - /> - </label> <span className="text-xs text-zinc-500"> - {row.kind === "local" ? "Local" : "Remote"} + {row.kind === "local" ? "Local · one slot" : "Remote · one slot"} </span> </div> @@ -273,6 +291,15 @@ function WorkerCard({ id={`${uid}-model`} /> )} + {app?.fields.remoteUrl && ( + <CardField + label="chough server URL" + value={cfg.remoteUrl ?? ""} + onChange={(v) => onPatchConfig({ remoteUrl: v })} + hint="chough only: transcribe via a chough --server (sets CHOUGH_URL). Point two copies of this worker at the same server for two togglable slots. Blank = run chough locally." + id={`${uid}-remoteurl`} + /> + )} {app?.fields.chunkSize && ( <CardField label="Chunk size (seconds)" diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx @@ -23,8 +23,7 @@ export type WorkerView = { kind: "local" | "remote"; appId?: string; priority: number; - slots: number; - inUse: number; + busy: boolean; state: WorkerRuntimeState; degraded: boolean; enabled: boolean; @@ -134,7 +133,7 @@ function stateBadge(w: WorkerView): { label: string; className: string } { switch (w.state) { case "enabled": return { - label: w.inUse > 0 ? "busy" : "idle", + label: w.busy ? "busy" : "idle", className: "bg-green-100 dark:bg-green-950 text-green-800 dark:text-green-200 border-green-300 dark:border-green-800", }; @@ -177,9 +176,6 @@ function WorkerCard({ {badge.label} </span> <span className="text-xs text-zinc-500">{engine}</span> - <span className="text-xs text-zinc-500" title="slots in use / total"> - {w.inUse}/{w.slots} slots - </span> <span className="text-xs text-zinc-400">priority {w.priority}</span> <span className="ml-auto flex items-center gap-1.5"> {(w.state === "disabled" || w.state === "draining" || w.degraded) && ( @@ -197,7 +193,7 @@ function WorkerCard({ <> <button type="button" - disabled={pending || w.inUse === 0} + disabled={pending || !w.busy} onClick={() => run(() => drainWorkerAction(w.id))} aria-label={`drain ${w.name}`} title="Stop taking new work; let in-flight transcriptions finish" diff --git a/editor/app/workers/page.tsx b/editor/app/workers/page.tsx @@ -22,11 +22,11 @@ export default async function WorkersPage() { </Link> </div> <p className="text-sm text-zinc-500"> - Transcription workers and their live slot usage. Enable, disable, or drain - a worker to free up a CPU/GPU for other programs, then turn it back on - when done. Disabling every worker pauses running batches (they wait for a - worker) instead of failing. Worker config (engines, slots, priority) lives - on the{" "} + Transcription workers — one slot each. Enable, disable, or drain a worker + to free up a CPU/GPU for other programs, then turn it back on when done. + Disabling every worker pauses running batches (they wait for a worker) + instead of failing. Worker config (engines, priority, copies) lives on + the{" "} <Link href="/settings" className="underline"> Settings </Link>{" "} diff --git a/editor/e2e/settings.spec.ts b/editor/e2e/settings.spec.ts @@ -56,20 +56,19 @@ test("saves global default social links", async ({ page }) => { expect(saved.socialLinks?.[0].svg).toContain("currentColor"); }); -test("shows the migrated default worker and persists an added worker", async ({ +test("shows the migrated default workers and persists an added worker", async ({ page, }) => { await page.goto("/settings"); - // An empty settings.json migrates to one enabled whisper-cpp worker. - const w1 = page.getByLabel("worker 1 name"); - await expect(w1).toHaveValue(/whisper/i); - await expect(page.getByLabel("worker 1 slots")).toHaveValue("2"); + // An empty settings.json migrates to one whisper-cpp worker PER slot — the old + // default parallelTranscriptions of 2 becomes two enabled whisper workers. + await expect(page.getByLabel("worker 1 name")).toHaveValue(/whisper/i); + await expect(page.getByLabel("worker 2 name")).toHaveValue(/whisper/i); - // Add a second (CPU) worker, name it, give it 1 slot. + // Add a third worker on a different engine. await page.getByRole("button", { name: "+ Local worker" }).click(); - await page.getByLabel("worker 2 name").fill("CPU chough"); - await page.getByLabel("worker 2 engine").selectOption("chough"); - await page.getByLabel("worker 2 slots").fill("1"); + await page.getByLabel("worker 3 name").fill("CPU chough"); + await page.getByLabel("worker 3 engine").selectOption("chough"); await page.getByRole("button", { name: /save settings/i }).click(); await expect( @@ -77,22 +76,49 @@ test("shows the migrated default worker and persists an added worker", async ({ ).toBeVisible(); const saved = await readJson<{ - workers?: { name: string; appId?: string; slots: number; priority: number }[]; + workers?: { name: string; appId?: string; priority: number }[]; }>("test-settings.json"); - expect(saved.workers).toHaveLength(2); - expect(saved.workers?.[1].name).toBe("CPU chough"); - expect(saved.workers?.[1].appId).toBe("chough"); - expect(saved.workers?.[1].slots).toBe(1); + expect(saved.workers).toHaveLength(3); + expect(saved.workers?.[2].name).toBe("CPU chough"); + expect(saved.workers?.[2].appId).toBe("chough"); // List order is priority. - expect(saved.workers?.[1].priority).toBe(1); + expect(saved.workers?.[2].priority).toBe(2); await page.reload(); - await expect(page.getByLabel("worker 2 name")).toHaveValue("CPU chough"); + await expect(page.getByLabel("worker 3 name")).toHaveValue("CPU chough"); +}); + +test("copy duplicates a worker into a new slot", async ({ page }) => { + await page.goto("/settings"); + // Configure worker 1 as a chough-server worker, then copy it — the pattern for + // running two togglable slots against one chough --server. + await page.getByLabel("worker 1 engine").selectOption("chough"); + await page.getByLabel("worker 1 name").fill("Chough server"); + await page.getByLabel(/chough server URL/i).fill("http://localhost:8165"); + await page.getByRole("button", { name: "copy worker 1" }).click(); + + await expect(page.getByLabel("worker 2 name")).toHaveValue("Chough server (copy)"); + + await page.getByRole("button", { name: /save settings/i }).click(); + await expect( + page.getByRole("status").filter({ hasText: "Saved" }), + ).toBeVisible(); + + const saved = await readJson<{ + workers?: { appId?: string; config?: { remoteUrl?: string } }[]; + }>("test-settings.json"); + const choughs = (saved.workers ?? []).filter((w) => w.appId === "chough"); + expect(choughs.length).toBe(2); + expect(choughs.every((w) => w.config?.remoteUrl === "http://localhost:8165")).toBe( + true, + ); }); test("rejects saving with no enabled workers", async ({ page }) => { await page.goto("/settings"); + // Migration produced two enabled workers; disable both. await page.getByLabel("worker 1 enabled").uncheck(); + await page.getByLabel("worker 2 enabled").uncheck(); await page.getByRole("button", { name: /save settings/i }).click(); await expect(page.locator("form").getByRole("alert")).toContainText( /at least one worker must be enabled/i, diff --git a/editor/e2e/worker-remote.spec.ts b/editor/e2e/worker-remote.spec.ts @@ -0,0 +1,149 @@ +// Remote-worker protocol. The test server runs with WORKER_TOKEN set (see +// package.json dev:test), so its /api/worker/* endpoints are live. Two angles: +// 1. The acceptor API directly (auth + a raw upload→poll→result round-trip). +// 2. A full end-to-end where a worker points at the instance's OWN base URL. +// A worker is one slot, so the outer lease holds the remote worker's slot +// while the acceptor's internal transcribe falls to the local worker — no +// delegation loop — exercising the whole client→HTTP→server→pool→result path. + +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { pathExists, resetData, resolvePath, writeSettings } from "./helpers"; +import { baseUrl } from "./baseUrl"; + +const TOKEN = "test-worker-token"; + +async function makeTranscribeChannel(slug: string, ids: string[]) { + const root = resolvePath(`test-transcripts/channels/${slug}`); + await mkdir(root, { recursive: true }); + await writeFile( + `${root}/config.json`, + JSON.stringify({ + handling: "transcribe", + name: slug, + url: "https://odysee.com/@example", + audioFormat: "mp3", + }), + ); + for (const id of ids) { + await mkdir(`${root}/data/${id}`, { recursive: true }); + await writeFile(`${root}/data/${id}/audio.mp3`, `fake audio ${id}\n`); + } +} + +test.beforeEach(async () => { + await resetData("empty"); +}); + +test("worker health endpoint enforces the bearer token", async ({ request }) => { + const noAuth = await request.get(`${baseUrl}/api/worker/health`); + expect(noAuth.status()).toBe(401); + + const wrong = await request.get(`${baseUrl}/api/worker/health`, { + headers: { authorization: "Bearer nope" }, + }); + expect(wrong.status()).toBe(401); + + const ok = await request.get(`${baseUrl}/api/worker/health`, { + headers: { authorization: `Bearer ${TOKEN}` }, + }); + expect(ok.status()).toBe(200); + const body = await ok.json(); + expect(body.ok).toBe(true); + expect(Array.isArray(body.workers)).toBe(true); +}); + +test("acceptor transcribes an uploaded audio and returns the result", async ({ + request, +}) => { + // One local worker so the acceptor has something to run the upload on. + await writeSettings({ + workers: [ + { id: "cpu", name: "CPU", kind: "local", enabled: true, priority: 0, appId: "whisper-cpp", config: {} }, + ], + }); + + const post = await request.post(`${baseUrl}/api/worker/transcribe`, { + headers: { + authorization: `Bearer ${TOKEN}`, + "content-type": "application/octet-stream", + "x-audio-name": "audio.mp3", + "x-job-label": "uploaded1", + }, + data: Buffer.from("fake audio uploaded1\n"), + }); + expect(post.status()).toBe(202); + const { remoteJobId } = await post.json(); + expect(remoteJobId).toBeTruthy(); + + // Poll events until done. + await expect + .poll( + async () => { + const r = await request.get( + `${baseUrl}/api/worker/transcribe/${remoteJobId}/events`, + { headers: { authorization: `Bearer ${TOKEN}` } }, + ); + const b = await r.json(); + return b.status; + }, + { timeout: 30_000 }, + ) + .toBe("done"); + + const result = await request.get( + `${baseUrl}/api/worker/transcribe/${remoteJobId}/result`, + { headers: { authorization: `Bearer ${TOKEN}` } }, + ); + expect(result.status()).toBe(200); + const transcript = await result.json(); + // fake-whisper emits a whisper-json doc with a transcription array. + expect(transcript).toHaveProperty("transcription"); + + // Cleanup removes the scratch; result then 404s. + const del = await request.delete( + `${baseUrl}/api/worker/transcribe/${remoteJobId}`, + { headers: { authorization: `Bearer ${TOKEN}` } }, + ); + expect(del.status()).toBe(200); + const gone = await request.get( + `${baseUrl}/api/worker/transcribe/${remoteJobId}/result`, + { headers: { authorization: `Bearer ${TOKEN}` } }, + ); + expect(gone.status()).toBe(404); +}); + +test("a batch dispatched to a remote worker transcribes via the HTTP round-trip", async ({ + page, +}) => { + test.setTimeout(60_000); + // Remote worker (priority 0) points at THIS instance; a local worker backs the + // acceptor. The remote worker's single slot is held by the outer lease while + // the acceptor runs, so its internal transcribe uses the local worker — no loop. + await writeSettings({ + workers: [ + { id: "self", name: "Self Remote", kind: "remote", enabled: true, priority: 0, remote: { baseUrl, token: TOKEN } }, + { id: "cpu", name: "CPU", kind: "local", enabled: true, priority: 1, appId: "whisper-cpp", config: {} }, + ], + }); + await makeTranscribeChannel("remote-chan", ["vidremote1"]); + + await page.goto("/channels/remote-chan"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + + // The job log proves the video was dispatched over the remote path. + await expect( + page.getByLabel("Transcribe missing output"), + ).toContainText("via remote", { timeout: 30_000 }); + + // And the transcript lands locally, pulled back from the remote. + await expect + .poll( + () => + pathExists( + "test-transcripts/channels/remote-chan/data/vidremote1/transcript.json", + ), + { timeout: 30_000 }, + ) + .toBe(true); +}); diff --git a/editor/e2e/workers.spec.ts b/editor/e2e/workers.spec.ts @@ -36,11 +36,12 @@ async function transcriptCount(slug: string, ids: string[]): Promise<number> { return n; } -// Two enabled local workers so disabling one still leaves the system running. +// Two enabled local workers (one slot each) so disabling one still leaves the +// system running. const TWO_WORKERS = { workers: [ - { id: "gpu", name: "GPU", kind: "local", enabled: true, priority: 0, slots: 1, appId: "whisper-cpp", config: {} }, - { id: "cpu", name: "CPU", kind: "local", enabled: true, priority: 1, slots: 2, appId: "whisper-cpp", config: {} }, + { id: "gpu", name: "GPU", kind: "local", enabled: true, priority: 0, appId: "whisper-cpp", config: {} }, + { id: "cpu", name: "CPU", kind: "local", enabled: true, priority: 1, appId: "whisper-cpp", config: {} }, ], }; @@ -48,13 +49,13 @@ test.beforeEach(async () => { await resetData("empty"); }); -test("lists configured workers with their slot counts", async ({ page }) => { +test("lists configured workers", async ({ page }) => { await writeSettings(TWO_WORKERS); await page.goto("/workers"); const gpu = page.getByRole("listitem").filter({ hasText: "GPU" }); const cpu = page.getByRole("listitem").filter({ hasText: "CPU" }); - await expect(gpu).toContainText("0/1 slots"); - await expect(cpu).toContainText("0/2 slots"); + await expect(gpu).toContainText("idle"); + await expect(cpu).toContainText("idle"); }); test("disable then enable a worker flips its state", async ({ page }) => { @@ -99,7 +100,7 @@ test("pausing all workers pauses a running batch instead of failing it; resume c // to observe. The fake whisper runs slowly for ids containing "slowop". await writeSettings({ workers: [ - { id: "only", name: "Only", kind: "local", enabled: true, priority: 0, slots: 1, appId: "whisper-cpp", config: {} }, + { id: "only", name: "Only", kind: "local", enabled: true, priority: 0, appId: "whisper-cpp", config: {} }, ], }); const ids = ["slowop1", "slowop2", "slowop3"]; diff --git a/editor/package.json b/editor/package.json @@ -5,8 +5,8 @@ "type": "module", "scripts": { "dev": "next dev --port ${EDITOR_PORT:-3001}", - "dev:test": "TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 next dev --port ${PORT:-3011}", - "start:test": "TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 next start --port ${PORT:-3011}", + "dev:test": "WORKER_TOKEN=test-worker-token TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 next dev --port ${PORT:-3011}", + "start:test": "WORKER_TOKEN=test-worker-token TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 next start --port ${PORT:-3011}", "build": "next build", "start": "next start --port ${EDITOR_PORT:-3001}", "lint": "eslint",