commit 92107aae467bfd3a2b8235e6a6c6b74486ca0e9f
parent f8206b9b317ce193b1f1809a245d1545e26bb964
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 13 Jun 2026 23:25:50 -0400
Transcription workers: data model + global pool scheduler (Phases 1–2)
Rework transcription from a single global engine into configurable workers.
Phase 1 — data model + migration:
- common/lib/workers.ts: Worker type (local/remote, priority, slots, per-worker
engine config or remote target), sanitizers, validation, and migration that
turns a pre-worker settings.json into a default enabled worker (slots = old
parallelTranscriptions) plus disabled workers for any other configured apps.
- settings.ts: add `workers` to SiteSettings; workers are the source of truth on
save, with transcriptionApp/transcriptionApps kept as a deprecated shadow for
rollback. Existing installs behave identically.
Phase 2 — global worker pool scheduler (local workers):
- common/jobs/workerPool.ts: process-wide singleton. Priority-ordered free-list
acquire/release (every enabled worker runs concurrently, highest priority
filled first), FIFO waiter parking, runtime enable/disable/drain, reconfigure
that never drops an in-use worker, and failure/degrade tracking.
- transcribeOne.ts: takes an explicit worker instead of reading global settings;
adds TranscribeError (transport vs transcription class) and a shared
transcribeWithWorker() that acquires a pool worker, retries across workers on
transport failures, and tracks the per-task progress.
- whisperBatch.ts: pool-driven (no pLimit). Worker slots are the sole throttle
for transcription parallelism; when all workers are disabled a batch pauses in
acquire() instead of failing, and resumes on re-enable.
- taskHooks/registry: per-task progress parser is chosen by the task's worker
engine; JobTask carries workerId.
- Single-video + inline-fallback transcribe callers route through the pool.
- Drop the now-redundant per-run Concurrency control from transcription actions
and UI (worker slots replace it); other batch ops keep their own concurrency.
Typechecks clean; whisper + transcription-app-migration e2e specs pass.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
16 files changed, 906 insertions(+), 253 deletions(-)
diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts
@@ -2,8 +2,10 @@ import path from "node:path";
import fs from "fs-extra";
import { execa } from "execa";
import type { Paths } from "../lib/paths";
-import { getSettings } from "../lib/settings";
+import type { Worker } from "../lib/workers";
import { getTranscriptionApp } from "../lib/transcriptionApps";
+import { getWorkerPool } from "../jobs/workerPool";
+import type { TaskTracker } from "../jobs/taskHooks";
import { normalizeTranscript } from "./normalizeTranscript";
import { isRealAudioFile } from "../lib/videoStatus";
@@ -33,12 +35,30 @@ export type TranscribeOneOptions = {
videoId: string;
audioFilename: string;
strictAudio?: boolean;
+ // The worker that runs this transcription. Callers obtain one from the global
+ // worker pool (getWorkerPool().acquire()). A local worker runs its engine here
+ // via execa; a remote worker delegates to another app instance.
+ worker: Worker;
onLog?: (msg: 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";
+ }
+}
+
export async function transcribeOneVideo(
opts: TranscribeOneOptions,
): Promise<TranscribeOneOutcome> {
@@ -67,38 +87,54 @@ export async function transcribeOneVideo(
log(`Using ${resolvedAudio} instead of ${opts.audioFilename}`);
}
- const settings = getSettings();
- const app = getTranscriptionApp(settings.transcriptionApp);
- const appConfig = settings.transcriptionApps[app.id] ?? {};
- const bin = appConfig.bin?.trim() || app.defaultBin();
- log(`Transcribe ${opts.videoId} start (${resolvedAudio}) via ${app.id}`);
+ const worker = opts.worker;
const start = Date.now();
// Each app writes "<outputBase><ext>" relative to cwd (the ext differs: whisper
// appends ".json", chough writes the exact -o path). Use a tmp base and rename
// the app-declared output file on success so a SIGTERM mid-write can't leave a
// half-baked transcript.json.
const tmpBase = `transcript.tmp-${process.pid}`;
- const build = app.build({
- audioFile: resolvedAudio,
- outputBase: tmpBase,
- config: appConfig,
- });
- const child = execa(bin, build.argv, {
- cwd: opts.videoDir,
- cancelSignal: opts.signal,
- all: true,
- buffer: false,
- env: build.env ? { ...process.env, ...build.env } : undefined,
- });
- child.all?.on("data", (c: Buffer) => log(c.toString("utf8")));
- await child;
- const producedPath = path.join(opts.videoDir, build.outputFile);
- if (!(await pathExists(producedPath))) {
- throw new Error(
- `transcription with ${app.id} produced no ${build.outputFile} in ${opts.videoDir}`,
+ 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",
);
+ } else {
+ const app = getTranscriptionApp(worker.appId);
+ const appConfig = worker.config ?? {};
+ const bin = appConfig.bin?.trim() || app.defaultBin();
+ log(
+ `Transcribe ${opts.videoId} start (${resolvedAudio}) via ${app.id} [${worker.id}]`,
+ );
+ const build = app.build({
+ audioFile: resolvedAudio,
+ outputBase: tmpBase,
+ config: appConfig,
+ });
+ const child = execa(bin, build.argv, {
+ cwd: opts.videoDir,
+ cancelSignal: opts.signal,
+ all: true,
+ buffer: false,
+ env: build.env ? { ...process.env, ...build.env } : undefined,
+ });
+ child.all?.on("data", (c: Buffer) => log(c.toString("utf8")));
+ await child;
+ const producedPath = path.join(opts.videoDir, build.outputFile);
+ if (!(await pathExists(producedPath))) {
+ throw new TranscribeError(
+ `transcription with ${app.id} produced no ${build.outputFile} in ${opts.videoDir}`,
+ "transcription",
+ );
+ }
+ await rename(producedPath, transcriptPath);
+ outputFormat = build.outputFormat;
}
- await rename(producedPath, transcriptPath);
log(
`Transcribe ${opts.videoId} done in ${((Date.now() - start) / 1000).toFixed(2)}s`,
);
@@ -109,7 +145,7 @@ export async function transcribeOneVideo(
await normalizeTranscript({
videoDir: opts.videoDir,
channelSlug,
- formatHint: build.outputFormat,
+ formatHint: outputFormat,
log,
});
} catch (err) {
@@ -128,3 +164,107 @@ function deriveChannelSlug(paths: Paths, videoDir: string): string | null {
const [slug] = rel.split(path.sep);
return slug || null;
}
+
+// How many workers a single video will try before giving up, when each attempt
+// fails with a transport-class error (remote unreachable, etc.). A
+// transcription-class failure stops immediately — retrying the same audio
+// elsewhere is pointless. Also bounds the loop if a worker keeps flapping.
+export const MAX_WORKER_ATTEMPTS = 3;
+
+export type TranscribeWithWorkerOptions = {
+ paths: Paths;
+ videoDir: string;
+ videoId: string;
+ audioFilename: string;
+ strictAudio?: boolean;
+ // Optional per-task progress tracker; when provided the acquired worker drives
+ // which progress parser is used and which worker the Workers page shows.
+ tracker?: TaskTracker;
+ taskId?: string;
+ taskLabel?: string;
+ onLog?: (msg: string) => void;
+ // Hard cancel. Also unblocks a parked acquire.
+ signal?: AbortSignal;
+ // Soft cancel (drain): when aborted, a parked acquire unblocks and this throws
+ // an AbortError so the caller skips the video without starting it.
+ drainSignal?: AbortSignal;
+};
+
+// Acquire a worker from the global pool and transcribe one video through it,
+// retrying on a different worker when a transport-class failure occurs (a remote
+// went away). This is the single entry point batches and single-video actions
+// share so pool accounting (lease release, success/failure marks) is consistent.
+// Throws TranscribeError on a transcription-class failure or exhausted retries,
+// and an AbortError when cancelled/drained while parked.
+export async function transcribeWithWorker(
+ opts: TranscribeWithWorkerOptions,
+): Promise<TranscribeOneOutcome> {
+ const log = opts.onLog ?? ((m: string) => console.log(m));
+ const pool = getWorkerPool();
+ // Unblock a parked acquire on either a hard cancel or a drain.
+ const acquireSignal =
+ opts.signal && opts.drainSignal
+ ? AbortSignal.any([opts.signal, opts.drainSignal])
+ : (opts.signal ?? opts.drainSignal);
+ let lastErr: unknown;
+ for (let attempt = 0; attempt < MAX_WORKER_ATTEMPTS; attempt++) {
+ if (opts.signal?.aborted || opts.drainSignal?.aborted) throw abortError();
+ let lease;
+ try {
+ lease = await pool.acquire(acquireSignal);
+ } catch {
+ throw abortError(); // cancelled or drained while parked
+ }
+ const worker = lease.worker;
+ const task = opts.tracker?.start({
+ id: opts.taskId ?? opts.videoId,
+ label: opts.taskLabel ?? opts.videoId,
+ kind: "transcribe",
+ appId: worker.appId,
+ workerId: worker.id,
+ });
+ try {
+ const outcome = await transcribeOneVideo({
+ paths: opts.paths,
+ videoDir: opts.videoDir,
+ videoId: opts.videoId,
+ audioFilename: opts.audioFilename,
+ strictAudio: opts.strictAudio,
+ worker,
+ onLog: task ? task.onLog : log,
+ signal: opts.signal,
+ });
+ pool.markSuccess(worker.id);
+ return outcome;
+ } catch (err) {
+ if (opts.signal?.aborted) throw err;
+ const failureClass =
+ err instanceof TranscribeError ? err.failureClass : "transcription";
+ if (failureClass === "transport") {
+ pool.markFailure(worker.id);
+ lastErr = err;
+ log(
+ `Worker ${worker.id} transport failure on ${opts.videoId}: ${String(err)} — retrying on another worker`,
+ );
+ continue;
+ }
+ throw err;
+ } finally {
+ task?.end();
+ lease.release();
+ }
+ }
+ throw (
+ lastErr ??
+ new TranscribeError(
+ `no worker could transcribe ${opts.videoId}`,
+ "transport",
+ )
+ );
+}
+
+function abortError(): Error {
+ const err = new Error("transcription aborted");
+ err.name = "AbortError";
+ return err;
+}
diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts
@@ -1,13 +1,12 @@
import path from "node:path";
import fs from "fs-extra";
-import pLimit from "p-limit";
import type { Paths } from "../lib/paths";
import type { AudioFormat } from "../lib/channelConfig";
import { isVideoTranscribed, readVideoFiles } from "../lib/videoStatus";
-import { transcribeOneVideo } from "./transcribeOne";
+import { transcribeWithWorker } from "./transcribeOne";
import { resolveShardItems } from "./shard";
import { countNotYetTranscribed } from "./channels";
-import { getSettings } from "../lib/settings";
+import { getWorkerPool } from "../jobs/workerPool";
import type { TaskTracker } from "../jobs/taskHooks";
import type { JobProgress } from "../jobs/registry";
@@ -16,7 +15,6 @@ const { pathExists, readdir, appendFile, readFile, ensureFile } = fs;
export type WhisperBatchOptions = {
channelSlug: string;
paths: Paths;
- concurrency?: number;
audioFilename?: string;
audioFormat?: AudioFormat;
strictAudioFormat?: boolean;
@@ -56,7 +54,6 @@ export type WhisperBatchResult = {
export async function runWhisperBatch({
channelSlug,
paths,
- concurrency,
audioFilename,
audioFormat,
strictAudioFormat = false,
@@ -76,7 +73,12 @@ export async function runWhisperBatch({
const channelDir = path.join(paths.channelsDir, channelSlug);
const dataDir = path.join(channelDir, "data");
const failureListFile = path.join(channelDir, "failed-transcriptions");
- const limit = pLimit(concurrency ?? getSettings().parallelTranscriptions);
+ // Refresh the global worker pool from current settings so this batch sees the
+ // latest worker config. Reconcile never drops an in-use worker, so this is safe
+ // even if another batch is somehow in flight. The pool's total free slots —
+ // not a pLimit — is what throttles this batch's parallelism.
+ const pool = getWorkerPool();
+ pool.reconfigure();
const resolvedAudioFilename =
audioFilename ?? (audioFormat ? `audio.${audioFormat}` : "audio.mp3");
const strict = strictAudioFormat && audioFormat !== undefined;
@@ -156,78 +158,75 @@ export async function runWhisperBatch({
let failed = 0;
let skipped = 0;
- await Promise.all(
- videoDirs.map((videoDir) =>
- limit(async () => {
- if (signal?.aborted) {
- skipped++;
- return;
- }
- // Drain (soft-cancel): don't start NEW transcriptions, but tasks that
- // already passed this gate keep running to completion.
- if (drainSignal?.aborted) {
- skipped++;
- return;
- }
- const videoPath = path.join(dataDir, videoDir);
- if (failedSet.has(videoDir)) {
- log(`Skipping previously failed transcription for ${videoDir}`);
- skipped++;
- return;
- }
- // A transcript already exists if either whisper has run
- // (transcript.json) OR yt-dlp wrote auto/manual English subs
- // (transcript.en.vtt, or a regional/auto fallback like en-US). Without
- // the VTT branch, videos on a youtube-handling channel with only VTT
- // auto-subs end up attempted -> fail with "no audio file found" -> get
- // added to failed-transcriptions on every run.
- if (isVideoTranscribed(await readVideoFiles(videoPath))) {
- log(`Transcription for ${videoDir} already exists`);
- skipped++;
- return;
- }
- if (
- strict &&
- !(await pathExists(path.join(videoPath, resolvedAudioFilename)))
- ) {
- log(
- `Skipping ${videoDir}: no ${resolvedAudioFilename} (strict mode)`,
- );
- skipped++;
- return;
- }
- attempted++;
- const task = tracker?.start({
- id: videoDir,
- label: videoDir,
- kind: "transcribe",
- });
- try {
- await transcribeOneVideo({
- paths,
- videoDir: videoPath,
- videoId: videoDir,
- audioFilename: resolvedAudioFilename,
- strictAudio: strict,
- onLog: task ? task.onLog : log,
- signal,
- });
- succeededCount++;
- } catch (err) {
- if (signal?.aborted) {
- log(`Transcribe ${videoDir} cancelled`);
- skipped++;
- return;
- }
- log(`FAILED TO TRANSCRIBE ${videoDir}: ${String(err)}`);
- await appendFile(failureListFile, `${videoDir}\n`);
- failed++;
- } finally {
- task?.end();
- }
- }),
- ),
- );
+ // No pLimit: every video starts concurrently and parks in pool.acquire() until
+ // a worker slot frees. The pool's total free slots throttle the batch, and a
+ // freed high-priority slot is handed to the oldest waiter. When all workers are
+ // disabled, acquire() blocks indefinitely — the batch pauses (it does not fail)
+ // and resumes when a worker is re-enabled.
+ const runOne = async (videoDir: string): Promise<void> => {
+ if (signal?.aborted || drainSignal?.aborted) {
+ skipped++;
+ return;
+ }
+ const videoPath = path.join(dataDir, videoDir);
+ if (failedSet.has(videoDir)) {
+ log(`Skipping previously failed transcription for ${videoDir}`);
+ skipped++;
+ return;
+ }
+ // A transcript already exists if either whisper has run (transcript.json) OR
+ // yt-dlp wrote auto/manual English subs (transcript.en.vtt, or a regional/auto
+ // fallback like en-US). Without the VTT branch, videos on a youtube-handling
+ // channel with only VTT auto-subs end up attempted -> fail with "no audio file
+ // found" -> get added to failed-transcriptions on every run.
+ if (isVideoTranscribed(await readVideoFiles(videoPath))) {
+ log(`Transcription for ${videoDir} already exists`);
+ skipped++;
+ return;
+ }
+ if (strict && !(await pathExists(path.join(videoPath, resolvedAudioFilename)))) {
+ log(`Skipping ${videoDir}: no ${resolvedAudioFilename} (strict mode)`);
+ skipped++;
+ return;
+ }
+
+ attempted++;
+ try {
+ // transcribeWithWorker acquires a pool worker (parking until one frees, or
+ // forever while all are disabled), retries across workers on transport
+ // failures, and tracks the per-video task.
+ await transcribeWithWorker({
+ paths,
+ videoDir: videoPath,
+ videoId: videoDir,
+ audioFilename: resolvedAudioFilename,
+ strictAudio: strict,
+ tracker,
+ taskId: videoDir,
+ taskLabel: videoDir,
+ onLog: log,
+ signal,
+ drainSignal,
+ });
+ succeededCount++;
+ } catch (err) {
+ // Cancelled or drained while parked/in-flight: a skip, not a failure.
+ if (
+ signal?.aborted ||
+ drainSignal?.aborted ||
+ (err as Error)?.name === "AbortError"
+ ) {
+ log(`Transcribe ${videoDir} cancelled`);
+ skipped++;
+ return;
+ }
+ log(`FAILED TO TRANSCRIBE ${videoDir}: ${String(err)}`);
+ await appendFile(failureListFile, `${videoDir}\n`);
+ failed++;
+ }
+ };
+
+ await Promise.all(videoDirs.map(runOne));
return { attempted, succeeded: succeededCount, failed, skipped };
}
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -28,6 +28,9 @@ export type JobTask = {
fraction?: number; // 0..1, undefined until first parseable progress
detail?: string; // e.g. "2.31MiB/s ETA 00:03:12" or "01:50:02 / 03:24:47"
startedAt: number;
+ // For transcribe tasks: the worker running it (id into the worker pool), so
+ // the Workers page can show which engine each task is on.
+ workerId?: string;
};
export type JobRecord = {
diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts
@@ -1,7 +1,6 @@
import type { JobTaskKind } from "./registry";
import type { JobRunContext } from "./streamCommand";
import { parseDownloadProgress } from "./progressParsers";
-import { getSettings } from "../lib/settings";
import { getTranscriptionApp } from "../lib/transcriptionApps";
// A handle for one in-flight sub-operation. `onLog` is a drop-in replacement
@@ -15,7 +14,17 @@ export type TaskHandle = {
};
export type TaskTracker = {
- start: (init: { id: string; label: string; kind: JobTaskKind }) => TaskHandle;
+ start: (init: {
+ id: string;
+ label: string;
+ kind: JobTaskKind;
+ // For transcribe tasks: the engine (app id) running this task, used to pick
+ // the right progress parser. A remote worker streams pre-parsed progress, so
+ // it passes no appId and the tracker uses a pass-through.
+ appId?: string;
+ // The worker that owns this task, surfaced on the Workers page.
+ workerId?: string;
+ }) => TaskHandle;
};
type TaskCtx = Pick<
@@ -31,17 +40,20 @@ export function makeTaskTracker(
forwardLog: (line: string) => void,
): TaskTracker {
return {
- start({ id, label, kind }) {
+ start({ id, label, kind, appId, workerId }) {
if (!ctx) {
return { onLog: forwardLog, end: () => {} };
}
const startedAt = Date.now();
- ctx.addTask({ id, label, kind, startedAt });
+ ctx.addTask({ id, label, kind, startedAt, workerId });
// Transcription progress output is app-specific (whisper's segment
- // timestamps vs chough's ETA bars), so pick the active app's parser.
+ // timestamps vs chough's ETA bars), so pick the parser for THIS task's
+ // worker engine. A transcribe task with no appId (e.g. a remote worker that
+ // streams pre-parsed progress) falls through to the download parser, which
+ // is harmless for non-matching lines.
const transcribeParser =
- kind === "transcribe"
- ? getTranscriptionApp(getSettings().transcriptionApp).makeProgressParser()
+ kind === "transcribe" && appId
+ ? getTranscriptionApp(appId).makeProgressParser()
: null;
let ended = false;
const onLog = (line: string) => {
diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts
@@ -0,0 +1,310 @@
+// Global transcription worker pool — the scheduler that distributes individual
+// transcriptions across configured workers.
+//
+// Like the JobRegistry (registry.ts), this is a process-wide singleton held on
+// globalThis so every request/server-action shares ONE view of which worker has
+// a free slot. That global-ness is the whole point: "the GPU worker has exactly
+// 1 slot across everything" can only be enforced in one shared place. A batch
+// (runWhisperBatch) no longer owns a pLimit; it acquires a lease per video from
+// this pool and releases it when the transcription finishes.
+//
+// Scheduling is a priority-ordered free-list (decided with the user): every
+// enabled worker runs concurrently, and when several have a free slot the
+// highest-priority one is handed out first. When NO worker is eligible (all
+// disabled/draining/full), acquire() simply parks on a FIFO queue — so a batch
+// pauses, neither failing nor busy-looping, and resumes the instant a worker is
+// re-enabled.
+
+import type { Worker } from "../lib/workers";
+import { getSettings } from "../lib/settings";
+
+// Consecutive transport/exec failures before a worker is auto-marked degraded
+// (skipped by the scheduler until re-enabled). See markFailure / Phase 5.
+export const WORKER_DEGRADE_THRESHOLD = 3;
+
+export type WorkerRuntimeState = "enabled" | "draining" | "disabled";
+
+// A granted slot. release() is idempotent; the controller calls it in a finally.
+export type Lease = {
+ worker: Worker;
+ release: () => void;
+};
+
+// Client-safe snapshot of one worker for the Workers page / active-jobs payload.
+export type WorkerSummary = {
+ id: string;
+ name: string;
+ kind: Worker["kind"];
+ appId?: string;
+ priority: number;
+ slots: number;
+ inUse: number;
+ state: WorkerRuntimeState;
+ degraded: boolean;
+ enabled: boolean; // persisted intent (config.enabled)
+};
+
+type PoolEntry = {
+ config: Worker;
+ inUse: number;
+ state: WorkerRuntimeState;
+ degraded: boolean;
+ consecutiveFailures: number;
+ // Set when a worker is removed from settings while still busy: drain then drop.
+ retireWhenIdle: boolean;
+};
+
+type Waiter = {
+ resolve: (lease: Lease) => void;
+ reject: (err: Error) => void;
+ onAbort?: () => void;
+};
+
+class WorkerPool {
+ // Insertion order is the tiebreak for equal priority, so use a Map (ordered).
+ private entries = new Map<string, PoolEntry>();
+ private waiters: Waiter[] = [];
+ private initialized = false;
+
+ // Lazily seed from settings on first use. Idempotent.
+ private ensureInit(): void {
+ if (this.initialized) return;
+ this.initialized = true;
+ this.reconfigure(getSettings().workers);
+ }
+
+ // Re-sync the pool with a worker list (defaults to current settings). Updates
+ // config in place, adds new workers, and retires removed ones — but NEVER drops
+ // a worker that still has in-flight leases (that would oversubscribe hardware
+ // the moment it's re-added). A busy removed worker is marked retireWhenIdle and
+ // dropped when its last lease releases.
+ reconfigure(workers?: Worker[]): void {
+ this.initialized = true;
+ const next = workers ?? getSettings().workers;
+ const seen = new Set<string>();
+ for (const w of next) {
+ seen.add(w.id);
+ const existing = this.entries.get(w.id);
+ if (existing) {
+ existing.config = w;
+ existing.retireWhenIdle = false;
+ // Reconcile runtime state with persisted intent. Disabling a busy worker
+ // drains it; (re-)enabling clears degraded.
+ if (w.enabled) {
+ existing.state = "enabled";
+ existing.degraded = false;
+ existing.consecutiveFailures = 0;
+ } else if (existing.inUse > 0) {
+ existing.state = "draining";
+ } else {
+ existing.state = "disabled";
+ }
+ } else {
+ this.entries.set(w.id, {
+ config: w,
+ inUse: 0,
+ state: w.enabled ? "enabled" : "disabled",
+ degraded: false,
+ consecutiveFailures: 0,
+ retireWhenIdle: false,
+ });
+ }
+ }
+ // 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) {
+ entry.retireWhenIdle = true;
+ entry.state = "draining";
+ } else {
+ this.entries.delete(id);
+ }
+ }
+ // A reconfigure can make new slots eligible (worker enabled/added).
+ this.pump();
+ }
+
+ private eligible(entry: PoolEntry): boolean {
+ return (
+ entry.state === "enabled" &&
+ !entry.degraded &&
+ entry.inUse < entry.config.slots
+ );
+ }
+
+ // Highest-priority free entry: (priority asc, insertion order asc).
+ private pickFree(): PoolEntry | null {
+ let best: PoolEntry | null = null;
+ for (const entry of this.entries.values()) {
+ if (!this.eligible(entry)) continue;
+ if (!best || entry.config.priority < best.config.priority) best = entry;
+ }
+ return best;
+ }
+
+ private grant(entry: PoolEntry): Lease {
+ entry.inUse++;
+ 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";
+ }
+ }
+ this.pump();
+ };
+ return { worker: entry.config, release };
+ }
+
+ // Try to satisfy parked waiters in FIFO order while free slots remain.
+ private pump(): void {
+ while (this.waiters.length > 0) {
+ const entry = this.pickFree();
+ if (!entry) break;
+ const waiter = this.waiters.shift()!;
+ waiter.onAbort?.();
+ waiter.resolve(this.grant(entry));
+ }
+ }
+
+ // Acquire a lease on the highest-priority free worker. Resolves immediately if
+ // one is free, otherwise parks until a slot frees or a worker is enabled. If
+ // `signal` aborts while parked (or before), rejects with an AbortError so the
+ // caller treats it like a cancel.
+ acquire(signal?: AbortSignal): Promise<Lease> {
+ this.ensureInit();
+ if (signal?.aborted) {
+ return Promise.reject(abortError());
+ }
+ const entry = this.pickFree();
+ if (entry && this.waiters.length === 0) {
+ return Promise.resolve(this.grant(entry));
+ }
+ return new Promise<Lease>((resolve, reject) => {
+ const waiter: Waiter = { resolve, reject };
+ if (signal) {
+ const onAbort = () => {
+ const idx = this.waiters.indexOf(waiter);
+ if (idx >= 0) this.waiters.splice(idx, 1);
+ reject(abortError());
+ };
+ waiter.onAbort = () => signal.removeEventListener("abort", onAbort);
+ signal.addEventListener("abort", onAbort, { once: true });
+ }
+ this.waiters.push(waiter);
+ // A new arrival can't create a free slot, but pump keeps FIFO honest if one
+ // freed between pickFree and here.
+ this.pump();
+ });
+ }
+
+ // --- Runtime controls (Workers page). These mutate the live pool; callers
+ // also persist the new `enabled` flag to settings so it survives a restart. ---
+
+ enableWorker(id: string): boolean {
+ const entry = this.entries.get(id);
+ if (!entry) return false;
+ entry.state = "enabled";
+ entry.degraded = false;
+ entry.consecutiveFailures = 0;
+ entry.config = { ...entry.config, enabled: true };
+ this.pump();
+ return true;
+ }
+
+ // Hard disable: stop granting new leases immediately. In-flight leases keep
+ // running until their controllers release (a true stop also hard-cancels the
+ // job). State goes straight to disabled even if busy.
+ disableWorker(id: string): boolean {
+ const entry = this.entries.get(id);
+ if (!entry) return false;
+ entry.state = "disabled";
+ entry.config = { ...entry.config, enabled: false };
+ return true;
+ }
+
+ // Drain: stop granting NEW leases but let in-flight transcriptions finish;
+ // flip to disabled once the last lease releases. The "free up the GPU" path.
+ drainWorker(id: string): boolean {
+ const entry = this.entries.get(id);
+ if (!entry) return false;
+ entry.config = { ...entry.config, enabled: false };
+ entry.state = entry.inUse > 0 ? "draining" : "disabled";
+ return true;
+ }
+
+ // --- Failure tracking (auto-disable, Phase 5) ---
+
+ markSuccess(id: string): void {
+ const entry = this.entries.get(id);
+ if (entry) entry.consecutiveFailures = 0;
+ }
+
+ // Returns true if this failure pushed the worker over the degrade threshold.
+ markFailure(id: string): boolean {
+ const entry = this.entries.get(id);
+ if (!entry) return false;
+ entry.consecutiveFailures++;
+ if (
+ !entry.degraded &&
+ entry.consecutiveFailures >= WORKER_DEGRADE_THRESHOLD
+ ) {
+ entry.degraded = true;
+ return true;
+ }
+ return false;
+ }
+
+ // True when at least one worker is enabled and not degraded (ignoring momentary
+ // fullness — a full but enabled pool is working at capacity, not paused). Its
+ // negation is what surfaces the "paused — waiting for a worker" state: a running
+ // batch whose every worker is disabled/draining/degraded, so nothing can start.
+ hasEligibleWorker(): boolean {
+ this.ensureInit();
+ for (const entry of this.entries.values()) {
+ if (entry.state === "enabled" && !entry.degraded) return true;
+ }
+ return false;
+ }
+
+ summary(): WorkerSummary[] {
+ this.ensureInit();
+ return Array.from(this.entries.values())
+ .map((e) => ({
+ id: e.config.id,
+ name: e.config.name,
+ kind: e.config.kind,
+ appId: e.config.appId,
+ priority: e.config.priority,
+ slots: e.config.slots,
+ inUse: e.inUse,
+ state: e.state,
+ degraded: e.degraded,
+ enabled: e.config.enabled,
+ }))
+ .sort((a, b) => a.priority - b.priority || a.id.localeCompare(b.id));
+ }
+}
+
+function abortError(): Error {
+ const err = new Error("worker acquire aborted");
+ err.name = "AbortError";
+ return err;
+}
+
+declare global {
+ // eslint-disable-next-line no-var
+ var __yttWorkerPool__: WorkerPool | undefined;
+}
+
+export function getWorkerPool(): WorkerPool {
+ if (!globalThis.__yttWorkerPool__) {
+ globalThis.__yttWorkerPool__ = new WorkerPool();
+ }
+ return globalThis.__yttWorkerPool__;
+}
diff --git a/common/lib/settings.ts b/common/lib/settings.ts
@@ -5,10 +5,17 @@ import {
type AppInstanceConfig,
DEFAULT_TRANSCRIBE_ARGS,
DEFAULT_TRANSCRIPTION_APP_ID,
- getTranscriptionApp,
TRANSCRIPTION_APPS,
- validateTranscribeArgs,
} from "./transcriptionApps";
+import {
+ type Worker,
+ defaultWorkersFromApps,
+ sanitizeWorkerConfig,
+ sanitizeWorkers,
+ validateWorkers,
+} from "./workers";
+
+export type { Worker } from "./workers";
// Transcribe placeholder/arg helpers now live with the whisper-cpp app in
// transcriptionApps.ts. Re-exported here so existing import sites keep working.
@@ -35,8 +42,15 @@ export type SiteSettings = {
// or "chough"). Selected globally; see common/lib/transcriptionApps.ts.
transcriptionApp: string;
// Per-app configuration, keyed by app id. Each app reads only its own block;
- // a missing block means "use the app's defaults".
+ // a missing block means "use the app's defaults". DEPRECATED in favor of
+ // `workers` (each local worker carries its own config); kept one release to
+ // drive migration and allow rollback. See common/lib/workers.ts.
transcriptionApps: Record<string, AppInstanceConfig>;
+ // Configured transcription workers (named processing slots). The scheduler
+ // distributes each video to the highest-priority free worker. A settings.json
+ // predating this field is migrated to a single enabled worker from the active
+ // app (see defaultWorkersFromApps). See common/lib/workers.ts.
+ workers: Worker[];
// Browser spec (e.g. "firefox", "chrome:Default") passed to
// `yt-dlp --cookies-from-browser` ONLY on the auth-retry attempt of the
// managed per-video download flow. Empty string = disabled.
@@ -131,6 +145,7 @@ function defaults(): SiteSettings {
maxTranscriptPageBytes: TRANSCRIPT_PAGE_DEFAULT_BYTES,
transcriptionApp: DEFAULT_TRANSCRIPTION_APP_ID,
transcriptionApps: {},
+ workers: [],
cookiesFromBrowser: "",
sleepBetweenDownloadsSeconds: SLEEP_BETWEEN_DOWNLOADS_DEFAULT_SECONDS,
parallelTranscriptions: PARALLEL_TRANSCRIPTIONS_DEFAULT,
@@ -288,6 +303,19 @@ export function getSettings(): SiteSettings {
merged.autoRefreshIntervalSeconds,
);
merged.socialLinks = parseSocialLinks(merged.socialLinks);
+ // Workers. When the file predates the worker model (no `workers` key),
+ // synthesize a default list from the (now-settled) active app + per-app
+ // configs so existing installs behave identically. Otherwise sanitize the
+ // stored list.
+ if (parsed.workers === undefined) {
+ merged.workers = defaultWorkersFromApps(
+ merged.transcriptionApp,
+ merged.transcriptionApps,
+ merged.parallelTranscriptions,
+ );
+ } else {
+ merged.workers = sanitizeWorkers(merged.workers);
+ }
return merged;
}
@@ -316,16 +344,7 @@ export function sanitizeTranscriptionApps(
const out: Record<string, AppInstanceConfig> = {};
for (const [id, raw] of Object.entries(value as Record<string, unknown>)) {
if (!raw || typeof raw !== "object") continue;
- const r = raw as Record<string, unknown>;
- const cfg: AppInstanceConfig = {};
- if (typeof r.bin === "string") cfg.bin = r.bin;
- if (typeof r.model === "string") cfg.model = r.model;
- if (typeof r.remoteUrl === "string") cfg.remoteUrl = r.remoteUrl;
- if (typeof r.chunkSize === "number" && Number.isFinite(r.chunkSize)) {
- cfg.chunkSize = Math.floor(r.chunkSize);
- }
- if (Array.isArray(r.customArgs)) cfg.customArgs = r.customArgs.map(String);
- out[id] = cfg;
+ out[id] = sanitizeWorkerConfig(raw);
}
return out;
}
@@ -374,25 +393,41 @@ function migrateLegacyTranscription(
}
export async function writeSettings(next: SiteSettings): Promise<void> {
+ // Workers are the source of truth. A caller that still sets only the
+ // deprecated transcriptionApp/transcriptionApps (no `workers`) gets a list
+ // synthesized from them, so old call sites keep working during the transition.
+ let workers = sanitizeWorkers(next.workers);
+ if (workers.length === 0) {
+ const fallbackAppId =
+ typeof next.transcriptionApp === "string" &&
+ TRANSCRIPTION_APPS[next.transcriptionApp]
+ ? next.transcriptionApp
+ : DEFAULT_TRANSCRIPTION_APP_ID;
+ workers = defaultWorkersFromApps(
+ fallbackAppId,
+ sanitizeTranscriptionApps(next.transcriptionApps),
+ clampParallelTranscriptions(next.parallelTranscriptions),
+ );
+ }
+ const workersErr = validateWorkers(workers);
+ if (workersErr) throw new Error(workersErr);
+
+ // Keep the deprecated transcriptionApp/transcriptionApps as a faithful shadow
+ // of the workers (used only for rollback to a pre-worker build — a file with a
+ // `workers` key is never re-migrated). The active app is the first enabled
+ // local worker; the per-app map mirrors each local worker's config.
+ const firstLocal =
+ workers.find((w) => w.enabled && w.kind === "local") ??
+ workers.find((w) => w.kind === "local");
const appId =
- typeof next.transcriptionApp === "string" &&
- TRANSCRIPTION_APPS[next.transcriptionApp]
- ? next.transcriptionApp
+ firstLocal?.appId && TRANSCRIPTION_APPS[firstLocal.appId]
+ ? firstLocal.appId
: DEFAULT_TRANSCRIPTION_APP_ID;
- const app = getTranscriptionApp(appId);
const transcriptionApps = sanitizeTranscriptionApps(next.transcriptionApps);
- const activeCfg = transcriptionApps[appId] ?? {};
- const resolvedBin = activeCfg.bin?.trim() || app.defaultBin();
- if (!resolvedBin) {
- throw new Error(`Transcribe binary is required for ${app.label}`);
- }
- if (
- app.fields.customArgs &&
- activeCfg.customArgs &&
- activeCfg.customArgs.length > 0
- ) {
- const argsErr = validateTranscribeArgs(activeCfg.customArgs);
- if (argsErr) throw new Error(argsErr);
+ for (const w of workers) {
+ if (w.kind === "local" && w.appId) {
+ transcriptionApps[w.appId] = w.config ?? {};
+ }
}
const socialLinks: SocialLink[] = [];
for (const link of parseSocialLinks(next.socialLinks)) {
@@ -415,6 +450,7 @@ export async function writeSettings(next: SiteSettings): Promise<void> {
maxTranscriptPageBytes: clampPageBytes(next.maxTranscriptPageBytes),
transcriptionApp: appId,
transcriptionApps,
+ workers,
cookiesFromBrowser:
typeof next.cookiesFromBrowser === "string"
? next.cookiesFromBrowser.trim()
diff --git a/common/lib/workers.ts b/common/lib/workers.ts
@@ -0,0 +1,238 @@
+// Transcription workers: named, independently-configured processing slots.
+//
+// A worker is one engine the scheduler can hand a video to. The three apps in
+// transcriptionApps.ts (whisper-cpp/chough/parakeet) are the *types* of a LOCAL
+// worker — a local worker is just an (appId, AppInstanceConfig) pair promoted to
+// a named, repeatable slot. A REMOTE worker delegates the transcription to the
+// same app running on another machine on the LAN (see common/controller/
+// remoteTranscribe.ts). Workers carry a priority and a slot count; the global
+// scheduler (common/jobs/workerPool.ts) distributes each video — per task — to
+// the highest-priority free worker.
+//
+// Pure definitions + sanitizers only. Dependency direction is one-way:
+// settings.ts -> workers.ts -> transcriptionApps.ts. Keep it that way.
+
+import {
+ type AppInstanceConfig,
+ DEFAULT_TRANSCRIPTION_APP_ID,
+ TRANSCRIPTION_APPS,
+ getTranscriptionApp,
+ validateTranscribeArgs,
+} from "./transcriptionApps";
+
+export type WorkerKind = "local" | "remote";
+
+// Where a remote worker delegates. The remote runs its OWN worker pool and picks
+// among ITS local workers — so the primary stores only how to reach it, not which
+// engine to use.
+export type RemoteWorkerConfig = {
+ baseUrl: string; // e.g. http://gpu-box.lan:3011
+ // Outbound bearer token sent with every /api/worker request to this remote.
+ // The accepting side validates against its own WORKER_TOKEN env, never this.
+ token?: string;
+ // When true the remote shares the transcripts mount, so we send
+ // {channelSlug, videoId} instead of uploading the audio bytes.
+ sharedFs?: boolean;
+};
+
+export type Worker = {
+ // Stable slug; used in settings, task ids, and logs.
+ id: string;
+ // Human label shown in the UI.
+ name: string;
+ kind: WorkerKind;
+ 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[];
+
+ // LOCAL: an instance of a TRANSCRIPTION_APPS entry + its per-worker config.
+ appId?: string;
+ config?: AppInstanceConfig;
+
+ // REMOTE: how to reach the delegate instance.
+ 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 {
+ const cfg: AppInstanceConfig = {};
+ if (!value || typeof value !== "object") return cfg;
+ const r = value as Record<string, unknown>;
+ if (typeof r.bin === "string") cfg.bin = r.bin;
+ if (typeof r.model === "string") cfg.model = r.model;
+ if (typeof r.remoteUrl === "string") cfg.remoteUrl = r.remoteUrl;
+ if (typeof r.chunkSize === "number" && Number.isFinite(r.chunkSize)) {
+ cfg.chunkSize = Math.floor(r.chunkSize);
+ }
+ if (Array.isArray(r.customArgs)) cfg.customArgs = r.customArgs.map(String);
+ return cfg;
+}
+
+function sanitizeRemoteConfig(value: unknown): RemoteWorkerConfig {
+ const cfg: RemoteWorkerConfig = { baseUrl: "" };
+ if (!value || typeof value !== "object") return cfg;
+ const r = value as Record<string, unknown>;
+ if (typeof r.baseUrl === "string") cfg.baseUrl = r.baseUrl.trim();
+ if (typeof r.token === "string" && r.token.trim()) cfg.token = r.token.trim();
+ if (r.sharedFs === true) cfg.sharedFs = true;
+ return cfg;
+}
+
+function slugify(input: string): string {
+ return input
+ .toLowerCase()
+ .replace(/[^a-z0-9]+/g, "-")
+ .replace(/^-+|-+$/g, "");
+}
+
+function uniqueId(base: string, seen: Set<string>): string {
+ let id = slugify(base) || "worker";
+ if (!seen.has(id)) return id;
+ let n = 2;
+ while (seen.has(`${id}-${n}`)) n++;
+ return `${id}-${n}`;
+}
+
+// Coerce a raw settings.workers value into a clean Worker[]. Generates stable
+// ids, fills defaults, and guarantees unique ids. Validation of "can this worker
+// actually run" (resolvable bin, valid baseUrl) is separate — see validateWorkers.
+export function sanitizeWorkers(value: unknown): Worker[] {
+ if (!Array.isArray(value)) return [];
+ const out: Worker[] = [];
+ const seenIds = new Set<string>();
+ value.forEach((raw, index) => {
+ if (!raw || typeof raw !== "object") return;
+ const r = raw as Record<string, unknown>;
+ const kind: WorkerKind = r.kind === "remote" ? "remote" : "local";
+ const rawId = typeof r.id === "string" ? r.id.trim() : "";
+ const rawName = typeof r.name === "string" ? r.name.trim() : "";
+ const id = uniqueId(rawId || rawName || `worker-${index + 1}`, seenIds);
+ seenIds.add(id);
+ const worker: Worker = {
+ id,
+ name: rawName || id,
+ kind,
+ enabled: r.enabled !== false,
+ priority:
+ 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");
+ if (tags.length > 0) worker.tags = tags;
+ }
+ if (kind === "local") {
+ worker.appId =
+ typeof r.appId === "string" && TRANSCRIPTION_APPS[r.appId]
+ ? r.appId
+ : DEFAULT_TRANSCRIPTION_APP_ID;
+ worker.config = sanitizeWorkerConfig(r.config);
+ } else {
+ worker.remote = sanitizeRemoteConfig(r.remote);
+ }
+ out.push(worker);
+ });
+ return out;
+}
+
+// Validate a worker list for saving. Returns an error string or null. Enforces:
+// at least one enabled worker, each enabled local worker has a resolvable binary
+// and (if it surfaces customArgs) a valid template, each remote has a valid URL.
+export function validateWorkers(workers: Worker[]): string | null {
+ if (!workers.some((w) => w.enabled)) {
+ return "At least one worker must be enabled";
+ }
+ for (const w of workers) {
+ const label = w.name || w.id;
+ if (w.kind === "local") {
+ if (!w.appId || !TRANSCRIPTION_APPS[w.appId]) {
+ return `Worker "${label}" has an unknown engine`;
+ }
+ if (!w.enabled) continue;
+ const app = getTranscriptionApp(w.appId);
+ const cfg = w.config ?? {};
+ const bin = cfg.bin?.trim() || app.defaultBin();
+ if (!bin) return `Worker "${label}" needs a binary`;
+ if (app.fields.customArgs && cfg.customArgs && cfg.customArgs.length > 0) {
+ const err = validateTranscribeArgs(cfg.customArgs);
+ if (err) return err;
+ }
+ } else {
+ const url = w.remote?.baseUrl?.trim();
+ if (!url) return `Remote worker "${label}" needs a base URL`;
+ try {
+ const u = new URL(url);
+ if (u.protocol !== "http:" && u.protocol !== "https:") {
+ return `Remote worker "${label}" base URL must be http(s)`;
+ }
+ } catch {
+ return `Remote worker "${label}" has an invalid base URL`;
+ }
+ }
+ }
+ return 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.
+export function defaultWorkersFromApps(
+ activeAppId: string,
+ appConfigs: Record<string, AppInstanceConfig>,
+ parallelSlots: number,
+): Worker[] {
+ const activeApp = getTranscriptionApp(activeAppId);
+ const workers: Worker[] = [
+ {
+ id: "default",
+ name: 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;
+ seen.add(id);
+ workers.push({
+ id,
+ name: TRANSCRIPTION_APPS[appId].label,
+ kind: "local",
+ enabled: false,
+ priority: priority++,
+ slots: 1,
+ appId,
+ config,
+ });
+ }
+ return workers;
+}
diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts
@@ -20,7 +20,7 @@ import { writeDownloadOutcome } from "../lib/downloadOutcome-server";
import { loadRawMetadata } from "../lib/transcripts-server";
import { evaluateDownloadFilters } from "../lib/downloadFilters";
import type { Paths } from "../lib/paths";
-import { transcribeOneVideo } from "../controller/transcribeOne";
+import { transcribeWithWorker } from "../controller/transcribeOne";
import { extractVideoId, outputArgsForUrl } from "./runYtdlp";
import { runAudioCheckedYtdlp } from "./audioCheckedDownload";
import { DOWNLOAD_PROGRESS_TEMPLATE } from "../jobs/progressParsers";
@@ -605,9 +605,11 @@ async function runManagedDownload(
if (attemptSucceeded(fallbackRes.exitCode)) {
fellBackToTranscribe = true;
if (opts.inlineTranscribeOnFallback) {
- // Inline whisper: matches whisperVideoAction's shape.
+ // Inline whisper: matches whisperVideoAction's shape. Routes through
+ // the worker pool so the inline transcription respects worker config
+ // and slot limits like any other.
try {
- await transcribeOneVideo({
+ await transcribeWithWorker({
paths: opts.paths,
videoDir,
videoId,
diff --git a/editor/app/actionable/components/InlineActionButton.tsx b/editor/app/actionable/components/InlineActionButton.tsx
@@ -51,7 +51,6 @@ async function runAction(variant: Variant): Promise<StreamActionResult> {
variant.slug,
undefined,
undefined,
- undefined,
variant.audioFormat,
);
}
diff --git a/editor/app/channels/[slug]/bulkVideoActions.ts b/editor/app/channels/[slug]/bulkVideoActions.ts
@@ -22,9 +22,8 @@ export async function bulkTranscribeAction(
slug: string,
videoIds: string[],
queueKey?: string,
- concurrency?: number,
): Promise<StreamActionResult> {
- return transcribeBucketAction(slug, videoIds, queueKey, concurrency);
+ return transcribeBucketAction(slug, videoIds, queueKey);
}
export async function bulkRetryDownloadAction(
diff --git a/editor/app/channels/[slug]/components/VideoListPane.tsx b/editor/app/channels/[slug]/components/VideoListPane.tsx
@@ -6,7 +6,6 @@ import { useMemo, useState, useTransition } from "react";
import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand";
import type { VideoRow, VideoFilter } from "../lib/videoRows";
import { filterRows } from "../lib/videoRows";
-import { ConcurrencyControl } from "../../../components/ConcurrencyControl";
import { QueueControl } from "../../../components/QueueControl";
import {
bulkMarkUntranscribableAction,
@@ -24,17 +23,8 @@ type Props = {
defaultTranscribeQueue: string;
defaultDownloadQueue: string;
existingQueues: string[];
- defaultConcurrency: number;
};
-function parseConcurrency(s: string): number | undefined {
- const trimmed = s.trim();
- if (trimmed === "") return undefined;
- const n = Number(trimmed);
- if (!Number.isFinite(n) || n < 1) return undefined;
- return Math.floor(n);
-}
-
const FILTER_OPTIONS: { value: VideoFilter; label: string }[] = [
{ value: "all", label: "All" },
{ value: "needs_action", label: "Needs action" },
@@ -69,7 +59,6 @@ export function VideoListPane({
defaultTranscribeQueue,
defaultDownloadQueue,
existingQueues,
- defaultConcurrency,
}: Props) {
const router = useRouter();
const searchParams = useSearchParams();
@@ -80,7 +69,6 @@ export function VideoListPane({
const [lastResult, setLastResult] = useState<BulkActionSummary | null>(null);
const [error, setError] = useState<string | null>(null);
const [transcribeQueue, setTranscribeQueue] = useState(defaultTranscribeQueue);
- const [transcribeConcurrency, setTranscribeConcurrency] = useState("");
const [retryQueue, setRetryQueue] = useState(defaultDownloadQueue);
const [retryAbortOnError, setRetryAbortOnError] = useState(false);
@@ -311,23 +299,12 @@ export function VideoListPane({
existingQueues={existingQueues}
actionLabel="bulk transcribe"
/>
- <ConcurrencyControl
- value={transcribeConcurrency}
- onChange={setTranscribeConcurrency}
- defaultLimit={defaultConcurrency}
- actionLabel="bulk transcribe"
- />
<button
type="button"
disabled={pending}
onClick={() =>
doStreamingBulk((s, ids) =>
- bulkTranscribeAction(
- s,
- ids,
- transcribeQueue,
- parseConcurrency(transcribeConcurrency),
- ),
+ bulkTranscribeAction(s, ids, transcribeQueue),
)
}
className="px-2 py-1 rounded border border-zinc-300 dark:border-zinc-700 text-xs hover:bg-zinc-100 dark:hover:bg-zinc-800 disabled:opacity-50"
diff --git a/editor/app/channels/[slug]/components/stages/TranscribeStage.tsx b/editor/app/channels/[slug]/components/stages/TranscribeStage.tsx
@@ -6,7 +6,6 @@ import {
AUDIO_FORMAT_VALUES,
type AudioFormat,
} from "yt-dlp-transcript-common/lib/channelConfig";
-import { ConcurrencyControl } from "../../../../components/ConcurrencyControl";
import { QueueControl } from "../../../../components/QueueControl";
import {
ShardControl,
@@ -24,20 +23,11 @@ import { VideoIdList } from "../VideoIdList";
type AudioFormatChoice = AudioFormat | "any";
-function parseConcurrency(s: string): number | undefined {
- const trimmed = s.trim();
- if (trimmed === "") return undefined;
- const n = Number(trimmed);
- if (!Number.isFinite(n) || n < 1) return undefined;
- return Math.floor(n);
-}
-
type Props = {
slug: string;
existingQueues: string[];
failedVideoIds: string[];
downloadedNoTranscriptIds: string[];
- defaultConcurrency: number;
defaultQueueKey: string;
missingShard: ShardConfigSummary | null;
};
@@ -47,7 +37,6 @@ export function TranscribeStage({
existingQueues,
failedVideoIds,
downloadedNoTranscriptIds,
- defaultConcurrency,
defaultQueueKey,
missingShard,
}: Props) {
@@ -61,14 +50,12 @@ export function TranscribeStage({
slug={slug}
ids={downloadedNoTranscriptIds}
existingQueues={existingQueues}
- defaultConcurrency={defaultConcurrency}
defaultQueueKey={defaultQueueKey}
/>
)}
<TranscribeMissingSection
slug={slug}
existingQueues={existingQueues}
- defaultConcurrency={defaultConcurrency}
defaultQueueKey={defaultQueueKey}
missingShard={missingShard}
demoted={bucketCount > 0}
@@ -87,17 +74,14 @@ function BucketTranscribeSection({
slug,
ids,
existingQueues,
- defaultConcurrency,
defaultQueueKey,
}: {
slug: string;
ids: string[];
existingQueues: string[];
- defaultConcurrency: number;
defaultQueueKey: string;
}) {
const [queue, setQueue] = useState(defaultQueueKey);
- const [concurrency, setConcurrency] = useState("");
const [audioFormat, setAudioFormat] = useState<AudioFormatChoice>("any");
const [strictFormat, setStrictFormat] = useState(false);
return (
@@ -120,7 +104,6 @@ function BucketTranscribeSection({
slug,
ids,
queue,
- parseConcurrency(concurrency),
audioFormat === "any" ? undefined : audioFormat,
audioFormat !== "any" && strictFormat,
)
@@ -137,12 +120,6 @@ function BucketTranscribeSection({
existingQueues={existingQueues}
actionLabel="Transcribe downloaded"
/>
- <ConcurrencyControl
- value={concurrency}
- onChange={setConcurrency}
- defaultLimit={defaultConcurrency}
- actionLabel="Transcribe downloaded"
- />
<AudioFormatControl
actionLabel="Transcribe downloaded"
value={audioFormat}
@@ -160,20 +137,17 @@ function BucketTranscribeSection({
function TranscribeMissingSection({
slug,
existingQueues,
- defaultConcurrency,
defaultQueueKey,
missingShard,
demoted,
}: {
slug: string;
existingQueues: string[];
- defaultConcurrency: number;
defaultQueueKey: string;
missingShard: ShardConfigSummary | null;
demoted: boolean;
}) {
const [missingQueue, setMissingQueue] = useState(defaultQueueKey);
- const [missingConcurrency, setMissingConcurrency] = useState("");
const [missingReverse, setMissingReverse] = useState(false);
const [missingAudioFormat, setMissingAudioFormat] =
useState<AudioFormatChoice>("any");
@@ -199,7 +173,6 @@ function TranscribeMissingSection({
transcribeMissingAction(
slug,
missingQueue,
- parseConcurrency(missingConcurrency),
missingReverse,
missingAudioFormat === "any" ? undefined : missingAudioFormat,
missingAudioFormat !== "any" && missingStrictFormat,
@@ -219,12 +192,6 @@ function TranscribeMissingSection({
existingQueues={existingQueues}
actionLabel="Transcribe missing"
/>
- <ConcurrencyControl
- value={missingConcurrency}
- onChange={setMissingConcurrency}
- defaultLimit={defaultConcurrency}
- actionLabel="Transcribe missing"
- />
<AudioFormatControl
actionLabel="Transcribe missing"
value={missingAudioFormat}
diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx
@@ -21,7 +21,6 @@ import {
} from "yt-dlp-transcript-common/controller/shard";
import { loadDownloadOutcome } from "yt-dlp-transcript-common/lib/downloadOutcome-server";
import { getPaths } from "yt-dlp-transcript-common/lib/paths";
-import { getSettings } from "yt-dlp-transcript-common/lib/settings";
import { resolvePrimaryVtt } from "yt-dlp-transcript-common/lib/videoStatus";
import {
platformQueueKey,
@@ -124,7 +123,6 @@ export default async function ChannelDetailPage({
const filter = parseFilter(rawFilter);
const queryRaw = typeof sp.q === "string" ? sp.q : "";
const paths = getPaths();
- const settings = getSettings();
const config = await readChannelConfig(paths, slug);
if (!config) notFound();
@@ -246,7 +244,6 @@ export default async function ChannelDetailPage({
existingQueues={existingQueues}
failedVideoIds={failedVideoIds}
downloadedNoTranscriptIds={actionableDownloadedNoTranscriptIds}
- defaultConcurrency={settings.parallelTranscriptions}
defaultQueueKey={TRANSCRIPTION_QUEUE}
missingShard={transcribeMissingShard}
/>
@@ -400,7 +397,6 @@ export default async function ChannelDetailPage({
defaultTranscribeQueue={TRANSCRIPTION_QUEUE}
defaultDownloadQueue={platformDefaultQueueKey}
existingQueues={existingQueues}
- defaultConcurrency={settings.parallelTranscriptions}
/>
}
right={
diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts
@@ -23,7 +23,7 @@ import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"
import { setDoNotClean } from "yt-dlp-transcript-common/lib/doNotClean-server";
import { pruneFailedTranscriptions } from "yt-dlp-transcript-common/controller/failedTranscriptions";
import { transcodeAudio } from "yt-dlp-transcript-common/controller/transcode";
-import { transcribeOneVideo } from "yt-dlp-transcript-common/controller/transcribeOne";
+import { transcribeWithWorker } from "yt-dlp-transcript-common/controller/transcribeOne";
import { findVideoSourceUrl } from "yt-dlp-transcript-common/controller/undownloadedVideos";
import { getSettings } from "yt-dlp-transcript-common/lib/settings";
import { downloadOneManaged } from "yt-dlp-transcript-common/ytdlp/downloadOneManaged";
@@ -105,24 +105,16 @@ export async function transcribeOneAction(
channelSlug: slug,
videoId,
fn: async (onLog, signal, _setProgress, ctx) => {
- const task = makeTaskTracker(ctx, onLog).start({
- id: videoId,
- label: videoId,
- kind: "transcribe",
+ await transcribeWithWorker({
+ paths,
+ videoDir,
+ videoId,
+ audioFilename,
+ tracker: makeTaskTracker(ctx, onLog),
+ onLog,
+ signal,
});
- try {
- await transcribeOneVideo({
- paths,
- videoDir,
- videoId,
- audioFilename,
- onLog: task.onLog,
- signal,
- });
- revalidatePath(`/channels/${slug}/videos/${videoId}`);
- } finally {
- task.end();
- }
+ revalidatePath(`/channels/${slug}/videos/${videoId}`);
},
});
}
@@ -218,23 +210,15 @@ export async function whisperVideoAction(
} else {
onLog(`Audio already on disk for ${videoId}; skipping download.`);
}
- const task = makeTaskTracker(ctx, onLog).start({
- id: videoId,
- label: videoId,
- kind: "transcribe",
+ await transcribeWithWorker({
+ paths,
+ videoDir,
+ videoId,
+ audioFilename: `audio.${audioFormat}`,
+ tracker: makeTaskTracker(ctx, onLog),
+ onLog,
+ signal,
});
- try {
- await transcribeOneVideo({
- paths,
- videoDir,
- videoId,
- audioFilename: `audio.${audioFormat}`,
- onLog: task.onLog,
- signal,
- });
- } finally {
- task.end();
- }
revalidatePath(`/channels/${slug}/videos/${videoId}`);
revalidatePath(`/channels/${slug}`);
},
diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts
@@ -33,12 +33,6 @@ import {
} from "yt-dlp-transcript-common/jobs/streamCommand";
import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks";
-function sanitizeConcurrency(c: number | undefined): number | undefined {
- if (c === undefined) return undefined;
- if (!Number.isFinite(c) || c < 1) return undefined;
- return Math.floor(c);
-}
-
function sanitizeAudioFormat(f: AudioFormat | undefined): AudioFormat | undefined {
if (f === undefined) return undefined;
return AUDIO_FORMAT_VALUES.includes(f) ? f : undefined;
@@ -47,7 +41,6 @@ function sanitizeAudioFormat(f: AudioFormat | undefined): AudioFormat | undefine
export async function transcribeMissingAction(
slug: string,
queueKey?: string,
- concurrency?: number,
reverse?: boolean,
audioFormat?: AudioFormat,
strictAudioFormat?: boolean,
@@ -55,7 +48,6 @@ export async function transcribeMissingAction(
shardIndex?: number,
): Promise<StreamActionResult> {
const paths = getPaths();
- const limit = sanitizeConcurrency(concurrency);
const fmt = sanitizeAudioFormat(audioFormat);
return runManagedFunction({
kind: "whisper-all",
@@ -74,7 +66,6 @@ export async function transcribeMissingAction(
const result = await runWhisperBatch({
channelSlug: slug,
paths,
- concurrency: limit,
reverse: reverse === true,
audioFormat: fmt,
strictAudioFormat: fmt !== undefined && strictAudioFormat === true,
@@ -99,12 +90,10 @@ export async function transcribeBucketAction(
slug: string,
ids: string[],
queueKey?: string,
- concurrency?: number,
audioFormat?: AudioFormat,
strictAudioFormat?: boolean,
): Promise<StreamActionResult> {
const paths = getPaths();
- const limit = sanitizeConcurrency(concurrency);
const fmt = sanitizeAudioFormat(audioFormat);
const cleaned = Array.from(new Set(ids.map((id) => id.trim()).filter(Boolean)));
if (cleaned.length === 0) {
@@ -128,7 +117,6 @@ export async function transcribeBucketAction(
const result = await runWhisperBatch({
channelSlug: slug,
paths,
- concurrency: limit,
audioFormat: fmt,
strictAudioFormat: fmt !== undefined && strictAudioFormat === true,
ids: cleaned,
diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts
@@ -202,6 +202,9 @@ export async function saveSettingsAction(
maxTranscriptPageBytes: parsed,
transcriptionApp,
transcriptionApps,
+ // Phase 3 replaces this action with a worker-list editor. Until then, leave
+ // workers empty: writeSettings() synthesizes them from the app fields above.
+ workers: [],
cookiesFromBrowser,
sleepBetweenDownloadsSeconds: sleepParsed,
parallelTranscriptions: parallelParsed,