Archilyzer · Source

archilyzer

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

commit e487ca6aaf9095a7bfb9c151c1eeebc01c6d23da
parent 673c5603a517e07f5fa4a9b126b535782b6b8f15
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Tue, 23 Jun 2026 21:40:48 -0400

Merge fix/auto-transcribe-priority-and-drain: manual transcribe priority + Drain hang fix

Diffstat:
Mcommon/controller/autoRunner.ts | 18++++++++++++++++++
Mcommon/controller/transcribeOne.ts | 5++++-
Mcommon/controller/transcribeOneFromQueue.ts | 5+++++
Acommon/jobs/workerPool.test.ts | 120+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/workerPool.ts | 54++++++++++++++++++++++++++++++++++++++++++++----------
Meditor/CHANGELOG.md | 2++
Meditor/app/auto-queue/components/AutoQueueView.tsx | 9+++++++++
Meditor/e2e/auto-queue.spec.ts | 82+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
8 files changed, 284 insertions(+), 11 deletions(-)

diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -425,6 +425,7 @@ async function runLoop( tracker, onLog, signal, + drainSignal: ctx.drainSignal, onChildJob: (jid) => childJobIds.set(pick.videoId, jid), }) .then((r) => { @@ -533,6 +534,13 @@ type LaunchArgs = { tracker: ReturnType<typeof makeTaskTracker>; onLog: (line: string) => void; signal: AbortSignal; + // Soft-cancel (drain): forwarded into the transcription unit's worker-pool + // acquire so a parked (not-yet-started) unit unblocks and skips on drain — + // otherwise the runner's drain loop waits forever for it to clear. A running + // transcription is unaffected (drainSignal never reaches the engine). Unused by + // the download branch (downloads serialize on the registry queue, not the pool, + // and a drain lets the in-flight download finish by design). + drainSignal?: AbortSignal; // Download only: called with the child download job's id once enqueued, so the // runner can hard-cancel it when stopped. No-op for transcription units. onChildJob?: (jobId: string) => void; @@ -564,6 +572,16 @@ async function launchUnit(args: LaunchArgs): Promise<UnitResult> { taskLabel: `${args.channelSlug}/${args.pick.videoId}`, onLog: args.onLog, signal: args.signal, + // Drain (soft cancel): unblocks this unit's worker-pool acquire if it's + // parked waiting for a slot, so the runner's drain can complete instead of + // hanging. A unit already transcribing finishes (the engine only sees the + // hard signal). + drainSignal: args.drainSignal, + // Acquire the worker pool slot at background priority: a manual (foreground) + // transcribe preempts queued auto units for the next freed slot. The running + // unit is never interrupted — it drains, then the slot goes to the manual + // waiter. Mirrors how auto-download units yield to a manual sync. + background: true, }); return { outcome: res.outcome }; } diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -270,6 +270,9 @@ export type TranscribeWithWorkerOptions = { // Soft cancel (drain): when aborted, a parked acquire unblocks and this throws // an AbortError so the caller skips the video without starting it. drainSignal?: AbortSignal; + // Auto-runner units pass true so they park BEHIND any manual (foreground) + // acquire in the worker pool — a manual transcribe preempts queued auto work. + background?: boolean; }; // Acquire a worker from the global pool and transcribe one video through it, @@ -293,7 +296,7 @@ export async function transcribeWithWorker( if (opts.signal?.aborted || opts.drainSignal?.aborted) throw abortError(); let lease; try { - lease = await pool.acquire(acquireSignal); + lease = await pool.acquire(acquireSignal, { background: opts.background }); } catch { throw abortError(); // cancelled or drained while parked } diff --git a/common/controller/transcribeOneFromQueue.ts b/common/controller/transcribeOneFromQueue.ts @@ -49,6 +49,9 @@ export type TranscribeOneOptions = { onLog?: (msg: string) => void; signal?: AbortSignal; drainSignal?: AbortSignal; + // Forwarded to the worker-pool acquire: auto-runner units pass true so they + // yield a free slot to any manual (foreground) transcription waiting on it. + background?: boolean; }; const skip = (): TranscribeOneResult => ({ attempted: false, outcome: "skipped" }); @@ -65,6 +68,7 @@ export async function transcribeOneFromQueue({ onLog, signal, drainSignal, + background, }: TranscribeOneOptions): Promise<TranscribeOneResult> { const log = onLog ?? ((m: string) => console.log(m)); const channelDir = path.join(paths.channelsDir, channelSlug); @@ -115,6 +119,7 @@ export async function transcribeOneFromQueue({ onLog: log, signal, drainSignal, + background, }); if (outcome === "paused") { // Parakeet was paused with progress cached; leave the video untranscribed diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts @@ -0,0 +1,120 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import type { Worker } from "../lib/workers"; +import { WorkerPool } from "./workerPool"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/workerPool.test.ts +// +// Covers the foreground/background waiter ordering added so a manual +// (foreground) transcription preempts queued auto-runner (background) units for +// the next freed worker slot — the running unit is never interrupted, it drains. +// Each test builds a fresh `new WorkerPool()` (not the global singleton) and +// seeds it with reconfigure(workers, {applyEnabled:true}), which avoids reading +// settings.json so the scheduler is fully isolated. + +function worker(id: string, priority = 0): Worker { + return { id, name: id, kind: "local", enabled: true, priority }; +} + +// A single-slot pool with its only worker already leased out (running), so every +// subsequent acquire parks as a waiter. Returns the pool and the running lease. +async function busyPool() { + const pool = new WorkerPool(); + pool.reconfigure([worker("w1")], { applyEnabled: true }); + const lease = await pool.acquire(); + return { pool, lease }; +} + +test("a foreground acquire is served before an earlier-parked background one", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + // Background parks FIRST, foreground SECOND — arrival order would serve B first. + const pB = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B"), l)); + const pF = pool.acquire().then((l) => (order.push("F"), l)); + + lease.release(); // one slot frees → highest-priority waiter (foreground) wins + const leaseF = await pF; + assert.deepEqual(order, ["F"]); + + leaseF.release(); // next slot → the background waiter + await pB; + assert.deepEqual(order, ["F", "B"]); +}); + +test("foreground waiters keep FIFO among themselves", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + const p1 = pool.acquire().then((l) => (order.push("1"), l)); + const p2 = pool.acquire().then((l) => (order.push("2"), l)); + + lease.release(); + const l1 = await p1; + assert.deepEqual(order, ["1"]); + l1.release(); + await p2; + assert.deepEqual(order, ["1", "2"]); +}); + +test("a foreground acquire slots ahead of several queued background ones", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + const b1 = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B1"), l)); + const b2 = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B2"), l)); + const b3 = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B3"), l)); + // Arrives last but is foreground → must be served first, then B1,B2,B3 in order. + const f = pool.acquire().then((l) => (order.push("F"), l)); + + let lease_ = lease; + for (const p of [f, b1, b2, b3]) { + lease_.release(); + lease_ = await p; + } + assert.deepEqual(order, ["F", "B1", "B2", "B3"]); +}); + +test("aborting a parked waiter removes it without disturbing the order", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + const ac = new AbortController(); + const p1 = pool.acquire().then((l) => (order.push("1"), l)); + const p2 = pool + .acquire(ac.signal) + .then((l) => (order.push("2"), l)) + .catch((e) => (order.push("2-aborted"), Promise.reject(e))); + const p3 = pool.acquire().then((l) => (order.push("3"), l)); + + ac.abort(); // cancel the middle waiter while parked + await assert.rejects(p2, /abort/i); + + lease.release(); + const l1 = await p1; + l1.release(); + await p3; + assert.deepEqual(order, ["2-aborted", "1", "3"]); +}); + +test("all-background acquires fall back to plain FIFO (backward compatible)", async () => { + const { pool, lease } = await busyPool(); + const order: string[] = []; + const b1 = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B1"), l)); + const b2 = pool + .acquire(undefined, { background: true }) + .then((l) => (order.push("B2"), l)); + + lease.release(); + const l1 = await b1; + assert.deepEqual(order, ["B1"]); + l1.release(); + await b2; + assert.deepEqual(order, ["B1", "B2"]); +}); diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts @@ -11,9 +11,18 @@ // Scheduling is a priority-ordered free-list (decided with the user): every // enabled worker runs concurrently, and when several have a free slot the // highest-priority one is handed out first. When NO worker is eligible (all -// disabled/draining/full), acquire() simply parks on a FIFO queue — so a batch +// disabled/draining/full), acquire() parks on the waiter queue — so a batch // pauses, neither failing nor busy-looping, and resumes the instant a worker is // re-enabled. +// +// Waiters are ordered foreground-before-background, not plain FIFO: a manual +// (foreground) acquire is served before any parked background one, while FIFO is +// preserved WITHIN each class. This mirrors the registry queue's foreground/ +// background insertion (see registry.ts enqueue) so the same "manual preempts +// auto" priority axis governs BOTH schedulers. The auto-runner marks its +// per-video transcription acquires `background: true`, so a manual transcribe +// jumps ahead of queued auto work — the running auto unit isn't interrupted, it +// drains, and the freed slot goes to the manual waiter. import type { Worker } from "../lib/workers"; import { getSettings } from "../lib/settings"; @@ -64,9 +73,12 @@ 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; }; -class WorkerPool { +export class WorkerPool { // Insertion order is the tiebreak for equal priority, so use a Map (ordered). private entries = new Map<string, PoolEntry>(); private waiters: Waiter[] = []; @@ -254,11 +266,33 @@ 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. + 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) { + insertAt = i; + break; + } + } + this.waiters.splice(insertAt, 0, waiter); + } + // Acquire a lease on the highest-priority free worker. Resolves immediately if - // one is free, otherwise parks until a slot frees or a worker is enabled. If - // `signal` aborts while parked (or before), rejects with an AbortError so the - // caller treats it like a cancel. - acquire(signal?: AbortSignal): Promise<Lease> { + // one is free, otherwise parks until a slot frees or a worker is enabled. A + // background acquire (auto-runner transcription unit) queues behind any + // 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> { this.ensureInit(); if (signal?.aborted) { return Promise.reject(abortError()); @@ -268,7 +302,7 @@ class WorkerPool { return Promise.resolve(this.grant(entry)); } return new Promise<Lease>((resolve, reject) => { - const waiter: Waiter = { resolve, reject }; + const waiter: Waiter = { resolve, reject, background: opts?.background }; if (signal) { const onAbort = () => { const idx = this.waiters.indexOf(waiter); @@ -278,9 +312,9 @@ class WorkerPool { waiter.onAbort = () => signal.removeEventListener("abort", onAbort); signal.addEventListener("abort", onAbort, { once: true }); } - this.waiters.push(waiter); - // A new arrival can't create a free slot, but pump keeps FIFO honest if one - // freed between pickFree and here. + this.enqueueWaiter(waiter); + // A new arrival can't create a free slot, but pump keeps the order honest if + // one freed between pickFree and here. this.pump(); }); } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,8 @@ # Changelog ## [Unreleased] +- **"Drain all" (and the per-runner Drain) no longer hangs the auto-transcribe runner.** Draining could leave the runner stuck "draining" forever — appearing to freeze the app — whenever one of its per-video transcription units was *parked* in the worker pool waiting for a free slot at the moment drain fired (much more likely now that auto units yield slots to manual transcriptions). The runner forwarded only its hard-cancel signal — never the soft `drainSignal` — to a parked unit's `pool.acquire()`, so a soft drain could never unblock it: the unit's promise never settled, the runner's in-flight count never reached zero, and its drain loop spun indefinitely (hard **Cancel**/**Stop** always worked, since that signal *was* forwarded). The runner now threads `ctx.drainSignal` into each transcription unit, so a parked (not-yet-started) unit unblocks and is skipped on drain while a unit already transcribing finishes normally — correct drain semantics, and the runner finalizes promptly. See `common/controller/autoRunner.ts` (the `launchUnit` drainSignal wiring) and the new parked-unit drain regression test in `editor/e2e/auto-queue.spec.ts`. +- **Manual transcriptions now preempt auto-queued ones — click "Transcribe missing" (or any non-auto transcribe) while auto-transcribe is running and yours goes next.** The auto-transcribe runner shares the global transcription **worker pool** with every manual transcribe (channel batch, bucket retry, single-video), so they already serialized — but the pool granted freed worker slots in plain arrival order, so a manual transcribe could wait behind the runner's next auto pick. The pool's waiter queue is now **priority-ordered**: a manual (foreground) acquire is served before any parked auto (background) one, with FIFO preserved within each class — the same foreground-vs-background priority the per-platform sync queue already uses, now applied to the pool too. The auto-runner acquires its per-video transcription slots at **background priority**, so a manual transcribe jumps ahead: the in-flight auto transcription is **not interrupted** (it finishes — "drains"), then every freed worker goes to the manual work until it's exhausted, then auto resumes. Worker N-way parallelism is untouched (only *parked* waiters are reordered). A future "auto-sync queues its own transcriptions" path gets the same yield for free by acquiring as background. See `common/jobs/workerPool.ts` (waiter priority), `common/controller/transcribeOne.ts`, `common/controller/transcribeOneFromQueue.ts`, and `common/controller/autoRunner.ts`. - **Auto-download now shares the same per-platform queue as a manual Sync, so you can click Sync while auto-download is running without risking a 429.** Each auto-download video is now launched as a real job on the channel's platform queue (e.g. `platform:youtube`) — the same queue Sync uses — instead of running out-of-band. The job registry serializes them, so the two never spawn yt-dlp against one platform at the same time (in either direction), and each auto-download unit appears as its own row in Active Jobs with progress. A manually clicked Sync **preempts** the runner's *queued* auto-download units on that platform (it doesn't interrupt one already downloading), and a Sync is **refused with a clear message** (showing remaining seconds) while that platform is in a rate-limit cooldown. A 429 hit by either path now records the shared cooldown, so manual and automatic downloads back off together. See `common/jobs/registry.ts` (queue priority), `common/controller/autoRunner.ts`, `common/jobs/downloadBackoff.ts`, and `common/jobs/streamCommand.ts`. - **Auto-queue downloads now back off per-platform on rate limits instead of hammering the source.** When a managed auto-download hits an HTTP 429 / "too many requests" or a network error (most often on Odysee), the runner pauses *that platform* for an exponential cooldown (1 min, doubling up to 30 min, with jitter) while other platforms keep flowing, and the affected video is retried after the cooldown rather than being burned as a false success. Previously the runner discarded the download outcome, counted the rate-limited video as done, and immediately re-hit the same platform — so it never actually downloaded and re-stormed the source on every restart. Failures are now classified against the full yt-dlp stderr (not just the last few lines), so a 429 that yt-dlp logs as a WARNING before failing with a different final error still triggers the backoff. Cooldowns persist across restarts (`.auto-queue` state), and a successful download clears the platform's backoff. To stop a runner immediately, use Cancel in Active Jobs (Drain still waits for the in-flight video to finish, by design). See `common/jobs/platformBackoff.ts` and `common/controller/autoRunner.ts`. - **Per-video progress bar now advances for subtitle-only (YouTube-handling) downloads.** YouTube-handling channels fetch only subtitles (`--skip-download`), which yt-dlp reports with no byte total — so the Active Jobs progress bar sat empty and jumped straight to 100%. It now steps once per subtitle track (e.g. `subs 1/2`) using the track list yt-dlp announces, while real media downloads (Odysee/transcribe) keep their byte-based bar. See `createDownloadProgressParser` in `common/jobs/progressParsers.ts`. diff --git a/editor/app/auto-queue/components/AutoQueueView.tsx b/editor/app/auto-queue/components/AutoQueueView.tsx @@ -166,6 +166,15 @@ function KindPanel({ </span> </div> + {kind === "transcription" && ( + <p className="text-xs text-zinc-500"> + Auto-transcribe runs at background priority: a manually-triggered + transcription preempts it for the next free worker — the in-flight auto + transcription finishes (drains), then the manual one runs, then auto + resumes. + </p> + )} + {status.cooldowns.length > 0 && ( <CooldownBanner cooldowns={status.cooldowns} /> )} diff --git a/editor/e2e/auto-queue.spec.ts b/editor/e2e/auto-queue.spec.ts @@ -517,6 +517,88 @@ test("Auto-queue page: the Drain button stops the runner gracefully", async ({ await expect.poll(() => runnerRunning(request), { timeout: 15_000 }).toBe(false); }); +// Regression: draining the runner while one of its units is PARKED in the worker +// pool (waiting for a slot a manual transcription is holding) used to hang +// forever — the runner never forwarded its drainSignal to the parked acquire, so +// the unit's promise never settled and the drain loop spun on `inFlight > 0`. +// Stage exactly that: one worker, a manual "Transcribe missing" on a slowop video +// holds it busy, the auto-runner then launches a unit that parks behind it. Drain +// must still complete. +test("Drain completes when an auto-transcribe unit is parked behind a busy worker", async ({ + page, + request, +}) => { + await resetData(null); + // "busy": a transcribe channel with a single slowop video. The fake-whisper + // fixture runs ~50s when the video id contains "slowop", so a manual transcribe + // holds the only worker long enough for the auto unit to park and for us to + // drain. "alpha": a normal pending video the auto-runner will pick. + const busy = resolvePath("test-transcripts/channels/busy"); + await mkdir(`${busy}/data/slowop1`, { recursive: true }); + await writeFile( + `${busy}/config.json`, + JSON.stringify({ handling: "transcribe", name: "busy", url: "https://example.com/@busy" }), + ); + await writeFile(`${busy}/data/slowop1/audio.mp3`, "fake audio slowop1\n"); + await makeChannel("alpha", ["a1"]); + + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + workers: ONE_WORKER, + autoQueue: transcriptionAutoQueue(ALPHA_ROOT), // runner targets ONLY alpha + }); + + // 1) Manual foreground transcription holds the single worker (slowop → ~50s). + await page.goto("/channels/busy"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + // Wait until it's actually running (worker leased), so the auto unit will park. + await expect + .poll( + async () => { + const res = await request.get(`${baseUrl}/api/jobs/active`); + const body = await res.json(); + const jobs = Array.isArray(body) ? body : (body.jobs ?? []); + return jobs.some( + (j: { kind?: string; status?: string }) => + j.kind === "whisper-all" && j.status === "running", + ); + }, + { timeout: 20_000 }, + ) + .toBe(true); + + // 2) Start the auto-runner; it picks a1 and parks behind the busy worker. + await startRunner(request); + await expect + .poll( + async () => + (await getStatus(request)).transcription.picks.some( + (p) => p.videoId === "a1", + ), + { timeout: 30_000 }, + ) + .toBe(true); + + // 3) Drain the runner. The parked a1 unit must unblock (skip) so the runner + // finalizes instead of hanging. + await page.goto("/auto-queue"); + const section = page.locator("section", { + has: page.getByRole("heading", { name: "Auto-transcribe" }), + }); + await section.getByRole("button", { name: "Drain Auto-transcribe" }).click(); + await expect.poll(() => runnerRunning(request), { timeout: 15_000 }).toBe(false); + + // Cleanup: hard-cancel the lingering slowop manual job so the worker frees and + // doesn't bleed into the next spec (afterEach only stops runners). + await page.goto("/jobs/active"); + for (const btn of await page.getByRole("button", { name: /^Cancel$/ }).all()) { + await btn.click().catch(() => {}); + } +}); + test("download: prioritizes channels across the per-platform queue", async ({ request, }) => {