commit c778ea20aa674e0b46ec22420ccfd818a8d11882
parent 5b8cb9c3403fe2476f78aed391c2458675df0bbd
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sun, 14 Jun 2026 01:58:57 -0400
Transcription workers: cross-reference Workers ↔ Active Jobs (Phase 6)
- /api/workers includes each running task's startedAt + channelSlug; the Workers
page now shows what each busy worker is transcribing: the video (linked to its
page), the channel, a live elapsed timer, percent, and progress detail.
- buildActiveJobs maps task.workerId → worker name; /jobs/active task bars now
read "Transcribing <id> on <worker>", so a batch's spread across workers is
visible at a glance.
e2e: a running transcription shows channel+video on the Workers page and "on
<worker>" on the Active Jobs page. Typechecks clean; workers spec green.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
6 files changed, 134 insertions(+), 23 deletions(-)
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -6,6 +6,7 @@
- **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 returns a transport error mid-job, the video is automatically retried on another worker; a genuine transcription failure on the remote is not retried. On a transport failure the remote's `GET /api/worker/health` is probed, and a remote confirmed **down** is auto-disabled (shown "degraded" on the Workers page, with **Enable** to retry once it's back) so neither the current video nor later ones keep burning attempts on it — they fail over to a healthy worker. A worker that racks up repeated failures while still reachable is auto-disabled after a few strikes.
- **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.
+- **The Workers page and Active Jobs page cross-reference each other.** Each busy worker on `/workers` now shows what it's transcribing right now — the video (linked), the channel it's in, a live elapsed timer, percent, and the engine's progress detail — not just a bare bar. Conversely, every in-flight transcription on `/jobs/active` now says which worker it's running **on** (e.g. "Transcribing <id> on GPU"), so you can see how a batch is spread across your workers at a glance.
- **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/workers/route.ts b/editor/app/api/workers/route.ts
@@ -11,10 +11,18 @@ export async function GET() {
const pool = getWorkerPool();
const summary = pool.summary();
- // Collect running transcribe tasks grouped by their worker.
+ // Collect running transcribe tasks grouped by their worker, with enough detail
+ // for the Workers page to show what each worker is doing right now.
const byWorker = new Map<
string,
- { id: string; label: string; fraction?: number; detail?: string }[]
+ {
+ id: string;
+ label: string;
+ fraction?: number;
+ detail?: string;
+ startedAt: number;
+ channelSlug?: string;
+ }[]
>();
for (const job of getRegistry().list()) {
if (job.status !== "running" || !job.tasks) continue;
@@ -26,6 +34,8 @@ export async function GET() {
label: t.label,
fraction: t.fraction,
detail: t.detail,
+ startedAt: t.startedAt,
+ channelSlug: job.channelSlug,
});
byWorker.set(t.workerId, list);
}
diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts
@@ -3,6 +3,7 @@ import {
isDrainableKind,
type JobRecord,
} from "yt-dlp-transcript-common/jobs/registry";
+import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool";
import {
readChannelStat,
type ChannelStat,
@@ -91,6 +92,11 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
if (stat) channelStats.set(channelSlugs[i], stat);
}
+ // Map worker id → display name so each transcribe task can show which worker
+ // it's running on.
+ const workerNames = new Map<string, string>();
+ for (const w of getWorkerPool().summary()) workerNames.set(w.id, w.name);
+
const now = Date.now();
const jobs: RunningJobsListItem[] = activeRecords.map((j) => ({
id: j.id,
@@ -110,6 +116,8 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
fraction: t.fraction,
detail: t.detail,
startedAt: t.startedAt,
+ workerId: t.workerId,
+ workerName: t.workerId ? workerNames.get(t.workerId) : undefined,
})),
draining: j.draining === true,
drainable:
diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx
@@ -15,6 +15,9 @@ export type RunningJobsTask = {
detail?: string;
// Epoch ms when this sub-operation started, for the live "running for" timer.
startedAt: number;
+ // The worker running this transcription (transcribe tasks only).
+ workerId?: string;
+ workerName?: string;
};
export type RunningJobsListItem = {
@@ -195,6 +198,12 @@ function TaskProgressBar({ task }: { task: RunningJobsTask }) {
<span className="truncate">
<span className="text-zinc-500">{verb} </span>
<span className="font-mono">{task.label}</span>
+ {task.workerName && (
+ <span className="text-zinc-400">
+ {" "}
+ on <span className="font-mono">{task.workerName}</span>
+ </span>
+ )}
</span>
<span className="font-mono text-zinc-500 shrink-0">{meta}</span>
</div>
diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx
@@ -1,6 +1,8 @@
"use client";
+import Link from "next/link";
import { useCallback, useEffect, useState, useTransition } from "react";
+import { formatDuration } from "yt-dlp-transcript-common/lib/format";
import type { WorkerRuntimeState } from "yt-dlp-transcript-common/jobs/workerPool";
import {
disableWorkerAction,
@@ -15,8 +17,22 @@ export type WorkerTask = {
label: string;
fraction?: number;
detail?: string;
+ startedAt: number;
+ channelSlug?: string;
};
+// Live wall-clock that re-renders once a second; null until mounted so SSR and
+// the first client render agree (no Date.now() hydration mismatch).
+function useNow(): number | null {
+ const [now, setNow] = useState<number | null>(null);
+ useEffect(() => {
+ setNow(Date.now());
+ const id = setInterval(() => setNow(Date.now()), 1000);
+ return () => clearInterval(id);
+ }, []);
+ return now;
+}
+
export type WorkerView = {
id: string;
name: string;
@@ -232,28 +248,10 @@ function WorkerCard({
</p>
)}
{w.tasks.length > 0 && (
- <ul className="flex flex-col gap-1">
+ <ul className="flex flex-col gap-1.5">
{w.tasks.map((t) => (
- <li key={t.id} className="flex items-center gap-2 text-xs">
- <span
- role="progressbar"
- aria-label={`Transcribing ${t.label}`}
- aria-valuenow={
- t.fraction !== undefined
- ? Math.round(t.fraction * 100)
- : undefined
- }
- className="relative h-2 w-32 overflow-hidden rounded bg-zinc-200 dark:bg-zinc-800"
- >
- <span
- className="absolute inset-y-0 left-0 bg-blue-500"
- style={{ width: `${Math.round((t.fraction ?? 0) * 100)}%` }}
- />
- </span>
- <span className="font-mono text-zinc-600 dark:text-zinc-400">
- {t.label}
- </span>
- {t.detail && <span className="text-zinc-400">{t.detail}</span>}
+ <li key={t.id}>
+ <TaskRow task={t} />
</li>
))}
</ul>
@@ -261,3 +259,58 @@ function WorkerCard({
</li>
);
}
+
+function TaskRow({ task }: { task: WorkerTask }) {
+ const now = useNow();
+ const hasFraction = typeof task.fraction === "number";
+ const pct = hasFraction ? Math.round((task.fraction as number) * 100) : 0;
+ const elapsed =
+ now === null
+ ? null
+ : formatDuration(Math.max(0, Math.round((now - task.startedAt) / 1000))) ||
+ "0:00";
+ const meta = [elapsed, hasFraction ? `${pct}%` : null, task.detail]
+ .filter(Boolean)
+ .join(" · ");
+ return (
+ <div className="flex flex-col gap-0.5">
+ <div className="flex items-baseline justify-between gap-2 text-xs">
+ <span className="truncate">
+ <span className="text-zinc-500">Transcribing </span>
+ {task.channelSlug ? (
+ <Link
+ href={`/channels/${task.channelSlug}/videos/${encodeURIComponent(task.label)}`}
+ className="font-mono underline hover:text-zinc-900 dark:hover:text-zinc-100"
+ >
+ {task.label}
+ </Link>
+ ) : (
+ <span className="font-mono">{task.label}</span>
+ )}
+ {task.channelSlug && (
+ <span className="text-zinc-400">
+ {" "}
+ in <span className="font-mono">{task.channelSlug}</span>
+ </span>
+ )}
+ </span>
+ <span className="font-mono text-zinc-500 shrink-0">{meta}</span>
+ </div>
+ <span
+ role="progressbar"
+ aria-label={`Transcribing ${task.label}`}
+ aria-valuenow={hasFraction ? pct : undefined}
+ className="relative block h-1.5 w-full overflow-hidden rounded bg-zinc-200 dark:bg-zinc-800"
+ >
+ {hasFraction ? (
+ <span
+ className="absolute inset-y-0 left-0 bg-blue-500"
+ style={{ width: `${pct}%` }}
+ />
+ ) : (
+ <span className="absolute inset-y-0 left-0 w-1/3 animate-pulse bg-blue-500" />
+ )}
+ </span>
+ </div>
+ );
+}
diff --git a/editor/e2e/workers.spec.ts b/editor/e2e/workers.spec.ts
@@ -144,3 +144,33 @@ test("pausing all workers pauses a running batch instead of failing it; resume c
.poll(() => transcriptCount("pause-batch", ids), { timeout: 60_000 })
.toBe(ids.length);
});
+
+test("Workers page and Active jobs cross-reference the running task", async ({
+ page,
+}) => {
+ test.setTimeout(60_000);
+ await writeSettings({
+ workers: [
+ { id: "only", name: "Only", kind: "local", enabled: true, priority: 0, appId: "whisper-cpp", config: {} },
+ ],
+ });
+ await makeTranscribeChannel("xref-chan", ["slowopx1"]);
+
+ await page.goto("/channels/xref-chan");
+ await page.getByRole("button", { name: "Transcribe missing" }).click();
+
+ // The Workers page shows the worker actively transcribing the video, with the
+ // channel it belongs to.
+ await page.goto("/workers");
+ const only = page.getByRole("listitem").filter({ hasText: "Only" });
+ await expect(only).toContainText("Transcribing", { timeout: 15_000 });
+ await expect(only).toContainText("slowopx1");
+ await expect(only).toContainText("xref-chan");
+
+ // The Active Jobs page shows that task running ON the worker "Only".
+ await page.goto("/jobs/active");
+ const section = page.locator(
+ "section[aria-label='Active jobs for xref-chan']",
+ );
+ await expect(section).toContainText(/on\s+Only/, { timeout: 15_000 });
+});