commit b6abf2b9fd6381add251a5887332ce54bed03e17
parent 61e2071a744642661d57cddb9b15e5ae6178d2a5
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 22 Jun 2026 16:16:14 -0400
Auto-queue: add Drain control; verify Drain/Cancel on runner jobs
The Active Jobs Drain and Cancel buttons already act on the runner jobs
(drainJobAction/cancelJobAction address by job id; auto-* kinds are in
DRAINABLE_KINDS and the loop honors drainSignal/signal). Added the missing
piece: a Drain control on the Auto-queue page next to Stop.
- drainAutoRunner(kind) -> registry.requestDrain(jobId) (graceful stop:
finish in-flight, start nothing new, then the runner ends "done").
- /api/auto-queue/control now accepts action "drain" (stop stays hard-cancel).
- Auto-queue page: Drain button alongside Stop.
New e2e: Active Jobs Drain and Cancel stop the runner; the page Drain button
stops it gracefully. Full auto-queue suite + production build green.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
5 files changed, 101 insertions(+), 5 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -565,9 +565,20 @@ export async function startAutoRunnersIfEnabled(
await startAutoRunner("download", paths);
}
-// Stop a runner (hard cancel its job). Returns true if one was running.
+// Stop a runner (hard cancel its job). In-flight units are aborted. Returns true
+// if one was running.
export function stopAutoRunner(kind: AutoQueueKind): boolean {
const live = getSingleton().runners.get(kind);
if (!live) return false;
return getRegistry().cancel(live.jobId);
}
+
+// Drain a runner (graceful stop): stop picking new work, let in-flight units
+// finish, then the loop exits and the job ends "done". Returns true if one was
+// running. Same soft-cancel the Active Jobs "Drain" button uses, just addressed
+// by kind instead of job id.
+export function drainAutoRunner(kind: AutoQueueKind): boolean {
+ const live = getSingleton().runners.get(kind);
+ if (!live) return false;
+ return getRegistry().requestDrain(live.jobId);
+}
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -3,7 +3,7 @@
## [Unreleased]
- **"Stop & keep progress" no longer mislabels the paused video as a failed transcription.** Using **Stop & keep progress** on a busy parakeet worker (or any partial-capable engine) sends the engine a graceful SIGTERM so it stops after the current window and the video resumes next run. But if the engine took longer than execa's 5-second force-kill window to exit — which a parakeet window routinely does, since finishing/stitching one ~480s window outlasts 5s — execa force-SIGKILLed it and the resulting "Command was killed with SIGTERM … forcefully terminated after 5000 milliseconds" error escaped the pause handling: it was treated as a genuine transcription failure and the video was written to the channel's `failed-transcriptions` file *permanently* (so even though its completed windows were cached for resume, it was skipped as "failed" on every later run). The transcribe path now recognizes that a force-killed **requested pause** is still a pause, not a failure — it returns the `paused` outcome (a skip, not a failure), so nothing lands in `failed-transcriptions` and the next "Transcribe missing" resumes it from the cached windows. Hard **Cancel** and **Drain** were never affected (their abort signal already classifies the kill as a skip). A video wrongly blacklisted by the old behavior won't auto-prune (it has real audio) — clear it with the channel's **Clear failed transcriptions** action to retry. See `common/controller/transcribeOne.ts`.
- **New Auto-queue: automatically transcribe (and download) across all channels by a configurable priority policy, instead of running one channel batch at a time.** Previously the only way to process pending work was to manually fire a per-channel batch (e.g. *Transcribe missing* on one channel), and since every transcription job serialized on a single queue, a batch ran to completion before any other channel got a turn — there was no way to say "do cornbreadman first, then fall back to hasanabi." The new **Auto-queue** page (under Pool → Auto-queue) adds two always-on runners, **auto-transcribe** and **auto-download**, each driven by a **policy tree**: order rules top-to-bottom for **strict** priority, or wrap rules in a group set to **round-robin** or **weighted-fair** (smooth weighted round-robin) to *alternate* between rulesets. A rule (leaf) matches a **channel**, a whole **platform**, or **all** channels, optionally narrowed to a snapshot **bucket** (e.g. prioritize `failedListed` retries over fresh `downloadedNoTranscript`), and any rule or group can carry a **max-workers** cap (a saturated subtree falls through to the next-priority sibling, like an HTB ceil). The highest-priority channel with available work claims the **next freed worker slot** — non-destructive, so a higher-priority video never kills an in-flight transcription, it just wins the next slot; when a channel's work runs out the runner falls back automatically. Transcription concurrency is bounded by the worker pool's eligible slots (so policy decides *which* video runs, the pool decides *how many*); downloads have no pool, so the runner gates to **one download per platform at a time**, matching the per-platform serial queue's politeness. Each runner is a real, drainable/cancellable job (visible on the Jobs pages), and the Auto-queue page shows live per-rule pending counts and a recent-pick log. Independent of the sync **Schedule** (which only decides *when* to re-fetch a channel) — manual batches keep working alongside it. Policies live in `settings.json` under `autoQueue` (defensively sanitized like `syncScheduler`); fairness cursors persist in `transcripts/.auto-queue/state.json`. See `common/jobs/autoQueuePolicy.ts` (pure selection engine + unit tests), `common/controller/autoRunner.ts`, `common/controller/transcribeOneFromQueue.ts` (shared per-video gating, also used by the existing whisper batch), and `editor/app/auto-queue/*`.
- - **Auto-queue refinements:** (1) the auto-transcribe runner now honors disabled workers — it no longer used CPU workers you'd turned off via Workers → *Set as default*. The runner started at boot and called the worker pool's `reconfigure()` before anything triggered the pool's lazy init, which set the pool's `initialized` flag *without* applying the saved `.worker-defaults.json` arrangement, so disabled workers came back enabled. The runner no longer pre-empts that init (it relies on the pool's own first-use init, which applies both settings and the saved default). (2) The runners no longer show up under a generic **"Other"** group on the **Active Jobs** screen — channel-less jobs are now grouped into their own labeled sections ("Auto-transcribe" / "Auto-download") via a shared `jobKindLabel` map (also used for friendlier kind labels in the running-jobs list), and each in-flight item is labeled `channel/videoId` so you can see which channel it's on. See `common/controller/autoRunner.ts`, `editor/app/jobs/jobKindLabels.ts`, and `editor/app/jobs/components/ActiveJobsLive.tsx`.
+ - **Auto-queue refinements:** (1) the auto-transcribe runner now honors disabled workers — it no longer used CPU workers you'd turned off via Workers → *Set as default*. The runner started at boot and called the worker pool's `reconfigure()` before anything triggered the pool's lazy init, which set the pool's `initialized` flag *without* applying the saved `.worker-defaults.json` arrangement, so disabled workers came back enabled. The runner no longer pre-empts that init (it relies on the pool's own first-use init, which applies both settings and the saved default). (2) The runners no longer show up under a generic **"Other"** group on the **Active Jobs** screen — channel-less jobs are now grouped into their own labeled sections ("Auto-transcribe" / "Auto-download") via a shared `jobKindLabel` map (also used for friendlier kind labels in the running-jobs list), and each in-flight item is labeled `channel/videoId` so you can see which channel it's on. See `common/controller/autoRunner.ts`, `editor/app/jobs/jobKindLabels.ts`, and `editor/app/jobs/components/ActiveJobsLive.tsx`. (3) The Auto-queue page now has a **Drain** control (graceful stop: finish in-flight items, start no new ones, then stop) alongside **Stop** (hard stop, aborts in-flight); both the **Drain** and **Cancel** buttons on the **Active Jobs** screen also act on the runner jobs.
- **Charts can now track content *added to the sites* over time, not just when creators uploaded it.** The chart engine previously only binned the time axis on a video's **upload date**. Two acquisition dates are now recorded per video — when *we* downloaded it and when *we* transcribed it — and the chart editor's X-axis gains a **Date field** selector (Uploaded / Downloaded / Transcribed) alongside the bin. A new **Content added** preset group ships three ready charts (cumulative *Library growth (added)*, cumulative *Transcribed over time*, and *Added per month* stacked by channel), and the default dashboard now includes the cumulative library-growth-by-acquisition chart so the progress view is present out of the box. Everything reuses the existing charts UI, so per-channel filtering, cumulative curves, CSV/PNG export, and shareable URLs all work unchanged. Acquisition dates come from the per-video `download-outcome.json` (`finishedAt`) and a new `transcribe-outcome.json` sidecar written when a transcript is finalized; the stats build falls back to file mtimes for content added before the sidecars existed. Requires a one-time data rebuild (`build:index` + `build:stats`) on the bumped `STATS_SCHEMA_VERSION`. See `common/lib/{stats,chartConfig,chartAggregate,chartShare,transcribeOutcome}.ts`, `common/controller/{buildStats,transcribeOne}.ts`, and `common/components/charts/ChartConfigEditor.tsx`.
- **The Schedule page is now a one-stop editor for per-channel sync cadence.** The `/scheduler` page used to be read-only — you could see each channel's interval, last sync, and next-due time, but to *change* a cadence you had to open that channel's editor (Source → Auto-sync), one channel at a time, and the headline global toggles lived only in Settings. Now each row's **Interval** cell is an inline editor: pick a preset (Default / Off / Every 10–30m / Hourly / 6h / 12h / Daily / Weekly) **or** choose **Custom (minutes)…** and type an exact minute count, then **Save** — writing just `syncIntervalMinutes` to that channel's `config.json` and leaving every other field untouched (it does *not* go through the full channel-form merge). The page also gained a **Global controls** block to toggle the master **enable**, the **default interval**, and the **internal heartbeat** right there (the advanced knobs — concurrency, quiet hours, backoff — still link out to Settings). The underlying due logic is unchanged: a channel auto-syncs on the next heartbeat once `now − lastSyncedAt ≥ its interval`. The live status columns keep polling every 5s, but each row's editor holds its own state seeded once from the stored value, so a refresh can't clobber an in-progress edit. The preset list is shared with the channel editor (`editor/app/scheduler/intervalPresets.ts`), and `GET /api/scheduler/status` now carries each channel's raw `configuredIntervalMinutes` so the editor can tell *inherit-default* from *explicit-off* from *explicit-minutes*. See `editor/app/scheduler/actions.ts` and `editor/app/scheduler/components/{ChannelIntervalEditor,SchedulerSettingsForm}.tsx`.
- **Scheduled sync can now run without an external cron job.** The sync scheduler previously only fired when an OS cron entry POSTed to `/api/scheduler/tick` (via `pnpm sync:tick`) — fine on a server, but a chore to set up just to call a function the editor already hosts in-process. The editor can now drive its own heartbeat through a **Next.js instrumentation hook** (`editor/instrumentation.ts`): on server startup it arms a single in-process timer that calls `runSchedulerTick()` directly — no HTTP, no cron, no token. Turn it on with **Settings → Sync scheduler → Internal heartbeat (seconds)**: `0` = off (keep using an external cron heartbeat), any positive value is clamped to `[15, 3600]`s and is the cadence the editor ticks itself at; the `SYNC_HEARTBEAT_SECONDS` env var overrides the setting at runtime. The timer is a self-rescheduling, `unref`'d `setTimeout` loop (so it never holds the process open and never overlaps a tick), re-reading the cadence each fire so a change takes effect on the next tick — though turning it on *from 0* needs a restart, since the timer is armed once at boot. It's modeled on the existing snapshot-scheduler timer and reuses the already-overlap-guarded `runSchedulerTick()`, so internal and external heartbeats are interchangeable and may even coexist. The **Schedule** page header now reports how ticks are driven ("internal heartbeat every N" vs. "external heartbeat (cron)"), and `GET /api/scheduler/status` carries the effective `heartbeatSeconds`. Defaults to off, so dev/test and existing cron installs are unchanged. One caveat for multi-instance deployments: the timer runs once *per server instance*, so a cluster against one data dir should set `SYNC_HEARTBEAT_SECONDS=0` on all but one — see `SCHEDULED_SYNC.md`.
diff --git a/editor/app/api/auto-queue/control/route.ts b/editor/app/api/auto-queue/control/route.ts
@@ -1,5 +1,6 @@
import { NextResponse } from "next/server";
import {
+ drainAutoRunner,
startAutoRunner,
stopAutoRunner,
} from "yt-dlp-transcript-common/controller/autoRunner";
@@ -13,7 +14,8 @@ export const dynamic = "force-dynamic";
// flipping the enable toggle. The editor admin surface is otherwise
// unauthenticated (trusted self-host), consistent with the rest of the app.
//
-// Body: { kind: "transcription" | "download", action: "start" | "stop" }.
+// Body: { kind: "transcription" | "download", action: "start" | "stop" | "drain" }.
+// "stop" hard-cancels (aborts in-flight); "drain" lets in-flight units finish.
export async function POST(req: Request) {
let body: { kind?: unknown; action?: unknown };
try {
@@ -37,8 +39,12 @@ export async function POST(req: Request) {
const stopped = stopAutoRunner(kind as AutoQueueKind);
return NextResponse.json({ ok: true, stopped });
}
+ if (action === "drain") {
+ const draining = drainAutoRunner(kind as AutoQueueKind);
+ return NextResponse.json({ ok: true, draining });
+ }
return NextResponse.json(
- { ok: false, error: "action must be 'start' or 'stop'" },
+ { ok: false, error: "action must be 'start', 'stop', or 'drain'" },
{ status: 400 },
);
}
diff --git a/editor/app/auto-queue/components/AutoQueueView.tsx b/editor/app/auto-queue/components/AutoQueueView.tsx
@@ -91,7 +91,7 @@ function KindPanel({
const running = status.runner.running;
const control = useCallback(
- async (action: "start" | "stop") => {
+ async (action: "start" | "stop" | "drain") => {
setBusy(true);
try {
await fetch("/api/auto-queue/control", {
@@ -141,7 +141,18 @@ function KindPanel({
</button>
<button
type="button"
+ aria-label={`Drain ${title}`}
+ title="Finish the in-flight items, start no new ones, then stop"
+ onClick={() => control("drain")}
+ disabled={busy || !running}
+ className="px-3 py-1.5 rounded-md border border-amber-300 dark:border-amber-800 text-sm font-medium text-amber-700 dark:text-amber-300 hover:bg-amber-50 dark:hover:bg-amber-950 disabled:opacity-50"
+ >
+ Drain
+ </button>
+ <button
+ type="button"
aria-label={`Stop ${title}`}
+ title="Stop now, aborting any in-flight item"
onClick={() => control("stop")}
disabled={busy || !running}
className="px-3 py-1.5 rounded-md border border-zinc-300 dark:border-zinc-700 text-sm font-medium hover:bg-zinc-100 dark:hover:bg-zinc-800 disabled:opacity-50"
diff --git a/editor/e2e/auto-queue.spec.ts b/editor/e2e/auto-queue.spec.ts
@@ -449,6 +449,74 @@ test("Active Jobs: the runner shows in its own labeled section, not Other", asyn
await expect(page.getByRole("heading", { name: "Other" })).toHaveCount(0);
});
+const IDLE_SETTINGS = (root: Group) => ({
+ adminTitle: "Test Admin",
+ maxTranscriptPageBytes: 8388608,
+ sleepBetweenDownloadsSeconds: 0,
+ minFreeDiskGB: 0,
+ workers: ONE_WORKER,
+ autoQueue: transcriptionAutoQueue(root),
+});
+
+const ALPHA_ROOT: Group = {
+ id: "root",
+ mode: "strict",
+ children: [{ id: "leaf-alpha", match: { type: "channel", value: "alpha" } }],
+};
+
+async function runnerRunning(
+ request: import("@playwright/test").APIRequestContext,
+): Promise<boolean> {
+ return (await getStatus(request)).transcription.runner.running;
+}
+
+test("Active Jobs: Drain and Cancel buttons work on the runner job", async ({
+ page,
+ request,
+}) => {
+ await resetData(null);
+ await makeChannel("alpha", []); // enabled but no pending work -> idle/running
+ await writeSettings(IDLE_SETTINGS(ALPHA_ROOT));
+ await startRunner(request);
+
+ await page.goto("/jobs/active");
+ const section = page.locator(
+ "section[aria-label='System jobs: Auto-transcribe']",
+ );
+ await expect(section.getByRole("button", { name: "Drain" })).toBeVisible();
+ await expect(section.getByRole("button", { name: /^Cancel$/ })).toBeVisible();
+
+ // Drain stops the (idle) runner — the job ends and leaves the active list.
+ await section.getByRole("button", { name: "Drain" }).click();
+ await expect.poll(() => runnerRunning(request), { timeout: 15_000 }).toBe(false);
+
+ // Start it again and Cancel from the row.
+ await startRunner(request);
+ await page.reload();
+ await expect(section.getByRole("button", { name: /^Cancel$/ })).toBeVisible();
+ await section.getByRole("button", { name: /^Cancel$/ }).click();
+ await expect.poll(() => runnerRunning(request), { timeout: 15_000 }).toBe(false);
+});
+
+test("Auto-queue page: the Drain button stops the runner gracefully", async ({
+ page,
+ request,
+}) => {
+ await resetData(null);
+ await makeChannel("alpha", []);
+ await writeSettings(IDLE_SETTINGS(ALPHA_ROOT));
+ await startRunner(request);
+
+ await page.goto("/auto-queue");
+ const section = page.locator("section", {
+ has: page.getByRole("heading", { name: "Auto-transcribe" }),
+ });
+ await expect.poll(() => runnerRunning(request), { timeout: 15_000 }).toBe(true);
+
+ await section.getByRole("button", { name: "Drain Auto-transcribe" }).click();
+ await expect.poll(() => runnerRunning(request), { timeout: 15_000 }).toBe(false);
+});
+
test("download: prioritizes channels across the per-platform queue", async ({
request,
}) => {