// Per-operation progress bars on /jobs, 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 { channelStage, generateReport, pathExists, resetData, resolvePath, } from "./helpers"; import { baseUrl } from "./baseUrl"; 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 { 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 generateReport(page, "tasks-one"); await page.goto(channelStage("tasks-one", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await page.goto("/jobs"); // 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); // The task leads with a live "running for" timer (m:ss) that ticks up. Read // the leading time token off the task row and assert it advances. const taskRow = taskBar.locator( "xpath=ancestor::div[contains(@class,'gap-0.5')][1]", ); const elapsedSeconds = async () => { const m = (await taskRow.innerText()).match(/(\d+):(\d{2})/); return m ? Number(m[1]) * 60 + Number(m[2]) : -1; }; const first = await elapsedSeconds(); expect(first).toBeGreaterThanOrEqual(0); await expect.poll(elapsedSeconds, { timeout: 10_000 }).toBeGreaterThan(first); }); test("active jobs shows an estimated time remaining once a task completes", async ({ page, }) => { test.setTimeout(90_000); await resetData(null); // More videos than the default concurrency of 4, so some tasks finish while // others are still in flight — i.e. completed > 0 and remaining > 0, the // window in which an ETA can be shown. const ids = ["slowope1", "slowope2", "slowope3", "slowope4", "slowope5", "slowope6"]; await makeTranscribeChannel("eta-one", "Eta One", ids); await invalidateCache(); await generateReport(page, "eta-one"); await page.goto(channelStage("eta-one", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await page.goto("/jobs"); const row = page .getByRole("row") .filter({ hasText: "eta-one" }) .filter({ hasText: "Transcribe all" }) .first(); // The batch-level progress bar is present (proves setProgress ran). await expect(row.getByText(/Transcripts: \d+ \/ 6/)).toBeVisible({ timeout: 15_000, }); // Before any task completes the average is unknown → "estimating…"; once one // finishes a concrete "~m:ss left" estimate appears. Poll for the concrete // form (it implies we passed through, or skipped straight past, estimating). await expect .poll(async () => (await row.innerText()).replace(/\s+/g, " "), { timeout: 60_000, }) .toMatch(/~\d+:\d{2} left/); }); 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 generateReport(page, "drain-a"); await page.goto(channelStage("drain-a", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await generateReport(page, "drain-b"); await page.goto(channelStage("drain-b", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await page.goto("/jobs"); // A bare slug filter would also match the channel's automatic refresh-report // row once it lands (see the same note further down), so pin the batch. const rowA = page .getByRole("row") .filter({ hasText: "drain-a" }) .filter({ hasText: "Transcribe all" }) .first(); await expect(rowA.getByText("running", { exact: true })).toBeVisible({ timeout: 15_000, }); // Drain A: let the in-flight transcriptions finish, start no new ones. await rowA.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("'Drain all' drains the running batch and cancels the queued one", async ({ page, }) => { test.setTimeout(90_000); await resetData(null); // A runs (6 videos > concurrency 4, so some stay in-flight when we drain); // B waits queued behind A on the shared transcription queue. const aIds = ["slowopd1", "slowopd2", "slowopd3", "slowopd4", "slowopd5", "slowopd6"]; const bIds = ["slowopd7"]; await makeTranscribeChannel("drainall-a", "DrainAll A", aIds); await makeTranscribeChannel("drainall-b", "DrainAll B", bIds); await invalidateCache(); await generateReport(page, "drainall-a"); await page.goto(channelStage("drainall-a", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await generateReport(page, "drainall-b"); await page.goto(channelStage("drainall-b", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await page.goto("/jobs"); const rowA = page .getByRole("row") .filter({ hasText: "drainall-a" }) .filter({ hasText: "Transcribe all" }) .first(); await expect(rowA.getByText("running", { exact: true })).toBeVisible({ timeout: 15_000, }); // Accept the confirm, then spin everything down in one click. page.on("dialog", (d) => d.accept()); await page.getByRole("button", { name: "Drain all" }).click(); // B (queued) ends cancelled and never transcribes — proof Drain all cancels // the queue rather than letting it run after A finishes (the per-job Drain // test above shows the opposite). A drains to "done" with a partial subset. await page.goto("/jobs"); const bRow = page .getByRole("row") .filter({ hasText: "drainall-b" }) .filter({ hasText: "Transcribe all" }) .first(); await expect(bRow).toContainText("cancelled", { timeout: 20_000 }); expect(await transcriptCount("drainall-b", bIds)).toBe(0); const aRow = page .getByRole("row") .filter({ hasText: "drainall-a" }) .filter({ hasText: "Transcribe all" }) .first(); await expect(aRow).toContainText("done", { timeout: 20_000 }); const aDone = await transcriptCount("drainall-a", aIds); expect(aDone).toBeGreaterThanOrEqual(1); expect(aDone).toBeLessThan(aIds.length); }); test("a queued job can be cancelled directly from its row without opening the log", async ({ page, }) => { test.setTimeout(90_000); await resetData(null); // A runs; B waits queued behind A on the shared transcription queue. const aIds = ["slowopq1", "slowopq2", "slowopq3", "slowopq4", "slowopq5", "slowopq6"]; const bIds = ["slowopq7"]; await makeTranscribeChannel("qcancel-a", "QCancel A", aIds); await makeTranscribeChannel("qcancel-b", "QCancel B", bIds); await invalidateCache(); await generateReport(page, "qcancel-a"); await page.goto(channelStage("qcancel-a", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await generateReport(page, "qcancel-b"); await page.goto(channelStage("qcancel-b", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await page.goto("/jobs"); const rowB = page .getByRole("row") .filter({ hasText: "qcancel-b" }) .filter({ hasText: "Transcribe all" }) .first(); // B's row shows as queued, with a Cancel button right there (no "Show log"). await expect(rowB.getByText("queued", { exact: true })).toBeVisible({ timeout: 15_000, }); await rowB.getByRole("button", { name: /^Cancel$/ }).click(); // B ends cancelled and never transcribes; A keeps running and finishes. await page.goto("/jobs"); const bRow = page .getByRole("row") .filter({ hasText: "qcancel-b" }) .filter({ hasText: "Transcribe all" }) .first(); await expect(bRow).toContainText("cancelled", { timeout: 20_000 }); expect(await transcriptCount("qcancel-b", bIds)).toBe(0); }); 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 generateReport(page, "cancel-drain"); await page.goto(channelStage("cancel-drain", "transcribe")); await page.getByRole("button", { name: "Transcribe missing" }).click(); await page.goto("/jobs"); const liveRow = page .getByRole("row") .filter({ hasText: "cancel-drain" }) .filter({ hasText: "Transcribe all" }) .first(); await expect(liveRow.getByText("running", { exact: true })).toBeVisible({ timeout: 15_000, }); // Drain first, then hard-cancel before the in-flight work finishes. await liveRow.getByRole("button", { name: "Drain" }).click(); await liveRow.getByRole("button", { name: /^Cancel$/ }).click(); // The job ends cancelled (hard cancel wins over the in-progress drain). // Target the whisper-all row specifically: finishing the batch now also // queues an automatic `refresh-report` job for the same channel (the global // debounced report refresh), which would otherwise be matched by a bare // slug filter. await page.goto("/jobs"); const row = page .getByRole("row") .filter({ hasText: "cancel-drain" }) .filter({ hasText: "Transcribe all" }) .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); });