Archilyzer · Source

archilyzer

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

commit 147f0527761b773c422ba47174cf99464df3c49c
parent 3afd0dc326ed91c040c27a76845be6265715b605
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 14 Jun 2026 00:31:39 -0400

Transcription workers: settings editor + Workers page (Phase 3)

Replace the single-engine settings UI with a worker-list editor and add a live
Workers page for runtime control.

- Settings: WorkersField (add/remove/reorder, per-engine fields, slots, enabled,
  remote baseUrl/token/sharedFs) serialized as workersJson; saveSettingsAction
  parses + validates via sanitizeWorkers/validateWorkers and reconfigures the
  live pool (applyEnabled). Drops the old App dropdown + Parallel transcriptions
  field (worker slots replace per-run/global concurrency throughout).
- Workers page (/workers): polls GET /api/workers (pool summary joined with
  in-flight transcribe tasks by workerId). Per-worker Enable / Disable / Drain,
  plus Pause all / Resume-to-previous-state. Runtime controls are transient
  operator overrides (not persisted); reconfigure gained applyEnabled so a
  batch-start reconfigure preserves them while a settings save applies new intent.
- Pool: pauseAll/resumeAll/isPaused with a pre-pause snapshot, sticky across
  reconfigure; runtime controls sync the snapshot so resume honors explicit
  per-worker choices made while paused.
- Nav: add Workers under Pool.
- e2e: settings worker-editor + no-enabled-worker rejection; workers.spec covers
  list/disable/enable/pause/resume and the key promise — pausing all workers
  parks a running batch (stays running, doesn't fail) and resume completes it.

CHANGELOG updated.

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

Diffstat:
Mcommon/jobs/workerPool.ts | 109+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Meditor/CHANGELOG.md | 3+++
Aeditor/app/api/workers/route.ts | 40++++++++++++++++++++++++++++++++++++++++
Meditor/app/layout.tsx | 1+
Meditor/app/settings/actions.ts | 116+++++++++++++++++++++++--------------------------------------------------------
Meditor/app/settings/components/SettingsForm.tsx | 111++++++++++---------------------------------------------------------------------
Aeditor/app/settings/components/WorkersField.tsx | 375+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/workers/actions.ts | 50++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/workers/components/WorkersView.tsx | 267+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/workers/page.tsx | 38++++++++++++++++++++++++++++++++++++++
Meditor/e2e/settings.spec.ts | 35++++++++++++++++++++++++-----------
Aeditor/e2e/workers.spec.ts | 145+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
12 files changed, 1078 insertions(+), 212 deletions(-)

diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts @@ -65,21 +65,35 @@ class WorkerPool { private entries = new Map<string, PoolEntry>(); private waiters: Waiter[] = []; private initialized = false; + // When set, the pool is in a temporary "pause all": every worker is forced + // disabled and this records each worker's pre-pause state so resumeAll can + // restore it. Runtime-only (not persisted) — a restart returns workers to + // their configured enabled state. See pauseAll / resumeAll. + private pausedSnapshot: Map<string, WorkerRuntimeState> | null = null; // Lazily seed from settings on first use. Idempotent. private ensureInit(): void { if (this.initialized) return; this.initialized = true; - this.reconfigure(getSettings().workers); + // Initial seed: apply the persisted enabled flags as the starting state. + this.reconfigure(getSettings().workers, { applyEnabled: true }); } // 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 { + // each worker's 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. + // + // `applyEnabled` controls whether an EXISTING worker's runtime state is reset + // from its persisted `enabled` flag. The settings-save path passes true (a new + // enabled intent should take effect). The batch-start path passes false (the + // default) so an operator's runtime enable/disable/drain from the Workers page + // is preserved across the reconfigure a batch does on start — otherwise every + // batch would silently re-enable a worker the operator just turned off. + reconfigure(workers?: Worker[], opts?: { applyEnabled?: boolean }): void { this.initialized = true; + const applyEnabled = opts?.applyEnabled === true; const next = workers ?? getSettings().workers; const seen = new Set<string>(); for (const w of next) { @@ -88,17 +102,20 @@ class WorkerPool { 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"; + if (applyEnabled) { + // Reconcile runtime state with the (just-saved) 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: leave runtime state untouched (preserve operator overrides). } else { this.entries.set(w.id, { config: w, @@ -120,10 +137,50 @@ class WorkerPool { this.entries.delete(id); } } + // Keep "pause all" sticky across a reconfigure (settings save / batch start): + // re-disable everything and snapshot any newly-added worker's intended state. + if (this.pausedSnapshot) { + for (const [id, entry] of this.entries) { + if (!this.pausedSnapshot.has(id)) { + this.pausedSnapshot.set(id, entry.state); + } + entry.state = "disabled"; + } + } // A reconfigure can make new slots eligible (worker enabled/added). this.pump(); } + // Temporary pause: disable every worker (stop granting new leases; in-flight + // leases finish on their own), recording each worker's pre-pause state so + // resumeAll restores it exactly. Idempotent — a second call is a no-op. + pauseAll(): void { + this.ensureInit(); + if (this.pausedSnapshot) return; + this.pausedSnapshot = new Map(); + for (const [id, entry] of this.entries) { + this.pausedSnapshot.set(id, entry.state); + entry.state = "disabled"; + } + } + + // Undo pauseAll: restore each worker to the state it had when paused (a + // transient "draining" restores to enabled). Wakes parked waiters. + resumeAll(): void { + this.ensureInit(); + if (!this.pausedSnapshot) return; + for (const [id, prev] of this.pausedSnapshot) { + const entry = this.entries.get(id); + if (entry) entry.state = prev === "draining" ? "enabled" : prev; + } + this.pausedSnapshot = null; + this.pump(); + } + + isPaused(): boolean { + return this.pausedSnapshot !== null; + } + private eligible(entry: PoolEntry): boolean { return ( entry.state === "enabled" && @@ -203,8 +260,18 @@ class WorkerPool { }); } - // --- Runtime controls (Workers page). These mutate the live pool; callers - // also persist the new `enabled` flag to settings so it survives a restart. --- + // --- Runtime controls (Workers page). Transient operator overrides: they + // mutate the live pool's runtime `state` only, never settings.json. A restart + // returns workers to their configured enabled state. They survive the + // config-only reconfigure a batch does on start (see reconfigure's + // applyEnabled). While a pause-all is active, they also update the snapshot so + // resumeAll honors the operator's explicit choice for that worker. --- + + // Reflect an operator override into the pause snapshot so resumeAll restores + // the operator's intent for that worker rather than its pre-pause state. + private syncSnapshot(id: string, state: WorkerRuntimeState): void { + if (this.pausedSnapshot) this.pausedSnapshot.set(id, state); + } enableWorker(id: string): boolean { const entry = this.entries.get(id); @@ -212,7 +279,7 @@ class WorkerPool { entry.state = "enabled"; entry.degraded = false; entry.consecutiveFailures = 0; - entry.config = { ...entry.config, enabled: true }; + this.syncSnapshot(id, "enabled"); this.pump(); return true; } @@ -224,7 +291,7 @@ class WorkerPool { const entry = this.entries.get(id); if (!entry) return false; entry.state = "disabled"; - entry.config = { ...entry.config, enabled: false }; + this.syncSnapshot(id, "disabled"); return true; } @@ -233,8 +300,8 @@ class WorkerPool { 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"; + this.syncSnapshot(id, "disabled"); return true; } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -3,6 +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. +- **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. - **Parakeet transcriptions show a per-video ETA.** While a video is being transcribed with the parakeet app, its per-task progress bar on `/jobs/active` now shows an estimate of how long that single video has left (e.g. `segment 3/12 · ETA 6:10`), alongside the existing segment count and percent. The wrapper (`scripts/parakeet-stitch.mjs`) measures each segment's real transcription wall-time and projects the remaining time as **average time per completed segment × remaining segments**, emitting it on its progress lines; the parser surfaces that ETA in the task detail. It appears from the second segment onward (the first segment has no average to project from yet). This is distinct from the batch-level "~4:30 left" estimate across all videos. diff --git a/editor/app/api/workers/route.ts b/editor/app/api/workers/route.ts @@ -0,0 +1,40 @@ +import { NextResponse } from "next/server"; +import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; +import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; + +export const dynamic = "force-dynamic"; + +// Live worker status for the Workers page. Polled ~1s by WorkersView. Joins the +// pool's per-worker slot/state summary with the registry's in-flight transcribe +// tasks (which carry workerId) so each card can show what it's currently running. +export async function GET() { + const pool = getWorkerPool(); + const summary = pool.summary(); + + // Collect running transcribe tasks grouped by their worker. + const byWorker = new Map< + string, + { id: string; label: string; fraction?: number; detail?: string }[] + >(); + for (const job of getRegistry().list()) { + if (job.status !== "running" || !job.tasks) continue; + for (const t of job.tasks) { + if (t.kind !== "transcribe" || !t.workerId) continue; + const list = byWorker.get(t.workerId) ?? []; + list.push({ + id: t.id, + label: t.label, + fraction: t.fraction, + detail: t.detail, + }); + byWorker.set(t.workerId, list); + } + } + + const workers = summary.map((w) => ({ + ...w, + tasks: byWorker.get(w.id) ?? [], + })); + + return NextResponse.json({ paused: pool.isPaused(), workers }); +} diff --git a/editor/app/layout.tsx b/editor/app/layout.tsx @@ -52,6 +52,7 @@ const NAV_GROUPS: NavGroup[] = [ links: [ { href: "/jobs", label: "Jobs", badgeKey: "jobs" }, { href: "/jobs/active", label: "Active", badgeKey: "running" }, + { href: "/workers", label: "Workers" }, { href: "/build", label: "Build" }, { href: "/actionable", label: "Actionable" }, ], diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts @@ -8,19 +8,20 @@ import { isReportDebouncePreset, normalizeSocialSvg, parseSocialLinks, - PARALLEL_TRANSCRIPTIONS_MAX, + PARALLEL_TRANSCRIPTIONS_DEFAULT, SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS, TRANSCRIPT_PAGE_HARD_CAP_BYTES, TRANSCRIPT_PAGE_MIN_BYTES, - validateTranscribeArgs, writeSettings, type SiteSettings, type SocialLink, } from "yt-dlp-transcript-common/lib/settings"; +import { DEFAULT_TRANSCRIPTION_APP_ID } from "yt-dlp-transcript-common/lib/transcriptionApps"; import { - TRANSCRIPTION_APPS, - type AppInstanceConfig, -} from "yt-dlp-transcript-common/lib/transcriptionApps"; + sanitizeWorkers, + validateWorkers, +} from "yt-dlp-transcript-common/lib/workers"; +import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; export type SaveResult = { ok: true } | { ok: false; error: string }; @@ -36,9 +37,6 @@ export async function saveSettingsAction( const sleepRaw = String( formData.get("sleepBetweenDownloadsSeconds") ?? "", ).trim(); - const parallelRaw = String( - formData.get("parallelTranscriptions") ?? "", - ).trim(); const autoRefreshRaw = String( formData.get("autoRefreshIntervalSeconds") ?? "", ).trim(); @@ -54,64 +52,17 @@ export async function saveSettingsAction( if (!adminTitle) return { ok: false, error: "Admin title is required" }; - const transcriptionApp = String(formData.get("transcriptionApp") ?? "").trim(); - const activeApp = TRANSCRIPTION_APPS[transcriptionApp]; - if (!activeApp) { - return { ok: false, error: "Unknown transcription app" }; - } - - // Collect per-app config from namespaced fields (app.<id>.<field>). Only the - // selected app's customArgs are validated; writeSettings re-checks the - // resolved binary. - const transcriptionApps: Record<string, AppInstanceConfig> = {}; - for (const app of Object.values(TRANSCRIPTION_APPS)) { - const cfg: AppInstanceConfig = {}; - const bin = String(formData.get(`app.${app.id}.bin`) ?? "").trim(); - if (bin) cfg.bin = bin; - if (app.fields.model) { - const model = String(formData.get(`app.${app.id}.model`) ?? "").trim(); - if (model) cfg.model = model; - } - if (app.fields.remoteUrl) { - const remoteUrl = String( - formData.get(`app.${app.id}.remoteUrl`) ?? "", - ).trim(); - if (remoteUrl) cfg.remoteUrl = remoteUrl; - } - if (app.fields.chunkSize) { - const chunkRaw = String( - formData.get(`app.${app.id}.chunkSize`) ?? "", - ).trim(); - if (chunkRaw) { - const n = Number.parseInt(chunkRaw, 10); - if (!Number.isFinite(n) || n <= 0) { - return { - ok: false, - error: `${app.label} chunk size must be a positive number`, - }; - } - cfg.chunkSize = n; - } - } - if (app.fields.customArgs) { - const customArgs = String(formData.get(`app.${app.id}.customArgs`) ?? "") - .split("\n") - .map((s) => s.trim()) - .filter(Boolean); - if (customArgs.length > 0) cfg.customArgs = customArgs; - } - if (Object.keys(cfg).length > 0) transcriptionApps[app.id] = cfg; - } - - const activeCfg = transcriptionApps[transcriptionApp]; - if ( - activeApp.fields.customArgs && - activeCfg?.customArgs && - activeCfg.customArgs.length > 0 - ) { - const argsErr = validateTranscribeArgs(activeCfg.customArgs); - if (argsErr) return { ok: false, error: argsErr }; + // Workers are submitted as a JSON array by the WorkersField client component. + // Sanitize + validate here for a friendly error; writeSettings re-validates. + let workersInput: unknown; + try { + workersInput = JSON.parse(String(formData.get("workersJson") ?? "[]")); + } catch { + return { ok: false, error: "Workers payload is malformed" }; } + const workers = sanitizeWorkers(workersInput); + const workersErr = validateWorkers(workers); + if (workersErr) return { ok: false, error: workersErr }; const parsed = Number.parseInt(maxBytesRaw, 10); if (!Number.isFinite(parsed)) { @@ -144,17 +95,6 @@ export async function saveSettingsAction( }; } - const parallelParsed = Number.parseInt(parallelRaw, 10); - if (!Number.isFinite(parallelParsed)) { - return { ok: false, error: "parallelTranscriptions must be a number" }; - } - if (parallelParsed < 1 || parallelParsed > PARALLEL_TRANSCRIPTIONS_MAX) { - return { - ok: false, - error: `parallelTranscriptions must be between 1 and ${PARALLEL_TRANSCRIPTIONS_MAX}`, - }; - } - // 0 disables passive refresh; any other value must land in the allowed window. const autoRefreshParsed = Number.parseInt(autoRefreshRaw, 10); if (!Number.isFinite(autoRefreshParsed)) { @@ -200,21 +140,31 @@ export async function saveSettingsAction( const next: SiteSettings = { adminTitle, 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: [], + // Workers are the source of truth; writeSettings derives the deprecated + // transcriptionApp/transcriptionApps shadow from them. parallelTranscriptions + // is vestigial (kept only for rollback) so pass the default. + transcriptionApp: DEFAULT_TRANSCRIPTION_APP_ID, + transcriptionApps: {}, + workers, cookiesFromBrowser, sleepBetweenDownloadsSeconds: sleepParsed, - parallelTranscriptions: parallelParsed, + parallelTranscriptions: PARALLEL_TRANSCRIPTIONS_DEFAULT, inlineTranscribeOnFallback, skipLiveDownloads, reportDebouncePreset, autoRefreshIntervalSeconds: autoRefreshParsed, socialLinks, }; - await writeSettings(next); + try { + await writeSettings(next); + } catch (e) { + return { ok: false, error: (e as Error).message }; + } + // Apply the saved worker list to the live pool immediately so the Workers page + // and the next batch see it without waiting for a batch-start reconfigure. + // applyEnabled: a saved Enabled change is an explicit intent and takes effect now. + getWorkerPool().reconfigure(next.workers, { applyEnabled: true }); revalidatePath("/settings"); + revalidatePath("/workers"); return { ok: true }; } diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx @@ -12,6 +12,7 @@ import { toSocialRow, type SocialRow, } from "../../components/SocialLinksField"; +import { WorkersField } from "./WorkersField"; type Props = { initial: SiteSettings; @@ -26,7 +27,6 @@ export function SettingsForm({ initial, apps }: Props) { const [social, setSocial] = useState<SocialRow[]>(() => initial.socialLinks.map(toSocialRow), ); - const [appId, setAppId] = useState<string>(initial.transcriptionApp); return ( <form action={formAction} className="flex flex-col gap-4 max-w-xl"> @@ -45,97 +45,21 @@ export function SettingsForm({ initial, apps }: Props) { type="number" /> <fieldset className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded p-3"> - <legend className="px-1 text-sm font-medium">Transcription</legend> + <legend className="px-1 text-sm font-medium"> + Transcription workers + </legend> <p className="text-xs text-zinc-500"> - The app used by the channel &ldquo;Transcribe missing&rdquo; / - &ldquo;Retry failures&rdquo; actions. Each app handles its own - arguments and output format. Per-app fields below; only the selected - app&apos;s settings are used. + 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{" "} + <a href="/workers" className="underline"> + Workers + </a>{" "} + page. </p> - <label className="flex flex-col gap-1 text-sm"> - <span className="font-medium">App</span> - <select - name="transcriptionApp" - value={appId} - onChange={(e) => setAppId(e.target.value)} - className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm" - > - {apps.map((app) => ( - <option key={app.id} value={app.id}> - {app.label} - </option> - ))} - </select> - </label> - {apps.map((app) => { - const cfg = initial.transcriptionApps[app.id] ?? {}; - return ( - <div - key={app.id} - hidden={app.id !== appId} - className="flex flex-col gap-3 border-t border-zinc-200 dark:border-zinc-800 pt-3" - > - <Field - label="Binary" - name={`app.${app.id}.bin`} - defaultValue={cfg.bin ?? ""} - hint="Path or name of the executable. Leave blank to use the default for this app." - /> - {app.fields.model && ( - <Field - label="Model" - name={`app.${app.id}.model`} - defaultValue={cfg.model ?? ""} - hint="whisper.cpp: substituted for {model}. chough: sets CHOUGH_MODEL. parakeet: the .gguf model path. Leave blank for the app default." - /> - )} - {app.fields.remoteUrl && ( - <Field - label="Remote URL" - name={`app.${app.id}.remoteUrl`} - defaultValue={cfg.remoteUrl ?? ""} - hint="chough only: transcribe via a remote chough --server (CHOUGH_URL). Leave blank for local transcription." - /> - )} - {app.fields.chunkSize && ( - <Field - label="Chunk size (seconds)" - name={`app.${app.id}.chunkSize`} - defaultValue={ - cfg.chunkSize !== undefined ? String(cfg.chunkSize) : "" - } - type="number" - hint={ - app.id === "parakeet" - ? "parakeet: per-window length in seconds for overlapping splitting. Leave blank for the default (480s)." - : "chough only (-c). Leave blank for the app default (60s)." - } - /> - )} - {app.fields.customArgs && ( - <label className="flex flex-col gap-1 text-sm"> - <span className="font-medium">Custom args (one per line)</span> - <textarea - name={`app.${app.id}.customArgs`} - defaultValue={(cfg.customArgs ?? []).join("\n")} - rows={6} - className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm font-mono" - /> - <span className="text-xs text-zinc-500"> - Leave blank for the default whisper.cpp call:{" "} - <code> - -ojf -l en -m {"{model}"} -of {"{outputBase}"}{" "} - {"{audioFile}"} - </code> - . Placeholders: <code>{"{audioFile}"}</code> (required),{" "} - <code>{"{outputBase}"}</code> (required),{" "} - <code>{"{model}"}</code> (optional). - </span> - </label> - )} - </div> - ); - })} + <WorkersField initial={initial.workers} apps={apps} name="workersJson" /> </fieldset> <Field label="Cookies from browser (retry only)" @@ -151,13 +75,6 @@ export function SettingsForm({ initial, apps }: Props) { hint="Pause inserted between per-video yt-dlp invocations in managed batch downloads (download-from-playlist, download-missing). 0 disables. Default 10s. Each channel can override this in its Advanced settings." /> <Field - label="Parallel transcriptions" - name="parallelTranscriptions" - defaultValue={String(initial.parallelTranscriptions)} - type="number" - hint="Default number of videos transcribed at once when a run doesn't set its own concurrency (1–16). Default 2. Higher values use more CPU/RAM. The per-run Concurrency input on a channel overrides this for a single run." - /> - <Field label="Auto-refresh interval (seconds)" name="autoRefreshIntervalSeconds" defaultValue={String(initial.autoRefreshIntervalSeconds)} diff --git a/editor/app/settings/components/WorkersField.tsx b/editor/app/settings/components/WorkersField.tsx @@ -0,0 +1,375 @@ +"use client"; + +import { useId, useState } from "react"; +import type { Worker } from "yt-dlp-transcript-common/lib/workers"; +import type { TranscriptionAppDescriptor } from "yt-dlp-transcript-common/lib/transcriptionApps"; + +// Editor for the transcription worker list. Manages an ordered array of workers +// (list order == priority, top = highest) and serializes it to a hidden JSON +// input the save action reads (mirrors SocialLinksField/socialLinksJson). The +// server re-sanitizes and validates, so this component only needs to produce a +// reasonable shape. + +type Row = Worker & { _key: number }; + +let keyCounter = 0; +function withKeys(workers: Worker[]): Row[] { + return workers.map((w) => ({ ...w, _key: keyCounter++ })); +} + +function blankLocal(appId: string): Row { + return { + _key: keyCounter++, + id: "", + name: "", + kind: "local", + enabled: true, + priority: 0, + slots: 1, + appId, + config: {}, + }; +} + +function blankRemote(): Row { + return { + _key: keyCounter++, + id: "", + name: "", + kind: "remote", + enabled: true, + priority: 0, + slots: 1, + remote: { baseUrl: "" }, + }; +} + +// Strip the React-only _key and set priority from list order before serializing. +function serialize(rows: Row[]): string { + return JSON.stringify( + rows.map(({ _key, ...w }, i) => { + void _key; + return { ...w, priority: i }; + }), + ); +} + +type Props = { + initial: Worker[]; + apps: TranscriptionAppDescriptor[]; + name: string; +}; + +export function WorkersField({ initial, apps, name }: Props) { + const [rows, setRows] = useState<Row[]>(() => withKeys(initial)); + const defaultAppId = apps[0]?.id ?? "whisper-cpp"; + + function update(next: Row[]) { + setRows(next); + } + function patch(key: number, change: Partial<Row>) { + update(rows.map((r) => (r._key === key ? { ...r, ...change } : r))); + } + function patchConfig(key: number, change: Partial<NonNullable<Worker["config"]>>) { + update( + rows.map((r) => + r._key === key ? { ...r, config: { ...r.config, ...change } } : r, + ), + ); + } + function patchRemote(key: number, change: Partial<NonNullable<Worker["remote"]>>) { + update( + rows.map((r) => + r._key === key + ? { ...r, remote: { baseUrl: "", ...r.remote, ...change } } + : r, + ), + ); + } + function remove(key: number) { + update(rows.filter((r) => r._key !== key)); + } + function move(index: number, dir: -1 | 1) { + const j = index + dir; + if (j < 0 || j >= rows.length) return; + const next = [...rows]; + [next[index], next[j]] = [next[j], next[index]]; + update(next); + } + + return ( + <div className="flex flex-col gap-3"> + <input type="hidden" name={name} value={serialize(rows)} readOnly /> + {rows.length === 0 && ( + <p className="text-xs text-red-700 dark:text-red-300"> + No workers defined — add at least one so transcription can run. + </p> + )} + {rows.map((row, i) => ( + <WorkerCard + key={row._key} + row={row} + index={i} + total={rows.length} + apps={apps} + onPatch={(c) => patch(row._key, c)} + onPatchConfig={(c) => patchConfig(row._key, c)} + onPatchRemote={(c) => patchRemote(row._key, c)} + onRemove={() => remove(row._key)} + onMove={(dir) => move(i, dir)} + /> + ))} + <div className="flex gap-2"> + <button + type="button" + onClick={() => update([...rows, blankLocal(defaultAppId)])} + 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" + > + + Local worker + </button> + <button + type="button" + onClick={() => update([...rows, blankRemote()])} + 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" + > + + Remote worker + </button> + </div> + </div> + ); +} + +function WorkerCard({ + row, + index, + total, + apps, + onPatch, + onPatchConfig, + onPatchRemote, + onRemove, + onMove, +}: { + row: Row; + index: number; + total: number; + apps: TranscriptionAppDescriptor[]; + onPatch: (c: Partial<Row>) => void; + onPatchConfig: (c: Partial<NonNullable<Worker["config"]>>) => void; + onPatchRemote: (c: Partial<NonNullable<Worker["remote"]>>) => void; + onRemove: () => void; + onMove: (dir: -1 | 1) => void; +}) { + const uid = useId(); + const app = apps.find((a) => a.id === row.appId); + const cfg = row.config ?? {}; + const remote = row.remote ?? { baseUrl: "" }; + return ( + <div + aria-label={`worker ${index + 1}`} + className="flex flex-col gap-2 rounded border border-zinc-200 dark:border-zinc-800 p-3" + > + <div className="flex flex-wrap items-center gap-2"> + <span className="text-xs text-zinc-500" title="priority (lower = preferred)"> + #{index + 1} + </span> + <input + type="text" + aria-label={`worker ${index + 1} name`} + placeholder="Worker name" + value={row.name} + onChange={(e) => onPatch({ name: e.target.value })} + className="flex-1 min-w-40 rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm" + /> + <label className="flex items-center gap-1 text-xs"> + <input + type="checkbox" + checked={row.enabled} + onChange={(e) => onPatch({ enabled: e.target.checked })} + aria-label={`worker ${index + 1} enabled`} + /> + Enabled + </label> + <div className="flex items-center gap-1"> + <button + type="button" + onClick={() => onMove(-1)} + disabled={index === 0} + aria-label={`move worker ${index + 1} up`} + className="px-1.5 py-0.5 rounded border border-zinc-300 dark:border-zinc-700 text-xs disabled:opacity-40" + > + ↑ + </button> + <button + type="button" + onClick={() => onMove(1)} + disabled={index === total - 1} + aria-label={`move worker ${index + 1} down`} + className="px-1.5 py-0.5 rounded border border-zinc-300 dark:border-zinc-700 text-xs disabled:opacity-40" + > + ↓ + </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" + > + Remove + </button> + </div> + </div> + + <div className="flex flex-wrap items-end gap-3"> + {row.kind === "local" && ( + <label className="flex flex-col gap-1 text-xs"> + <span className="font-medium">Engine</span> + <select + value={row.appId ?? ""} + onChange={(e) => onPatch({ appId: e.target.value })} + aria-label={`worker ${index + 1} engine`} + className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm" + > + {apps.map((a) => ( + <option key={a.id} value={a.id}> + {a.label} + </option> + ))} + </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"} + </span> + </div> + + {row.kind === "local" ? ( + <div className="flex flex-col gap-2 border-t border-zinc-100 dark:border-zinc-800 pt-2"> + <CardField + label="Binary" + value={cfg.bin ?? ""} + onChange={(v) => onPatchConfig({ bin: v })} + hint="Path/name of the executable. Blank = the engine's default." + id={`${uid}-bin`} + /> + {app?.fields.model && ( + <CardField + label="Model" + value={cfg.model ?? ""} + onChange={(v) => onPatchConfig({ model: v })} + hint="whisper.cpp: {model}. chough: CHOUGH_MODEL. parakeet: .gguf path. Blank = default." + id={`${uid}-model`} + /> + )} + {app?.fields.chunkSize && ( + <CardField + label="Chunk size (seconds)" + type="number" + value={cfg.chunkSize !== undefined ? String(cfg.chunkSize) : ""} + onChange={(v) => + onPatchConfig({ + chunkSize: v.trim() === "" ? undefined : Number(v), + }) + } + hint="parakeet: per-window seconds. chough: -c. Blank = default." + id={`${uid}-chunk`} + /> + )} + {app?.fields.customArgs && ( + <label className="flex flex-col gap-1 text-xs"> + <span className="font-medium">Custom args (one per line)</span> + <textarea + value={(cfg.customArgs ?? []).join("\n")} + onChange={(e) => + onPatchConfig({ + customArgs: e.target.value + .split("\n") + .map((s) => s.trim()) + .filter(Boolean), + }) + } + rows={4} + className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm font-mono" + /> + <span className="text-zinc-500"> + Blank = default. Placeholders: <code>{"{audioFile}"}</code>,{" "} + <code>{"{outputBase}"}</code> (both required),{" "} + <code>{"{model}"}</code>. + </span> + </label> + )} + </div> + ) : ( + <div className="flex flex-col gap-2 border-t border-zinc-100 dark:border-zinc-800 pt-2"> + <CardField + label="Base URL" + value={remote.baseUrl} + onChange={(v) => onPatchRemote({ baseUrl: v })} + hint="Another instance of this app, e.g. http://gpu-box.lan:3001. It runs its own worker pool." + id={`${uid}-url`} + /> + <CardField + label="Token" + value={remote.token ?? ""} + onChange={(v) => onPatchRemote({ token: v })} + hint="Bearer token sent to the remote (must match its WORKER_TOKEN). Blank if the remote has no token." + id={`${uid}-token`} + /> + <label className="flex items-center gap-2 text-xs"> + <input + type="checkbox" + checked={remote.sharedFs === true} + onChange={(e) => onPatchRemote({ sharedFs: e.target.checked })} + /> + <span> + Shared filesystem — send the video path instead of uploading audio + (only if the remote mounts the same transcripts dir). + </span> + </label> + </div> + )} + </div> + ); +} + +function CardField({ + label, + value, + onChange, + hint, + type = "text", + id, +}: { + label: string; + value: string; + onChange: (v: string) => void; + hint?: string; + type?: string; + id: string; +}) { + return ( + <label className="flex flex-col gap-1 text-xs" htmlFor={id}> + <span className="font-medium">{label}</span> + <input + id={id} + type={type} + value={value} + onChange={(e) => onChange(e.target.value)} + className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm" + /> + {hint && <span className="text-zinc-500">{hint}</span>} + </label> + ); +} diff --git a/editor/app/workers/actions.ts b/editor/app/workers/actions.ts @@ -0,0 +1,50 @@ +"use server"; + +import { revalidatePath } from "next/cache"; +import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; + +// Runtime worker controls for the Workers page. These are transient operator +// overrides on the live pool — they are NOT written to settings.json (a restart +// returns workers to their configured enabled state). The settings form is where +// the persisted default-enabled state lives. + +export type WorkerActionResult = { ok: boolean; error?: string }; + +function refresh() { + revalidatePath("/workers"); +} + +export async function enableWorkerAction(id: string): Promise<WorkerActionResult> { + const ok = getWorkerPool().enableWorker(id); + refresh(); + return ok ? { ok: true } : { ok: false, error: `Unknown worker "${id}"` }; +} + +export async function disableWorkerAction(id: string): Promise<WorkerActionResult> { + const ok = getWorkerPool().disableWorker(id); + refresh(); + return ok ? { ok: true } : { ok: false, error: `Unknown worker "${id}"` }; +} + +// Drain: stop giving the worker new transcriptions but let in-flight ones finish +// (the "free up the GPU when done with the current video" path). +export async function drainWorkerAction(id: string): Promise<WorkerActionResult> { + const ok = getWorkerPool().drainWorker(id); + refresh(); + return ok ? { ok: true } : { ok: false, error: `Unknown worker "${id}"` }; +} + +// Temporary pause-all: disable every worker, remembering each one's current +// state. Running batches pause (they wait for a worker) rather than failing. +export async function pauseAllWorkersAction(): Promise<WorkerActionResult> { + getWorkerPool().pauseAll(); + refresh(); + return { ok: true }; +} + +// Resume: restore every worker to the state it had when pauseAll ran. +export async function resumeAllWorkersAction(): Promise<WorkerActionResult> { + getWorkerPool().resumeAll(); + refresh(); + return { ok: true }; +} diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx @@ -0,0 +1,267 @@ +"use client"; + +import { useCallback, useEffect, useState, useTransition } from "react"; +import type { WorkerRuntimeState } from "yt-dlp-transcript-common/jobs/workerPool"; +import { + disableWorkerAction, + drainWorkerAction, + enableWorkerAction, + pauseAllWorkersAction, + resumeAllWorkersAction, +} from "../actions"; + +export type WorkerTask = { + id: string; + label: string; + fraction?: number; + detail?: string; +}; + +export type WorkerView = { + id: string; + name: string; + kind: "local" | "remote"; + appId?: string; + priority: number; + slots: number; + inUse: number; + state: WorkerRuntimeState; + degraded: boolean; + enabled: boolean; + tasks: WorkerTask[]; +}; + +export type WorkersPayload = { + paused: boolean; + workers: WorkerView[]; +}; + +const POLL_MS = 1000; + +export function WorkersView({ initial }: { initial: WorkersPayload }) { + const [payload, setPayload] = useState<WorkersPayload>(initial); + const [pending, startTransition] = useTransition(); + + const refetch = useCallback(async () => { + try { + const res = await fetch("/api/workers", { cache: "no-store" }); + if (res.ok) setPayload((await res.json()) as WorkersPayload); + } catch { + // transient — keep polling + } + }, []); + + useEffect(() => { + let cancelled = false; + let timer: ReturnType<typeof setTimeout> | null = null; + async function tick() { + if (!cancelled) await refetch(); + if (!cancelled) timer = setTimeout(tick, POLL_MS); + } + timer = setTimeout(tick, POLL_MS); + return () => { + cancelled = true; + if (timer) clearTimeout(timer); + }; + }, [refetch]); + + function run(action: () => Promise<unknown>) { + startTransition(async () => { + await action(); + await refetch(); + }); + } + + const { workers, paused } = payload; + + return ( + <div className="flex flex-col gap-3"> + <div className="flex items-center gap-3"> + {paused ? ( + <button + type="button" + disabled={pending} + onClick={() => run(resumeAllWorkersAction)} + className="px-3 py-1.5 rounded-md bg-amber-500 text-white text-sm font-medium hover:bg-amber-600 disabled:opacity-50" + > + Resume all workers + </button> + ) : ( + <button + type="button" + disabled={pending || workers.length === 0} + onClick={() => run(pauseAllWorkersAction)} + className="px-3 py-1.5 rounded-md border border-zinc-300 dark:border-zinc-700 text-sm hover:bg-zinc-100 dark:hover:bg-zinc-800 disabled:opacity-50" + > + Pause all + </button> + )} + {paused && ( + <span role="status" className="text-sm text-amber-700 dark:text-amber-300"> + All workers paused — batches are waiting. Resume restores each + worker&apos;s previous state. + </span> + )} + </div> + + {workers.length === 0 ? ( + <p className="text-sm text-zinc-500 border border-dashed border-zinc-300 dark:border-zinc-700 rounded p-4"> + No workers configured. Add one on the{" "} + <a href="/settings" className="underline"> + Settings + </a>{" "} + page. + </p> + ) : ( + <ul className="flex flex-col gap-2"> + {workers.map((w) => ( + <WorkerCard key={w.id} worker={w} pending={pending} run={run} /> + ))} + </ul> + )} + </div> + ); +} + +function stateBadge(w: WorkerView): { label: string; className: string } { + if (w.degraded) { + return { + label: "degraded", + className: + "bg-red-100 dark:bg-red-950 text-red-800 dark:text-red-200 border-red-300 dark:border-red-800", + }; + } + switch (w.state) { + case "enabled": + return { + label: w.inUse > 0 ? "busy" : "idle", + className: + "bg-green-100 dark:bg-green-950 text-green-800 dark:text-green-200 border-green-300 dark:border-green-800", + }; + case "draining": + return { + label: "draining", + className: + "bg-amber-100 dark:bg-amber-950 text-amber-800 dark:text-amber-200 border-amber-300 dark:border-amber-800", + }; + default: + return { + label: "disabled", + className: + "bg-zinc-100 dark:bg-zinc-800 text-zinc-600 dark:text-zinc-300 border-zinc-300 dark:border-zinc-700", + }; + } +} + +function WorkerCard({ + worker: w, + pending, + run, +}: { + worker: WorkerView; + pending: boolean; + run: (action: () => Promise<unknown>) => void; +}) { + const badge = stateBadge(w); + const engine = w.kind === "remote" ? "remote" : (w.appId ?? "local"); + return ( + <li + aria-label={`worker ${w.name}`} + className="flex flex-col gap-2 rounded border border-zinc-200 dark:border-zinc-800 p-3 bg-white dark:bg-zinc-900" + > + <div className="flex flex-wrap items-center gap-2"> + <span className="font-medium text-sm">{w.name}</span> + <span + className={`text-[11px] px-1.5 py-0.5 rounded-full border ${badge.className}`} + > + {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) && ( + <button + type="button" + disabled={pending} + onClick={() => run(() => enableWorkerAction(w.id))} + aria-label={`enable ${w.name}`} + className="px-2 py-1 rounded border border-green-300 dark:border-green-800 text-xs text-green-700 dark:text-green-300 hover:bg-green-50 dark:hover:bg-green-950 disabled:opacity-50" + > + Enable + </button> + )} + {w.state === "enabled" && !w.degraded && ( + <> + <button + type="button" + disabled={pending || w.inUse === 0} + onClick={() => run(() => drainWorkerAction(w.id))} + aria-label={`drain ${w.name}`} + title="Stop taking new work; let in-flight transcriptions finish" + className="px-2 py-1 rounded border border-amber-300 dark:border-amber-800 text-xs text-amber-700 dark:text-amber-300 hover:bg-amber-50 dark:hover:bg-amber-950 disabled:opacity-50" + > + Drain + </button> + <button + type="button" + disabled={pending} + onClick={() => run(() => disableWorkerAction(w.id))} + aria-label={`disable ${w.name}`} + title="Stop immediately (in-flight transcriptions keep running)" + 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" + > + Disable + </button> + </> + )} + {w.state === "draining" && ( + <button + type="button" + disabled={pending} + onClick={() => run(() => disableWorkerAction(w.id))} + aria-label={`force disable ${w.name}`} + 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" + > + Disable now + </button> + )} + </span> + </div> + {w.degraded && ( + <p className="text-xs text-red-700 dark:text-red-300"> + Auto-disabled after repeated failures. Enable to retry. + </p> + )} + {w.tasks.length > 0 && ( + <ul className="flex flex-col gap-1"> + {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> + ))} + </ul> + )} + </li> + ); +} diff --git a/editor/app/workers/page.tsx b/editor/app/workers/page.tsx @@ -0,0 +1,38 @@ +import type { Metadata } from "next"; +import Link from "next/link"; +import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; +import { WorkersView, type WorkersPayload } from "./components/WorkersView"; + +export const dynamic = "force-dynamic"; + +export const metadata: Metadata = { title: "Workers" }; + +export default async function WorkersPage() { + const pool = getWorkerPool(); + const initial: WorkersPayload = { + paused: pool.isPaused(), + workers: pool.summary().map((w) => ({ ...w, tasks: [] })), + }; + return ( + <div className="flex flex-col gap-4"> + <div className="flex items-center justify-between"> + <h1 className="text-2xl font-semibold">Workers</h1> + <Link href="/settings" className="text-sm underline"> + Configure workers + </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{" "} + <Link href="/settings" className="underline"> + Settings + </Link>{" "} + page. + </p> + <WorkersView initial={initial} /> + </div> + ); +} diff --git a/editor/e2e/settings.spec.ts b/editor/e2e/settings.spec.ts @@ -56,33 +56,46 @@ test("saves global default social links", async ({ page }) => { expect(saved.socialLinks?.[0].svg).toContain("currentColor"); }); -test("defaults parallel transcriptions to 2 and persists a new value", async ({ +test("shows the migrated default worker and persists an added worker", async ({ page, }) => { await page.goto("/settings"); - await expect(page.getByLabel(/parallel transcriptions/i)).toHaveValue("2"); + // 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"); + + // Add a second (CPU) worker, name it, give it 1 slot. + 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(/parallel transcriptions/i).fill("3"); await page.getByRole("button", { name: /save settings/i }).click(); await expect( page.getByRole("status").filter({ hasText: "Saved" }), ).toBeVisible(); - const saved = await readJson<{ parallelTranscriptions?: number }>( - "test-settings.json", - ); - expect(saved.parallelTranscriptions).toBe(3); + const saved = await readJson<{ + workers?: { name: string; appId?: string; slots: number; 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); + // List order is priority. + expect(saved.workers?.[1].priority).toBe(1); await page.reload(); - await expect(page.getByLabel(/parallel transcriptions/i)).toHaveValue("3"); + await expect(page.getByLabel("worker 2 name")).toHaveValue("CPU chough"); }); -test("rejects out-of-range parallel transcriptions", async ({ page }) => { +test("rejects saving with no enabled workers", async ({ page }) => { await page.goto("/settings"); - await page.getByLabel(/parallel transcriptions/i).fill("99"); + await page.getByLabel("worker 1 enabled").uncheck(); await page.getByRole("button", { name: /save settings/i }).click(); await expect(page.locator("form").getByRole("alert")).toContainText( - "between", + /at least one worker must be enabled/i, ); }); diff --git a/editor/e2e/workers.spec.ts b/editor/e2e/workers.spec.ts @@ -0,0 +1,145 @@ +// The Workers page: live per-worker status plus the runtime enable/disable/drain +// controls and the temporary pause-all/resume-all. These are transient operator +// overrides on the global worker pool — they don't touch settings.json. + +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { pathExists, resetData, resolvePath, writeSettings } from "./helpers"; + +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`); + } +} + +async function transcriptCount(slug: string, ids: string[]): Promise<number> { + let n = 0; + for (const id of ids) { + if ( + await pathExists(`test-transcripts/channels/${slug}/data/${id}/transcript.json`) + ) { + n++; + } + } + return n; +} + +// Two enabled local workers 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: {} }, + ], +}; + +test.beforeEach(async () => { + await resetData("empty"); +}); + +test("lists configured workers with their slot counts", 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"); +}); + +test("disable then enable a worker flips its state", async ({ page }) => { + await writeSettings(TWO_WORKERS); + await page.goto("/workers"); + const gpu = page.getByRole("listitem").filter({ hasText: "GPU" }); + + await gpu.getByRole("button", { name: "disable GPU" }).click(); + await expect(gpu).toContainText("disabled"); + + await gpu.getByRole("button", { name: "enable GPU" }).click(); + await expect(gpu).toContainText(/idle|busy/); +}); + +test("pause all disables every worker; resume restores them", async ({ + page, +}) => { + await writeSettings(TWO_WORKERS); + await page.goto("/workers"); + + // Put one worker into a non-default state first so resume must restore it. + const gpu = page.getByRole("listitem").filter({ hasText: "GPU" }); + await gpu.getByRole("button", { name: "disable GPU" }).click(); + await expect(gpu).toContainText("disabled"); + + await page.getByRole("button", { name: "Pause all" }).click(); + await expect(page.getByRole("status")).toContainText(/paused/i); + const cpu = page.getByRole("listitem").filter({ hasText: "CPU" }); + await expect(cpu).toContainText("disabled"); + + await page.getByRole("button", { name: "Resume all workers" }).click(); + // CPU returns to enabled (idle); GPU stays disabled — its pre-pause state. + await expect(cpu).toContainText(/idle|busy/); + await expect(gpu).toContainText("disabled"); +}); + +test("pausing all workers pauses a running batch instead of failing it; resume completes it", async ({ + page, +}) => { + test.setTimeout(90_000); + // One worker, one slot, so videos transcribe one at a time and parking is easy + // 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: {} }, + ], + }); + const ids = ["slowop1", "slowop2", "slowop3"]; + await makeTranscribeChannel("pause-batch", ids); + + await page.goto("/channels/pause-batch"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + + // Wait until the batch is running with its first transcription in flight. + await page.goto("/jobs/active"); + await expect( + page.getByRole("progressbar", { name: /Transcribing slowop/ }).first(), + ).toBeVisible({ timeout: 15_000 }); + + // Pause all workers. The in-flight one finishes (disable lets it complete) but + // no new transcription starts — the batch parks, it does NOT fail. + await page.goto("/workers"); + await page.getByRole("button", { name: "Pause all" }).click(); + await expect(page.getByRole("status")).toContainText(/paused/i); + + // Give it time: at most the in-flight transcript lands, then progress stalls. + await page.waitForTimeout(6_000); + const whilePaused = await transcriptCount("pause-batch", ids); + expect(whilePaused).toBeLessThan(ids.length); + + // The batch is still running (paused), not failed/done: the Active Jobs page + // only lists non-terminal jobs, so the pause-batch section is still present. + await page.goto("/jobs/active"); + const section = page.locator( + "section[aria-label='Active jobs for pause-batch']", + ); + await expect(section.getByText("running", { exact: true })).toBeVisible({ + timeout: 10_000, + }); + + // Resume → the remaining videos transcribe and the batch finishes. + await page.goto("/workers"); + await page.getByRole("button", { name: "Resume all workers" }).click(); + + await expect + .poll(() => transcriptCount("pause-batch", ids), { timeout: 60_000 }) + .toBe(ids.length); +});