Archilyzer · Source

archilyzer

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

commit 40349e80dd2522384c2eef4c22bd0293b8c59ecf
parent 92988bd46de8d84b2c4b56d1b907f59eb0796b0b
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Tue, 26 May 2026 18:40:03 -0400

per-operation queue

Diffstat:
Mcommon/controller/whisperBatch.ts | 24+++++++++++++++++++++++-
Acommon/jobs/progressParsers.ts | 88+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/registry.ts | 87+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/streamCommand.ts | 28+++++++++++++++++++++++++++-
Acommon/jobs/taskHooks.ts | 60++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/ytdlp/runYtdlp.ts | 39+++++++++++++++++++++++++++++----------
Meditor/CHANGELOG.md | 2++
Aeditor/app/api/jobs/active/route.ts | 11+++++++++++
Meditor/app/channels/[slug]/pipelineActions.ts | 5++++-
Meditor/app/channels/[slug]/videos/[id]/videoActions.ts | 88++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------------
Meditor/app/channels/[slug]/whisperActions.ts | 9+++++++--
Meditor/app/jobs/actions.ts | 8++++++++
Aeditor/app/jobs/active/buildActiveJobs.ts | 98+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/active/page.tsx | 132+++----------------------------------------------------------------------------
Aeditor/app/jobs/components/ActiveJobsLive.tsx | 95+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/DrainJobButton.tsx | 41+++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/components/RunningJobsList.tsx | 94+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
Meditor/e2e/fixtures/bin/fake-whisper.mjs | 21+++++++++++++++++++++
Meditor/e2e/fixtures/bin/fake-ytdlp.mjs | 13+++++++++++++
Aeditor/e2e/jobs-batch-tasks-drain.spec.ts | 166+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
20 files changed, 927 insertions(+), 182 deletions(-)

diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts @@ -6,6 +6,7 @@ import type { AudioFormat } from "../lib/channelConfig"; import { VTT_FILENAME, WHISPER_FILENAME } from "../lib/videoStatus"; import { transcribeOneVideo } from "./transcribeOne"; import { resolveShardItems } from "./shard"; +import type { TaskTracker } from "../jobs/taskHooks"; const { pathExists, readdir, appendFile, readFile, ensureFile } = fs; @@ -25,6 +26,12 @@ export type WhisperBatchOptions = { ids?: string[]; onLog?: (msg: string) => void; signal?: AbortSignal; + // Soft-cancel: when aborted, stop starting NEW transcriptions but let + // in-flight ones finish. Separate from the hard `signal`. + drainSignal?: AbortSignal; + // When provided, each transcription is tracked as a per-operation task with + // its own parsed progress bar on the Active Jobs screen. + tracker?: TaskTracker; }; export type WhisperBatchResult = { @@ -47,6 +54,8 @@ export async function runWhisperBatch({ ids, onLog, signal, + drainSignal, + tracker, }: WhisperBatchOptions): Promise<WhisperBatchResult> { const log = onLog ?? ((m: string) => console.log(m)); const channelDir = path.join(paths.channelsDir, channelSlug); @@ -99,6 +108,12 @@ export async function runWhisperBatch({ skipped++; return; } + // Drain (soft-cancel): don't start NEW transcriptions, but tasks that + // already passed this gate keep running to completion. + if (drainSignal?.aborted) { + skipped++; + return; + } const videoPath = path.join(dataDir, videoDir); if (failedSet.has(videoDir)) { log(`Skipping previously failed transcription for ${videoDir}`); @@ -130,6 +145,11 @@ export async function runWhisperBatch({ return; } attempted++; + const task = tracker?.start({ + id: videoDir, + label: videoDir, + kind: "transcribe", + }); try { await transcribeOneVideo({ paths, @@ -137,7 +157,7 @@ export async function runWhisperBatch({ videoId: videoDir, audioFilename: resolvedAudioFilename, strictAudio: strict, - onLog: log, + onLog: task ? task.onLog : log, signal, }); succeededCount++; @@ -150,6 +170,8 @@ export async function runWhisperBatch({ log(`FAILED TO TRANSCRIBE ${videoDir}: ${String(err)}`); await appendFile(failureListFile, `${videoDir}\n`); failed++; + } finally { + task?.end(); } }), ), diff --git a/common/jobs/progressParsers.ts b/common/jobs/progressParsers.ts @@ -0,0 +1,88 @@ +// Pure parsers that turn a line of yt-dlp / whisper shell output into a +// per-operation progress update. No I/O — kept dependency-free so they can be +// unit-tested directly if a unit runner is ever added. Used by taskHooks.ts to +// drive the per-task progress bars on the Active Jobs screen. + +export type ProgressUpdate = { fraction?: number; detail?: string }; + +function clamp01(n: number): number { + return Math.max(0, Math.min(1, n)); +} + +function formatClock(totalSeconds: number): string { + const s = Math.max(0, Math.floor(totalSeconds)); + const hh = Math.floor(s / 3600); + const mm = Math.floor((s % 3600) / 60); + const ss = s % 60; + const pad = (n: number) => String(n).padStart(2, "0"); + return hh > 0 ? `${hh}:${pad(mm)}:${pad(ss)}` : `${mm}:${pad(ss)}`; +} + +// yt-dlp download progress, e.g. +// [download] 0.1% of ~ 1.04GiB at 2.31MiB/s ETA Unknown (frag 3/2599) +// [download] 100% of 123.95KiB in 00:00:01 +// Returns null for non-progress lines (Destination:, Downloading item, …). +export function parseDownloadProgress(line: string): ProgressUpdate | null { + if (!line.includes("[download]")) return null; + const pctMatch = line.match(/\[download\]\s+([\d.]+)%/); + const fragMatch = line.match(/\(frag\s+(\d+)\/(\d+)\)/); + const rateMatch = line.match(/at\s+([\d.]+\s*[KMG]?i?B\/s)/); + const etaMatch = line.match(/ETA\s+([\d:]+)/); + + let fraction: number | undefined; + if (pctMatch) fraction = clamp01(Number.parseFloat(pctMatch[1]) / 100); + + const detailParts: string[] = []; + if (rateMatch) detailParts.push(rateMatch[1].replace(/\s+/g, "")); + if (etaMatch) detailParts.push(`ETA ${etaMatch[1]}`); + if (fragMatch) { + detailParts.push(`frag ${fragMatch[1]}/${fragMatch[2]}`); + // When the total size is unknown there's no usable percent; fall back to + // the fragment ratio so the bar still advances. + if (fraction === undefined) { + const done = Number.parseInt(fragMatch[1], 10); + const total = Number.parseInt(fragMatch[2], 10); + if (total > 0) fraction = clamp01(done / total); + } + } + + if (fraction === undefined && detailParts.length === 0) return null; + const out: ProgressUpdate = {}; + if (fraction !== undefined) out.fraction = fraction; + if (detailParts.length > 0) out.detail = detailParts.join(" · "); + return out; +} + +// whisper-cli output is stateful: a header line declares the audio's total +// length, and subsequent `[hh:mm:ss --> hh:mm:ss]` lines report the segment +// just transcribed. Progress = latest end-timestamp / total length. +export function createTranscribeProgressParser(): { + feed: (line: string) => ProgressUpdate | null; +} { + let totalSec: number | undefined; + return { + feed(line: string): ProgressUpdate | null { + const headerMatch = line.match(/\(\d+\s+samples,\s+([\d.]+)\s+sec\)/); + if (headerMatch) { + totalSec = Number.parseFloat(headerMatch[1]); + return null; + } + const tsMatch = line.match( + /^\s*\[(\d{2}):(\d{2}):(\d{2})\.\d+\s*-->\s*(\d{2}):(\d{2}):(\d{2})\.\d+\]/, + ); + if (!tsMatch) return null; + const endSec = + Number.parseInt(tsMatch[4], 10) * 3600 + + Number.parseInt(tsMatch[5], 10) * 60 + + Number.parseInt(tsMatch[6], 10); + if (totalSec === undefined || totalSec <= 0) { + // No header yet — show position but not a fraction. + return { detail: formatClock(endSec) }; + } + return { + fraction: clamp01(endSec / totalSec), + detail: `${formatClock(endSec)} / ${formatClock(totalSec)}`, + }; + }, + }; +} diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -15,6 +15,21 @@ export type JobProgress = { target: number; }; +export type JobTaskKind = "download" | "transcribe"; + +// A single in-flight sub-operation within a job (one video download or one +// transcription). Only currently-running tasks are kept on the record — they +// are added on start and pruned on completion — so memory stays bounded by the +// job's concurrency regardless of how many videos a batch spans. +export type JobTask = { + id: string; // videoId (or URL-derived id) + label: string; // human label + kind: JobTaskKind; + fraction?: number; // 0..1, undefined until first parseable progress + detail?: string; // e.g. "2.31MiB/s ETA 00:03:12" or "01:50:02 / 03:24:47" + startedAt: number; +}; + export type JobRecord = { id: string; kind: string; @@ -29,7 +44,13 @@ export type JobRecord = { logPath: string; child?: ChildProcess; abortController?: AbortController; + // Separate from abortController: a drain (soft-cancel) aborts this signal so + // a batch stops starting NEW sub-operations but lets in-flight ones finish, + // then completes normally (releasing the queue). Hard cancel aborts both. + drainController?: AbortController; + draining?: boolean; progress?: JobProgress; + tasks?: JobTask[]; }; export type QueueSnapshot = { @@ -120,6 +141,9 @@ class JobRegistry { job.endedAt = Date.now(); if (typeof exitCode === "number") job.exitCode = exitCode; } + // Release per-task and drain references on terminal jobs. + job.tasks = []; + job.drainController = undefined; const q = this.queues.get(job.queueKey); if (!q) return; const idx = q.indexOf(job); @@ -153,6 +177,8 @@ class JobRegistry { // idempotent and will start the next queued job. job.status = "cancelled"; job.abortController?.abort(); + // Also abort the drain signal so any code awaiting a drain unblocks. + job.drainController?.abort(); if (job.child) { job.child.kill("SIGTERM"); setTimeout(() => { @@ -164,6 +190,50 @@ class JobRegistry { return false; } + // Soft-cancel ("drain"): stop the job from starting NEW sub-operations but + // let in-flight ones run to completion. The batch's fn returns normally, + // which finalizes it "done" and promotes the next queued job. A queued job + // has nothing in flight, so draining it is equivalent to a plain cancel. + requestDrain(id: string): boolean { + const job = this.jobs.get(id); + if (!job) return false; + if (job.status === "queued") return this.cancel(id); + if (job.status === "running") { + job.draining = true; + job.drainController?.abort(); + return true; + } + return false; + } + + // Per-task tracking. No-ops if the job is gone. tasks is created lazily and + // only ever holds currently-running sub-operations. + addTask(jobId: string, task: JobTask): void { + const job = this.jobs.get(jobId); + if (!job) return; + if (!job.tasks) job.tasks = []; + job.tasks.push(task); + } + + updateTask( + jobId: string, + taskId: string, + patch: { fraction?: number; detail?: string }, + ): void { + const job = this.jobs.get(jobId); + const task = job?.tasks?.find((t) => t.id === taskId); + if (!task) return; + if (patch.fraction !== undefined) task.fraction = patch.fraction; + if (patch.detail !== undefined) task.detail = patch.detail; + } + + removeTask(jobId: string, taskId: string): void { + const job = this.jobs.get(jobId); + if (!job?.tasks) return; + const idx = job.tasks.findIndex((t) => t.id === taskId); + if (idx >= 0) job.tasks.splice(idx, 1); + } + // Snapshot of queues for UI. Sorted by queue name. listQueues(): QueueSnapshot[] { const out: QueueSnapshot[] = []; @@ -207,3 +277,20 @@ export function getRegistry(): JobRegistry { export function newJobId(): string { return `${Date.now().toString(36)}-${Math.random().toString(36).slice(2, 8)}`; } + +// Batch kinds whose controllers honor the drain signal (stop starting new +// sub-operations, let in-flight ones finish). The Drain button is only offered +// for running jobs of these kinds. +const DRAINABLE_KINDS = new Set<string>([ + "whisper-all", + "whisper-bucket-downloaded-no-transcript", + "download-from-playlist", + "download-missing", + "download-missing-subs", + "sync", + "retry-bucket", +]); + +export function isDrainableKind(kind: string): boolean { + return DRAINABLE_KINDS.has(kind); +} diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -7,6 +7,7 @@ import { newJobId, type JobProgress, type JobRecord, + type JobTask, } from "./registry"; import type { Paths } from "../lib/paths"; import { makeSafeController } from "../lib/safeStreamController"; @@ -30,11 +31,26 @@ export type RunManagedCommandOpts = CommonOpts & { env?: Record<string, string | undefined>; }; +// Passed as the optional 4th argument to a managed function. Existing callers +// declare three-arg arrows and silently ignore it; batch / single-op callers +// opt in to per-task progress reporting and drain (soft-cancel) support. +export type JobRunContext = { + jobId: string; + drainSignal: AbortSignal; + addTask: (task: JobTask) => void; + updateTask: ( + taskId: string, + patch: { fraction?: number; detail?: string }, + ) => void; + removeTask: (taskId: string) => void; +}; + export type RunManagedFunctionOpts = CommonOpts & { fn: ( onLog: (line: string) => void, signal: AbortSignal, setProgress: (snap: JobProgress) => void, + ctx: JobRunContext, ) => Promise<void>; }; @@ -197,6 +213,8 @@ export async function runManagedFunction( ignoreFileStreamErrors(fileStream); const abort = new AbortController(); record.abortController = abort; + const drain = new AbortController(); + record.drainController = drain; const onLog = (line: string) => { const text = line.endsWith("\n") ? line : `${line}\n`; @@ -208,8 +226,16 @@ export async function runManagedFunction( record.progress = snap; }; + const ctx: JobRunContext = { + jobId: id, + drainSignal: drain.signal, + addTask: (task) => registry.addTask(id, task), + updateTask: (taskId, patch) => registry.updateTask(id, taskId, patch), + removeTask: (taskId) => registry.removeTask(id, taskId), + }; + opts - .fn(onLog, abort.signal, setProgress) + .fn(onLog, abort.signal, setProgress, ctx) .then(() => { if (record.status === "cancelled") { registry.finalize(id, "cancelled"); diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts @@ -0,0 +1,60 @@ +import type { JobTaskKind } from "./registry"; +import type { JobRunContext } from "./streamCommand"; +import { + parseDownloadProgress, + createTranscribeProgressParser, +} from "./progressParsers"; + +// A handle for one in-flight sub-operation. `onLog` is a drop-in replacement +// for the controller's existing `onLog`: it forwards every line to the shared +// job log AND feeds it to the matching progress parser, updating the task's +// per-operation progress bar. `end()` prunes the task (call it in a finally so +// failures don't leak a task into record.tasks). +export type TaskHandle = { + onLog: (line: string) => void; + end: () => void; +}; + +export type TaskTracker = { + start: (init: { id: string; label: string; kind: JobTaskKind }) => TaskHandle; +}; + +type TaskCtx = Pick<JobRunContext, "addTask" | "updateTask" | "removeTask">; + +// Build a tracker from a job's run context. `forwardLog` is the shared `onLog` +// the controller already writes to. When `ctx` is undefined (a caller that +// doesn't opt in) the tracker degrades to plain forwarding. +export function makeTaskTracker( + ctx: TaskCtx | undefined, + forwardLog: (line: string) => void, +): TaskTracker { + return { + start({ id, label, kind }) { + if (!ctx) { + return { onLog: forwardLog, end: () => {} }; + } + ctx.addTask({ id, label, kind, startedAt: Date.now() }); + const transcribeParser = + kind === "transcribe" ? createTranscribeProgressParser() : null; + let ended = false; + const onLog = (line: string) => { + forwardLog(line); + // yt-dlp rewrites the progress line in place with carriage returns, so + // split on both \r and \n to see each discrete update. + for (const part of line.split(/[\r\n]+/)) { + if (!part) continue; + const update = transcribeParser + ? transcribeParser.feed(part) + : parseDownloadProgress(part); + if (update) ctx.updateTask(id, update); + } + }; + const end = () => { + if (ended) return; + ended = true; + ctx.removeTask(id); + }; + return { onLog, end }; + }, + }; +} diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -27,6 +27,7 @@ import { resolveEffectiveAvailability } from "../lib/availability-server"; import { backfillAvailabilityFromMetadata } from "../controller/backfillAvailability"; import { resolveShardItems } from "../controller/shard"; import { downloadOneManaged } from "./downloadOneManaged"; +import type { TaskTracker } from "../jobs/taskHooks"; export type YtdlpMode = | "store-playlist" @@ -44,6 +45,12 @@ export type RunYtdlpOpts = { paths: Paths; onLog: (s: string) => void; signal: AbortSignal; + // Soft-cancel: when aborted, finish the current download but start no new + // ones. Separate from the hard `signal`. + drainSignal?: AbortSignal; + // When provided, each per-video download is tracked as a per-operation task + // with its own parsed progress bar on the Active Jobs screen. + tracker?: TaskTracker; // download-missing only: drop --download-archive so videos already listed // in the archive still get re-fetched (recovery from a wrong/stale archive). ignoreArchive?: boolean; @@ -440,18 +447,30 @@ async function downloadPlaylistManaged( items.map((url) => limit(async () => { if (opts.signal.aborted) return; + // Drain (soft-cancel): finish the in-flight download, start no more. + if (opts.drainSignal?.aborted) return; if (firstFailure && abortOnError) return; - const outcome = await downloadOneManaged({ - channelSlug: opts.channelSlug, - channelConfig: effectiveChannelConfig, - paths: opts.paths, - videoUrl: url, - onLog: opts.onLog, - signal: opts.signal, - globalCookiesFromBrowser: globalCookies || undefined, - appendArchive: !opts.ignoreArchive, - inlineTranscribeOnFallback, + const task = opts.tracker?.start({ + id: extractVideoId(url) ?? url, + label: extractVideoId(url) ?? url, + kind: "download", }); + let outcome; + try { + outcome = await downloadOneManaged({ + channelSlug: opts.channelSlug, + channelConfig: effectiveChannelConfig, + paths: opts.paths, + videoUrl: url, + onLog: task ? task.onLog : opts.onLog, + signal: opts.signal, + globalCookiesFromBrowser: globalCookies || undefined, + appendArchive: !opts.ignoreArchive, + inlineTranscribeOnFallback, + }); + } finally { + task?.end(); + } if (outcome.status === "failed") { failedCount++; if (abortOnError && !firstFailure) { diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,8 @@ # Changelog ## [Unreleased] +- **Per-operation progress bars on `/jobs/active`.** Each running download/transcription now shows its own live progress bar parsed from the tool's shell output — yt-dlp's download percent (and fragment count) and whisper's transcribed position against the audio length. Batch jobs (Transcribe missing, Download from playlist, Sync, …) list a bar per in-flight video underneath the batch's overall bar; standalone single-video jobs get one too. The screen now polls about once a second so the bars advance live. +- **"Drain" (soft-cancel) for batch jobs.** Alongside the existing Cancel (which immediately kills everything), running batches now offer **Drain**: it lets the in-flight operations finish, starts no new ones, then completes the job normally and releases the queue so the next batch can run. Hard Cancel still works at any time, including mid-drain. - **"Build static export" moved from the Build page to the top of the Deploy page**, co-locating it with the export changelog preview, "Cut release", and "Deploy static export" so the whole build → review → cut → deploy sequence lives on one page. - **The embedded single-video view on the channel page is now collapsed by default**, keeping the channel view compact. Clicking a video in the list expands the panel and scrolls to it; navigating away from a video (no video selected) collapses it again. A manual collapse/expand toggle is available on the panel, mirroring the video list's existing collapse. - **Unlisted videos now get their own clickable list in the channel availability diagnostics**, alongside Deleted/Private/Members-only/Needs-auth/Error. Previously unlisted videos were only shown as a count, so there was no way to jump to the specific videos. diff --git a/editor/app/api/jobs/active/route.ts b/editor/app/api/jobs/active/route.ts @@ -0,0 +1,11 @@ +import { NextResponse } from "next/server"; +import { buildActiveJobsPayload } from "../../../jobs/active/buildActiveJobs"; + +export const dynamic = "force-dynamic"; + +// Backs the ~1s client poll on the Active Jobs screen so per-task progress +// bars advance live without a full RSC refresh. +export async function GET() { + const payload = await buildActiveJobsPayload(); + return NextResponse.json(payload); +} diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -22,6 +22,7 @@ import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; function defaultQueueKey(config: ChannelConfig): string { return platformQueueKey(config.platform ?? detectPlatform(config.url)); @@ -61,7 +62,7 @@ async function runPipelineAction( queueKey === undefined ? defaultQueueKey(channelConfig) : queueKey.trim(), paths, channelSlug: slug, - fn: async (onLog, signal, setProgress) => { + fn: async (onLog, signal, setProgress, ctx) => { if (mode !== "store-playlist") { const stat = await readChannelStat(paths, slug); if (stat) { @@ -94,6 +95,8 @@ async function runPipelineAction( paths, onLog, signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), ignoreArchive: options?.ignoreArchive, abortOnError: options?.abortOnError, shardTotal: options?.shardTotal, diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts @@ -26,6 +26,7 @@ import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; function videoQueueKey(config: ChannelConfig, override: string | undefined): string { if (override === undefined) { @@ -99,16 +100,25 @@ export async function transcribeOneAction( paths, channelSlug: slug, videoId, - fn: async (onLog, signal) => { - await transcribeOneVideo({ - paths, - videoDir, - videoId, - audioFilename, - onLog, - signal, + fn: async (onLog, signal, _setProgress, ctx) => { + const task = makeTaskTracker(ctx, onLog).start({ + id: videoId, + label: videoId, + kind: "transcribe", }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); + try { + await transcribeOneVideo({ + paths, + videoDir, + videoId, + audioFilename, + onLog: task.onLog, + signal, + }); + revalidatePath(`/channels/${slug}/videos/${videoId}`); + } finally { + task.end(); + } }, }); } @@ -136,20 +146,29 @@ export async function downloadVideoPipelineAction( paths, channelSlug: slug, videoId, - fn: async (onLog, signal) => { - await downloadOneManaged({ - channelSlug: slug, - channelConfig: r.config, - paths, - videoUrl: url, - onLog, - signal, - globalCookiesFromBrowser: settings.cookiesFromBrowser || undefined, - inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, - appendArchive: true, + fn: async (onLog, signal, _setProgress, ctx) => { + const task = makeTaskTracker(ctx, onLog).start({ + id: videoId, + label: videoId, + kind: "download", }); - revalidatePath(`/channels/${slug}/videos/${videoId}`); - revalidatePath(`/channels/${slug}`); + try { + await downloadOneManaged({ + channelSlug: slug, + channelConfig: r.config, + paths, + videoUrl: url, + onLog: task.onLog, + signal, + globalCookiesFromBrowser: settings.cookiesFromBrowser || undefined, + inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, + appendArchive: true, + }); + revalidatePath(`/channels/${slug}/videos/${videoId}`); + revalidatePath(`/channels/${slug}`); + } finally { + task.end(); + } }, }); } @@ -169,7 +188,7 @@ export async function whisperVideoAction( paths, channelSlug: slug, videoId, - fn: async (onLog, signal) => { + fn: async (onLog, signal, _setProgress, ctx) => { const audioFormat = r.config.audioFormat ?? "mp3"; const entries = await readdir(videoDir).catch(() => [] as string[]); const hasAudio = entries.some((e: string) => e.startsWith("audio.")); @@ -194,14 +213,23 @@ export async function whisperVideoAction( } else { onLog(`Audio already on disk for ${videoId}; skipping download.`); } - await transcribeOneVideo({ - paths, - videoDir, - videoId, - audioFilename: `audio.${audioFormat}`, - onLog, - signal, + const task = makeTaskTracker(ctx, onLog).start({ + id: videoId, + label: videoId, + kind: "transcribe", }); + try { + await transcribeOneVideo({ + paths, + videoDir, + videoId, + audioFilename: `audio.${audioFormat}`, + onLog: task.onLog, + signal, + }); + } finally { + task.end(); + } revalidatePath(`/channels/${slug}/videos/${videoId}`); revalidatePath(`/channels/${slug}`); }, diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts @@ -27,6 +27,7 @@ import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; function defaultQueueKey(slug: string): string { return `channel:${slug}`; @@ -61,7 +62,7 @@ export async function transcribeMissingAction( queueKey: queueKey === undefined ? TRANSCRIPTION_QUEUE : queueKey.trim(), paths, channelSlug: slug, - fn: async (onLog, signal, setProgress) => { + fn: async (onLog, signal, setProgress, ctx) => { const stat = await readChannelStat(paths, slug); if (stat) { setProgress({ @@ -81,6 +82,8 @@ export async function transcribeMissingAction( shardIndex, onLog, signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), }); onLog( `Whisper batch: ${result.succeeded} succeeded, ${result.failed} failed, ${result.skipped} skipped, ${result.attempted} attempted.`, @@ -110,7 +113,7 @@ export async function transcribeBucketAction( queueKey: queueKey === undefined ? TRANSCRIPTION_QUEUE : queueKey.trim(), paths, channelSlug: slug, - fn: async (onLog, signal, setProgress) => { + fn: async (onLog, signal, setProgress, ctx) => { const stat = await readChannelStat(paths, slug); if (stat) { const remaining = await countNotYetTranscribed(paths, slug, cleaned); @@ -129,6 +132,8 @@ export async function transcribeBucketAction( ids: cleaned, onLog, signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), }); onLog( `Whisper bucket: ${result.succeeded} succeeded, ${result.failed} failed, ${result.skipped} skipped, ${result.attempted} attempted.`, diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts @@ -11,6 +11,14 @@ export async function cancelJobAction(id: string): Promise<{ ok: boolean }> { return { ok }; } +// Soft-cancel: let the batch's in-flight sub-operations finish, start no new +// ones, then complete and release the queue for the next job. +export async function drainJobAction(id: string): Promise<{ ok: boolean }> { + const ok = getRegistry().requestDrain(id); + revalidatePath("/jobs"); + return { ok }; +} + export async function clearArchivedAction(): Promise<{ deleted: number }> { const deleted = await clearArchivedLogs(getPaths()); revalidatePath("/jobs"); diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts @@ -0,0 +1,98 @@ +import { + getRegistry, + isDrainableKind, + type JobProgress, +} from "yt-dlp-transcript-common/jobs/registry"; +import { + readChannelStat, + type ChannelStat, +} from "yt-dlp-transcript-common/controller/channels"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import type { RunningJobsListItem } from "../components/RunningJobsList"; + +export type ActiveJobsPayload = { + jobs: RunningJobsListItem[]; + channels: { slug: string; displayName: string }[]; +}; + +function computeJobProgressView( + snap: JobProgress | undefined, + stat: ChannelStat | undefined, +): RunningJobsListItem["progress"] { + if (!snap || !stat) return undefined; + const current = + snap.metric === "downloads" ? stat.downloadCount : stat.transcriptCount; + const range = Math.max(0, snap.target - snap.initial); + const advance = Math.max(0, current - snap.initial); + const pct = + range === 0 ? 100 : Math.max(0, Math.min(100, (advance / range) * 100)); + return { + metric: snap.metric, + initial: snap.initial, + current, + target: snap.target, + pct, + }; +} + +// Builds the Active Jobs screen payload from the in-memory registry plus +// on-disk channel stats. Shared by the server-rendered page (initial state) +// and the /api/jobs/active poll route (live updates). +export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { + const registry = getRegistry(); + const activeRecords = registry + .list() + .filter((j) => j.status === "running" || j.status === "queued") + .sort((a, b) => { + if (a.status === b.status) return 0; + return a.status === "running" ? -1 : 1; + }); + + const channelSlugs = Array.from( + new Set( + activeRecords + .map((j) => j.channelSlug) + .filter((s): s is string => Boolean(s)), + ), + ); + + const paths = getPaths(); + const statResults = await Promise.all( + channelSlugs.map((slug) => readChannelStat(paths, slug)), + ); + const channelStats = new Map<string, ChannelStat>(); + for (let i = 0; i < channelSlugs.length; i++) { + const stat = statResults[i]; + if (stat) channelStats.set(channelSlugs[i], stat); + } + + const jobs: RunningJobsListItem[] = activeRecords.map((j) => ({ + id: j.id, + kind: j.kind, + status: j.status as "running" | "queued", + queueKey: j.queueKey, + channelSlug: j.channelSlug, + videoId: j.videoId, + progress: + j.status === "running" && j.channelSlug + ? computeJobProgressView(j.progress, channelStats.get(j.channelSlug)) + : undefined, + tasks: j.tasks?.map((t) => ({ + id: t.id, + label: t.label, + kind: t.kind, + fraction: t.fraction, + detail: t.detail, + })), + draining: j.draining === true, + drainable: + j.status === "running" && isDrainableKind(j.kind) && j.draining !== true, + })); + + const channels = channelSlugs.map((slug) => ({ + slug, + displayName: channelStats.get(slug)?.config.name ?? slug, + })); + + return { jobs, channels }; +} diff --git a/editor/app/jobs/active/page.tsx b/editor/app/jobs/active/page.tsx @@ -1,143 +1,19 @@ import type { Metadata } from "next"; -import Link from "next/link"; -import { - getRegistry, - type JobProgress, -} from "yt-dlp-transcript-common/jobs/registry"; -import { - readChannelStat, - type ChannelStat, -} from "yt-dlp-transcript-common/controller/channels"; -import { getPaths } from "yt-dlp-transcript-common/lib/paths"; -import { JobsAutoRefresh } from "../components/JobsAutoRefresh"; -import { - RunningJobsList, - type RunningJobsListItem, -} from "../components/RunningJobsList"; +import { buildActiveJobsPayload } from "./buildActiveJobs"; +import { ActiveJobsLive } from "../components/ActiveJobsLive"; export const dynamic = "force-dynamic"; export const metadata: Metadata = { title: "Active jobs" }; -function computeJobProgressView( - snap: JobProgress | undefined, - stat: ChannelStat | undefined, -): RunningJobsListItem["progress"] { - if (!snap || !stat) return undefined; - const current = - snap.metric === "downloads" ? stat.downloadCount : stat.transcriptCount; - const range = Math.max(0, snap.target - snap.initial); - const advance = Math.max(0, current - snap.initial); - const pct = - range === 0 ? 100 : Math.max(0, Math.min(100, (advance / range) * 100)); - return { - metric: snap.metric, - initial: snap.initial, - current, - target: snap.target, - pct, - }; -} - export default async function ActiveJobsPage() { - const registry = getRegistry(); - const activeRecords = registry - .list() - .filter((j) => j.status === "running" || j.status === "queued") - .sort((a, b) => { - if (a.status === b.status) return 0; - return a.status === "running" ? -1 : 1; - }); - - const channelSlugs = Array.from( - new Set( - activeRecords - .map((j) => j.channelSlug) - .filter((s): s is string => Boolean(s)), - ), - ); - - const paths = getPaths(); - const statResults = await Promise.all( - channelSlugs.map((slug) => readChannelStat(paths, slug)), - ); - const channelStats = new Map<string, ChannelStat>(); - for (let i = 0; i < channelSlugs.length; i++) { - const stat = statResults[i]; - if (stat) channelStats.set(channelSlugs[i], stat); - } - - const activeJobs: RunningJobsListItem[] = activeRecords.map((j) => ({ - id: j.id, - kind: j.kind, - status: j.status as "running" | "queued", - queueKey: j.queueKey, - channelSlug: j.channelSlug, - videoId: j.videoId, - progress: - j.status === "running" && j.channelSlug - ? computeJobProgressView(j.progress, channelStats.get(j.channelSlug)) - : undefined, - })); - - const jobsBySlug = new Map<string, RunningJobsListItem[]>(); - const unassigned: RunningJobsListItem[] = []; - for (const job of activeJobs) { - if (job.channelSlug) { - const list = jobsBySlug.get(job.channelSlug) ?? []; - list.push(job); - jobsBySlug.set(job.channelSlug, list); - } else { - unassigned.push(job); - } - } - + const initial = await buildActiveJobsPayload(); return ( <div className="flex flex-col gap-4"> - <JobsAutoRefresh hasNonTerminal={activeJobs.length > 0} /> <div className="flex items-center justify-between"> <h1 className="text-2xl font-semibold">Active jobs</h1> </div> - {activeJobs.length === 0 ? ( - <p className="text-sm text-zinc-500 border border-dashed border-zinc-300 dark:border-zinc-700 rounded p-4"> - No active jobs. - </p> - ) : ( - <div className="flex flex-col gap-4"> - {channelSlugs.map((slug) => { - const jobs = jobsBySlug.get(slug) ?? []; - const stat = channelStats.get(slug); - const displayName = stat?.config.name ?? slug; - return ( - <section - key={slug} - aria-label={`Active jobs for ${displayName}`} - className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded-md p-3 bg-white dark:bg-zinc-900" - > - <div className="flex flex-wrap items-baseline justify-between gap-2"> - <Link - href={`/channels/${slug}`} - className="text-base font-medium underline hover:text-zinc-900 dark:hover:text-zinc-100" - > - {displayName} - </Link> - <span className="font-mono text-xs text-zinc-500">{slug}</span> - </div> - <RunningJobsList jobs={jobs} hideChannelSlug /> - </section> - ); - })} - {unassigned.length > 0 && ( - <section - aria-label="Active jobs without a channel" - className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded-md p-3 bg-white dark:bg-zinc-900" - > - <h2 className="text-base font-medium">Other</h2> - <RunningJobsList jobs={unassigned} /> - </section> - )} - </div> - )} + <ActiveJobsLive initial={initial} /> </div> ); } diff --git a/editor/app/jobs/components/ActiveJobsLive.tsx b/editor/app/jobs/components/ActiveJobsLive.tsx @@ -0,0 +1,95 @@ +"use client"; + +import Link from "next/link"; +import { useEffect, useState } from "react"; +import type { ActiveJobsPayload } from "../active/buildActiveJobs"; +import { RunningJobsList, type RunningJobsListItem } from "./RunningJobsList"; + +// Polls /api/jobs/active so per-task progress bars advance live (the server +// component only provides the initial paint). Replaces the coarser 2.5s +// router.refresh for this screen. +const POLL_MS = 1000; + +export function ActiveJobsLive({ initial }: { initial: ActiveJobsPayload }) { + const [payload, setPayload] = useState<ActiveJobsPayload>(initial); + + useEffect(() => { + let cancelled = false; + let timer: ReturnType<typeof setTimeout> | null = null; + async function tick() { + try { + const res = await fetch("/api/jobs/active", { cache: "no-store" }); + if (res.ok && !cancelled) { + setPayload((await res.json()) as ActiveJobsPayload); + } + } catch { + // transient — keep polling + } finally { + if (!cancelled) timer = setTimeout(tick, POLL_MS); + } + } + timer = setTimeout(tick, POLL_MS); + return () => { + cancelled = true; + if (timer) clearTimeout(timer); + }; + }, []); + + const { jobs, channels } = payload; + + const jobsBySlug = new Map<string, RunningJobsListItem[]>(); + const unassigned: RunningJobsListItem[] = []; + for (const job of jobs) { + if (job.channelSlug) { + const list = jobsBySlug.get(job.channelSlug) ?? []; + list.push(job); + jobsBySlug.set(job.channelSlug, list); + } else { + unassigned.push(job); + } + } + + if (jobs.length === 0) { + return ( + <p className="text-sm text-zinc-500 border border-dashed border-zinc-300 dark:border-zinc-700 rounded p-4"> + No active jobs. + </p> + ); + } + + return ( + <div className="flex flex-col gap-4"> + {channels.map(({ slug, displayName }) => { + const channelJobs = jobsBySlug.get(slug) ?? []; + if (channelJobs.length === 0) return null; + return ( + <section + key={slug} + aria-label={`Active jobs for ${displayName}`} + className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded-md p-3 bg-white dark:bg-zinc-900" + > + <div className="flex flex-wrap items-baseline justify-between gap-2"> + <Link + href={`/channels/${slug}`} + className="text-base font-medium underline hover:text-zinc-900 dark:hover:text-zinc-100" + > + {displayName} + </Link> + <span className="font-mono text-xs text-zinc-500">{slug}</span> + </div> + <RunningJobsList jobs={channelJobs} hideChannelSlug /> + </section> + ); + })} + {unassigned.length > 0 && ( + <section + aria-label="Active jobs without a channel" + className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded-md p-3 bg-white dark:bg-zinc-900" + > + <h2 className="text-base font-medium">Other</h2> + <RunningJobsList jobs={unassigned} /> + </section> + )} + </div> + ); +} diff --git a/editor/app/jobs/components/DrainJobButton.tsx b/editor/app/jobs/components/DrainJobButton.tsx @@ -0,0 +1,41 @@ +"use client"; + +import { useState } from "react"; +import { useRouter } from "next/navigation"; +import { drainJobAction } from "../actions"; + +type Props = { + jobId: string; + draining?: boolean; +}; + +export function DrainJobButton({ jobId, draining }: Props) { + const [busy, setBusy] = useState(false); + const router = useRouter(); + if (draining) { + return ( + <span className="px-2 py-1 rounded border border-amber-300 dark:border-amber-800 text-xs font-medium text-amber-700 dark:text-amber-300"> + Draining… + </span> + ); + } + return ( + <button + type="button" + onClick={async () => { + setBusy(true); + try { + await drainJobAction(jobId); + router.refresh(); + } finally { + setBusy(false); + } + }} + disabled={busy} + title="Finish the in-flight operations, start no new ones, then release the queue" + className="px-2 py-1 rounded border border-amber-300 dark:border-amber-800 text-xs font-medium text-amber-700 dark:text-amber-300 hover:bg-amber-50 dark:hover:bg-amber-950 disabled:opacity-50" + > + {busy ? "Draining…" : "Drain"} + </button> + ); +} diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx @@ -3,6 +3,16 @@ import Link from "next/link"; import { useState } from "react"; import { JobLogTail } from "../[id]/components/JobLogTail"; +import { DrainJobButton } from "./DrainJobButton"; +import { CancelJobButton } from "./CancelJobButton"; + +export type RunningJobsTask = { + id: string; + label: string; + kind: "download" | "transcribe"; + fraction?: number; + detail?: string; +}; export type RunningJobsListItem = { id: string; @@ -18,6 +28,9 @@ export type RunningJobsListItem = { target: number; pct: number; }; + tasks?: RunningJobsTask[]; + draining?: boolean; + drainable?: boolean; }; type Props = { @@ -100,22 +113,85 @@ function JobRow({ queue <code className="font-mono">{job.queueKey || "—"}</code> </span> )} - <button - type="button" - onClick={() => setShowLog((s) => !s)} - aria-expanded={showLog} - aria-label={`${showLog ? "hide" : "show"} log for job ${job.id}`} - className="ml-auto px-2 py-1 rounded-md bg-zinc-200 dark:bg-zinc-800 text-zinc-900 dark:text-zinc-100 text-xs font-medium hover:opacity-90" - > - {showLog ? "Hide log" : "Show log"} - </button> + <div className="ml-auto flex items-center gap-2"> + {(job.drainable || job.draining) && ( + <DrainJobButton jobId={job.id} draining={job.draining} /> + )} + {job.status === "running" && <CancelJobButton jobId={job.id} />} + <button + type="button" + onClick={() => setShowLog((s) => !s)} + aria-expanded={showLog} + aria-label={`${showLog ? "hide" : "show"} log for job ${job.id}`} + className="px-2 py-1 rounded-md bg-zinc-200 dark:bg-zinc-800 text-zinc-900 dark:text-zinc-100 text-xs font-medium hover:opacity-90" + > + {showLog ? "Hide log" : "Show log"} + </button> + </div> </div> {job.progress && <JobProgressBar progress={job.progress} />} + {job.tasks && job.tasks.length > 0 && ( + <ul + aria-label="In-progress operations" + className="flex flex-col gap-1.5 pl-2 border-l-2 border-zinc-200 dark:border-zinc-800" + > + {job.tasks.map((t) => ( + <li key={t.id}> + <TaskProgressBar task={t} /> + </li> + ))} + </ul> + )} {showLog && <JobLogTail jobId={job.id} initiallyRunning />} </li> ); } +function TaskProgressBar({ task }: { task: RunningJobsTask }) { + const verb = task.kind === "download" ? "Downloading" : "Transcribing"; + const hasFraction = typeof task.fraction === "number"; + const pct = hasFraction ? Math.round((task.fraction as number) * 100) : 0; + const fillClass = + task.kind === "download" + ? "bg-emerald-300 dark:bg-emerald-700" + : "bg-emerald-500"; + return ( + <div className="flex flex-col gap-0.5"> + <div className="flex items-baseline justify-between gap-2 text-xs"> + <span className="truncate"> + <span className="text-zinc-500">{verb} </span> + <span className="font-mono">{task.label}</span> + </span> + <span className="font-mono text-zinc-500 shrink-0"> + {hasFraction ? `${pct}%` : ""} + {task.detail ? ` ${task.detail}` : ""} + </span> + </div> + <div + role="progressbar" + aria-valuemin={0} + aria-valuemax={100} + aria-valuenow={hasFraction ? pct : undefined} + aria-label={`${verb} ${task.label}`} + className="relative h-1.5 w-full overflow-hidden rounded bg-zinc-200 dark:bg-zinc-800" + > + {hasFraction ? ( + <div + className={`absolute inset-y-0 left-0 ${fillClass}`} + style={{ width: `${pct}%` }} + /> + ) : ( + // Indeterminate: no parseable progress yet (e.g. before whisper's + // header line or yt-dlp's first percent). + <div + className={`absolute inset-y-0 left-0 w-1/3 animate-pulse ${fillClass}`} + /> + )} + </div> + </div> + ); +} + function JobProgressBar({ progress, }: { diff --git a/editor/e2e/fixtures/bin/fake-whisper.mjs b/editor/e2e/fixtures/bin/fake-whisper.mjs @@ -18,6 +18,27 @@ if (!tmpBase) { process.exit(2); } +const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); +const pad = (n) => String(n).padStart(2, "0"); +const clock = (s) => `00:${pad(Math.floor(s / 60))}:${pad(s % 60)}.000`; + +// Videos whose dir (== video id) contains the SLOWOP marker emit whisper-style +// progress output with real delays, so the per-operation progress bars and +// drain (soft-cancel) behaviour can be observed on the Active Jobs screen. +// Videos without the marker stay instant so the rest of the suite is fast. +if (process.cwd().toLowerCase().includes("slowop")) { + const totalSec = 50; + process.stdout.write( + `main: processing '${audio}' (800000 samples, ${totalSec}.0 sec), 4 threads, 1 processors, 5 beams + best of 5, lang = en, task = transcribe, timestamps = 1 ...\n`, + ); + for (let t = 5; t <= totalSec; t += 5) { + process.stdout.write( + `[${clock(t - 5)} --> ${clock(t)}] slow whisper segment ${t}\n`, + ); + await sleep(700); + } +} + const doc = { systeminfo: "fake", model: { type: "fake" }, diff --git a/editor/e2e/fixtures/bin/fake-ytdlp.mjs b/editor/e2e/fixtures/bin/fake-ytdlp.mjs @@ -112,6 +112,19 @@ async function downloadOne(url, opts = {}) { } await ensureDir(videoDir); process.stdout.write(`[download] Fetching ${id}\n`); + // Videos whose id contains the SLOWOP marker emit yt-dlp-style progress + // lines with real delays so the per-operation progress bar (and drain) can + // be observed on the Active Jobs screen. Unmarked ids stay instant. + if (id.toLowerCase().includes("slowop")) { + const fragTotal = 6; + for (let f = 1; f <= fragTotal; f++) { + const pct = ((f / fragTotal) * 100).toFixed(1); + process.stdout.write( + `[download] ${pct}% of ~ 10.00MiB at 1.50MiB/s ETA 00:00:0${Math.max(0, fragTotal - f)} (frag ${f}/${fragTotal})\n`, + ); + await sleep(500); + } + } if (!alreadyHasMetadata) await writeMetadata(videoDir, id); if (audioFmt) { await writeFile( diff --git a/editor/e2e/jobs-batch-tasks-drain.spec.ts b/editor/e2e/jobs-batch-tasks-drain.spec.ts @@ -0,0 +1,166 @@ +// Per-operation progress bars on the Active Jobs screen, and the "Drain" +// (soft-cancel) control: a drained batch finishes its in-flight sub-operations +// without starting new ones, then completes (done, not cancelled) and releases +// the queue for the next batch. The fake whisper/yt-dlp binaries emit slow, +// progress-style output for any video id containing "slowop". + +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { pathExists, resetData, resolvePath } from "./helpers"; + +const baseUrl = "http://localhost:3011"; + +async function invalidateCache() { + await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {}); +} + +// Build a transcribe-handling channel with `count` slow videos (each has only +// audio.mp3 on disk, so "Transcribe missing" will run whisper on all of them). +async function makeTranscribeChannel(slug: string, name: 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, + 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; +} + +test("active jobs shows a per-operation progress bar that advances", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData(null); + const ids = ["slowop1", "slowop2"]; + await makeTranscribeChannel("tasks-one", "Tasks One", ids); + await invalidateCache(); + + await page.goto("/channels/tasks-one"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + await page.goto("/jobs/active"); + + // A per-task "Transcribing" bar shows up for the in-flight operation(s). + const taskBar = page + .getByRole("progressbar", { name: /Transcribing slowop/ }) + .first(); + await expect(taskBar).toBeVisible({ timeout: 15_000 }); + + // Its parsed position advances as whisper emits more timestamp lines. + await expect + .poll( + async () => { + const v = await taskBar.getAttribute("aria-valuenow"); + return v ? Number(v) : 0; + }, + { timeout: 20_000 }, + ) + .toBeGreaterThan(0); +}); + +test("draining a batch finishes in-flight work, skips the rest, and releases the queue", async ({ + page, +}) => { + test.setTimeout(90_000); + await resetData(null); + // 6 videos with the default concurrency of 4 guarantees some are still + // queued inside the batch when we drain, so they get skipped. + const aIds = ["slowopa1", "slowopa2", "slowopa3", "slowopa4", "slowopa5", "slowopa6"]; + const bIds = ["slowopb1"]; + await makeTranscribeChannel("drain-a", "Drain A", aIds); + await makeTranscribeChannel("drain-b", "Drain B", bIds); + await invalidateCache(); + + // Both default to the shared transcription queue, so B waits behind A. + await page.goto("/channels/drain-a"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + await page.goto("/channels/drain-b"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + + await page.goto("/jobs/active"); + const sectionA = page.locator("section[aria-label='Active jobs for Drain A']"); + await expect(sectionA.getByText("running", { exact: true })).toBeVisible({ + timeout: 15_000, + }); + + // Drain A: let the in-flight transcriptions finish, start no new ones. + await sectionA.getByRole("button", { name: "Drain" }).click(); + + // B (previously queued) eventually runs to completion — proof the queue was + // released by A finishing normally. + await expect + .poll(() => transcriptCount("drain-b", bIds), { timeout: 60_000 }) + .toBe(1); + + // A finished only its in-flight subset, not all six — proof drain stopped it + // from starting new operations. + const aDone = await transcriptCount("drain-a", aIds); + expect(aDone).toBeGreaterThanOrEqual(1); + expect(aDone).toBeLessThan(aIds.length); + + // A finalized as "done" (a soft-cancel), not "cancelled". + await page.goto("/jobs"); + const aRow = page.getByRole("row").filter({ hasText: "drain-a" }).first(); + await expect(aRow).toContainText("done", { timeout: 15_000 }); +}); + +test("hard Cancel during a drain ends the job cancelled without stream errors", async ({ + page, +}) => { + test.setTimeout(90_000); + await resetData(null); + const ids = ["slowopc1", "slowopc2", "slowopc3", "slowopc4", "slowopc5", "slowopc6"]; + await makeTranscribeChannel("cancel-drain", "Cancel Drain", ids); + await invalidateCache(); + await fetch(`${baseUrl}/api/test/uncaught-count`, { method: "DELETE" }).catch( + () => {}, + ); + + await page.goto("/channels/cancel-drain"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + + await page.goto("/jobs/active"); + const section = page.locator( + "section[aria-label='Active jobs for Cancel Drain']", + ); + await expect(section.getByText("running", { exact: true })).toBeVisible({ + timeout: 15_000, + }); + + // Drain first, then hard-cancel before the in-flight work finishes. + await section.getByRole("button", { name: "Drain" }).click(); + await section.getByRole("button", { name: /^Cancel$/ }).click(); + + // The job ends cancelled (hard cancel wins over the in-progress drain). + await page.goto("/jobs"); + const row = page.getByRole("row").filter({ hasText: "cancel-drain" }).first(); + await expect(row).toContainText("cancelled", { timeout: 20_000 }); + + await page.waitForTimeout(1_000); + const counts = (await ( + await fetch(`${baseUrl}/api/test/uncaught-count`) + ).json()) as { uncaught: number; messages: string[] }; + expect( + counts.messages.filter((m) => m.includes("Controller is already closed")), + ).toEqual([]); + expect(counts.uncaught).toBe(0); +});