Archilyzer · Source

archilyzer

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

commit 6b6e79cfdc79a17c4c2859d560db451bb69482d3
parent abcfb4233f72ce9f921e78ebd16618afef19f738
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 26 Jun 2026 16:26:00 -0400

Merge refactor/job-queue-scheduler: unified scheduler + runPool, reorder & retry

Foundational job-queue refactor (8 phases), each independently green:
- One job-kind metadata table (jobKinds.ts) replacing 4 scattered data sites;
  replay dispatch is now a data-driven lookup (adding a kind ~1 file vs ~5).
- One scheduler (scheduler.ts) owning queue ordering via a single compareTier
  comparator, shared by the registry and the worker pool (priority tiers +
  per-queue concurrency + reorder/promote). Deletes the duplicated priority
  logic.
- One concurrency primitive (runPool) replacing the auto-runner's hand-written
  loop AND whisperBatch's Promise.all; the no-event-loop-spin wait is encoded
  once and unit-tested, eliminating the Drain-all-hangs bug class by construction.
- Granular control UX: reorder/promote queued jobs, and one-click Retry +
  Retry-all for failed jobs (re-run from spec, promoted to run next).

Also fixes a pre-existing stale e2e locator (jobs-batch-tasks-drain filtered
rows by raw kind 'whisper-all' instead of the rendered 'Transcribe all' label).

49 unit tests + the queue/drain/reorder/retry e2e suites green; typecheck
(common+editor) clean.

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

Diffstat:
Mcommon/controller/autoRunner.ts | 299++++++++++++++++++++++++++++++++++++-------------------------------------------
Mcommon/controller/whisperBatch.ts | 34+++++++++++++++++++++++++---------
Acommon/jobs/concurrentRunner.test.ts | 175+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/concurrentRunner.ts | 134+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/jobKinds.test.ts | 77+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/jobKinds.ts | 232+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/registry.ts | 181++++++++++++++++++++++++++++++++++---------------------------------------------
Acommon/jobs/scheduler.test.ts | 182+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/scheduler.ts | 213+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/workerPool.test.ts | 42++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/workerPool.ts | 36+++++++++++++++++++++---------------
Meditor/CHANGELOG.md | 3+++
Meditor/app/api/test/invalidate-cache/route.ts | 6++++++
Meditor/app/jobs/actions.ts | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/active/buildActiveJobs.ts | 21+++++++++++++++++++++
Meditor/app/jobs/components/JobsTable.tsx | 4++++
Aeditor/app/jobs/components/ReorderJobButtons.tsx | 74++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/RetryAllFailedButton.tsx | 30++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/RetryJobButton.tsx | 45+++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/components/RunningJobsList.tsx | 13+++++++++++++
Meditor/app/jobs/jobKindLabels.ts | 34++++++----------------------------
Aeditor/app/jobs/jobReplayRegistry.ts | 187+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/page.tsx | 9++++++++-
Meditor/app/jobs/runJobSpec.ts | 147++++++-------------------------------------------------------------------------
Meditor/e2e/jobs-batch-tasks-drain.spec.ts | 8++++----
Aeditor/e2e/jobs-reorder.spec.ts | 103+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/e2e/jobs-retry.spec.ts | 90+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
27 files changed, 1993 insertions(+), 459 deletions(-)

diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -6,6 +6,7 @@ import { detectPlatform } from "../lib/platform"; import { getWorkerPool } from "../jobs/workerPool"; import { getRegistry } from "../jobs/registry"; import { runManagedFunction } from "../jobs/streamCommand"; +import { runPool } from "../jobs/concurrentRunner"; import { requestChannelSnapshot } from "../jobs/snapshotScheduler"; import { makeTaskTracker } from "../jobs/taskHooks"; import type { JobRunContext } from "../jobs/streamCommand"; @@ -220,6 +221,10 @@ function eligibleSlots(): number { // --- The fill-to-capacity loop --------------------------------------------- +// One unit of work selected by the policy, carried from the runPool `next` +// (selection + reservation) into `run` (execution + release). +type Picked = { pick: WorkPick; channelSlug: string; unitPlatform: string | null }; + async function runLoop( kind: AutoQueueKind, paths: Paths, @@ -239,56 +244,7 @@ async function runLoop( pruneExpired(kindState.platformBackoff, Date.now()); const completed = new Set<string>(); const completedOrder: string[] = []; - let counters = { transcribed: 0, failed: 0, skipped: 0 }; - - // Wake coordination: a finishing unit (or an abort) wakes the loop out of its - // wait. pendingWake guards against a wake that fires between waits. - let waiter: (() => void) | null = null; - let pendingWake = false; - const wake = () => { - if (waiter) { - const w = waiter; - waiter = null; - w(); - } else { - pendingWake = true; - } - }; - // Wait up to maxMs, woken early by wake() (a finishing unit) or by a hard - // cancel / drain firing during the wait (the abort listeners below). An - // ALREADY-aborted signal must NOT short-circuit to an immediately-resolved - // promise: after a drain (or hard cancel) the signal stays aborted for the - // rest of the run, so the loops that wait for in-flight units to clear - // (`while (inFlight) await waitNext()`, and the at-capacity wait once - // target == 0 while draining) would become a timer-less microtask spin that - // starves the event loop — a still-running unit's child-process exit (a - // macrotask) would then never be delivered, inFlight would never reach zero, - // and a CPU core would peg forever (the "Drain all hangs the app" bug). So we - // only fast-path a real pending wake here; an already-aborted signal falls - // through to a real timer, and a finishing unit still wakes us promptly. - const waitNext = (maxMs: number): Promise<void> => - new Promise<void>((resolve) => { - if (pendingWake) { - pendingWake = false; - resolve(); - return; - } - let done = false; - const finish = () => { - if (done) return; - done = true; - clearTimeout(timer); - signal.removeEventListener("abort", finish); - ctx.drainSignal.removeEventListener("abort", finish); - if (waiter === finish) waiter = null; - resolve(); - }; - const timer = setTimeout(finish, maxMs); - timer.unref?.(); - waiter = finish; - signal.addEventListener("abort", finish, { once: true }); - ctx.drainSignal.addEventListener("abort", finish, { once: true }); - }); + const counters = { transcribed: 0, failed: 0, skipped: 0 }; const markCompleted = (videoId: string) => { if (completed.has(videoId)) return; @@ -322,6 +278,12 @@ async function runLoop( let metaCache: ChannelMeta[] = []; let metaAt = 0; + // A graceful stop (the job record is gone after an e2e reset, or the policy was + // disabled) is modeled as a soft drain: stop picking, let in-flight finish. + // Combined with the job's real drain signal so either one ends the run. + const stopController = new AbortController(); + const drainSignal = AbortSignal.any([ctx.drainSignal, stopController.signal]); + // NOTE: do NOT call getWorkerPool().reconfigure() here. At server boot the // runner is the first thing to touch the pool, and reconfigure() sets the // pool's `initialized` flag WITHOUT applying the saved "Set as default" @@ -332,36 +294,45 @@ async function runLoop( // exactly once. Live settings changes are pushed by saveSettingsAction's own // reconfigure(applyEnabled:true), so nothing is lost by not refreshing here. - onLog(`Auto-${kind} runner started.`); + // On a HARD cancel, actively cancel each in-flight/queued child download job: a + // running one's yt-dlp is SIGTERM'd (via downloadOneManaged's signal), a + // still-queued one is dropped — either way launchUnit's `await res.done` + // resolves and the unit's run() settles so runPool's in-flight clears. (Drain + // is soft: it leaves them to finish. Transcription units settle via the hard + // signal threaded into the unit, so they need no explicit cancel here.) + const onHardAbort = () => { + for (const jid of childJobIds.values()) getRegistry().cancel(jid); + }; + if (signal.aborted) onHardAbort(); + else signal.addEventListener("abort", onHardAbort, { once: true }); + + // Transcription is throttled by the worker pool's slot count; download has no + // pool, so its real limiter is the per-platform gate in next() — leave the + // numeric cap at the policy max (or unbounded). runPool drops the target to 0 + // itself while draining, so this never needs to check the drain signal. + const limit = (): number => { + const policy = getSettings().autoQueue[kind]; + return kind === "transcription" + ? Math.min(eligibleSlots(), policy.maxWorkers ?? Number.POSITIVE_INFINITY) + : (policy.maxWorkers ?? Number.POSITIVE_INFINITY); + }; - while (!signal.aborted) { + // Select the next unit and RESERVE its slot (so the next selection sees it), + // or return null when there's nothing to pick right now / the runner should + // stop. runPool calls this only while below the limit and not draining. + const next = async (): Promise<Picked | null> => { // Stop if the job record is gone (e2e registry reset) or no longer running - // (hard cancel flips it to "cancelled"). Drain keeps it "running" and is - // handled via ctx.drainSignal below. + // (a hard cancel flips it to "cancelled" and also fires `signal`). const rec = getRegistry().get(live.jobId); - if (!rec || rec.status !== "running") break; - + if (!rec || rec.status !== "running") { + stopController.abort(); + return null; + } const policy: AutoQueuePolicy = getSettings().autoQueue[kind]; if (!policy.enabled) { onLog(`Auto-${kind} disabled — stopping runner.`); - break; - } - const drained = ctx.drainSignal.aborted; - // Transcription is throttled by the worker pool's slot count; download has no - // pool, so its real limiter is the per-platform gate below — leave the - // numeric cap at the policy max (or unbounded). - const target = drained - ? 0 - : kind === "transcription" - ? Math.min(eligibleSlots(), policy.maxWorkers ?? Number.POSITIVE_INFINITY) - : (policy.maxWorkers ?? Number.POSITIVE_INFINITY); - - if (live.inFlight.size >= target) { - // At capacity (or paused/drained). When draining/cancelling and nothing is - // left in flight, the run is complete. - if (live.inFlight.size === 0 && (drained || signal.aborted)) break; - await waitNext(IDLE_POLL_MS); - continue; + stopController.abort(); + return null; } // Refresh the (rarely-changing) channel list/platforms on a TTL. @@ -413,93 +384,19 @@ async function runLoop( } const pick = selectNextWork(policy.root, pending, runtime, live.active); - if (!pick) { - if (live.inFlight.size === 0 && (drained || signal.aborted)) break; - await waitNext(IDLE_POLL_MS); - continue; - } + if (!pick) return null; const channelSlug = owner.get(pick.videoId); if (!channelSlug) { // Shouldn't happen (every pending id came from a channel) — guard anyway. markCompleted(pick.videoId); - continue; + return null; } const unitPlatform = kind === "download" ? platformKey(channelSlug, slugToPlatform) : null; - let result: UnitResult = { outcome: "failed" }; - launchUnit({ - kind, - paths, - pick, - channelSlug, - tracker, - onLog, - signal, - drainSignal: ctx.drainSignal, - onChildJob: (jid) => childJobIds.set(pick.videoId, jid), - }) - .then((r) => { - result = r; - }) - .catch((err) => { - result = { outcome: "failed" }; - onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`); - }) - .finally(() => { - childJobIds.delete(pick.videoId); - live.inFlight.delete(pick.videoId); - for (const id of pick.path) { - live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1); - } - if (unitPlatform) { - platformInFlight.set( - unitPlatform, - Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1), - ); - } - - if (result.outcome === "transcribed") counters.transcribed++; - else if (result.outcome === "failed") counters.failed++; - else counters.skipped++; - - // Per-platform backoff: a rate-limit / network failure pauses the whole - // platform with an exponential cooldown, and the video is deliberately - // NOT marked completed so it's retried once the cooldown lapses. Any - // other outcome marks the video done-for-session; a success also clears - // the platform's cooldown. - const backoffHit = - unitPlatform !== null && - (result.failureClass === "rate_limit" || - result.failureClass === "network"); - if (backoffHit && unitPlatform) { - const entry = nextBackoff( - kindState.platformBackoff[unitPlatform], - Date.now(), - ); - kindState.platformBackoff[unitPlatform] = entry; - const secs = Math.round((entry.until - Date.now()) / 1000); - onLog( - `Auto-download: ${unitPlatform} ${result.failureClass} — backing off ${secs}s (attempt ${entry.fails}). ${pick.videoId} will retry after cooldown.`, - ); - } else { - if (unitPlatform && result.outcome === "transcribed") { - clearBackoff(kindState.platformBackoff, unitPlatform); - } - markCompleted(pick.videoId); - } - // Persist so the cooldown (and reset) survive a restart. - persist(); - // The runner job has no channelSlug, so the registry's own - // snapshot-on-finish hook doesn't fire — request a regen here so the - // channel's pending counts refresh and a restart sees fresh buckets. - requestChannelSnapshot(paths, channelSlug); - wake(); - }); - - // Reserve the slot and account active counts up the path BEFORE looping, so - // the next pick sees the reservation (and ancestor caps reflect it). + // Reserve the slot and account active counts up the path BEFORE returning, + // so the next selection sees the reservation (and ancestor caps reflect it). live.inFlight.set(pick.videoId, { videoId: pick.videoId, leafId: pick.leafId, @@ -521,17 +418,95 @@ async function runLoop( }); persist(); onLog(`Auto-${kind}: ${channelSlug}/${pick.videoId} (leaf ${pick.leafId})`); - } + return { pick, channelSlug, unitPlatform }; + }; + + // Run one reserved unit, then release its reservation and fold its outcome + // into counters/backoff. Mirrors the old `.finally` accounting exactly. + const run = async (picked: Picked, runSignal: AbortSignal): Promise<void> => { + const { pick, channelSlug, unitPlatform } = picked; + let result: UnitResult = { outcome: "failed" }; + try { + result = await launchUnit({ + kind, + paths, + pick, + channelSlug, + tracker, + onLog, + signal: runSignal, + drainSignal: ctx.drainSignal, + onChildJob: (jid) => childJobIds.set(pick.videoId, jid), + }); + } catch (err) { + result = { outcome: "failed" }; + onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`); + } finally { + childJobIds.delete(pick.videoId); + live.inFlight.delete(pick.videoId); + for (const id of pick.path) { + live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1); + } + if (unitPlatform) { + platformInFlight.set( + unitPlatform, + Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1), + ); + } + + if (result.outcome === "transcribed") counters.transcribed++; + else if (result.outcome === "failed") counters.failed++; + else counters.skipped++; + + // Per-platform backoff: a rate-limit / network failure pauses the whole + // platform with an exponential cooldown, and the video is deliberately NOT + // marked completed so it's retried once the cooldown lapses. Any other + // outcome marks the video done-for-session; a success also clears the + // platform's cooldown. + const backoffHit = + unitPlatform !== null && + (result.failureClass === "rate_limit" || + result.failureClass === "network"); + if (backoffHit && unitPlatform) { + const entry = nextBackoff( + kindState.platformBackoff[unitPlatform], + Date.now(), + ); + kindState.platformBackoff[unitPlatform] = entry; + const secs = Math.round((entry.until - Date.now()) / 1000); + onLog( + `Auto-download: ${unitPlatform} ${result.failureClass} — backing off ${secs}s (attempt ${entry.fails}). ${pick.videoId} will retry after cooldown.`, + ); + } else { + if (unitPlatform && result.outcome === "transcribed") { + clearBackoff(kindState.platformBackoff, unitPlatform); + } + markCompleted(pick.videoId); + } + // Persist so the cooldown (and reset) survive a restart. + persist(); + // The runner job has no channelSlug, so the registry's own + // snapshot-on-finish hook doesn't fire — request a regen here so the + // channel's pending counts refresh and a restart sees fresh buckets. + requestChannelSnapshot(paths, channelSlug); + } + }; + + onLog(`Auto-${kind} runner started.`); + + // One audited fill-to-capacity primitive replaces the old hand-written + // wake/waitNext loop AND its terminal drain-wait. runPool stops pulling on a + // drain (or stop), lets in-flight units finish, and never spins the event + // loop on an already-aborted signal — the "Drain all hangs the app" bug class. + await runPool<Picked>({ + next, + run, + limit, + signal, + drainSignal, + idlePollMs: IDLE_POLL_MS, + }); - // Drain remaining in-flight units (graceful stop / drain / cancel). On a HARD - // cancel, actively cancel each in-flight/queued child download job: a running - // one's yt-dlp is SIGTERM'd (via downloadOneManaged's signal), a still-queued - // one is dropped — either way launchUnit's `await res.done` resolves and the - // unit's .finally clears the slot. Drain (soft) leaves them to finish. - if (signal.aborted) { - for (const jid of childJobIds.values()) getRegistry().cancel(jid); - } - while (live.inFlight.size > 0) await waitNext(IDLE_POLL_MS); persist(); onLog( `Auto-${kind} runner finished: ${counters.transcribed} done, ${counters.failed} failed, ${counters.skipped} skipped.`, diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts @@ -8,6 +8,7 @@ import { readVideoFiles, } from "../lib/videoStatus"; import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; +import { runPool } from "../jobs/concurrentRunner"; import { pruneFailedTranscriptions } from "./failedTranscriptions"; import { resolveShardItems } from "./shard"; import { countNotYetTranscribed } from "./channels"; @@ -190,13 +191,17 @@ export async function runWhisperBatch({ let failed = 0; let skipped = 0; - // No pLimit: every video starts concurrently and parks in pool.acquire() until - // a worker slot frees. The pool's total free slots throttle the batch, and a - // freed high-priority slot is handed to the oldest waiter. When all workers are - // disabled, acquire() blocks indefinitely — the batch pauses (it does not fail) - // and resumes when a worker is re-enabled. The per-video gating + transcribe is - // shared with the cross-channel auto-runner via transcribeOneFromQueue. - const runOne = async (videoDir: string): Promise<void> => { + // The same runPool primitive the cross-channel auto-runner uses, instead of + // Promise.all over every video at once. It keeps up to `limit()` units in + // flight (the pool's eligible slot count), pulling the next as each finishes — + // so the batch fills the pool without dumping a backlog of parked acquires, and + // a drain stops pulling while letting in-flight units finish (no event-loop + // spin). When all workers are disabled the limit is 0 and the batch PAUSES (it + // does not fail), resuming when a worker is re-enabled. Per-video gating + + // transcribe is shared with the auto-runner via transcribeOneFromQueue. + const NEVER = new AbortController().signal; + let idx = 0; + const runOne = async (videoDir: string, runSignal: AbortSignal): Promise<void> => { const res = await transcribeOneFromQueue({ paths, channelSlug, @@ -206,7 +211,7 @@ export async function runWhisperBatch({ failedSet, tracker, onLog: log, - signal, + signal: runSignal, drainSignal, }); if (res.attempted) attempted++; @@ -215,7 +220,18 @@ export async function runWhisperBatch({ else skipped++; }; - await Promise.all(videoDirs.map(runOne)); + await runPool<string>({ + next: async () => (idx < videoDirs.length ? videoDirs[idx++] : null), + run: runOne, + // Eligible (enabled, non-degraded) worker slots — the real throughput + // ceiling. 0 ⇒ paused until a worker is re-enabled (runPool idle-waits). + limit: () => + pool.summary().filter((s) => s.state === "enabled" && !s.degraded).length, + signal: signal ?? NEVER, + drainSignal: drainSignal ?? NEVER, + finite: true, + idlePollMs: 3000, + }); return { attempted, succeeded: succeededCount, failed, skipped }; } diff --git a/common/jobs/concurrentRunner.test.ts b/common/jobs/concurrentRunner.test.ts @@ -0,0 +1,175 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { runPool } from "./concurrentRunner"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/concurrentRunner.test.ts +// +// runPool is the single audited fill-to-capacity primitive. The headline test is +// the already-aborted-signal no-spin case: it reproduces the "Drain all hangs" +// bug class, where a microtask-only wait starves the macrotask that delivers a +// running item's completion. With the correct wait, these resolve promptly; a +// regression would hang, which the timeout guard surfaces as a failure. + +const never = new AbortController().signal; +const sleep = (ms: number) => new Promise<void>((r) => setTimeout(r, ms)); + +function withTimeout<T>(p: Promise<T>, ms: number, msg: string): Promise<T> { + return Promise.race([ + p, + sleep(ms).then(() => { + throw new Error(`timeout: ${msg}`); + }), + ]) as Promise<T>; +} + +test("respects the dynamic limit(): never exceeds the ceiling", async () => { + let i = 0; + const items = Array.from({ length: 8 }, (_, n) => n); + let inFlight = 0; + let maxInFlight = 0; + let ran = 0; + await runPool<number>({ + next: async () => (i < items.length ? items[i++] : null), + run: async () => { + inFlight++; + maxInFlight = Math.max(maxInFlight, inFlight); + await sleep(5); + ran++; + inFlight--; + }, + limit: () => 2, + signal: never, + drainSignal: new AbortController().signal, + idlePollMs: 50, + finite: true, + }); + assert.equal(ran, 8, "all items ran"); + assert.ok(maxInFlight <= 2, `max concurrent was ${maxInFlight}, expected <= 2`); +}); + +test("finite mode stops once next() is exhausted and nothing is in flight", async () => { + let i = 0; + let ran = 0; + await withTimeout( + runPool<number>({ + next: async () => (i < 3 ? i++ : null), + run: async () => { + await sleep(1); + ran++; + }, + limit: () => 4, + signal: never, + drainSignal: new AbortController().signal, + finite: true, + }), + 2000, + "finite pool should terminate", + ); + assert.equal(ran, 3, "all 3 items ran, then the pool terminated on exhaustion"); +}); + +test("drain stops pulling new work but lets in-flight finish (no-spin)", async () => { + const drainAc = new AbortController(); + let pulled = 0; + let ranAfterDrain = 0; + let drained = false; + await withTimeout( + runPool<number>({ + // Daemon: would pull forever, but drain fires from inside the first run. + next: async () => { + if (drained) { + // Should never be pulled after drain (target drops to 0). + ranAfterDrain++; + return null; + } + return pulled++; + }, + run: async () => { + // Simulate "drain fires while a unit is actively running", then the unit + // completes via a macrotask (setTimeout). If the post-drain wait spun on + // microtasks, this timer would be starved and the pool would hang. + drainAc.abort(); + drained = true; + await sleep(20); + }, + limit: () => 1, + signal: never, + drainSignal: drainAc.signal, + idlePollMs: 1000, + }), + 2000, + "drain with a running unit must complete, not hang", + ); + assert.equal(pulled, 1, "exactly one item was pulled before drain"); + assert.equal(ranAfterDrain, 0, "no work pulled after drain"); +}); + +test("already-aborted drain at submit time: in-flight unit still completes", async () => { + // Pre-abort the drain, but stage one unit that is launched in the SAME tick it + // observes the abort. The classic starvation case: signal already aborted, a + // real timer (not a microtask) must pace the wait for the running unit. + const drainAc = new AbortController(); + let launched = false; + let finished = false; + await withTimeout( + runPool<number>({ + next: async () => { + if (launched) return null; + launched = true; + // Abort the drain right as we hand back the first (and only) item, so the + // very next loop turn is already-drained with a unit about to run. + drainAc.abort(); + return 0; + }, + run: async () => { + await sleep(20); // completes via a macrotask + finished = true; + }, + limit: () => 1, + signal: never, + drainSignal: drainAc.signal, + idlePollMs: 1000, + }), + 2000, + "already-aborted drain must not starve the running unit", + ); + // The unit was pulled before the loop re-checked the (now aborted) drain, so it + // runs to completion under soft drain. + assert.ok(finished, "the in-flight unit completed under drain"); +}); + +test("hard cancel aborts in-flight work and returns", async () => { + const cancelAc = new AbortController(); + let observedAbort = false; + await withTimeout( + runPool<number>({ + next: async () => 0, + run: async (_item, signal) => { + await new Promise<void>((resolve) => { + if (signal.aborted) { + observedAbort = true; + resolve(); + return; + } + signal.addEventListener( + "abort", + () => { + observedAbort = true; + resolve(); + }, + { once: true }, + ); + // Fire the hard cancel shortly after the unit starts. + setTimeout(() => cancelAc.abort(), 10); + }); + }, + limit: () => 1, + signal: cancelAc.signal, + drainSignal: new AbortController().signal, + idlePollMs: 1000, + }), + 2000, + "hard cancel must terminate the pool", + ); + assert.ok(observedAbort, "in-flight run() observed the hard-cancel signal"); +}); diff --git a/common/jobs/concurrentRunner.ts b/common/jobs/concurrentRunner.ts @@ -0,0 +1,134 @@ +// One event-driven fill-to-capacity primitive, replacing both the auto-runner's +// hand-written waitNext/wake loop and whisperBatch's Promise.all. Pulls work +// from `next()` up to a dynamic `limit()`, runs items concurrently, and stops on +// hard cancel (signal) or soft drain (drainSignal). Encoding the wait logic ONCE +// here — with the no-spin invariant below — removes the whole class of +// event-loop-starvation bugs (the "Drain all hangs the app" bug) that came from +// hand-rolling this loop in two places. + +export type RunPoolOptions<T> = { + // Pull the next item to run, or null/undefined if nothing is available right + // now. For a finite batch (`finite: true`) null means "exhausted". + next: () => Promise<T | null | undefined>; + // Process one item. Receives the HARD-cancel signal so it can abort promptly; + // a soft drain intentionally lets in-flight items finish. + run: (item: T, signal: AbortSignal) => Promise<void>; + // Dynamic concurrency ceiling (e.g. eligible worker slots). Re-read each loop. + limit: () => number; + signal: AbortSignal; // hard cancel: stop pulling, abort in-flight, return + drainSignal: AbortSignal; // soft drain: stop pulling, let in-flight finish, return + // Poll interval when waiting for capacity / new work. The wait is woken early + // by a finishing item or an abort, so this is just a safety ceiling. + idlePollMs?: number; + // true for a fixed work set that should stop once `next()` is exhausted and + // nothing is in flight (batch). false (default) for a long-lived daemon that + // idles waiting for new work and only stops on drain/cancel. + finite?: boolean; +}; + +export async function runPool<T>(opts: RunPoolOptions<T>): Promise<void> { + const { signal, drainSignal } = opts; + const idle = opts.idlePollMs ?? 3000; + const inFlight = new Set<Promise<void>>(); + + // Wait up to maxMs, woken early by wake() (a finishing item) or by a hard + // cancel / drain firing DURING the wait (the abort listeners below). + // + // The invariant that prevents event-loop starvation: an ALREADY-aborted signal + // must NOT short-circuit to an immediately-resolved promise. After a drain (or + // hard cancel) the signal stays aborted for the rest of the run, so a loop that + // waits for in-flight items to clear would otherwise become a timer-less + // microtask spin — starving the macrotask/timer/I/O phases, so a still-running + // item's child-process exit (a macrotask) never arrives, inFlight never empties, + // and a CPU core pegs forever. So we only fast-path a real pending wake; an + // already-aborted signal falls through to a real timer, and a finishing item + // still wakes us promptly. + let pendingWake = false; + let waiter: (() => void) | null = null; + const wake = (): void => { + if (waiter) { + const w = waiter; + waiter = null; + w(); + } else { + pendingWake = true; + } + }; + const waitNext = (maxMs: number): Promise<void> => + new Promise<void>((resolve) => { + if (pendingWake) { + pendingWake = false; + resolve(); + return; + } + let done = false; + const finish = (): void => { + if (done) return; + done = true; + clearTimeout(timer); + signal.removeEventListener("abort", finish); + drainSignal.removeEventListener("abort", finish); + if (waiter === finish) waiter = null; + resolve(); + }; + const timer = setTimeout(finish, maxMs); + timer.unref?.(); + waiter = finish; + // Registering on an ALREADY-aborted signal does not fire — the timer paces + // the wait once draining, which is exactly what we want. + signal.addEventListener("abort", finish, { once: true }); + drainSignal.addEventListener("abort", finish, { once: true }); + }); + + while (true) { + if (signal.aborted) break; + const drained = drainSignal.aborted; + const target = drained ? 0 : Math.max(0, opts.limit()); + + if (inFlight.size >= target) { + // At/over capacity, or draining (target 0). Only stop once nothing is in + // flight AND we've been told to (drain/cancel). A transiently-zero limit + // while neither draining nor cancelling just waits — it does NOT terminate. + if (inFlight.size === 0 && (drainSignal.aborted || signal.aborted)) break; + await waitNext(idle); + continue; + } + + const item = await opts.next(); + // Re-read the abort state after the async next(): a drain/cancel may have + // fired while it ran (using the live signals, not the stale top-of-loop + // snapshot, makes a stop that arrives mid-next() take effect immediately). + if (item === null || item === undefined) { + // Nothing available. A finite batch with nothing in flight is done; a + // daemon idles unless draining/cancelling. + if ( + inFlight.size === 0 && + (drainSignal.aborted || signal.aborted || opts.finite) + ) { + break; + } + await waitNext(idle); + continue; + } + + // async thunk so a synchronous throw from run() becomes a contained + // rejection (settled by Promise.allSettled) rather than escaping the pool. + const p = (async () => opts.run(item, signal))(); + inFlight.add(p); + // Remove from the in-flight set once it settles and wake the loop so it can + // pull more / observe drain completion. Use then(settle, settle) rather than + // finally(): the rejection handler CONSUMES a rejecting run() (no leaked + // unhandledRejection), while Promise.allSettled below still awaits `p`. `p` + // is declared before `settle`, so the closure reference is safe. + const settle = () => { + inFlight.delete(p); + wake(); + }; + void p.then(settle, settle); + } + + // Let in-flight items settle: a soft drain lets them finish; a hard cancel + // lets them observe `signal` and abort. Either way we don't return until the + // event loop has delivered their completions. + await Promise.allSettled([...inFlight]); +} diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts @@ -0,0 +1,77 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { isDrainableKind, jobKindLabel, getJobKind } from "./jobKinds"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/jobKinds.test.ts +// +// Phase 1 of the queue refactor moved the drainable-kinds set and the label map +// into one metadata table (jobKinds.ts). These snapshots pin the PRE-refactor +// values verbatim, so the consolidation is provably behavior-preserving: every +// kind that drained before still drains, every label is byte-identical, and +// kinds that previously had no label still fall back to their raw kind. + +// Verbatim copy of the old DRAINABLE_KINDS set from registry.ts. +const OLD_DRAINABLE = [ + "whisper-all", + "whisper-bucket-downloaded-no-transcript", + "download-from-playlist", + "download-missing", + "download-missing-subs", + "sync", + "retry-bucket", + "auto-transcribe", + "auto-download", +]; + +// Verbatim copy of the old JOB_KIND_LABELS map from editor jobKindLabels.ts. +const OLD_LABELS: Record<string, string> = { + "auto-transcribe": "Auto-transcribe", + "auto-download": "Auto-download runner", + "auto-download-unit": "Auto-download", + "whisper-all": "Transcribe all", + "whisper-bucket-downloaded-no-transcript": "Transcribe downloaded audio", + "download-from-playlist": "Download from playlist", + "download-missing": "Download missing", + "download-missing-subs": "Download missing subs", + "import-one": "Import video", + "redownload-archive": "Archive source video", + "retry-bucket": "Retry", + "clean-audio-transcribed": "Clean audio", + "check-kept-deleted": "Check kept videos", + "persist-kept": "Persist kept videos", + "backup-saved-videos": "Back up saved videos", + "verify-saved-video-backup": "Verify saved-video backup", + sync: "Sync", +}; + +test("drainable kinds match the old DRAINABLE_KINDS set exactly", () => { + const drainSet = new Set(OLD_DRAINABLE); + for (const kind of OLD_DRAINABLE) { + assert.equal(isDrainableKind(kind), true, `${kind} should drain`); + } + // No kind in the table drains unless it was in the old set. + for (const kind of Object.keys(OLD_LABELS)) { + assert.equal( + isDrainableKind(kind), + drainSet.has(kind), + `${kind} drainability changed`, + ); + } + // Unknown kinds are not drainable. + assert.equal(isDrainableKind("no-such-kind"), false); +}); + +test("labels match the old JOB_KIND_LABELS map exactly", () => { + for (const [kind, label] of Object.entries(OLD_LABELS)) { + assert.equal(jobKindLabel(kind), label, `label for ${kind}`); + } +}); + +test("label-less and unknown kinds fall back to the raw kind", () => { + // Replayable kinds that never had a label entry must still fall back. + assert.equal(jobKindLabel("store-playlist"), "store-playlist"); + assert.equal(jobKindLabel("transcode-failures"), "transcode-failures"); + // Genuinely unknown kind. + assert.equal(jobKindLabel("totally-unknown"), "totally-unknown"); + assert.equal(getJobKind("totally-unknown"), undefined); +}); diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -0,0 +1,232 @@ +// The single source of truth for what a job `kind` IS. Before this table, the +// same per-kind knowledge was scattered across four places: the DRAINABLE_KINDS +// set (registry.ts), the JOB_KIND_LABELS map (editor jobKindLabels.ts), the +// switch in editor runJobSpec.ts, and the spec-presence bookmark check. Phase 1 +// of the queue refactor consolidates the drainable + label data here; later +// phases consume `defaultTier` (scheduler) and lean on `bookmarkable`. +// +// Adding a new job kind should mean adding ONE entry here (plus its replay +// handler in editor/app/jobs/jobReplayRegistry.ts if it is bookmarkable). + +// How a kind picks its registry queueKey. Descriptive only — the real key is +// still computed by the action that creates the job (channelQueueKey / +// downloadQueueKey / ""). "parallel" means queueKey === "" (runs unserialized). +export type QueueKeyStrategy = "parallel" | "platform" | "custom"; + +// Scheduler priority tiers. A higher-priority tier's queued jobs run before a +// lower one's; FIFO within a tier. Consumed starting Phase 4 (the scheduler); +// today every non-background job is effectively "foreground". +export type SchedulerTier = "urgent" | "foreground" | "background"; + +export type JobKindMeta = { + kind: string; + // Human-readable label for the UI. Optional: a kind without one falls back to + // its raw machine kind (preserving the prior JOB_KIND_LABELS fallback), so a + // new kind is never invisible. + label?: string; + // The Drain (soft-cancel) button is only offered for running jobs of a + // drainable kind — i.e. one whose controller honors the drain signal (stops + // starting new sub-operations, lets in-flight ones finish). + drainable: boolean; + // Whether this kind's action attaches a replayable JobSpec (and thus has a + // replay handler). Actual bookmark-ability still keys off spec PRESENCE on the + // record at runtime — this flag just says the kind CAN be bookmarked/retried. + bookmarkable: boolean; + // Descriptive: how the action derives its queueKey. Not consumed as logic. + queueKeyStrategy: QueueKeyStrategy; + // Fallback scheduler tier when a record is not explicitly background. Left + // undefined today so Phase 4 maps it to "foreground" (unchanged behavior). + defaultTier?: SchedulerTier; +}; + +// One entry per kind known to the system. `label` is included only where the +// old JOB_KIND_LABELS had one, so jobKindLabel() behavior is byte-identical. +const JOB_KINDS: Record<string, JobKindMeta> = { + "auto-transcribe": { + kind: "auto-transcribe", + label: "Auto-transcribe", + drainable: true, + bookmarkable: false, + queueKeyStrategy: "parallel", + }, + "auto-download": { + kind: "auto-download", + label: "Auto-download runner", + drainable: true, + bookmarkable: false, + queueKeyStrategy: "parallel", + }, + "auto-download-unit": { + kind: "auto-download-unit", + label: "Auto-download", + drainable: false, + bookmarkable: false, + queueKeyStrategy: "platform", + }, + "whisper-all": { + kind: "whisper-all", + label: "Transcribe all", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "whisper-bucket-downloaded-no-transcript": { + kind: "whisper-bucket-downloaded-no-transcript", + label: "Transcribe downloaded audio", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "download-from-playlist": { + kind: "download-from-playlist", + label: "Download from playlist", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "platform", + }, + "download-missing": { + kind: "download-missing", + label: "Download missing", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "platform", + }, + "download-missing-subs": { + kind: "download-missing-subs", + label: "Download missing subs", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "platform", + }, + "import-one": { + kind: "import-one", + label: "Import video", + drainable: false, + bookmarkable: false, + queueKeyStrategy: "custom", + }, + "redownload-archive": { + kind: "redownload-archive", + label: "Archive source video", + drainable: false, + bookmarkable: false, + queueKeyStrategy: "custom", + }, + "retry-bucket": { + kind: "retry-bucket", + label: "Retry", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "platform", + }, + "clean-audio-transcribed": { + kind: "clean-audio-transcribed", + label: "Clean audio", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "check-kept-deleted": { + kind: "check-kept-deleted", + label: "Check kept videos", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "persist-kept": { + kind: "persist-kept", + label: "Persist kept videos", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "backup-saved-videos": { + kind: "backup-saved-videos", + label: "Back up saved videos", + drainable: false, + bookmarkable: false, + queueKeyStrategy: "custom", + }, + "verify-saved-video-backup": { + kind: "verify-saved-video-backup", + label: "Verify saved-video backup", + drainable: false, + bookmarkable: false, + queueKeyStrategy: "custom", + }, + sync: { + kind: "sync", + label: "Sync", + drainable: true, + bookmarkable: true, + queueKeyStrategy: "platform", + }, + // Replayable kinds that never had a JOB_KIND_LABELS entry: label omitted so + // jobKindLabel() keeps falling back to the raw kind (unchanged behavior). + "store-playlist": { + kind: "store-playlist", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "clear-failed-transcriptions": { + kind: "clear-failed-transcriptions", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "transcode-failures": { + kind: "transcode-failures", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "transcode-untranscoded": { + kind: "transcode-untranscoded", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "remove-failed-transcodings": { + kind: "remove-failed-transcodings", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "clear-failed-transcodings": { + kind: "clear-failed-transcodings", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "clean-extra-audio-formats": { + kind: "clean-extra-audio-formats", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, + "remove-wrong-format-audio": { + kind: "remove-wrong-format-audio", + drainable: false, + bookmarkable: true, + queueKeyStrategy: "custom", + }, +}; + +export function getJobKind(kind: string): JobKindMeta | undefined { + return JOB_KINDS[kind]; +} + +// Human-readable label, falling back to the raw kind for unknown/label-less +// kinds (a new kind is never invisible). +export function jobKindLabel(kind: string): string { + return JOB_KINDS[kind]?.label ?? kind; +} + +export function isDrainableKind(kind: string): boolean { + return JOB_KINDS[kind]?.drainable ?? false; +} + +export function isBookmarkableKind(kind: string): boolean { + return JOB_KINDS[kind]?.bookmarkable ?? false; +} diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -1,5 +1,7 @@ import type { ChildProcess } from "node:child_process"; import type { JobSpec } from "./jobSpec"; +import { getScheduler } from "./scheduler"; +import { getJobKind, type SchedulerTier } from "./jobKinds"; export type JobStatus = | "queued" @@ -88,11 +90,14 @@ export type QueueSnapshot = { type StartFn = () => void; type CancelFn = () => void; +// The registry owns job LIFECYCLE/STATE (the JobRecord, meta sidecars, tasks, +// terminal transitions). Queue ORDERING — which job runs vs. waits, and the +// foreground-before-background priority — is delegated to the shared Scheduler +// (scheduler.ts), the single source of truth that the worker pool also shares. +// Every non-"" job is submitted with concurrency 1 (strict per-queueKey serial), +// preserving the long-standing one-running-job-per-queue behavior. class JobRegistry { private jobs = new Map<string, JobRecord>(); - private queues = new Map<string, JobRecord[]>(); - private starts = new WeakMap<JobRecord, StartFn>(); - private cancels = new WeakMap<JobRecord, CancelFn>(); register(record: JobRecord): void { this.jobs.set(record.id, record); @@ -118,65 +123,45 @@ class JobRegistry { ); } - // Add a job to its queue. If the queue is empty, start it immediately and - // mark "running". Otherwise leave it "queued" and remember the start/cancel - // callbacks for when it becomes head-of-queue or is cancelled. + // Submit a job for scheduling. The scheduler decides whether it runs now (its + // queueKey has a free slot) or waits, and invokes `start` when it becomes head + // of its queue — immediately, or later when an earlier job finalizes. // - // Priority: a foreground (default) job is inserted AHEAD of any queued - // background jobs but never displaces the running head (q[0]) — so a manually - // clicked Sync jumps ahead of the auto-download runner's queued units without - // interrupting one that's already mid-download. Foreground jobs keep FIFO - // order among themselves; background jobs always append. finalize() still just - // promotes q[0], which this ordering keeps as the highest-priority job. + // Priority: a foreground (default) job runs ahead of any queued background + // jobs but never displaces a running one — so a manually clicked Sync jumps + // ahead of the auto-download runner's queued units without interrupting one + // that's already mid-download. This is the scheduler's compareTier ordering + // (background -> "background" tier, everything else -> "foreground"). // - // queueKey === "" is special: the job bypasses queue serialization entirely - // and runs immediately, parallel to any other work. Cancellation still - // works (running → SIGTERM/abort) but there's no queue map entry to clean - // up in finalize. + // queueKey === "" is special: the job bypasses serialization entirely and runs + // immediately, parallel to any other work, untracked by the scheduler. + // Cancellation still works (running -> SIGTERM/abort). enqueue( record: JobRecord, callbacks: { start: StartFn; onCancel: CancelFn }, ): { willRunNow: boolean; position: number } { - this.starts.set(record, callbacks.start); - this.cancels.set(record, callbacks.onCancel); - if (record.queueKey === "") { + // Setting status/startedAt is coupled with actually launching, so wrap the + // caller's start: the scheduler runs this exactly when the job starts. + const start: StartFn = () => { record.status = "running"; record.startedAt = Date.now(); callbacks.start(); - return { willRunNow: true, position: 0 }; - } - let q = this.queues.get(record.queueKey); - if (!q) { - q = []; - this.queues.set(record.queueKey, q); - } - if (q.length === 0) { - q.push(record); - record.status = "running"; - record.startedAt = Date.now(); - callbacks.start(); - return { willRunNow: true, position: 0 }; - } - if (record.background) { - // Background work always queues at the very back. - q.push(record); - } else { - // Foreground work slots in before the first QUEUED background job (never - // before the running head at index 0), preserving FIFO among foreground. - let insertAt = q.length; - for (let i = 1; i < q.length; i++) { - if (q[i].background) { - insertAt = i; - break; - } - } - q.splice(insertAt, 0, record); - } - return { willRunNow: false, position: q.indexOf(record) }; + }; + const tier: SchedulerTier = record.background + ? "background" + : getJobKind(record.kind)?.defaultTier ?? "foreground"; + return getScheduler().submit({ + id: record.id, + queueKey: record.queueKey, + tier, + concurrency: 1, + start, + onCancel: callbacks.onCancel, + }); } - // Idempotent: marks the job terminal (if not already), splices it out of - // its queue, and starts the next queued job in that queue. + // Idempotent: marks the job terminal (if not already), then tells the + // scheduler it is done so the next queued job on that queueKey is promoted. finalize( id: string, status: "done" | "failed" | "cancelled", @@ -192,29 +177,16 @@ class JobRegistry { // 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); - if (idx >= 0) q.splice(idx, 1); - if (q.length === 0) { - this.queues.delete(job.queueKey); - return; - } - const next = q[0]; - if (next.status === "queued") { - next.status = "running"; - next.startedAt = Date.now(); - const startFn = this.starts.get(next); - if (startFn) startFn(); - } + getScheduler().complete(id); } cancel(id: string): boolean { const job = this.jobs.get(id); if (!job) return false; if (job.status === "queued") { - const onCancel = this.cancels.get(job); - if (onCancel) onCancel(); + // Remove from the scheduler (fires the queued job's onCancel), then mark + // it terminal. finalize's scheduler.complete is a no-op by then. + getScheduler().cancel(id); this.finalize(id, "cancelled"); return true; } @@ -307,31 +279,48 @@ class JobRegistry { job.completedTaskMs = (job.completedTaskMs ?? 0) + durationMs; } - // Snapshot of queues for UI. Sorted by queue name. + // Snapshot of queues for UI, resolved from the scheduler's id views back to + // JobRecords. Already sorted by queue name (the scheduler sorts). Running/ + // queued ids always resolve: the 100-job eviction never drops non-terminal + // jobs, but filter defensively anyway. listQueues(): QueueSnapshot[] { - const out: QueueSnapshot[] = []; - for (const [name, q] of this.queues.entries()) { - const running = q.find((j) => j.status === "running"); - const queued = q.filter((j) => j.status === "queued"); - out.push({ name, running, queued }); - } - return out.sort((a, b) => a.name.localeCompare(b.name)); + return getScheduler() + .queues() + .map((v) => ({ + name: v.name, + running: v.running + .map((id) => this.jobs.get(id)) + .find((j): j is JobRecord => Boolean(j)), + queued: v.queued + .map((id) => this.jobs.get(id)) + .filter((j): j is JobRecord => Boolean(j)), + })); } // Names of queues that currently have any non-terminal jobs. Used by the // UI to populate the QueuePicker dropdown. activeQueueNames(): string[] { - return Array.from(this.queues.keys()).sort(); + return getScheduler() + .queues() + .map((v) => v.name); } - // Position of the job in its queue (0 = currently running). Returns -1 if - // the job is no longer in any queue (terminal). + // Position of the job in its queue (0 = currently running). Returns -1 if the + // job is no longer in any queue (terminal or unserialized ""). positionInQueue(id: string): number { - const job = this.jobs.get(id); - if (!job) return -1; - const q = this.queues.get(job.queueKey); - if (!q) return -1; - return q.indexOf(job); + return getScheduler().positionInQueue(id); + } + + // Granular manual control: move a QUEUED job one slot toward (-1) or away from + // (+1) the head, or to the front of the queued section. Never reorders a + // running job. Returns false if the move isn't possible. Delegates to the + // scheduler, which owns ordering. + reorder(id: string, dir: -1 | 1): boolean { + return getScheduler().reorder(id, dir); + } + + promote(id: string): boolean { + return getScheduler().promote(id); } } @@ -351,23 +340,7 @@ 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", - // The auto-queue runners: drain stops launching new units, lets in-flight - // finish, then the loop returns and the job completes. See autoRunner.ts. - "auto-transcribe", - "auto-download", -]); - -export function isDrainableKind(kind: string): boolean { - return DRAINABLE_KINDS.has(kind); -} +// Which kinds honor the drain signal now lives in the job-kind metadata table +// (jobKinds.ts), the single source of truth. Re-exported here so existing +// importers of `isDrainableKind` from the registry keep resolving unchanged. +export { isDrainableKind } from "./jobKinds"; diff --git a/common/jobs/scheduler.test.ts b/common/jobs/scheduler.test.ts @@ -0,0 +1,182 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { createScheduler, compareTier, type ScheduledTask } from "./scheduler"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/scheduler.test.ts +// +// The scheduler is the single source of truth for queue ordering, generalizing +// registry.enqueue's foreground-before-background insertion to priority tiers. +// Each test uses a fresh createScheduler() so there is no shared global state. + +let seq = 0; +function task( + queueKey: string, + tier: ScheduledTask["tier"], + concurrency = 1, +): { task: ScheduledTask; started: () => boolean; cancelled: () => boolean } { + let started = false; + let cancelled = false; + const t: ScheduledTask = { + id: `t${seq++}`, + queueKey, + tier, + concurrency, + start: () => { + started = true; + }, + onCancel: () => { + cancelled = true; + }, + }; + return { task: t, started: () => started, cancelled: () => cancelled }; +} + +test("compareTier orders urgent < foreground < background", () => { + assert.ok(compareTier("urgent", "foreground") < 0); + assert.ok(compareTier("foreground", "background") < 0); + assert.ok(compareTier("background", "urgent") > 0); + assert.equal(compareTier("foreground", "foreground"), 0); +}); + +test("first task on a queue runs immediately", () => { + const s = createScheduler(); + const a = task("k", "foreground"); + const r = s.submit(a.task); + assert.deepEqual(r, { willRunNow: true, position: 0 }); + assert.ok(a.started()); + assert.deepEqual(s.queues(), [{ name: "k", running: [a.task.id], queued: [] }]); +}); + +test("tiers order the queued section: urgent before foreground before background", () => { + const s = createScheduler(); + const head = task("k", "foreground"); // runs + const bg = task("k", "background"); + const fg = task("k", "foreground"); + const urg = task("k", "urgent"); + s.submit(head.task); + s.submit(bg.task); + s.submit(fg.task); + s.submit(urg.task); + const view = s.queues()[0]; + assert.deepEqual(view.running, [head.task.id]); + assert.deepEqual(view.queued, [urg.task.id, fg.task.id, bg.task.id]); + // Only the running head actually started. + assert.ok(head.started()); + assert.ok(!fg.started() && !bg.started() && !urg.started()); +}); + +test("FIFO within a tier", () => { + const s = createScheduler(); + const head = task("k", "foreground"); + const a = task("k", "foreground"); + const b = task("k", "foreground"); + s.submit(head.task); + s.submit(a.task); + s.submit(b.task); + assert.deepEqual(s.queues()[0].queued, [a.task.id, b.task.id]); +}); + +test("concurrency > 1 runs that many at once, queues the rest", () => { + const s = createScheduler(); + const ts = Array.from({ length: 5 }, () => task("k", "foreground", 3)); + ts.forEach((t) => s.submit(t.task)); + const view = s.queues()[0]; + assert.equal(view.running.length, 3); + assert.equal(view.queued.length, 2); + assert.ok(ts.slice(0, 3).every((t) => t.started())); + assert.ok(ts.slice(3).every((t) => !t.started())); +}); + +test("complete promotes the next queued task and fires its start", () => { + const s = createScheduler(); + const head = task("k", "foreground"); + const next = task("k", "foreground"); + s.submit(head.task); + s.submit(next.task); + assert.ok(!next.started()); + s.complete(head.task.id); + assert.ok(next.started()); + assert.deepEqual(s.queues(), [{ name: "k", running: [next.task.id], queued: [] }]); +}); + +test("completing the last task removes the queue", () => { + const s = createScheduler(); + const a = task("k", "foreground"); + s.submit(a.task); + s.complete(a.task.id); + assert.deepEqual(s.queues(), []); + assert.equal(s.positionInQueue(a.task.id), -1); +}); + +test("cancel removes a queued task and fires onCancel; running/unknown return false", () => { + const s = createScheduler(); + const head = task("k", "foreground"); + const q = task("k", "foreground"); + s.submit(head.task); + s.submit(q.task); + assert.equal(s.cancel(q.task.id), true); + assert.ok(q.cancelled()); + assert.deepEqual(s.queues()[0].queued, []); + // Running head can't be cancelled via the scheduler (caller hard-aborts). + assert.equal(s.cancel(head.task.id), false); + assert.equal(s.cancel("nope"), false); +}); + +test("promote moves a queued task to the front of the queued section", () => { + const s = createScheduler(); + const head = task("k", "foreground"); + const a = task("k", "background"); + const b = task("k", "background"); + s.submit(head.task); + s.submit(a.task); + s.submit(b.task); + assert.deepEqual(s.queues()[0].queued, [a.task.id, b.task.id]); + assert.equal(s.promote(b.task.id), true); + assert.deepEqual(s.queues()[0].queued, [b.task.id, a.task.id]); + // Running task can't be promoted. + assert.equal(s.promote(head.task.id), false); +}); + +test("reorder swaps adjacent queued tasks but never moves the running head", () => { + const s = createScheduler(); + const head = task("k", "foreground"); + const a = task("k", "foreground"); + const b = task("k", "foreground"); + s.submit(head.task); + s.submit(a.task); + s.submit(b.task); + // Move b up past a. + assert.equal(s.reorder(b.task.id, -1), true); + assert.deepEqual(s.queues()[0].queued, [b.task.id, a.task.id]); + // b is now at the front of the queued section: can't move up into running. + assert.equal(s.reorder(b.task.id, -1), false); + // a is last: can't move down past the end. + assert.equal(s.reorder(a.task.id, 1), false); + // Running head can't be reordered. + assert.equal(s.reorder(head.task.id, 1), false); +}); + +test('queueKey "" runs immediately in parallel and is not tracked', () => { + const s = createScheduler(); + const a = task("", "foreground"); + const b = task("", "foreground"); + assert.deepEqual(s.submit(a.task), { willRunNow: true, position: 0 }); + assert.deepEqual(s.submit(b.task), { willRunNow: true, position: 0 }); + assert.ok(a.started() && b.started()); + assert.deepEqual(s.queues(), []); + assert.equal(s.positionInQueue(a.task.id), -1); + s.complete(a.task.id); // no-op, no throw +}); + +test("positionInQueue reports 0 for running head, then queued offsets", () => { + const s = createScheduler(); + const head = task("k", "foreground"); + const a = task("k", "foreground"); + const b = task("k", "foreground"); + s.submit(head.task); + s.submit(a.task); + s.submit(b.task); + assert.equal(s.positionInQueue(head.task.id), 0); + assert.equal(s.positionInQueue(a.task.id), 1); + assert.equal(s.positionInQueue(b.task.id), 2); +}); diff --git a/common/jobs/scheduler.ts b/common/jobs/scheduler.ts @@ -0,0 +1,213 @@ +import type { SchedulerTier } from "./jobKinds"; + +// The unified queue-ordering core. Before this, "foreground preempts background" +// was implemented twice — once in registry.enqueue (job ordering) and again in +// workerPool.acquire (slot waiters). Both are now expressible through ONE +// comparator (compareTier) and one ordering policy here, generalized from a +// background boolean to priority tiers. +// +// A scheduler owns N independent queues keyed by queueKey. Within a queue, up to +// `concurrency` tasks run at once (the "running" prefix); the rest wait, ordered +// by tier (urgent < foreground < background) with FIFO inside a tier. queueKey +// === "" is special: such tasks bypass serialization and run immediately in +// parallel (preserving the registry's long-standing behavior), and are not +// tracked in any queue. + +// Lower number = higher priority = runs sooner. Negative when `a` outranks `b`. +const TIER_RANK: Record<SchedulerTier, number> = { + urgent: 0, + foreground: 1, + background: 2, +}; + +export function compareTier(a: SchedulerTier, b: SchedulerTier): number { + return TIER_RANK[a] - TIER_RANK[b]; +} + +export type ScheduledTask = { + id: string; + queueKey: string; + tier: SchedulerTier; + // Max simultaneously-running tasks on this queueKey (1 = strict serial, today's + // behavior for every non-"" queue). The first submitter's value wins for a + // given live queue. + concurrency: number; + // Becomes-head-of-queue callback: invoked exactly once when the task starts + // running (immediately on submit if there's room, else later via complete()). + start: () => void; + // Invoked if the task is cancelled while still queued (never started). + onCancel: () => void; +}; + +export type QueueView = { name: string; running: string[]; queued: string[] }; + +export interface Scheduler { + submit(t: ScheduledTask): { willRunNow: boolean; position: number }; + // Mark a task done (normal completion OR post-cancel cleanup) and promote the + // next eligible queued task(s). Idempotent; no-op if the id is unknown. + complete(id: string): void; + // Remove a QUEUED task, firing its onCancel. Returns false for a running or + // unknown task (the caller hard-aborts a running one, then calls complete()). + cancel(id: string): boolean; + // Granular control: move a QUEUED task one slot toward (-1) or away from (+1) + // the head, within the queued section only. Never moves a running task. + reorder(id: string, dir: -1 | 1): boolean; + // Move a QUEUED task to the front of the queued section. Never preempts a + // running task. + promote(id: string): boolean; + // 0 = running head, 1.. = queued position. -1 if gone or untracked (""). + positionInQueue(id: string): number; + queues(): QueueView[]; +} + +type Entry = ScheduledTask & { status: "running" | "queued" }; + +// Invariant for every queue array: all "running" entries form a contiguous +// prefix, followed by all "queued" entries. submit/complete/promote/reorder all +// preserve this. +class InMemoryScheduler implements Scheduler { + private queuesByKey = new Map<string, Entry[]>(); + private byId = new Map<string, Entry>(); + + submit(t: ScheduledTask): { willRunNow: boolean; position: number } { + // Unserialized: run immediately, never tracked. + if (t.queueKey === "") { + t.start(); + return { willRunNow: true, position: 0 }; + } + let q = this.queuesByKey.get(t.queueKey); + if (!q) { + q = []; + this.queuesByKey.set(t.queueKey, q); + } + const runningCount = q.filter((e) => e.status === "running").length; + if (runningCount < t.concurrency) { + // Room to run now: append to the end of the running prefix. + const entry: Entry = { ...t, status: "running" }; + q.splice(runningCount, 0, entry); + this.byId.set(t.id, entry); + entry.start(); + return { willRunNow: true, position: runningCount }; + } + // At capacity: queue by tier, inserting before the first queued entry of a + // strictly lower priority (preserving FIFO within a tier). + const entry: Entry = { ...t, status: "queued" }; + let insertAt = q.length; + for (let i = runningCount; i < q.length; i++) { + if (compareTier(t.tier, q[i].tier) < 0) { + insertAt = i; + break; + } + } + q.splice(insertAt, 0, entry); + this.byId.set(t.id, entry); + return { willRunNow: false, position: insertAt }; + } + + complete(id: string): void { + const entry = this.byId.get(id); + if (!entry) return; + this.byId.delete(id); + const q = this.queuesByKey.get(entry.queueKey); + if (!q) return; + const idx = q.indexOf(entry); + if (idx >= 0) q.splice(idx, 1); + this.fill(entry.queueKey, q); + } + + cancel(id: string): boolean { + const entry = this.byId.get(id); + if (!entry || entry.status !== "queued") return false; + this.byId.delete(id); + const q = this.queuesByKey.get(entry.queueKey); + if (!q) return true; + const idx = q.indexOf(entry); + if (idx >= 0) q.splice(idx, 1); + // Removing a queued entry frees no slot, so no promotion is needed. + if (q.length === 0) this.queuesByKey.delete(entry.queueKey); + entry.onCancel(); + return true; + } + + reorder(id: string, dir: -1 | 1): boolean { + const entry = this.byId.get(id); + if (!entry || entry.status !== "queued") return false; + const q = this.queuesByKey.get(entry.queueKey); + if (!q) return false; + const i = q.indexOf(entry); + const j = i + dir; + // Stay within the queued section: never swap with a running entry or past + // the array bounds. + if (j < 0 || j >= q.length || q[j].status !== "queued") return false; + [q[i], q[j]] = [q[j], q[i]]; + return true; + } + + promote(id: string): boolean { + const entry = this.byId.get(id); + if (!entry || entry.status !== "queued") return false; + const q = this.queuesByKey.get(entry.queueKey); + if (!q) return false; + const i = q.indexOf(entry); + const runningCount = q.filter((e) => e.status === "running").length; + if (i <= runningCount) return false; // already at the front of the queue + q.splice(i, 1); + q.splice(runningCount, 0, entry); + return true; + } + + positionInQueue(id: string): number { + const entry = this.byId.get(id); + if (!entry) return -1; + const q = this.queuesByKey.get(entry.queueKey); + if (!q) return -1; + return q.indexOf(entry); + } + + queues(): QueueView[] { + const out: QueueView[] = []; + for (const [name, q] of this.queuesByKey.entries()) { + out.push({ + name, + running: q.filter((e) => e.status === "running").map((e) => e.id), + queued: q.filter((e) => e.status === "queued").map((e) => e.id), + }); + } + return out.sort((a, b) => a.name.localeCompare(b.name)); + } + + // Promote queued entries to running until the concurrency ceiling is reached. + // The queued entries are already in priority order, and sit immediately after + // the running prefix, so promoting the boundary entry preserves the invariant. + private fill(queueKey: string, q: Entry[]): void { + if (q.length === 0) { + this.queuesByKey.delete(queueKey); + return; + } + let runningCount = q.filter((e) => e.status === "running").length; + const cap = q[0]?.concurrency ?? 1; + while (runningCount < cap && runningCount < q.length) { + const next = q[runningCount]; + if (next.status !== "queued") break; + next.status = "running"; + runningCount++; + next.start(); + } + } +} + +declare global { + // eslint-disable-next-line no-var + var __yttScheduler__: Scheduler | undefined; +} + +export function createScheduler(): Scheduler { + return new InMemoryScheduler(); +} + +export function getScheduler(): Scheduler { + if (!globalThis.__yttScheduler__) { + globalThis.__yttScheduler__ = new InMemoryScheduler(); + } + return globalThis.__yttScheduler__; +} diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts @@ -118,3 +118,45 @@ test("all-background acquires fall back to plain FIFO (backward compatible)", as await b2; assert.deepEqual(order, ["B1", "B2"]); }); + +test("tier ordering: urgent jumps ahead of foreground, which jumps ahead of background", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + // Park in arrival order background, foreground, urgent — served in tier order. + const b = pool + .acquire(undefined, { tier: "background" }) + .then((l) => (order.push("B"), l)); + const f = pool + .acquire(undefined, { tier: "foreground" }) + .then((l) => (order.push("F"), l)); + const u = pool + .acquire(undefined, { tier: "urgent" }) + .then((l) => (order.push("U"), l)); + + let lease_ = lease; + for (const p of [u, f, b]) { + lease_.release(); + lease_ = await p; + } + assert.deepEqual(order, ["U", "F", "B"]); +}); + +test("legacy { background: true } maps to the background tier", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + // A background (legacy flag) waiter parks first; a tier:foreground waiter must + // still preempt it, proving the boolean maps onto the same comparator. + const b = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B"), l)); + const f = pool + .acquire(undefined, { tier: "foreground" }) + .then((l) => (order.push("F"), l)); + + lease.release(); + const lf = await f; + assert.deepEqual(order, ["F"]); + lf.release(); + await b; + assert.deepEqual(order, ["F", "B"]); +}); diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts @@ -28,6 +28,8 @@ import type { Worker } from "../lib/workers"; import { getSettings } from "../lib/settings"; import { getPaths } from "../lib/paths"; import { readWorkerDefaults } from "./workerDefaults"; +import { compareTier } from "./scheduler"; +import type { SchedulerTier } from "./jobKinds"; // Consecutive transport/exec failures before a worker is auto-marked degraded // (skipped by the scheduler until re-enabled). See markFailure / Phase 5. @@ -73,9 +75,10 @@ type Waiter = { resolve: (lease: Lease) => void; reject: (err: Error) => void; onAbort?: () => void; - // Background waiters queue behind foreground ones (see enqueueWaiter). Auto - // (background) transcriptions yield to manual (foreground) ones for free slots. - background?: boolean; + // Priority tier (see enqueueWaiter). Background waiters (auto transcriptions) + // queue behind foreground (manual) ones for free slots; an urgent waiter would + // jump ahead of both. Same compareTier ordering the registry's scheduler uses. + tier: SchedulerTier; }; export class WorkerPool { @@ -277,19 +280,15 @@ export class WorkerPool { } } - // Insert a parked waiter respecting priority: a foreground waiter slots in - // BEFORE the first background waiter (so manual work jumps ahead of queued auto - // work), preserving FIFO within each class; a background waiter always appends. - // Mirrors the registry queue's enqueue (registry.ts) so one priority model - // governs both schedulers. + // Insert a parked waiter respecting priority: a waiter slots in BEFORE the + // first waiter of strictly lower priority (so manual/foreground work jumps + // ahead of queued auto/background work), preserving FIFO within a tier. Uses + // the SAME compareTier comparator as the registry's scheduler — one priority + // model governs both, instead of two hand-rolled copies. private enqueueWaiter(waiter: Waiter): void { - if (waiter.background) { - this.waiters.push(waiter); - return; - } let insertAt = this.waiters.length; for (let i = 0; i < this.waiters.length; i++) { - if (this.waiters[i].background) { + if (compareTier(waiter.tier, this.waiters[i].tier) < 0) { insertAt = i; break; } @@ -303,7 +302,12 @@ export class WorkerPool { // foreground (manual) waiter — see enqueueWaiter. If `signal` aborts while // parked (or before), rejects with an AbortError so the caller treats it like a // cancel. - acquire(signal?: AbortSignal, opts?: { background?: boolean }): Promise<Lease> { + // opts.tier sets the priority tier; the legacy opts.background is still + // accepted and maps to the "background" tier (default "foreground"). + acquire( + signal?: AbortSignal, + opts?: { tier?: SchedulerTier; background?: boolean }, + ): Promise<Lease> { this.ensureInit(); if (signal?.aborted) { return Promise.reject(abortError()); @@ -312,8 +316,10 @@ export class WorkerPool { if (entry && this.waiters.length === 0) { return Promise.resolve(this.grant(entry)); } + const tier: SchedulerTier = + opts?.tier ?? (opts?.background ? "background" : "foreground"); return new Promise<Lease>((resolve, reject) => { - const waiter: Waiter = { resolve, reject, background: opts?.background }; + const waiter: Waiter = { resolve, reject, tier }; if (signal) { const onAbort = () => { const idx = this.waiters.indexOf(waiter); diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,9 @@ # Changelog ## [Unreleased] +- **One-click Retry for failed jobs (plus "Retry all failed").** A failed job that carries a replay descriptor (any bookmarkable kind — sync, download-missing, transcribe-all, retry-bucket, …) now shows a **Retry** button on the Jobs history table, and the page header gains a **Retry all failed** button whenever at least one such job is listed. Retry re-runs the job from its stored spec exactly like a bookmark re-run (so bucket jobs re-derive from the channel's *current* state), and the re-run **jumps ahead of other queued work** (it's promoted to the front of its queue, reusing the new reorder machinery) so a fix-and-retry runs next rather than at the back of the line. The spec is resolved from the live registry or, for an evicted/archived job, from its on-disk `<id>.meta.json` sidecar — so even a failure the 100-job cap has dropped is still retryable. Kinds with no replay descriptor (e.g. `import-one`) intentionally offer no Retry. See `editor/app/jobs/actions.ts` (`retryJobAction` / `retryAllFailedAction`), the new `RetryJobButton` / `RetryAllFailedButton`, and `editor/e2e/jobs-retry.spec.ts`. +- **Reorder and promote queued jobs from Active Jobs.** A queued job's row now carries **Promote / ↑ / ↓** controls (mirroring the auto-queue policy editor's move buttons) to change its order within its queue — Promote sends it to the front so it runs next, ↑/↓ nudge it one slot. Only actionable moves render (the first-queued job shows no up/promote, the last no down), and a running job is never displaced. See `editor/app/jobs/components/ReorderJobButtons.tsx`, the `reorderJobAction` / `promoteJobAction` server actions, and `editor/e2e/jobs-reorder.spec.ts`. +- **Job-queue internals unified onto one scheduler + one concurrency primitive (foundational refactor; no behavior change beyond the two features above).** The "foreground preempts background" priority was previously implemented twice (once for job queue ordering, once for worker-slot waiters); both now share a single `compareTier` comparator in a new `common/jobs/scheduler.ts` (priority tiers urgent/foreground/background, per-queue concurrency, and the reorder/promote operations the UI uses), which the registry delegates its queue ordering to. The auto-transcribe/-download runner's hand-written fill-to-capacity loop **and** the whisper batch's `Promise.all` are both replaced by one audited `runPool()` primitive (`common/jobs/concurrentRunner.ts`) that encodes the no-event-loop-spin wait once — structurally eliminating the "Drain all hangs" bug class rather than patching each loop. Per-kind metadata (label, drainability, bookmarkability) is consolidated into one `common/jobs/jobKinds.ts` table, and the bookmark/replay dispatch is now a data-driven lookup, so adding a job kind touches ~1 file instead of ~5. Covered by new unit tests (`scheduler`, `concurrentRunner`, `jobKinds`) and the existing queue/drain e2e suites. - **"Drain all" no longer hangs the server when an auto-transcribe/-download unit is actively running.** A second, distinct drain hang remained after the earlier parked-unit fix: the runner's internal `waitNext()` helper short-circuited to an *immediately-resolved* promise whenever a signal was *already* aborted — so once a soft drain fired (and `drainSignal` stays aborted for the rest of the run), every wait while in-flight units were still finishing returned with no delay. That turned the runner's poll loops into a timer-less **microtask spin** that starved the Node event loop (the spin was reached first in the main fill loop's at-capacity wait once the drain target dropped to 0, before the terminal drain-wait was ever hit). A transcription that drain deliberately lets finish completes via a child-process `exit` event — a *macrotask* — which the spin never let run, so the in-flight count never reached zero, a CPU core pegged, and the whole app appeared frozen. (The earlier fix only covered *parked* units, which settle via microtasks and so cleared even under the spin; a genuinely *running* unit depends on a macrotask and didn't.) `waitNext()` now only fast-paths a real pending `wake()`; an already-aborted signal falls through to a real timer, so all three wait sites pace instead of spinning while still being woken promptly by a finishing unit or by an abort firing mid-wait. The `whisper-all` batch was never affected (it `await Promise.all(...)` with no manual poll loop). See `common/controller/autoRunner.ts` (`waitNext`) and the new running-unit drain regression test in `editor/e2e/auto-queue.spec.ts`. - **First-class video-persistence UI (phase 5, the final phase): a Saved Videos area, per-channel retention controls, and per-video persist/unpersist.** The video-persistence subsystem built up over phases 1–4 is now driveable end to end from the editor. A new top-level **Saved videos** page (`/saved-videos`, in the Pool nav) summarizes the whole saved-video store — total count and size, per-channel breakdown (count, size, how many carry a backup checksum), the default store location, and the last backup time — and hosts the **backup configuration** (destination, scheduled on/off, interval) plus **Back up now** / **Verify backup** buttons. Each channel's **Cleanup stage** gains a **Retention & persistence** section (shown whenever keep-latest is on or the channel has saved videos) with live counts and three buttons: **Check kept videos** (re-probe the window for deleted-from-source videos and pin them), **Persist kept now** (a new bulk catch-up pass that re-fetches the source container for any in-window video whose source isn't saved yet — `persistKeptAction` / `common/controller/persistKept.ts`), and **Back up saved videos**. The **channel settings form** adds a Retention & persistence section: **keep latest** (window size), **extraction mode** (yt-dlp vs app-side ffmpeg), and a per-channel **saved-video store dir** override. Each **video page** gains a **Source video** card showing persisted status (file, size, stored time, keep reason, sha256, location) with an **Unpersist** control that moves the container back into the data dir, or a **Persist source video** button (re-fetch + archive) when it isn't saved yet. New job-kind label for `persist-kept`; `persist-kept` is re-runnable from bookmarks. The saved-video actions moved from `editor/app/savedVideos/` to `editor/app/saved-videos/` to match the route. Covered by `editor/e2e/saved-videos.spec.ts` and `common/controller/persistKept.test.ts`. (Deferred: surfacing kept-check/persist as Actionable-page rows, and streaming the player directly from the store — unpersist brings the container back to the data dir to play it.) - **Backups for the saved-video store: rsync mirror + per-backup manifest + drift verification (phase 4 of the video-persistence subsystem; backend + scheduler, UI lands later).** The (large, often irreplaceable) saved source videos can now be **backed up to a configured destination**. A backup walks every saved-video pointer across all channels (so per-channel store overrides are covered automatically) and **`rsync`-mirrors each container** into `<dest>/<slug>/<videoId>/` — incremental and resumable (`-a --partial`), additive (no deletes), so re-running only transfers changed or new files. It writes a **`backup-manifest.json`** at the destination root recording each container's canonical location, byte size, and a **streamed sha256**, and caches that hash back onto the live pointer. A **verify** step reads the manifest back and reports drift in four buckets — `missing`, `sizeMismatch`, `checksumMismatch` (re-hashing each present file), and `extra` (containers at the destination the manifest doesn't know about). New global settings block **`savedVideoBackup`** (`{ enabled, dest, intervalMinutes }`; a blank `dest` forces `enabled` off) plus a **`RSYNC_BIN`** env override. When enabled with a destination, the **sync scheduler** runs the backup automatically on its own cadence (a global, not per-channel, job — suppressed during quiet hours, tracked via `lastSavedVideoBackupAt`). Backups can also be run/verified manually via `backupSavedVideosAction` / `verifySavedVideoBackupAction` (managed jobs on a dedicated `saved-videos` queue). The destination is treated as a local filesystem path (a mounted backup disk). See the new `common/lib/savedVideoBackup.ts` (manifest types/parse), `common/controller/{backupSavedVideos,savedVideoInventory}.ts` (+ tests), `common/lib/paths.ts` (`rsyncBin`), `common/lib/settings.ts` (`savedVideoBackup`), `common/jobs/syncSchedulerState.ts`, `editor/app/savedVideos/backupActions.ts`, and `editor/app/scheduler/runTick.ts`. diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts @@ -30,6 +30,12 @@ function invalidate() { // use. (In production the pool legitimately persists; this is e2e-only.) // eslint-disable-next-line @typescript-eslint/no-explicit-any (globalThis as any).__yttWorkerPool__ = undefined; + // Reset the shared scheduler: it owns the per-queueKey ordering for the + // (now-wiped) registry, so a prior spec's still-"running" queue entry would + // otherwise make the next spec's same-named queue see a phantom busy slot and + // queue forever. Cleared in lockstep with the registry it serves. + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (globalThis as any).__yttScheduler__ = undefined; // Reset the auto-queue runner singleton too: its live entries reference job // ids in the (now-wiped) registry. A still-executing runner loop from a prior // spec detects its job is gone (registry.get -> undefined) and exits on its diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts @@ -4,6 +4,10 @@ import { revalidatePath } from "next/cache"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { clearArchivedLogs } from "yt-dlp-transcript-common/jobs/listJobs"; import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; +import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta"; +import type { JobSpec } from "yt-dlp-transcript-common/jobs/jobSpec"; +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { runJobSpec } from "./runJobSpec"; export async function cancelJobAction(id: string): Promise<{ ok: boolean }> { const ok = getRegistry().cancel(id); @@ -36,6 +40,75 @@ export async function drainAllAction(): Promise<{ count: number }> { return { count }; } +// Granular control: nudge a QUEUED job up (-1) or down (+1) within its queue. +// Only queued jobs move; the running head is never displaced. +export async function reorderJobAction( + id: string, + dir: -1 | 1, +): Promise<{ ok: boolean }> { + const ok = getRegistry().reorder(id, dir); + revalidatePath("/jobs"); + return { ok }; +} + +// Granular control: jump a QUEUED job to the front of its queue (it runs next +// when the current job releases the queue). +export async function promoteJobAction(id: string): Promise<{ ok: boolean }> { + const ok = getRegistry().promote(id); + revalidatePath("/jobs"); + return { ok }; +} + +// Resolve a job's replay descriptor: prefer the live registry record, fall back +// to the on-disk meta sidecar so a failed job the 100-job cap evicted (or one +// from before a restart) can still be retried. Null if the kind isn't +// replayable (no spec was ever attached). +async function resolveSpec(id: string): Promise<JobSpec | null> { + const live = getRegistry().get(id)?.spec; + if (live) return live; + const meta = await readJobMeta(getPaths(), id); + return meta?.spec ?? null; +} + +// One-click retry: re-run a (typically failed) job from its stored spec. Like a +// bookmark re-run, bucket jobs re-derive from the channel's CURRENT state. The +// re-run jumps ahead of other queued work (run next) by reusing the same +// promote the reorder buttons use — a no-op if it starts immediately. +export async function retryJobAction(id: string): Promise<StreamActionResult> { + const spec = await resolveSpec(id); + if (!spec) { + return { + ok: false, + error: "This job can't be retried — it has no replay descriptor.", + }; + } + const res = await runJobSpec(spec); + if (res.ok) getRegistry().promote(res.jobId); + revalidatePath("/jobs"); + return res; +} + +// Retry every currently-listed failed job that has a spec. Returns how many were +// re-launched (kinds without a replay descriptor are skipped). +export async function retryAllFailedAction(): Promise<{ count: number }> { + const registry = getRegistry(); + const failed = registry + .list() + .filter((j): j is typeof j & { spec: JobSpec } => + Boolean(j.status === "failed" && j.spec), + ); + let count = 0; + for (const j of failed) { + const res = await runJobSpec(j.spec); + if (res.ok) { + registry.promote(res.jobId); + count++; + } + } + revalidatePath("/jobs"); + return { count }; +} + 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 @@ -115,6 +115,19 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { const workerNames = new Map<string, string>(); for (const w of getWorkerPool().summary()) workerNames.set(w.id, w.name); + // Queue positions, to derive each queued job's reorder bounds. position 0 is + // the running head; queued jobs occupy 1..max. A parallel ("") job has + // position -1 and is never reorderable. + const positions = new Map<string, number>(); + const queueMaxPos = new Map<string, number>(); + for (const j of activeRecords) { + const pos = registry.positionInQueue(j.id); + positions.set(j.id, pos); + if (pos >= 0) { + queueMaxPos.set(j.queueKey, Math.max(queueMaxPos.get(j.queueKey) ?? 0, pos)); + } + } + const now = Date.now(); const jobs: RunningJobsListItem[] = activeRecords.map((j) => ({ id: j.id, @@ -145,6 +158,14 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { j.status === "running" && isDrainableKind(j.kind) && j.draining !== true, bookmarkable: Boolean(j.spec), background: j.background === true, + // A queued job can move up/promote if it isn't the first queued (position > + // 1, since the running head is at 0), and down if it isn't the last in its + // queue. Running and parallel ("") jobs are never reorderable. + canMoveUp: j.status === "queued" && (positions.get(j.id) ?? -1) > 1, + canMoveDown: + j.status === "queued" && + (positions.get(j.id) ?? -1) >= 1 && + (positions.get(j.id) ?? -1) < (queueMaxPos.get(j.queueKey) ?? 0), })); const channels = channelSlugs.map((slug) => ({ diff --git a/editor/app/jobs/components/JobsTable.tsx b/editor/app/jobs/components/JobsTable.tsx @@ -5,6 +5,7 @@ import Link from "next/link"; import type { JobListEntry } from "yt-dlp-transcript-common/jobs/listJobs"; import { CancelJobButton } from "./CancelJobButton"; import { BookmarkJobButton } from "./BookmarkJobButton"; +import { RetryJobButton } from "./RetryJobButton"; import { jobKindLabel } from "../jobKindLabels"; import { clearJobsFilters, @@ -280,6 +281,9 @@ export function JobsTable({ jobs }: { jobs: JobListEntry[] }) { </td> <td className="px-3 py-2 text-right"> <div className="flex items-center justify-end gap-2"> + {j.status === "failed" && j.bookmarkable && ( + <RetryJobButton jobId={j.id} /> + )} {j.bookmarkable && <BookmarkJobButton jobId={j.id} />} {(j.status === "running" || j.status === "queued") && ( <CancelJobButton jobId={j.id} /> diff --git a/editor/app/jobs/components/ReorderJobButtons.tsx b/editor/app/jobs/components/ReorderJobButtons.tsx @@ -0,0 +1,74 @@ +"use client"; + +import { useState } from "react"; +import { useRouter } from "next/navigation"; +import { promoteJobAction, reorderJobAction } from "../actions"; + +type Props = { + jobId: string; + // Whether the job can move toward / away from the head of its queue. The + // running head is never displaced, so the first queued job can't move up. + canMoveUp?: boolean; + canMoveDown?: boolean; +}; + +// Granular manual control for a QUEUED job: nudge it up/down one slot in its +// queue, or promote it to the front (it runs next when the queue frees). Mirrors +// the move buttons in the auto-queue PolicyTreeEditor. +export function ReorderJobButtons({ jobId, canMoveUp, canMoveDown }: Props) { + const [busy, setBusy] = useState(false); + const router = useRouter(); + const run = async (fn: () => Promise<unknown>) => { + setBusy(true); + try { + await fn(); + router.refresh(); + } finally { + setBusy(false); + } + }; + const btn = + "px-1.5 py-1 rounded border border-zinc-300 dark:border-zinc-700 text-xs font-medium text-zinc-700 dark:text-zinc-300 hover:bg-zinc-100 dark:hover:bg-zinc-800 disabled:opacity-40"; + // Only render an action that's actually possible: a first-queued job can't move + // up or promote; a last-queued job can't move down. + return ( + <div className="flex items-center gap-1"> + {canMoveUp && ( + <button + type="button" + disabled={busy} + onClick={() => run(() => promoteJobAction(jobId))} + aria-label={`promote job ${jobId} to front of queue`} + title="Run next (move to front of queue)" + className={btn} + > + Promote + </button> + )} + {canMoveUp && ( + <button + type="button" + disabled={busy} + onClick={() => run(() => reorderJobAction(jobId, -1))} + aria-label={`move job ${jobId} up`} + title="Move up one slot" + className={btn} + > + ↑ + </button> + )} + {canMoveDown && ( + <button + type="button" + disabled={busy} + onClick={() => run(() => reorderJobAction(jobId, 1))} + aria-label={`move job ${jobId} down`} + title="Move down one slot" + className={btn} + > + ↓ + </button> + )} + </div> + ); +} diff --git a/editor/app/jobs/components/RetryAllFailedButton.tsx b/editor/app/jobs/components/RetryAllFailedButton.tsx @@ -0,0 +1,30 @@ +"use client"; + +import { useState } from "react"; +import { useRouter } from "next/navigation"; +import { retryAllFailedAction } from "../actions"; + +// Re-run every currently-listed failed job that has a replay descriptor. Shown +// only when there's at least one such job (the page decides). +export function RetryAllFailedButton() { + const [busy, setBusy] = useState(false); + const router = useRouter(); + return ( + <button + type="button" + disabled={busy} + onClick={async () => { + setBusy(true); + try { + await retryAllFailedAction(); + router.refresh(); + } finally { + setBusy(false); + } + }} + className="px-3 py-1.5 rounded border border-blue-300 dark:border-blue-800 text-sm font-medium text-blue-700 dark:text-blue-300 hover:bg-blue-50 dark:hover:bg-blue-950 disabled:opacity-50" + > + {busy ? "Retrying…" : "Retry all failed"} + </button> + ); +} diff --git a/editor/app/jobs/components/RetryJobButton.tsx b/editor/app/jobs/components/RetryJobButton.tsx @@ -0,0 +1,45 @@ +"use client"; + +import { useState } from "react"; +import { useRouter } from "next/navigation"; +import { retryJobAction } from "../actions"; + +type Props = { + jobId: string; +}; + +// Re-run a failed job from its stored spec (bucket jobs re-derive from the +// channel's current state). The re-run jumps ahead of other queued work. +export function RetryJobButton({ jobId }: Props) { + const [busy, setBusy] = useState(false); + const [note, setNote] = useState<string | null>(null); + const router = useRouter(); + return ( + <span className="inline-flex items-center gap-1"> + <button + type="button" + disabled={busy} + onClick={async () => { + setBusy(true); + setNote(null); + try { + const res = await retryJobAction(jobId); + if (!res.ok) setNote(res.error); + router.refresh(); + } finally { + setBusy(false); + } + }} + aria-label={`retry job ${jobId}`} + className="px-2 py-1 rounded border border-blue-300 dark:border-blue-800 text-xs font-medium text-blue-700 dark:text-blue-300 hover:bg-blue-50 dark:hover:bg-blue-950 disabled:opacity-50" + > + {busy ? "Retrying…" : "Retry"} + </button> + {note && ( + <span className="text-xs text-zinc-500" title={note}> + {note} + </span> + )} + </span> + ); +} diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx @@ -8,6 +8,7 @@ import { jobKindLabel } from "../jobKindLabels"; import { DrainJobButton } from "./DrainJobButton"; import { CancelJobButton } from "./CancelJobButton"; import { BookmarkJobButton } from "./BookmarkJobButton"; +import { ReorderJobButtons } from "./ReorderJobButtons"; export type RunningJobsTask = { id: string; @@ -54,6 +55,11 @@ export type RunningJobsListItem = { // the same platform queue. Shown as an "auto" badge, and — when queued — as a // hint that a foreground job (a clicked Sync) is being let through first. background?: boolean; + // Reorder affordances for a QUEUED job: whether it can move toward (up/promote) + // or away from (down) the head of its queue. Both false for running jobs and + // for parallel ("") jobs that aren't serialized. + canMoveUp?: boolean; + canMoveDown?: boolean; }; type Props = { @@ -149,6 +155,13 @@ function JobRow({ )} <div className="ml-auto flex items-center gap-2"> {job.bookmarkable && <BookmarkJobButton jobId={job.id} />} + {job.status === "queued" && (job.canMoveUp || job.canMoveDown) && ( + <ReorderJobButtons + jobId={job.id} + canMoveUp={job.canMoveUp} + canMoveDown={job.canMoveDown} + /> + )} {job.status === "running" && (job.drainable || job.draining) && ( <DrainJobButton jobId={job.id} draining={job.draining} /> )} diff --git a/editor/app/jobs/jobKindLabels.ts b/editor/app/jobs/jobKindLabels.ts @@ -1,28 +1,6 @@ -// Human-readable labels for job `kind` values. The registry stores terse -// machine kinds (e.g. "whisper-all", "auto-transcribe"); the UI shows these. -// Unknown kinds fall back to the raw value so a new kind is never invisible. -const JOB_KIND_LABELS: Record<string, string> = { - "auto-transcribe": "Auto-transcribe", - // The long-lived orchestrator loop (one per kind) vs. the per-video download - // it now launches as a real job on the channel's platform queue. - "auto-download": "Auto-download runner", - "auto-download-unit": "Auto-download", - "whisper-all": "Transcribe all", - "whisper-bucket-downloaded-no-transcript": "Transcribe downloaded audio", - "download-from-playlist": "Download from playlist", - "download-missing": "Download missing", - "download-missing-subs": "Download missing subs", - "import-one": "Import video", - "redownload-archive": "Archive source video", - "retry-bucket": "Retry", - "clean-audio-transcribed": "Clean audio", - "check-kept-deleted": "Check kept videos", - "persist-kept": "Persist kept videos", - "backup-saved-videos": "Back up saved videos", - "verify-saved-video-backup": "Verify saved-video backup", - sync: "Sync", -}; - -export function jobKindLabel(kind: string): string { - return JOB_KIND_LABELS[kind] ?? kind; -} +// Human-readable labels for job `kind` values now live in the shared job-kind +// metadata table (common/jobs/jobKinds.ts), the single source of truth. This +// module re-exports the lookup so existing editor imports stay unchanged. The +// registry stores terse machine kinds (e.g. "whisper-all"); the UI shows these, +// falling back to the raw value so a new kind is never invisible. +export { jobKindLabel } from "yt-dlp-transcript-common/jobs/jobKinds"; diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts @@ -0,0 +1,187 @@ +// The data-driven replay table: maps a stored JobSpec's `kind` to the server +// action that re-runs it. Previously this was a hard-coded switch in +// runJobSpec.ts; collapsing it to a lookup means adding a bookmarkable kind is +// one colocated entry here (plus its jobKinds.ts metadata) instead of a new +// switch case. runJobSpec.ts is now a thin dispatcher over this table, and the +// retry actions (Phase 8) reuse the same seam. +// +// Server-only: reached through bookmarkActions ("use server"), so it is never +// bundled to the client. +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { readChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot"; +import type { JobSpec, ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; +import type { AudioFormat } from "yt-dlp-transcript-common/lib/channelConfig"; +import { + downloadAction, + downloadMissingAction, + downloadMissingSubsAction, + retryBucketAction, + storePlaylistAction, + syncAction, +} from "../channels/[slug]/pipelineActions"; +import { + checkKeptDeletedAction, + cleanAudioAction, + cleanExtraAudioFormatsAction, + clearFailedTranscodingsAction, + clearFailedTranscriptionsAction, + removeFailedTranscodingsAction, + removeWrongFormatAudioAction, + transcodeFailuresAction, + transcodeUntranscodedAction, + transcribeBucketAction, + transcribeMissingAction, +} from "../channels/[slug]/whisperActions"; +import { persistKeptAction } from "../channels/[slug]/persistActions"; + +export type ReplayHandler = (spec: JobSpec) => Promise<StreamActionResult>; + +// Bucket jobs re-derive their work from the channel's CURRENT snapshot (not a +// frozen id list), so re-running "retry partial downloads" always acts on +// whatever is partial right now. +async function idsForBucket( + slug: string, + bucket: ReplayBucket, +): Promise<string[]> { + const snap = await readChannelSnapshot(getPaths(), slug); + const ids = snap?.buckets?.[bucket]; + return Array.isArray(ids) ? ids : []; +} + +const str = (v: unknown): string | undefined => + typeof v === "string" ? v : undefined; +const bool = (v: unknown): boolean | undefined => + typeof v === "boolean" ? v : undefined; +const num = (v: unknown): number | undefined => + typeof v === "number" ? v : undefined; + +// Flag-style params + the captured queueKey for a spec. +function params(spec: JobSpec): { + p: Record<string, unknown>; + queueKey: string | undefined; +} { + const p = spec.params ?? {}; + return { p, queueKey: str(p.queueKey) }; +} + +export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = { + "whisper-all": (spec) => { + const { p, queueKey } = params(spec); + return transcribeMissingAction( + spec.slug, + queueKey, + bool(p.reverse), + str(p.audioFormat) as AudioFormat | undefined, + bool(p.strictAudioFormat), + num(p.shardTotal), + num(p.shardIndex), + ); + }, + "whisper-bucket-downloaded-no-transcript": async (spec) => { + const { p, queueKey } = params(spec); + if (!spec.bucket) return { ok: false, error: "Bookmark is missing its bucket." }; + const ids = await idsForBucket(spec.slug, spec.bucket); + if (ids.length === 0) { + return { ok: false, error: "Nothing to transcribe right now.", info: true }; + } + return transcribeBucketAction( + spec.slug, + ids, + queueKey, + str(p.audioFormat) as AudioFormat | undefined, + bool(p.strictAudioFormat), + spec.bucket, + ); + }, + "retry-bucket": async (spec) => { + const { p, queueKey } = params(spec); + if (!spec.bucket) return { ok: false, error: "Bookmark is missing its bucket." }; + const ids = await idsForBucket(spec.slug, spec.bucket); + if (ids.length === 0) { + return { ok: false, error: "Nothing to retry right now.", info: true }; + } + return retryBucketAction( + spec.slug, + ids, + queueKey, + bool(p.abortOnError), + str(p.handlingOverride), + spec.bucket, + ); + }, + "download-from-playlist": (spec) => { + const { p, queueKey } = params(spec); + return downloadAction(spec.slug, queueKey, bool(p.abortOnError)); + }, + "download-missing": (spec) => { + const { p, queueKey } = params(spec); + return downloadMissingAction( + spec.slug, + queueKey, + bool(p.ignoreArchive), + bool(p.abortOnError), + num(p.shardTotal), + num(p.shardIndex), + bool(p.keepSourceVideoOverride), + bool(p.extractImmediately), + str(p.audioFormatOverride) as AudioFormat | undefined, + ); + }, + sync: (spec) => { + const { queueKey } = params(spec); + return syncAction(spec.slug, queueKey); + }, + "download-missing-subs": (spec) => { + const { p, queueKey } = params(spec); + return downloadMissingSubsAction(spec.slug, queueKey, bool(p.abortOnError)); + }, + "store-playlist": (spec) => { + const { queueKey } = params(spec); + return storePlaylistAction(spec.slug, queueKey); + }, + "clear-failed-transcriptions": (spec) => { + const { queueKey } = params(spec); + return clearFailedTranscriptionsAction(spec.slug, queueKey); + }, + "transcode-failures": (spec) => { + const { p, queueKey } = params(spec); + return transcodeFailuresAction( + spec.slug, + str(p.targetFormat) as AudioFormat, + queueKey, + ); + }, + "transcode-untranscoded": (spec) => { + const { queueKey } = params(spec); + return transcodeUntranscodedAction(spec.slug, queueKey); + }, + "remove-failed-transcodings": (spec) => { + const { queueKey } = params(spec); + return removeFailedTranscodingsAction(spec.slug, queueKey); + }, + "clear-failed-transcodings": (spec) => { + const { queueKey } = params(spec); + return clearFailedTranscodingsAction(spec.slug, queueKey); + }, + "clean-extra-audio-formats": (spec) => { + const { queueKey } = params(spec); + return cleanExtraAudioFormatsAction(spec.slug, queueKey); + }, + "remove-wrong-format-audio": (spec) => { + const { queueKey } = params(spec); + return removeWrongFormatAudioAction(spec.slug, queueKey); + }, + "clean-audio-transcribed": (spec) => { + const { queueKey } = params(spec); + return cleanAudioAction(spec.slug, queueKey); + }, + "check-kept-deleted": (spec) => { + const { queueKey } = params(spec); + return checkKeptDeletedAction(spec.slug, queueKey); + }, + "persist-kept": (spec) => { + const { queueKey } = params(spec); + return persistKeptAction(spec.slug, queueKey); + }, +}; diff --git a/editor/app/jobs/page.tsx b/editor/app/jobs/page.tsx @@ -3,6 +3,7 @@ import { listAllJobs } from "yt-dlp-transcript-common/jobs/listJobs"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { ClearArchivedButton } from "./components/ClearArchivedButton"; import { JobsTable } from "./components/JobsTable"; +import { RetryAllFailedButton } from "./components/RetryAllFailedButton"; import { BookmarksMenu } from "./components/BookmarksMenu"; import { loadBookmarksView } from "./loadBookmarks"; @@ -15,11 +16,17 @@ export default async function JobsPage() { listAllJobs(getPaths()), loadBookmarksView(), ]); + const hasRetryableFailed = jobs.some( + (j) => j.status === "failed" && j.bookmarkable, + ); return ( <div className="flex flex-col gap-4"> <div className="flex items-center justify-between"> <h1 className="text-2xl font-semibold">Jobs</h1> - <ClearArchivedButton /> + <div className="flex items-center gap-2"> + {hasRetryableFailed && <RetryAllFailedButton />} + <ClearArchivedButton /> + </div> </div> <BookmarksMenu bookmarks={bookmarks} missingSlugs={missingSlugs} /> {jobs.length === 0 ? ( diff --git a/editor/app/jobs/runJobSpec.ts b/editor/app/jobs/runJobSpec.ts @@ -1,143 +1,18 @@ // Server-only dispatcher: reached only through bookmarkActions ("use server"), // so it is never bundled to the client. -import { getPaths } from "yt-dlp-transcript-common/lib/paths"; -import { readChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot"; -import type { JobSpec, ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; +// +// The per-kind replay logic lives in jobReplayRegistry.ts (a data-driven +// lookup); this is just the dispatch over it. Re-running RE-DERIVES the work +// from the channel's CURRENT state rather than replaying a frozen list — see +// jobReplayRegistry.ts and jobSpec.ts. +import type { JobSpec } from "yt-dlp-transcript-common/jobs/jobSpec"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; -import type { AudioFormat } from "yt-dlp-transcript-common/lib/channelConfig"; -import { - downloadAction, - downloadMissingAction, - downloadMissingSubsAction, - retryBucketAction, - storePlaylistAction, - syncAction, -} from "../channels/[slug]/pipelineActions"; -import { - checkKeptDeletedAction, - cleanAudioAction, - cleanExtraAudioFormatsAction, - clearFailedTranscodingsAction, - clearFailedTranscriptionsAction, - removeFailedTranscodingsAction, - removeWrongFormatAudioAction, - transcodeFailuresAction, - transcodeUntranscodedAction, - transcribeBucketAction, - transcribeMissingAction, -} from "../channels/[slug]/whisperActions"; -import { persistKeptAction } from "../channels/[slug]/persistActions"; - -// The single place that maps a stored JobSpec back to the server action that -// runs it. Bucket jobs re-derive their work from the channel's CURRENT snapshot -// (not a frozen id list) via idsForBucket, so re-running "retry partial -// downloads" always acts on whatever is partial right now. - -async function idsForBucket( - slug: string, - bucket: ReplayBucket, -): Promise<string[]> { - const snap = await readChannelSnapshot(getPaths(), slug); - const ids = snap?.buckets?.[bucket]; - return Array.isArray(ids) ? ids : []; -} - -const str = (v: unknown): string | undefined => - typeof v === "string" ? v : undefined; -const bool = (v: unknown): boolean | undefined => - typeof v === "boolean" ? v : undefined; -const num = (v: unknown): number | undefined => - typeof v === "number" ? v : undefined; +import { JOB_REPLAY_HANDLERS } from "./jobReplayRegistry"; export async function runJobSpec(spec: JobSpec): Promise<StreamActionResult> { - const p = spec.params ?? {}; - const queueKey = str(p.queueKey); - switch (spec.kind) { - case "whisper-all": - return transcribeMissingAction( - spec.slug, - queueKey, - bool(p.reverse), - str(p.audioFormat) as AudioFormat | undefined, - bool(p.strictAudioFormat), - num(p.shardTotal), - num(p.shardIndex), - ); - case "whisper-bucket-downloaded-no-transcript": { - if (!spec.bucket) return { ok: false, error: "Bookmark is missing its bucket." }; - const ids = await idsForBucket(spec.slug, spec.bucket); - if (ids.length === 0) { - return { ok: false, error: "Nothing to transcribe right now.", info: true }; - } - return transcribeBucketAction( - spec.slug, - ids, - queueKey, - str(p.audioFormat) as AudioFormat | undefined, - bool(p.strictAudioFormat), - spec.bucket, - ); - } - case "retry-bucket": { - if (!spec.bucket) return { ok: false, error: "Bookmark is missing its bucket." }; - const ids = await idsForBucket(spec.slug, spec.bucket); - if (ids.length === 0) { - return { ok: false, error: "Nothing to retry right now.", info: true }; - } - return retryBucketAction( - spec.slug, - ids, - queueKey, - bool(p.abortOnError), - str(p.handlingOverride), - spec.bucket, - ); - } - case "download-from-playlist": - return downloadAction(spec.slug, queueKey, bool(p.abortOnError)); - case "download-missing": - return downloadMissingAction( - spec.slug, - queueKey, - bool(p.ignoreArchive), - bool(p.abortOnError), - num(p.shardTotal), - num(p.shardIndex), - bool(p.keepSourceVideoOverride), - bool(p.extractImmediately), - str(p.audioFormatOverride) as AudioFormat | undefined, - ); - case "sync": - return syncAction(spec.slug, queueKey); - case "download-missing-subs": - return downloadMissingSubsAction(spec.slug, queueKey, bool(p.abortOnError)); - case "store-playlist": - return storePlaylistAction(spec.slug, queueKey); - case "clear-failed-transcriptions": - return clearFailedTranscriptionsAction(spec.slug, queueKey); - case "transcode-failures": - return transcodeFailuresAction( - spec.slug, - str(p.targetFormat) as AudioFormat, - queueKey, - ); - case "transcode-untranscoded": - return transcodeUntranscodedAction(spec.slug, queueKey); - case "remove-failed-transcodings": - return removeFailedTranscodingsAction(spec.slug, queueKey); - case "clear-failed-transcodings": - return clearFailedTranscodingsAction(spec.slug, queueKey); - case "clean-extra-audio-formats": - return cleanExtraAudioFormatsAction(spec.slug, queueKey); - case "remove-wrong-format-audio": - return removeWrongFormatAudioAction(spec.slug, queueKey); - case "clean-audio-transcribed": - return cleanAudioAction(spec.slug, queueKey); - case "check-kept-deleted": - return checkKeptDeletedAction(spec.slug, queueKey); - case "persist-kept": - return persistKeptAction(spec.slug, queueKey); - default: - return { ok: false, error: `Cannot re-run job kind: ${spec.kind}` }; + const handler = JOB_REPLAY_HANDLERS[spec.kind]; + if (!handler) { + return { ok: false, error: `Cannot re-run job kind: ${spec.kind}` }; } + return handler(spec); } diff --git a/editor/e2e/jobs-batch-tasks-drain.spec.ts b/editor/e2e/jobs-batch-tasks-drain.spec.ts @@ -204,7 +204,7 @@ test("'Drain all' drains the running batch and cancels the queued one", async ({ const bRow = page .getByRole("row") .filter({ hasText: "drainall-b" }) - .filter({ hasText: "whisper-all" }) + .filter({ hasText: "Transcribe all" }) .first(); await expect(bRow).toContainText("cancelled", { timeout: 20_000 }); expect(await transcriptCount("drainall-b", bIds)).toBe(0); @@ -212,7 +212,7 @@ test("'Drain all' drains the running batch and cancels the queued one", async ({ const aRow = page .getByRole("row") .filter({ hasText: "drainall-a" }) - .filter({ hasText: "whisper-all" }) + .filter({ hasText: "Transcribe all" }) .first(); await expect(aRow).toContainText("done", { timeout: 20_000 }); const aDone = await transcriptCount("drainall-a", aIds); @@ -252,7 +252,7 @@ test("a queued job can be cancelled directly from its row without opening the lo const bRow = page .getByRole("row") .filter({ hasText: "qcancel-b" }) - .filter({ hasText: "whisper-all" }) + .filter({ hasText: "Transcribe all" }) .first(); await expect(bRow).toContainText("cancelled", { timeout: 20_000 }); expect(await transcriptCount("qcancel-b", bIds)).toBe(0); @@ -294,7 +294,7 @@ test("hard Cancel during a drain ends the job cancelled without stream errors", const row = page .getByRole("row") .filter({ hasText: "cancel-drain" }) - .filter({ hasText: "whisper-all" }) + .filter({ hasText: "Transcribe all" }) .first(); await expect(row).toContainText("cancelled", { timeout: 20_000 }); diff --git a/editor/e2e/jobs-reorder.spec.ts b/editor/e2e/jobs-reorder.spec.ts @@ -0,0 +1,103 @@ +// Granular queue control: the ↑/↓/Promote buttons on a QUEUED job's row in +// Active Jobs reorder it within its queue. The visual list is sorted by recency, +// not queue position, so this asserts the EFFECT: a promoted job runs next when +// the head is cancelled, ahead of the job that was FIFO-first. + +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { resetData, resolvePath } from "./helpers"; +import { baseUrl } from "./baseUrl"; + +async function invalidateCache() { + await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {}); +} + +// A channel whose Sync runs slowly (fake-ytdlp honors --test-slow), so it stays +// "running"/"queued" long enough to reorder. +async function makeSlowChannel(slug: string, name: string) { + const root = resolvePath(`test-transcripts/channels/${slug}`); + await mkdir(root, { recursive: true }); + await writeFile( + `${root}/config.json`, + JSON.stringify({ + handling: "youtube", + name, + url: `https://www.youtube.com/@${slug}/videos`, + ytdlpExtraArgs: ["--test-slow"], + }), + ); +} + +async function setQueueViaCustom( + page: import("@playwright/test").Page, + actionLabel: string, + value: string, +) { + await page + .getByLabel(`queue for ${actionLabel}`, { exact: true }) + .selectOption({ label: "Custom…" }); + await page + .getByLabel(`custom queue for ${actionLabel}`, { exact: true }) + .fill(value); + await page + .getByRole("button", { name: `save custom queue for ${actionLabel}` }) + .click(); +} + +async function syncToQueue( + page: import("@playwright/test").Page, + slug: string, + queue: string, +) { + await page.goto(`/channels/${slug}`); + await setQueueViaCustom(page, "Sync", queue); + await page.getByRole("button", { name: "Sync" }).click(); +} + +test("Promote makes a queued job run next, ahead of the FIFO-first one", async ({ + page, +}) => { + test.setTimeout(90_000); + await resetData(null); + await makeSlowChannel("reorder-a", "Reorder A"); + await makeSlowChannel("reorder-b", "Reorder B"); + await makeSlowChannel("reorder-c", "Reorder C"); + await invalidateCache(); + + // A runs; B (FIFO-first) then C queue behind it on the shared queue. + await syncToQueue(page, "reorder-a", "qReorder"); + await expect(page.getByLabel("Sync output")).toContainText("test-slow", { + timeout: 15_000, + }); + await syncToQueue(page, "reorder-b", "qReorder"); + await expect(page.getByLabel("Sync queue status")).toBeVisible({ + timeout: 10_000, + }); + await syncToQueue(page, "reorder-c", "qReorder"); + await expect(page.getByLabel("Sync queue status")).toBeVisible({ + timeout: 10_000, + }); + + await page.goto("/jobs/active"); + const sectionA = page.locator("section[aria-label='Active jobs for Reorder A']"); + const sectionB = page.locator("section[aria-label='Active jobs for Reorder B']"); + const sectionC = page.locator("section[aria-label='Active jobs for Reorder C']"); + await expect(sectionA.getByText("running", { exact: true })).toBeVisible({ + timeout: 15_000, + }); + await expect(sectionB.getByText("queued", { exact: true })).toBeVisible(); + await expect(sectionC.getByText("queued", { exact: true })).toBeVisible(); + + // C is last in the queue, so its row offers Promote; B (first queued) does not. + await expect( + sectionB.getByRole("button", { name: /^promote job/ }), + ).toHaveCount(0); + await sectionC.getByRole("button", { name: /^promote job/ }).click(); + + // Cancel the running A; the PROMOTED C runs next, not the FIFO-first B. + await sectionA.getByRole("button", { name: /^Cancel$/ }).click(); + await expect(sectionC.getByText("running", { exact: true })).toBeVisible({ + timeout: 15_000, + }); + await expect(sectionB.getByText("queued", { exact: true })).toBeVisible(); +}); diff --git a/editor/e2e/jobs-retry.spec.ts b/editor/e2e/jobs-retry.spec.ts @@ -0,0 +1,90 @@ +// One-click retry: a failed job that carries a replay descriptor shows a Retry +// button that re-runs it (resolving the spec from the on-disk meta sidecar, so +// even an evicted/archived job is retryable). A failed job WITHOUT a spec offers +// no Retry. Failed jobs are staged as sidecar + log files so the scenario is +// deterministic and exercises the sidecar-fallback path directly. + +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { resetData, resolvePath } from "./helpers"; +import { baseUrl } from "./baseUrl"; + +async function invalidateCache() { + await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {}); +} + +async function makeSlowChannel(slug: string, name: string) { + const root = resolvePath(`test-transcripts/channels/${slug}`); + await mkdir(root, { recursive: true }); + await writeFile( + `${root}/config.json`, + JSON.stringify({ + handling: "youtube", + name, + url: `https://www.youtube.com/@${slug}/videos`, + ytdlpExtraArgs: ["--test-slow"], + }), + ); +} + +// Stage an archived job on disk: listAllJobs iterates `.log` files and merges +// the `.meta.json` sidecar, so both are needed for the row to appear. +async function writeArchivedJob( + id: string, + meta: Record<string, unknown>, +): Promise<void> { + const dir = resolvePath("test-transcripts/.jobs"); + await mkdir(dir, { recursive: true }); + await writeFile(`${dir}/${id}.log`, `archived log for ${id}\n`); + await writeFile(`${dir}/${id}.meta.json`, JSON.stringify({ id, ...meta })); +} + +test("Retry re-runs a failed job from its spec; spec-less failures offer none", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData(null); + await makeSlowChannel("retry-ch", "Retry Ch"); + + const ts = 1_700_000_000_000; + // A failed SYNC carrying a replay spec → retryable. + await writeArchivedJob("failedsync1", { + kind: "sync", + channelSlug: "retry-ch", + queueKey: "platform:youtube", + status: "failed", + queuedAt: ts, + startedAt: ts, + endedAt: ts + 1000, + exitCode: 1, + spec: { kind: "sync", slug: "retry-ch" }, + }); + // A failed job with NO spec → not retryable. + await writeArchivedJob("failednospec1", { + kind: "import-one", + channelSlug: "retry-ch", + status: "failed", + queuedAt: ts, + startedAt: ts, + endedAt: ts + 1000, + exitCode: 1, + }); + await invalidateCache(); + + await page.goto("/jobs"); + + const failedRow = page.getByRole("row").filter({ hasText: "failedsync1" }); + await expect(failedRow.getByRole("button", { name: "Retry" })).toBeVisible(); + // The spec-less failed job offers no Retry. + const noSpecRow = page.getByRole("row").filter({ hasText: "failednospec1" }); + await expect(noSpecRow.getByRole("button", { name: "Retry" })).toHaveCount(0); + + // Retry re-runs the sync (slow channel → it shows up running). + await failedRow.getByRole("button", { name: "Retry" }).click(); + + const runningSync = page + .getByRole("row") + .filter({ hasText: "retry-ch" }) + .filter({ hasText: "running" }); + await expect(runningSync).toBeVisible({ timeout: 15_000 }); +});