Archilyzer · Source

archilyzer

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

commit d21d6f927a6e2e851b3e053b93f25e4252fa1bf2
parent c778ea20aa674e0b46ec22420ccfd818a8d11882
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 14 Jun 2026 02:30:45 -0400

Transcription workers: resumable parakeet + pause + device (Phase 7)

Parakeet is now resumable, yt-dlp style, and pausable mid-run:
- scripts/parakeet-stitch.mjs writes each window's raw parakeet-cli JSON to a
  per-audio work dir (".<audio>.parakeet/") as it completes, validates the cache
  against the windowing (meta.json), and stitches the final transcript only once
  ALL windows are done — then removes the work dir. Re-running resumes from the
  cached windows. SIGTERM finishes the in-flight window and exits WITHOUT
  stitching (paused); completed windows stay cached for the next run.
- transcribeOne: a partial-capable engine that exits with no output after a
  pause request returns the new "paused" outcome (not a failure). whisperBatch
  counts it as skipped so the video stays untranscribed and the next run resumes.
- Workers page: a busy parakeet worker shows "Stop & keep progress" (pool
  activeStop registry + partialStopWorker + stopWorkerPartialAction), which
  finishes the current window, frees the worker, and leaves progress cached.

Selectable parakeet device:
- per-worker Device config (AppInstanceConfig.device) → parakeet-cli --device /
  PARAKEET_DEVICE; surfaced in the worker editor for the parakeet engine.

e2e: pause caches the window and writes no transcript, then a re-run resumes to a
complete transcript and cleans the cache; device persists. Typechecks clean;
parakeet + whisper + workers + remote specs green.

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

Diffstat:
Mcommon/controller/transcribeOne.ts | 41++++++++++++++++++++++++++++++++++++++++-
Mcommon/controller/whisperBatch.ts | 11+++++++++--
Mcommon/jobs/workerPool.ts | 30++++++++++++++++++++++++++++++
Mcommon/lib/transcriptionApps.ts | 15++++++++++++++-
Mcommon/lib/workers.ts | 1+
Meditor/CHANGELOG.md | 2++
Meditor/app/api/workers/route.ts | 3+++
Meditor/app/settings/components/WorkersField.tsx | 9+++++++++
Meditor/app/workers/actions.ts | 12++++++++++++
Meditor/app/workers/components/WorkersView.tsx | 15+++++++++++++++
Meditor/app/workers/page.tsx | 6+++++-
Meditor/e2e/fixtures/bin/fake-parakeet-stitch.mjs | 105+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Aeditor/e2e/parakeet-partial.spec.ts | 87+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mscripts/parakeet-stitch.mjs | 122+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
14 files changed, 407 insertions(+), 52 deletions(-)

diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -48,9 +48,15 @@ export type TranscribeOneOptions = { // onLog → the engine's parser instead). onProgress?: (patch: { fraction?: number; detail?: string }) => void; signal?: AbortSignal; + // Graceful "pause & keep progress": when aborted and the engine supports it + // (parakeet), the child is SIGTERM'd so it finishes the current window and + // exits without producing a transcript — its completed windows are cached for + // resume. transcribeOneVideo returns "paused" so the batch leaves the video + // untranscribed (next run resumes). Ignored by engines without partial support. + partialSignal?: AbortSignal; }; -export type TranscribeOneOutcome = "transcribed" | "already-exists"; +export type TranscribeOneOutcome = "transcribed" | "already-exists" | "paused"; // TranscribeError lives in its own module so the remote client can throw it // without an import cycle; re-exported here for existing import sites. @@ -137,9 +143,31 @@ export async function transcribeOneVideo( env: build.env ? { ...process.env, ...build.env } : undefined, }); child.all?.on("data", (c: Buffer) => log(c.toString("utf8"))); + // "Pause & keep progress": for an engine that handles SIGTERM by stopping + // gracefully and caching completed work (parakeet), send it on partialSignal. + // execa's cancelSignal is the separate hard-cancel `opts.signal`, so this + // SIGTERM lets the child exit 0 and `await child` resolves normally. The + // engine writes NO transcript on pause (it stitches only when fully done), + // so a missing output + a requested pause is a resumable pause, not a failure. + let pauseRequested = false; + if (app.supportsPartialStop && opts.partialSignal) { + const stop = () => { + pauseRequested = true; + log(`Pausing ${worker.id} on ${opts.videoId} (progress kept for resume)…`); + child.kill("SIGTERM"); + }; + if (opts.partialSignal.aborted) stop(); + else opts.partialSignal.addEventListener("abort", stop, { once: true }); + } await child; const producedPath = path.join(opts.videoDir, build.outputFile); if (!(await pathExists(producedPath))) { + if (pauseRequested) { + log( + `Transcribe ${opts.videoId} paused; completed windows cached — re-run to resume.`, + ); + return "paused"; + } throw new TranscribeError( `transcription with ${app.id} produced no ${build.outputFile} in ${opts.videoDir}`, "transcription", @@ -236,6 +264,15 @@ export async function transcribeWithWorker( appId: worker.appId, workerId: worker.id, }); + // For a partial-capable local engine, register a graceful "stop & keep + // partial" handle the Workers page can trigger for this worker. + const partialCapable = + worker.kind === "local" && + getTranscriptionApp(worker.appId).supportsPartialStop === true; + const partialController = partialCapable ? new AbortController() : undefined; + if (partialController) { + pool.setActiveStop(worker.id, () => partialController.abort()); + } try { const outcome = await transcribeOneVideo({ paths: opts.paths, @@ -247,6 +284,7 @@ export async function transcribeWithWorker( onLog: task ? task.onLog : log, onProgress: task ? task.update : undefined, signal: opts.signal, + partialSignal: partialController?.signal, }); pool.markSuccess(worker.id); return outcome; @@ -278,6 +316,7 @@ export async function transcribeWithWorker( } throw err; } finally { + if (partialController) pool.clearActiveStop(worker.id); task?.end(); lease.release(); } diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts @@ -195,7 +195,7 @@ export async function runWhisperBatch({ // 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({ + const outcome = await transcribeWithWorker({ paths, videoDir: videoPath, videoId: videoDir, @@ -208,7 +208,14 @@ export async function runWhisperBatch({ signal, drainSignal, }); - succeededCount++; + if (outcome === "paused") { + // Parakeet was paused with progress cached; leave the video untranscribed + // (not failed) so the next run resumes it. + log(`Transcribe ${videoDir} paused (will resume next run)`); + skipped++; + } else { + succeededCount++; + } } catch (err) { // Cancelled or drained while parked/in-flight: a skip, not a failure. if ( diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts @@ -53,6 +53,9 @@ type PoolEntry = { consecutiveFailures: number; // Set when a worker is removed from settings while still busy: drain then drop. retireWhenIdle: boolean; + // Registered by the in-flight transcription (if its engine supports it) to + // request a graceful "stop & keep partial". Cleared when the lease releases. + activeStop?: () => void; }; type Waiter = { @@ -300,6 +303,33 @@ class WorkerPool { return true; } + // --- Partial-stop registry: the in-flight transcription on a worker can + // register a callback to be gracefully stopped with a partial result. --- + + setActiveStop(id: string, fn: () => void): void { + const entry = this.entries.get(id); + if (entry) entry.activeStop = fn; + } + + clearActiveStop(id: string): void { + const entry = this.entries.get(id); + if (entry) entry.activeStop = undefined; + } + + // True if the worker had an active transcription that supports partial-stop and + // was asked to stop. + partialStopWorker(id: string): boolean { + const entry = this.entries.get(id); + if (!entry?.activeStop) return false; + entry.activeStop(); + return true; + } + + // Whether a worker currently has a partial-stoppable transcription in flight. + canStopPartial(id: string): boolean { + return !!this.entries.get(id)?.activeStop; + } + // --- Failure tracking (auto-disable, Phase 5) --- markSuccess(id: string): void { diff --git a/common/lib/transcriptionApps.ts b/common/lib/transcriptionApps.ts @@ -35,6 +35,9 @@ export type AppInstanceConfig = { // whisper.cpp custom argv template using the {audioFile}/{outputBase}/{model} // placeholders. Undefined = DEFAULT_TRANSCRIBE_ARGS. customArgs?: string[]; + // parakeet compute device passed to parakeet-cli (--device / PARAKEET_DEVICE), + // e.g. "cuda:0", "cpu". Undefined = parakeet-cli's default device. + device?: string; }; export type TranscribeBuild = { @@ -65,11 +68,16 @@ export type TranscriptionApp = { remoteUrl?: boolean; chunkSize?: boolean; customArgs?: boolean; + device?: boolean; }; // Default binary when the per-app `bin` override is empty. defaultBin: () => string; build: (input: TranscribeBuildInput) => TranscribeBuild; makeProgressParser: () => { feed: (line: string) => ProgressUpdate | null }; + // True when the engine can be stopped mid-run and still produce a usable + // partial transcript (parakeet stitches completed windows on SIGTERM). Surfaces + // a "Stop & keep partial" control on the Workers page. + supportsPartialStop?: boolean; }; // --------------------------------------------------------------------------- @@ -208,7 +216,10 @@ const chough: TranscriptionApp = { const parakeet: TranscriptionApp = { id: "parakeet", label: "parakeet.cpp (overlapping segments)", - fields: { model: true, chunkSize: true }, + fields: { model: true, chunkSize: true, device: true }, + // The overlapping-segment wrapper can stop after the current window and stitch + // a partial transcript on SIGTERM — so a long run is interruptible. + supportsPartialStop: true, defaultBin: () => getPaths().parakeetBin, build({ audioFile, outputBase, config }) { const paths = getPaths(); @@ -217,6 +228,8 @@ const parakeet: TranscriptionApp = { if (typeof config.chunkSize === "number" && config.chunkSize > 0) { argv.push("--segment", String(Math.floor(config.chunkSize))); } + const device = config.device?.trim(); + if (device) argv.push("--device", device); argv.push(audioFile); return { argv, diff --git a/common/lib/workers.ts b/common/lib/workers.ts @@ -73,6 +73,7 @@ export function sanitizeWorkerConfig(value: unknown): AppInstanceConfig { cfg.chunkSize = Math.floor(r.chunkSize); } if (Array.isArray(r.customArgs)) cfg.customArgs = r.customArgs.map(String); + if (typeof r.device === "string") cfg.device = r.device; return cfg; } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -6,6 +6,8 @@ - **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. +- **parakeet.cpp transcriptions are resumable, and can be paused mid-run.** The overlapping-segment wrapper now writes each window's raw parakeet-cli JSON to a per-audio work dir (`.<audio>.parakeet/`) as it finishes, and stitches the final transcript only once *all* windows are done (then removes the work dir). Re-running the same transcription picks up the cached windows and only does what's missing — yt-dlp-style resume, so a crash, cancel, or pause never loses completed windows. A busy parakeet worker on the Workers page shows a **Stop & keep progress** button: it finishes the in-flight window, stops and frees the worker (no transcript written yet), and the next "Transcribe missing" resumes from the cached windows and completes. Handy to reclaim a GPU mid-run. (whisper.cpp/chough run as a single pass and don't offer this.) +- **Selectable compute device for parakeet.cpp.** A parakeet worker gained a **Device** field (e.g. `cuda:0`, `cpu`) passed through to `parakeet-cli` (`--device`, also honored as the `PARAKEET_DEVICE` env). Combined with one-worker-per-slot and Copy, you can pin different parakeet workers to different GPUs. - **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`. diff --git a/editor/app/api/workers/route.ts b/editor/app/api/workers/route.ts @@ -44,6 +44,9 @@ export async function GET() { const workers = summary.map((w) => ({ ...w, tasks: byWorker.get(w.id) ?? [], + // True when the in-flight transcription can be stopped into a partial result + // (parakeet). Drives the "Stop & keep partial" button on the Workers page. + canStopPartial: pool.canStopPartial(w.id), })); return NextResponse.json({ paused: pool.isPaused(), workers }); diff --git a/editor/app/settings/components/WorkersField.tsx b/editor/app/settings/components/WorkersField.tsx @@ -300,6 +300,15 @@ function WorkerCard({ id={`${uid}-remoteurl`} /> )} + {app?.fields.device && ( + <CardField + label="Device" + value={cfg.device ?? ""} + onChange={(v) => onPatchConfig({ device: v })} + hint="parakeet compute device passed to parakeet-cli (--device), e.g. cuda:0 or cpu. Blank = parakeet-cli's default." + id={`${uid}-device`} + /> + )} {app?.fields.chunkSize && ( <CardField label="Chunk size (seconds)" diff --git a/editor/app/workers/actions.ts b/editor/app/workers/actions.ts @@ -48,3 +48,15 @@ export async function resumeAllWorkersAction(): Promise<WorkerActionResult> { refresh(); return { ok: true }; } + +// Gracefully stop the worker's in-flight transcription and keep the partial +// result (parakeet finishes the current window, stitches, and writes a partial). +export async function stopWorkerPartialAction( + id: string, +): Promise<WorkerActionResult> { + const ok = getWorkerPool().partialStopWorker(id); + refresh(); + return ok + ? { ok: true } + : { ok: false, error: `Worker "${id}" has nothing to stop` }; +} diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx @@ -10,6 +10,7 @@ import { enableWorkerAction, pauseAllWorkersAction, resumeAllWorkersAction, + stopWorkerPartialAction, } from "../actions"; export type WorkerTask = { @@ -43,6 +44,8 @@ export type WorkerView = { state: WorkerRuntimeState; degraded: boolean; enabled: boolean; + // True when the in-flight transcription can be stopped into a partial result. + canStopPartial: boolean; tasks: WorkerTask[]; }; @@ -205,6 +208,18 @@ function WorkerCard({ Enable </button> )} + {w.busy && w.canStopPartial && ( + <button + type="button" + disabled={pending} + onClick={() => run(() => stopWorkerPartialAction(w.id))} + aria-label={`stop ${w.name} keep progress`} + title="Finish the current window, then stop and free the worker; completed windows are cached and resume on the next run" + 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" + > + Stop &amp; keep progress + </button> + )} {w.state === "enabled" && !w.degraded && ( <> <button diff --git a/editor/app/workers/page.tsx b/editor/app/workers/page.tsx @@ -11,7 +11,11 @@ export default async function WorkersPage() { const pool = getWorkerPool(); const initial: WorkersPayload = { paused: pool.isPaused(), - workers: pool.summary().map((w) => ({ ...w, tasks: [] })), + workers: pool.summary().map((w) => ({ + ...w, + tasks: [], + canStopPartial: pool.canStopPartial(w.id), + })), }; return ( <div className="flex flex-col gap-4"> diff --git a/editor/e2e/fixtures/bin/fake-parakeet-stitch.mjs b/editor/e2e/fixtures/bin/fake-parakeet-stitch.mjs @@ -1,12 +1,25 @@ #!/usr/bin/env node // E2E fake for scripts/parakeet-stitch.mjs (the parakeet overlapping-segment // wrapper). Mimics the args the parakeet app builds (transcriptionApps.ts): -// parakeet-stitch.mjs --model <m> --output <tmpBase> [--segment <n>] <audioFile> -// and the wrapper's output contract: write chough-native JSON to EXACTLY the -// --output path (no extension appended), and print "segment X/Y" progress lines -// to stderr (consumed by createParakeetProgressParser). Run from cwd == the -// video dir. Stays instant — it never shells out to ffmpeg/parakeet-cli. -import { writeFile } from "node:fs/promises"; +// parakeet-stitch.mjs --model <m> --output <tmpBase> [--segment <n>] +// [--device <d>] <audioFile> +// Output contract: write chough-native JSON to EXACTLY the --output path; print +// "segment X/Y" progress lines to stderr. Run from cwd == the video dir. +// +// Mirrors the real wrapper's resumable design: each window's payload is written +// to a per-audio work dir (".<audio>.parakeet/") as it completes; the stitched +// transcript is produced only when ALL windows are done (then the work dir is +// removed); re-running resumes from cached windows. On SIGTERM it stops without +// producing output (paused) — completed windows stay cached. +// +// Two modes by video dir (cwd): +// * default: do both windows instantly and stitch. +// * "slowop" dirs: window 0 is instant; window 1 waits — on SIGTERM it pauses +// (no output, window 0 cached). A *resume* run (window 0 already cached) +// finishes window 1 immediately and stitches. +import { mkdir, writeFile, rm } from "node:fs/promises"; +import { existsSync } from "node:fs"; +import path from "node:path"; const argv = process.argv.slice(2); function arg(flag) { @@ -15,31 +28,71 @@ function arg(flag) { } const out = arg("--output") ?? arg("-o"); +const device = arg("--device"); const audio = argv[argv.length - 1]; if (!out) { process.stderr.write(`[fake-parakeet-stitch] missing --output\n`); process.exit(2); } +if (device) process.stderr.write(`parakeet-stitch: device ${device}\n`); -// Two synthetic windows so the stitched output and progress lines are exercised. -// Segment 1 has no ETA (no average yet); segment 2 carries the per-video ETA -// (avg time per segment × remaining segments), matching the real wrapper. -process.stderr.write(`parakeet-stitch: ${audio}: 10s -> 2 window(s)\n`); -process.stderr.write(`parakeet-stitch: segment 1/2 @0s — transcribing\n`); -process.stderr.write( - `parakeet-stitch: segment 2/2 @5s — transcribing (avg 1s/seg, ETA 0:01)\n`, -); +const slow = process.cwd().toLowerCase().includes("slowop"); +const workDir = path.resolve(`.${path.basename(audio)}.parakeet`); +const win = (i) => path.join(workDir, `win-${String(i).padStart(4, "0")}.json`); +const cue = (i) => ({ + start_time: i * 5, + end_time: i * 5 + 5, + text: `Synthetic parakeet output for ${audio} #${i + 1}`, +}); -const doc = { - duration_seconds: 10, - chunks: 2, - text: `Synthetic parakeet output for ${audio} #1 Synthetic parakeet output for ${audio} #2`, - chunk_data: [ - { start_time: 0, end_time: 5, text: `Synthetic parakeet output for ${audio} #1` }, - { start_time: 5, end_time: 10, text: `Synthetic parakeet output for ${audio} #2` }, - ], -}; +async function stitchAndFinish() { + const chunk_data = [cue(0), cue(1)]; + await writeFile( + out, + JSON.stringify({ + duration_seconds: 10, + chunks: 2, + text: chunk_data.map((c) => c.text).join(" "), + chunk_data, + }), + ); + await rm(workDir, { recursive: true, force: true }); + process.stderr.write(`parakeet-stitch: wrote 2 cues -> ${out}\n`); + process.exit(0); +} + +await mkdir(workDir, { recursive: true }); +const resuming = existsSync(win(0)); -// The wrapper writes EXACTLY the --output path (same as chough's -o). -await writeFile(out, JSON.stringify(doc)); -process.stderr.write(`parakeet-stitch: wrote 2 cues -> ${out}\n`); +// Window 0. +if (existsSync(win(0))) { + process.stderr.write(`parakeet-stitch: segment 1/2 @0s — cached\n`); +} else { + process.stderr.write(`parakeet-stitch: segment 1/2 @0s — transcribing\n`); + await writeFile(win(0), JSON.stringify({ start: 0, cue: cue(0) })); +} + +// Window 1. +if (existsSync(win(1))) { + process.stderr.write(`parakeet-stitch: segment 2/2 @5s — cached\n`); + await stitchAndFinish(); +} else if (!slow || resuming) { + // Fresh-and-fast, or resuming a paused run: finish window 1 now. + process.stderr.write(`parakeet-stitch: segment 2/2 @5s — transcribing\n`); + await writeFile(win(1), JSON.stringify({ start: 5, cue: cue(1) })); + await stitchAndFinish(); +} else { + // Slow first run: window 1 is interruptible. On SIGTERM, pause (no output) — + // window 0 stays cached for the resume run. Backstop completes if never stopped. + process.stderr.write(`parakeet-stitch: segment 2/2 @5s — transcribing\n`); + const onPause = () => { + process.stderr.write(`parakeet-stitch: paused after 1/2 window(s) — re-run to resume\n`); + process.exit(0); + }; + process.on("SIGTERM", onPause); + process.on("SIGINT", onPause); + setTimeout(async () => { + await writeFile(win(1), JSON.stringify({ start: 5, cue: cue(1) })); + await stitchAndFinish(); + }, 30_000); +} diff --git a/editor/e2e/parakeet-partial.spec.ts b/editor/e2e/parakeet-partial.spec.ts @@ -0,0 +1,87 @@ +// Phase 7: parakeet.cpp workers can be stopped into a partial transcript, and +// their compute device is configurable. The fake parakeet wrapper +// (fixtures/bin/fake-parakeet-stitch.mjs) stays alive for "slowop" ids and writes +// a partial (partial:true, one window) on SIGTERM. + +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { pathExists, readJson, 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`); + } +} + +test.beforeEach(async () => { + await resetData("empty"); +}); + +test("Stop & keep progress pauses a parakeet run, caches the window, and resumes", async ({ + page, +}) => { + test.setTimeout(60_000); + await writeSettings({ + workers: [ + { id: "gpu", name: "GPU parakeet", kind: "local", enabled: true, priority: 0, appId: "parakeet", config: {} }, + ], + }); + await makeTranscribeChannel("partial-chan", ["slowoppk1"]); + const dir = "test-transcripts/channels/partial-chan/data/slowoppk1"; + const cachedWindow = `${dir}/.audio.mp3.parakeet/win-0000.json`; + const transcript = `${dir}/transcript.json`; + + await page.goto("/channels/partial-chan"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + + // While the parakeet worker runs, the Workers page offers "Stop & keep progress". + await page.goto("/workers"); + const gpu = page.getByRole("listitem").filter({ hasText: "GPU parakeet" }); + const stop = gpu.getByRole("button", { name: /stop GPU parakeet keep progress/i }); + await expect(stop).toBeVisible({ timeout: 15_000 }); + await stop.click(); + + // Pause caches the completed window but writes no transcript yet. + await expect.poll(() => pathExists(cachedWindow), { timeout: 30_000 }).toBe(true); + expect(await pathExists(transcript)).toBe(false); + + // Re-running resumes from the cached window and completes the transcript. + await page.goto("/channels/partial-chan"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + await expect.poll(() => pathExists(transcript), { timeout: 30_000 }).toBe(true); + const doc = await readJson<{ chunk_data?: unknown[] }>(transcript); + expect(doc.chunk_data?.length).toBe(2); + // The work dir is cleaned once stitching succeeds. + expect(await pathExists(cachedWindow)).toBe(false); +}); + +test("a parakeet worker's device is configurable and persists", async ({ + page, +}) => { + await page.goto("/settings"); + // Migration makes worker 1 whisper-cpp; switch it to parakeet to reveal Device. + await page.getByLabel("worker 1 engine").selectOption("parakeet"); + await page.getByLabel(/^Device/).fill("cuda:0"); + await page.getByRole("button", { name: /save settings/i }).click(); + await expect( + page.getByRole("status").filter({ hasText: "Saved" }), + ).toBeVisible(); + + const saved = await readJson<{ + workers?: { appId?: string; config?: { device?: string } }[]; + }>("test-settings.json"); + expect(saved.workers?.[0].appId).toBe("parakeet"); + expect(saved.workers?.[0].config?.device).toBe("cuda:0"); +}); diff --git a/scripts/parakeet-stitch.mjs b/scripts/parakeet-stitch.mjs @@ -33,6 +33,15 @@ // --overlap <sec> window overlap (PARAKEET_OVERLAP_SEC, default 6) // --decoder <ctc|tdt> passed through to parakeet-cli (PARAKEET_DECODER) // --lang <locale> passed through to parakeet-cli (PARAKEET_LANG) +// --device <dev> compute device passed to parakeet-cli (PARAKEET_DEVICE) +// +// Resumable, yt-dlp style: each window's raw parakeet-cli JSON is written to a +// per-audio work dir (".<audio>.parakeet/" beside the audio) as it completes, +// and the final transcript is stitched only once ALL windows are done (then the +// work dir is removed). Re-running the same command picks up the cached window +// JSONs and only transcribes what's missing. On SIGTERM/SIGINT it finishes the +// in-flight window, leaves the completed window JSONs in place, and exits WITHOUT +// stitching — so a long run can be paused and resumed later with no lost work. // --gap <sec> start a new cue when the inter-word gap exceeds this (default 0.8) // --max-cue <sec> cap a single cue's duration (default 8) // -o, --output <f> output file (alternative to the positional arg) @@ -41,7 +50,7 @@ import { execFile } from "node:child_process"; import { promisify } from "node:util"; -import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { existsSync } from "node:fs"; import os from "node:os"; import path from "node:path"; @@ -77,6 +86,7 @@ function parseArgs(argv) { overlap: numEnv(process.env.PARAKEET_OVERLAP_SEC, 6), decoder: process.env.PARAKEET_DECODER ?? "", lang: process.env.PARAKEET_LANG ?? "", + device: process.env.PARAKEET_DEVICE ?? "", gap: 0.8, maxCue: 8, keepTemp: false, @@ -101,6 +111,7 @@ function parseArgs(argv) { case "--overlap": opts.overlap = Number(next()); break; case "--decoder": opts.decoder = next(); break; case "--lang": opts.lang = next(); break; + case "--device": opts.device = next(); break; case "--gap": opts.gap = Number(next()); break; case "--max-cue": opts.maxCue = Number(next()); break; case "-o": case "--output": opts.output = next(); break; @@ -182,6 +193,7 @@ async function transcribeWindow(opts, wav) { const args = ["transcribe", "--model", opts.model, "--timestamps", "--json", "--input", wav]; if (opts.decoder) args.push("--decoder", opts.decoder); if (opts.lang) args.push("--lang", opts.lang); + if (opts.device) args.push("--device", opts.device); const { stdout } = await execFileP(opts.cli, args, { maxBuffer: 256 * 1024 * 1024 }); return parseCliJson(stdout); } @@ -218,6 +230,18 @@ function groupCues(words, { gap, maxCue }) { })); } +// Resumable per-audio work dir: ".<audio>.parakeet/" beside the audio file. Holds +// one raw parakeet-cli JSON per completed window plus a meta.json describing the +// windowing, so a re-run reuses finished windows and a pause loses nothing. +function workDirFor(absAudio) { + return path.join( + path.dirname(absAudio), + `.${path.basename(absAudio)}.parakeet`, + ); +} + +const PARTIAL_PAUSE_CODE = 0; + async function main() { const opts = parseArgs(process.argv.slice(2)); if (!opts.audio) fail("missing <audioFile>\n\n" + HELP); @@ -242,19 +266,64 @@ async function main() { const N = offsets.length; progress(`${path.basename(opts.audio)}: ${round3(duration)}s -> ${N} window(s) of ${opts.segment}s (overlap ${overlap}s)`); + // Resumable work dir. Validate cached windows against the current windowing + // (duration/segment/overlap); if anything changed, start fresh. + const absAudio = path.resolve(opts.audio); + const workDir = workDirFor(absAudio); + const metaPath = path.join(workDir, "meta.json"); + const meta = { v: 1, duration: round3(duration), segment: opts.segment, overlap, n: N }; + await mkdir(workDir, { recursive: true }); + let cachedMeta = null; + try { + cachedMeta = JSON.parse(await readFile(metaPath, "utf8")); + } catch { + cachedMeta = null; + } + const compatible = + cachedMeta && + cachedMeta.v === meta.v && + cachedMeta.duration === meta.duration && + cachedMeta.segment === meta.segment && + cachedMeta.overlap === meta.overlap && + cachedMeta.n === meta.n; + if (!compatible) { + await rm(workDir, { recursive: true, force: true }); + await mkdir(workDir, { recursive: true }); + } + await writeFile(metaPath, JSON.stringify(meta)); + const winPath = (i) => path.join(workDir, `win-${String(i).padStart(4, "0")}.json`); + + // Graceful pause: on SIGTERM/SIGINT, finish the in-flight window then stop. + // Completed window JSONs stay in the work dir; we do NOT stitch (the final + // transcript is produced only when ALL windows are done). Re-running resumes. + let stopRequested = false; + const onStop = () => { + if (!stopRequested) progress("pause requested — finishing the current window, then stopping (re-run to resume)"); + stopRequested = true; + }; + process.on("SIGTERM", onStop); + process.on("SIGINT", onStop); + const tmp = await mkdtemp(path.join(os.tmpdir(), "parakeet-stitch-")); - // Per-window absolute-timestamped word lists. - const windowWords = []; - // Measured transcription wall-time (ms) of each completed segment, for ETA. - const segTimes = []; + const segTimes = []; // measured wall-time (ms) per freshly-transcribed window, for ETA + let done = 0; // windows present (cached or freshly done) try { for (let i = 0; i < N; i += 1) { const start = offsets[i]; + // Resume: reuse a window already transcribed in a prior run. + if (existsSync(winPath(i))) { + progress(`segment ${i + 1}/${N} @${round3(start)}s — cached`); + done += 1; + continue; + } + if (stopRequested) { + progress(`paused after ${done}/${N} window(s) — re-run to resume`); + process.exit(PARTIAL_PAUSE_CODE); + } const isLast = i === N - 1; // Last window runs to EOF (no -t) so we never miss a trailing fragment. const len = isLast ? undefined : opts.segment; const wav = path.join(tmp, `win-${String(i).padStart(4, "0")}.wav`); - // remaining = this segment + the ones not yet started. const eta = etaSuffix(segTimes, N - i); progress(`segment ${i + 1}/${N} @${round3(start)}s — slicing${eta}`); await sliceWav(opts.ffmpeg, opts.audio, start, len, wav); @@ -262,29 +331,38 @@ async function main() { const t0 = Date.now(); const doc = await transcribeWindow(opts, wav); segTimes.push(Date.now() - t0); - const words = Array.isArray(doc.words) ? doc.words : []; - // Shift window-relative times into absolute timeline. - windowWords.push( - words.map((w) => ({ - w: w.w ?? w.text ?? "", - start: (Number(w.start) || 0) + start, - end: (Number(w.end) || 0) + start, - conf: w.conf, - })), - ); + // Persist the raw parakeet-cli payload for this window (the "original json"), + // tagged with its absolute start offset. This is the durable, resumable + // artifact; stitching happens only at the end. + await writeFile(winPath(i), JSON.stringify({ start, doc })); + done += 1; if (!opts.keepTemp) await rm(wav, { force: true }); } } finally { if (!opts.keepTemp) await rm(tmp, { recursive: true, force: true }); } - // Boundary cut points: the midpoint of each overlap region. Window i owns - // words whose start is in [cut[i-1], cut[i]); the first/last windows are - // open-ended. This keeps exactly one copy of every word across the seams. + // All windows present — load the raw payloads and stitch the final transcript. + const windowWords = []; + for (let i = 0; i < N; i += 1) { + const { start, doc } = JSON.parse(await readFile(winPath(i), "utf8")); + const words = Array.isArray(doc.words) ? doc.words : []; + windowWords.push( + words.map((w) => ({ + w: w.w ?? w.text ?? "", + start: (Number(w.start) || 0) + start, + end: (Number(w.end) || 0) + start, + conf: w.conf, + })), + ); + } + + // Boundary cut points: the midpoint of each overlap region. Window i owns words + // whose start is in [cut[i-1], cut[i]); the first/last windows are open-ended. + // This keeps exactly one copy of every word across the seams. const cuts = []; for (let i = 0; i < N - 1; i += 1) { - const nextStart = offsets[i + 1]; - cuts.push(nextStart + overlap / 2); + cuts.push(offsets[i + 1] + overlap / 2); } const stitched = []; for (let i = 0; i < N; i += 1) { @@ -313,6 +391,8 @@ async function main() { process.stdout.write(json + "\n"); progress(`emitted ${chunkData.length} cues (${stitched.length} words) to stdout`); } + // Stitch succeeded — drop the per-window cache. + if (!opts.keepTemp) await rm(workDir, { recursive: true, force: true }); } main().catch((err) => {