commit aea1282c3245e3e3027a1cb251293431d3f5f0cd
parent a303efc1a364c232654338db9bc4f6809eaa45d9
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 23 Jun 2026 14:36:15 -0400
Auto-queue: share the per-platform queue with manual Sync
Auto-download units now run as real jobs on the channel's platform queue
(downloadQueueKey) — the same queue a manual Sync uses — instead of running
out-of-band. The job registry serializes them, so the two never spawn yt-dlp
against one platform at once (in either direction), eliminating the 429 you'd
hit by clicking Sync while auto-download was running.
- streamCommand: runManagedFunction returns a `done` completion promise + a
`background` flag; resolved from every terminal path.
- registry: `background` jobs queue behind foreground ones on the same key, so
a clicked Sync preempts queued auto-download units (never interrupts a running
one).
- autoRunner: launchUnit enqueues a managed `auto-download-unit` job and awaits
it; child jobs are hard-cancelled on stop; sync-written cooldowns are merged
in each iteration.
- downloadBackoff (new): shared per-platform 429 cooldown read/write.
- runYtdlp/pipelineActions: Sync is refused on a cooling-down platform with
remaining time; any download path records the shared cooldown on a 429.
- UI: auto-download-unit label everywhere, an "auto" badge + "waiting behind a
manual job" hint, `background` serialized to the client, and a live 429
cooldown banner on the Auto-queue page.
Tests: +registry priority-insertion and downloadBackoff unit tests; +an e2e
cooldown-refusal test. common suite 45/45 and the queues/auto-queue/jobs e2e
specs are green.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Diffstat:
20 files changed, 697 insertions(+), 78 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -32,6 +32,8 @@ import {
pruneExpired,
} from "../jobs/platformBackoff";
import { type DownloadFailureClass } from "../lib/availability";
+import { type DownloadOutcomeStatus } from "../lib/downloadOutcome";
+import { downloadQueueKey } from "../lib/queueKeys";
import { readChannelConfig } from "./channels";
import { readChannelSnapshot } from "./channelSnapshot";
import { transcribeOneFromQueue } from "./transcribeOneFromQueue";
@@ -53,6 +55,12 @@ import { downloadOneManaged } from "../ytdlp/downloadOneManaged";
export const AUTO_TRANSCRIBE_KIND = "auto-transcribe";
export const AUTO_DOWNLOAD_KIND = "auto-download";
+// One per-video download launched by the auto-download runner. Distinct from
+// the long-lived runner job (AUTO_DOWNLOAD_KIND): each unit is a real managed
+// job on the channel's platform queue (downloadQueueKey), so it serializes
+// against a manual sync. Marked background so a clicked sync preempts queued
+// units. NOT in DRAINABLE_KINDS — a single-video download has nothing to drain.
+export const AUTO_DOWNLOAD_UNIT_KIND = "auto-download-unit";
export function autoRunnerJobKind(kind: AutoQueueKind): string {
return kind === "transcription" ? AUTO_TRANSCRIBE_KIND : AUTO_DOWNLOAD_KIND;
@@ -294,6 +302,11 @@ async function runLoop(
const platformKey = (slug: string, slugToPlatform: Map<string, string>) =>
slugToPlatform.get(slug) ?? "unknown";
+ // videoId -> child download job id, for the units the runner currently has on
+ // the registry's platform queues. Used to hard-cancel them (kill yt-dlp or
+ // drop a still-queued unit) when the runner itself is cancelled. Download-only.
+ const childJobIds = new Map<string, string>();
+
let metaCache: ChannelMeta[] = [];
let metaAt = 0;
@@ -356,6 +369,20 @@ async function runLoop(
// download (busy) OR is in a rate-limit/network backoff window (cooling
// down), so a busy/throttled platform yields to the next-priority free one.
if (kind === "download") {
+ // Merge in any cooldown a manual sync/import wrote to the shared state
+ // (read-modify-write from outside the runner) since our last persist.
+ // Take the later `until` so a sync-set 429 cooldown is honored even though
+ // our own next persist would otherwise clobber it.
+ try {
+ const persisted = (await readAutoQueueState(paths)).download
+ .platformBackoff;
+ for (const [pf, e] of Object.entries(persisted)) {
+ const cur = kindState.platformBackoff[pf];
+ if (!cur || e.until > cur.until) kindState.platformBackoff[pf] = e;
+ }
+ } catch {
+ // Best-effort: a transient read failure just skips this iteration's merge.
+ }
const skip = new Set<string>();
for (const [pf, n] of platformInFlight) {
if (n >= PER_PLATFORM_CAP) skip.add(pf);
@@ -390,7 +417,16 @@ async function runLoop(
const unitPlatform =
kind === "download" ? platformKey(channelSlug, slugToPlatform) : null;
let result: UnitResult = { outcome: "failed" };
- launchUnit({ kind, paths, pick, channelSlug, tracker, onLog, signal })
+ launchUnit({
+ kind,
+ paths,
+ pick,
+ channelSlug,
+ tracker,
+ onLog,
+ signal,
+ onChildJob: (jid) => childJobIds.set(pick.videoId, jid),
+ })
.then((r) => {
result = r;
})
@@ -399,6 +435,7 @@ async function runLoop(
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);
@@ -473,7 +510,14 @@ async function runLoop(
onLog(`Auto-${kind}: ${channelSlug}/${pick.videoId} (leaf ${pick.leafId})`);
}
- // Drain remaining in-flight units (graceful stop / drain / cancel).
+ // 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(
@@ -489,6 +533,9 @@ type LaunchArgs = {
tracker: ReturnType<typeof makeTaskTracker>;
onLog: (line: string) => void;
signal: 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;
};
// The outcome of one unit. `failureClass` is set only for download failures
@@ -521,10 +568,14 @@ async function launchUnit(args: LaunchArgs): Promise<UnitResult> {
return { outcome: res.outcome };
}
// Download: resolve the video's source URL and run the same managed
- // single-video download the import/download-missing paths use, so archive
- // appends, auth/no-subs fallbacks and audio-integrity checks all apply. The
- // per-platform in-flight gate in the loop keeps this to one download per
- // platform at a time (matching the per-platform serial queue's politeness).
+ // single-video download the import/download-missing paths use (so archive
+ // appends, auth/no-subs fallbacks and audio-integrity checks all apply) — but
+ // as its OWN managed job on the channel's platform queue (downloadQueueKey),
+ // the SAME queue a manual sync uses. The registry then serializes this unit
+ // against any sync on that platform, so the two never spawn yt-dlp at once.
+ // Marked `background` so a clicked sync preempts queued units (without
+ // interrupting a running one). The runner's per-platform in-flight gate still
+ // keeps this to one outstanding unit per platform so the queue isn't flooded.
const config = await readChannelConfig(args.paths, args.channelSlug);
if (!config) return { outcome: "skipped" };
const url = await findVideoSourceUrl(
@@ -538,42 +589,74 @@ async function launchUnit(args: LaunchArgs): Promise<UnitResult> {
return { outcome: "skipped" };
}
const settings = getSettings();
- const task = args.tracker.start({
- id: args.pick.videoId,
- label: `${args.channelSlug}/${args.pick.videoId}`,
- kind: "download",
+
+ // The child job's fn sets these synchronously before it returns; the runner
+ // reads them after `await res.done` resolves (single-threaded → no race).
+ let unitStatus: DownloadOutcomeStatus | undefined;
+ let unitFailureClass: DownloadFailureClass | undefined;
+
+ const res = await runManagedFunction({
+ kind: AUTO_DOWNLOAD_UNIT_KIND,
+ queueKey: downloadQueueKey(config),
+ paths: args.paths,
+ channelSlug: args.channelSlug,
+ videoId: args.pick.videoId,
+ background: true,
+ fn: async (onLog, signal, _setProgress, ctx) => {
+ // Per-video progress lives on THIS child job's row (its own ctx), so the
+ // unit shows up in Active Jobs like a manual single-video download.
+ const task = makeTaskTracker(ctx, onLog).start({
+ id: args.pick.videoId,
+ label: `${args.channelSlug}/${args.pick.videoId}`,
+ kind: "download",
+ });
+ try {
+ // downloadOneManaged does NOT throw on a 429/network failure — it
+ // returns a record whose last attempt carries the stderr tail +
+ // availability class, already classified into failureClass.
+ const record = await downloadOneManaged({
+ channelSlug: args.channelSlug,
+ channelConfig: config,
+ paths: args.paths,
+ videoUrl: url,
+ onLog: task.onLog,
+ // The child job's own abort signal: registry.cancel(childJobId) aborts
+ // it (→ kills yt-dlp) when the runner is hard-cancelled.
+ signal,
+ globalCookiesFromBrowser: settings.cookiesFromBrowser || undefined,
+ inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback,
+ globalSkipLiveDownloads: settings.skipLiveDownloads,
+ appendArchive: true,
+ });
+ unitStatus = record.status;
+ unitFailureClass = record.failureClass;
+ // Land a failed download as a failed JOB (red row) rather than a silent
+ // "done"; the runner reads the captured status/class regardless.
+ if (
+ record.status === "failed" ||
+ record.status === "failed-corrupt-source"
+ ) {
+ throw new Error(`download failed (${record.failureClass ?? "unknown"})`);
+ }
+ } finally {
+ task.end();
+ }
+ },
});
- try {
- // downloadOneManaged does NOT throw on a 429/network failure — it returns
- // a record whose last attempt carries the stderr tail + availability class.
- // Classify it the same way the batch path does (runYtdlp.ts) so the runner
- // can back off per-platform instead of immediately re-hammering the source.
- const record = await downloadOneManaged({
- channelSlug: args.channelSlug,
- channelConfig: config,
- paths: args.paths,
- videoUrl: url,
- onLog: task.onLog,
- signal: args.signal,
- globalCookiesFromBrowser: settings.cookiesFromBrowser || undefined,
- inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback,
- globalSkipLiveDownloads: settings.skipLiveDownloads,
- appendArchive: true,
- });
- if (record.status === "skipped-filtered") return { outcome: "skipped" };
- if (record.status.startsWith("ok")) return { outcome: "transcribed" };
- // failureClass is computed inside downloadOneManaged against the full stderr
- // tail (catches a rate-limit WARNING masked by a later error line).
- return { outcome: "failed", failureClass: record.failureClass };
- } catch (err) {
- if (args.signal.aborted || (err as Error)?.name === "AbortError") {
- return { outcome: "skipped" };
- }
- args.onLog(`Auto-download failed for ${args.pick.videoId}: ${String(err)}`);
+ if (!res.ok) {
+ args.onLog(
+ `Auto-download: could not enqueue ${args.pick.videoId}: ${res.error}`,
+ );
return { outcome: "failed" };
- } finally {
- task.end();
}
+ args.onChildJob?.(res.jobId);
+ const term = await res.done;
+ // Cancelled (runner stopped, or dropped while still queued) → not a failure.
+ if (term.status === "cancelled") return { outcome: "skipped" };
+ if (unitStatus === "skipped-filtered") return { outcome: "skipped" };
+ if (unitStatus && unitStatus.startsWith("ok")) return { outcome: "transcribed" };
+ // failureClass drives the runner's per-platform backoff (rate_limit/network).
+ return { outcome: "failed", failureClass: unitFailureClass };
}
// --- Lifecycle -------------------------------------------------------------
diff --git a/common/jobs/downloadBackoff.test.ts b/common/jobs/downloadBackoff.test.ts
@@ -0,0 +1,77 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { mkdtemp, readFile, rm } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import path from "node:path";
+import type { Paths } from "../lib/paths";
+import {
+ platformCooldownRemainingMs,
+ recordDownloadBackoff,
+} from "./downloadBackoff";
+import { BACKOFF_BASE_MS } from "./platformBackoff";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/downloadBackoff.test.ts
+//
+// downloadBackoff only touches paths.autoQueueStateFile, so a stub Paths with
+// just that field exercises the full read-modify-write against a temp file.
+
+async function withTempPaths(
+ fn: (paths: Paths) => Promise<void>,
+): Promise<void> {
+ const dir = await mkdtemp(path.join(tmpdir(), "dl-backoff-"));
+ const paths = {
+ autoQueueStateFile: path.join(dir, ".auto-queue", "state.json"),
+ } as Paths;
+ try {
+ await fn(paths);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+}
+
+test("no cooldown on an untouched platform", async () => {
+ await withTempPaths(async (paths) => {
+ assert.equal(await platformCooldownRemainingMs("odysee", paths), 0);
+ });
+});
+
+test("recordDownloadBackoff opens a cooldown window", async () => {
+ await withTempPaths(async (paths) => {
+ await recordDownloadBackoff("odysee", paths);
+ const remaining = await platformCooldownRemainingMs("odysee", paths);
+ // First failure ≈ BACKOFF_BASE_MS (±10% jitter); just assert it's positive
+ // and within a sane upper bound.
+ assert.ok(remaining > 0, "platform is cooling down");
+ assert.ok(remaining <= BACKOFF_BASE_MS * 1.2, "within base window + jitter");
+ // It is recorded under the download kind only.
+ const state = JSON.parse(await readFile(paths.autoQueueStateFile, "utf8"));
+ assert.ok(state.download.platformBackoff.odysee, "stored under download kind");
+ assert.deepEqual(
+ state.transcription.platformBackoff,
+ {},
+ "transcription untouched",
+ );
+ });
+});
+
+test("consecutive failures escalate the fails counter", async () => {
+ await withTempPaths(async (paths) => {
+ await recordDownloadBackoff("youtube", paths);
+ await recordDownloadBackoff("youtube", paths);
+ const state = JSON.parse(await readFile(paths.autoQueueStateFile, "utf8"));
+ assert.equal(state.download.platformBackoff.youtube.fails, 2);
+ });
+});
+
+test("remaining is 0 once the window lapses", async () => {
+ await withTempPaths(async (paths) => {
+ // Hand-write a lapsed entry to simulate an expired cooldown.
+ const { writeAutoQueueState, emptyAutoQueueState } = await import(
+ "./autoQueueState"
+ );
+ const s = emptyAutoQueueState();
+ s.download.platformBackoff.odysee = { until: Date.now() - 1000, fails: 3 };
+ await writeAutoQueueState(paths, s);
+ assert.equal(await platformCooldownRemainingMs("odysee", paths), 0);
+ });
+});
diff --git a/common/jobs/downloadBackoff.ts b/common/jobs/downloadBackoff.ts
@@ -0,0 +1,43 @@
+// Shared accessors for the per-platform download rate-limit cooldown that both
+// the auto-download runner and manual sync/download paths respect.
+//
+// The runner already owns this state in memory and persists it to
+// `.auto-queue/state.json` (the "download" kind's platformBackoff). These
+// helpers let code OUTSIDE the runner — a clicked Sync, an import — observe and
+// extend the same cooldown via a fresh read-modify-write, so a 429 hit by
+// either path pauses the other. Because auto-download units now serialize on the
+// same per-platform job queue as sync (see autoRunner.ts / registry.ts), a
+// momentarily racy write here can at most delay politeness by an iteration — it
+// can no longer cause two concurrent yt-dlp processes to hammer one source.
+
+import { getPaths, type Paths } from "../lib/paths";
+import { readAutoQueueState, writeAutoQueueState } from "./autoQueueState";
+import { nextBackoff, pruneExpired } from "./platformBackoff";
+
+// Milliseconds remaining in the platform's current cooldown window, or 0 if it
+// is not cooling down. Reads the persisted "download" backoff state fresh.
+export async function platformCooldownRemainingMs(
+ platform: string,
+ paths: Paths = getPaths(),
+): Promise<number> {
+ const state = await readAutoQueueState(paths);
+ const entry = state.download.platformBackoff[platform];
+ const now = Date.now();
+ return entry && entry.until > now ? entry.until - now : 0;
+}
+
+// Record a rate-limit/network failure against a platform, escalating its
+// exponential cooldown. Read-modify-write of the shared download backoff state.
+export async function recordDownloadBackoff(
+ platform: string,
+ paths: Paths = getPaths(),
+): Promise<void> {
+ const state = await readAutoQueueState(paths);
+ const now = Date.now();
+ pruneExpired(state.download.platformBackoff, now);
+ state.download.platformBackoff[platform] = nextBackoff(
+ state.download.platformBackoff[platform],
+ now,
+ );
+ await writeAutoQueueState(paths, state);
+}
diff --git a/common/jobs/registry.test.ts b/common/jobs/registry.test.ts
@@ -0,0 +1,128 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { getRegistry, newJobId, type JobRecord } from "./registry";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/registry.test.ts
+//
+// Covers the per-queue priority insertion added so a manually clicked Sync
+// (foreground) preempts the auto-download runner's queued units (background)
+// without interrupting one already running. Each test uses a UNIQUE queueKey so
+// the shared global registry singleton doesn't leak state between tests.
+
+function rec(queueKey: string, background?: boolean): JobRecord {
+ return {
+ id: newJobId(),
+ kind: "test",
+ queueKey,
+ background,
+ status: "queued",
+ queuedAt: 0,
+ logPath: "",
+ };
+}
+
+// enqueue() needs start/onCancel callbacks; track which fired. In production
+// makeJob() registers the record before enqueue (so cancel/finalize/position
+// can look it up by id), so mirror that here.
+function enqueue(record: JobRecord) {
+ const calls = { started: false, cancelled: false };
+ getRegistry().register(record);
+ const result = getRegistry().enqueue(record, {
+ start: () => {
+ calls.started = true;
+ },
+ onCancel: () => {
+ calls.cancelled = true;
+ },
+ });
+ return { result, calls };
+}
+
+test("first job on a queue runs immediately", () => {
+ const q = `test:first:${newJobId()}`;
+ const a = rec(q);
+ const { result, calls } = enqueue(a);
+ assert.equal(result.willRunNow, true);
+ assert.equal(result.position, 0);
+ assert.equal(a.status, "running");
+ assert.equal(calls.started, true);
+});
+
+test("foreground inserts AHEAD of a queued background job", () => {
+ const q = `test:preempt:${newJobId()}`;
+ const head = rec(q); // runs immediately, never moves
+ enqueue(head);
+ const bg = rec(q, true); // queued background
+ const bgRes = enqueue(bg);
+ assert.equal(bgRes.result.position, 1);
+ const fg = rec(q); // foreground arrives after bg but should jump ahead of it
+ const fgRes = enqueue(fg);
+ assert.equal(fgRes.result.position, 1, "foreground slots in at index 1");
+ assert.equal(
+ getRegistry().positionInQueue(bg.id),
+ 2,
+ "background pushed back to index 2",
+ );
+ assert.equal(getRegistry().positionInQueue(head.id), 0, "running head fixed");
+});
+
+test("foreground jobs keep FIFO among themselves, before background", () => {
+ const q = `test:fifo:${newJobId()}`;
+ const head = rec(q);
+ enqueue(head);
+ const bg = rec(q, true);
+ enqueue(bg);
+ const fg1 = rec(q);
+ enqueue(fg1);
+ const fg2 = rec(q);
+ enqueue(fg2);
+ // Order: head(0), fg1(1), fg2(2), bg(3)
+ assert.equal(getRegistry().positionInQueue(fg1.id), 1);
+ assert.equal(getRegistry().positionInQueue(fg2.id), 2);
+ assert.equal(getRegistry().positionInQueue(bg.id), 3);
+});
+
+test("background jobs always append to the back", () => {
+ const q = `test:bgback:${newJobId()}`;
+ const head = rec(q);
+ enqueue(head);
+ const bg1 = rec(q, true);
+ enqueue(bg1);
+ const bg2 = rec(q, true);
+ const bg2Res = enqueue(bg2);
+ assert.equal(bg2Res.result.position, 2, "second background lands at the tail");
+});
+
+test("finalize promotes the foreground job before the background one", () => {
+ const q = `test:promote:${newJobId()}`;
+ const head = rec(q);
+ enqueue(head);
+ const bg = rec(q, true);
+ const bgStart = enqueue(bg);
+ const fg = rec(q);
+ const fgStart = enqueue(fg);
+ assert.equal(bgStart.calls.started, false);
+ assert.equal(fgStart.calls.started, false);
+ // Head finishes → q[0] (the foreground job) is promoted and started.
+ getRegistry().finalize(head.id, "done");
+ assert.equal(fg.status, "running", "foreground promoted first");
+ assert.equal(fgStart.calls.started, true);
+ assert.equal(bg.status, "queued", "background still waiting");
+});
+
+test("cancelling a queued background job fires onCancel and removes it", () => {
+ const q = `test:cancelbg:${newJobId()}`;
+ const head = rec(q);
+ enqueue(head);
+ const bg = rec(q, true);
+ const bgStart = enqueue(bg);
+ const ok = getRegistry().cancel(bg.id);
+ assert.equal(ok, true);
+ assert.equal(bgStart.calls.cancelled, true);
+ assert.equal(bg.status, "cancelled");
+ assert.equal(
+ getRegistry().positionInQueue(bg.id),
+ -1,
+ "removed from its queue",
+ );
+});
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -52,6 +52,10 @@ export type JobRecord = {
// A serializable replay descriptor, set for bookmarkable job kinds. Its
// presence is what makes a job "bookmarkable" in the UI. See jobSpec.ts.
spec?: JobSpec;
+ // Background jobs queue BEHIND foreground (default) jobs on the same queueKey
+ // (see enqueue). The auto-download runner marks its per-video units background
+ // so a manually-clicked Sync on that platform preempts the queued units.
+ background?: boolean;
status: JobStatus;
queuedAt: number;
startedAt?: number;
@@ -114,9 +118,16 @@ class JobRegistry {
);
}
- // Append a job to its queue. If it's the only entry, start it immediately
- // and mark "running". Otherwise leave it as "queued" and remember the
- // start/cancel callbacks for when it becomes head-of-queue or is cancelled.
+ // 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.
+ //
+ // 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.
//
// queueKey === "" is special: the job bypasses queue serialization entirely
// and runs immediately, parallel to any other work. Cancellation still
@@ -139,14 +150,29 @@ class JobRegistry {
q = [];
this.queues.set(record.queueKey, q);
}
- q.push(record);
- if (q.length === 1) {
+ if (q.length === 0) {
+ q.push(record);
record.status = "running";
record.startedAt = Date.now();
callbacks.start();
return { willRunNow: true, position: 0 };
}
- return { willRunNow: false, position: q.length - 1 };
+ 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) };
}
// Idempotent: marks the job terminal (if not already), splices it out of
diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts
@@ -7,6 +7,7 @@ import {
newJobId,
type JobProgress,
type JobRecord,
+ type JobStatus,
type JobTask,
} from "./registry";
import type { Paths } from "../lib/paths";
@@ -31,8 +32,19 @@ function requestSnapshotOnFinish(jobId: string, opts: CommonOpts): void {
requestChannelSnapshot(opts.paths, opts.channelSlug);
}
+// Resolves when a started job reaches a terminal state. Lets a server-side
+// caller (e.g. the auto-download runner) enqueue a job on a serialized queue
+// and await its completion — the queue may hold it behind other work, so the
+// `{ jobId, stream }` return alone can't tell the caller when it actually ran.
+export type JobDoneResult = { status: JobStatus; jobId: string };
+
export type StreamActionResult =
- | { ok: true; jobId: string; stream: ReadableStream<string> }
+ | {
+ ok: true;
+ jobId: string;
+ stream: ReadableStream<string>;
+ done: Promise<JobDoneResult>;
+ }
// `info: true` marks a non-error outcome that started no job (e.g. a
// re-derived bucket that's currently empty) so the UI can show it neutrally
// rather than as a red failure.
@@ -47,6 +59,11 @@ type CommonOpts = {
// When set, recorded on the job (and its meta sidecar) so the job can be
// bookmarked and re-launched later. See common/jobs/jobSpec.ts.
spec?: JobSpec;
+ // Background jobs queue BEHIND any foreground (default) job on the same
+ // queueKey: a manually-triggered job (sync, manual download) jumps ahead of
+ // queued background work without interrupting a running one. The auto-download
+ // runner marks its per-video units background so a clicked Sync preempts them.
+ background?: boolean;
};
export type RunManagedCommandOpts = CommonOpts & {
@@ -96,6 +113,7 @@ function makeJob(
channelSlug?: string,
videoId?: string,
spec?: JobSpec,
+ background?: boolean,
): { id: string; logPath: string; record: JobRecord } {
const id = newJobId();
const logPath = path.join(paths.jobsDir, `${id}.log`);
@@ -106,6 +124,7 @@ function makeJob(
channelSlug,
videoId,
spec,
+ background,
status: "queued",
queuedAt: Date.now(),
logPath,
@@ -120,6 +139,27 @@ function ignoreFileStreamErrors(stream: WriteStream): void {
stream.on("error", () => {});
}
+// A one-shot completion deferred attached to a managed job. `settle` is called
+// from every terminal path (finalized done/failed/cancelled, or cancelled
+// before start); the first call wins so callers awaiting `done` see the true
+// terminal status exactly once.
+function makeDoneDeferred(jobId: string): {
+ done: Promise<JobDoneResult>;
+ settle: (status: JobStatus) => void;
+} {
+ let resolveDone!: (r: JobDoneResult) => void;
+ const done = new Promise<JobDoneResult>((res) => {
+ resolveDone = res;
+ });
+ let settled = false;
+ const settle = (status: JobStatus) => {
+ if (settled) return;
+ settled = true;
+ resolveDone({ status, jobId });
+ };
+ return { done, settle };
+}
+
export async function runManagedCommand(
opts: RunManagedCommandOpts,
): Promise<StreamActionResult> {
@@ -132,8 +172,11 @@ export async function runManagedCommand(
opts.channelSlug,
opts.videoId,
opts.spec,
+ opts.background,
);
+ const { done, settle } = makeDoneDeferred(id);
+
const safe = makeSafeController<string>();
let fileStream: WriteStream | null = null;
let cancelledBeforeStart = false;
@@ -198,19 +241,21 @@ export async function runManagedCommand(
requestSnapshotOnFinish(id, opts);
// Persist terminal state (status/endedAt/exitCode now set by finalize).
void writeJobMeta(opts.paths, record);
+ settle(record.status);
});
};
const onCancel = () => {
cancelledBeforeStart = true;
safe.safeClose();
+ settle("cancelled");
};
registry.enqueue(record, { start, onCancel });
// Persist queued/running identity up front so a mid-run crash still leaves a
// sidecar; the .finally above rewrites it with the terminal state.
void writeJobMeta(opts.paths, record);
- return { ok: true, jobId: id, stream };
+ return { ok: true, jobId: id, stream, done };
}
export async function runManagedFunction(
@@ -225,8 +270,11 @@ export async function runManagedFunction(
opts.channelSlug,
opts.videoId,
opts.spec,
+ opts.background,
);
+ const { done, settle } = makeDoneDeferred(id);
+
const safe = makeSafeController<string>();
let fileStream: WriteStream | null = null;
let cancelledBeforeStart = false;
@@ -306,17 +354,19 @@ export async function runManagedFunction(
requestSnapshotOnFinish(id, opts);
// Persist terminal state (status/endedAt/exitCode now set by finalize).
void writeJobMeta(opts.paths, record);
+ settle(record.status);
});
};
const onCancel = () => {
cancelledBeforeStart = true;
safe.safeClose();
+ settle("cancelled");
};
registry.enqueue(record, { start, onCancel });
// Persist queued/running identity up front so a mid-run crash still leaves a
// sidecar; the .finally above rewrites it with the terminal state.
void writeJobMeta(opts.paths, record);
- return { ok: true, jobId: id, stream };
+ return { ok: true, jobId: id, stream, done };
}
diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts
@@ -82,6 +82,14 @@ export type RunYtdlpOpts = {
// mutating the channel config on disk. Useful for retrying old "youtube"
// videos as "transcribe".
handlingOverride?: ChannelHandling;
+ // Called when a per-video download fails with a rate-limit (HTTP 429) or
+ // network error, so the caller can record the SHARED per-platform cooldown
+ // (see common/jobs/downloadBackoff.ts). This entangles manual sync/download
+ // with the auto-download runner: a 429 hit by either path backs off the other.
+ // Kept as a callback so this module stays decoupled from the auto-queue state.
+ onPlatformBackoff?: (
+ failureClass: "rate_limit" | "network",
+ ) => void | Promise<void>;
};
export async function runYtdlp(opts: RunYtdlpOpts): Promise<void> {
@@ -593,25 +601,34 @@ async function runManagedDownloads(
skippedCount++;
} else if (outcome.status === "failed") {
failedCount++;
- if (abortOnError && !firstFailure) {
- // downloadOneManaged classifies against the full stderr tail; fall
- // back to the truncated attempt error only if it's somehow absent.
- const lastAttempt = outcome.attempts[outcome.attempts.length - 1];
- const failureClass =
- outcome.failureClass ??
- classifyDownloadFailure(
- lastAttempt?.error ?? "",
- lastAttempt?.availabilityClass,
- );
- // Per-video failures (private, members-only, age-gated, deleted)
- // are properties of that video, not signals the batch should
- // stop. Only abort on rate-limit / network / unclassified-fatal.
- if (failureClass !== "per_video") {
- firstFailure = new Error(
- `Managed download aborted (${failureClass}) for ${url}`,
- );
+ // downloadOneManaged classifies against the full stderr tail; fall
+ // back to the truncated attempt error only if it's somehow absent.
+ const lastAttempt = outcome.attempts[outcome.attempts.length - 1];
+ const failureClass =
+ outcome.failureClass ??
+ classifyDownloadFailure(
+ lastAttempt?.error ?? "",
+ lastAttempt?.availabilityClass,
+ );
+ // A rate-limit / network failure pauses the WHOLE platform via the
+ // shared cooldown, so the auto-download runner (and a later manual
+ // sync) back off too. Best-effort: never let bookkeeping break the
+ // download loop. Fires regardless of abortOnError / firstFailure.
+ if (failureClass === "rate_limit" || failureClass === "network") {
+ try {
+ await opts.onPlatformBackoff?.(failureClass);
+ } catch {
+ /* shared-state write is best-effort */
}
}
+ // Per-video failures (private, members-only, age-gated, deleted) are
+ // properties of that video, not signals the batch should stop. Only
+ // abort on rate-limit / network / unclassified-fatal.
+ if (abortOnError && !firstFailure && failureClass !== "per_video") {
+ firstFailure = new Error(
+ `Managed download aborted (${failureClass}) for ${url}`,
+ );
+ }
}
processedCount++;
const isLast = processedCount >= urls.length;
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **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`.
- **New Archilyzer homepage (hub): a standalone marketing/docs site, configurable here.** A fourth workspace package, `homepage`, builds a single instance-level static site (`output: "export"`) that sits above the per-content export sites — for info/docs pages and cross-site charts that don't belong on any one content site. It's managed from the new **Manage → Homepage** page: edit branding (title/header/description/tagline/public URL/Cloudflare project) and author **Markdown pages** (a slug + nav label + body, with a live markdown-to-jsx preview), which the hub renders server-side at build (the `index` page is the home body; every other slug gets a `/<slug>` route). Page content lives in the data dir (`sites/_homepage/`, a reserved id `listSiteIds()` ignores), so copy changes need no code deploy. The hub's **/stats** dashboard charts downloads/transcriptions completed across **all** content sites, leading with a per-site breakdown — the chart engine gains a **Site** grouping option (`groupBy: "site"`) that fans each video out to every site exposing its channel, resolved through a channel→sites map the hub supplies; combined whole-pool totals remain one series. Build with `pnpm build:homepage` (a `compose-homepage` step stages whole-pool stats + the channel→sites map ahead of `next build`). See `common/lib/{homepage,homepagePages}.ts`, `common/bin/compose-homepage.ts`, `common/components/charts/channelSites.tsx`, the `homepage/` package, and `editor/app/homepage/*`.
diff --git a/editor/app/actionable/components/InlineActionButton.tsx b/editor/app/actionable/components/InlineActionButton.tsx
@@ -65,7 +65,13 @@ async function runAction(variant: Variant): Promise<StreamActionResult> {
if (result && "error" in result) {
return { ok: false, error: result.error };
}
- return { ok: true, jobId: "", stream: new ReadableStream<string>() };
+ // refreshReport runs no managed job, so synthesize an already-complete result.
+ return {
+ ok: true,
+ jobId: "",
+ stream: new ReadableStream<string>(),
+ done: Promise.resolve({ status: "done", jobId: "" } as const),
+ };
}
export function InlineActionButton({
diff --git a/editor/app/auto-queue/components/AutoQueueView.tsx b/editor/app/auto-queue/components/AutoQueueView.tsx
@@ -2,7 +2,11 @@
import { useCallback, useEffect, useRef, useState } from "react";
import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState";
-import type { AutoQueueStatusPayload, AutoQueueKindStatus } from "../status";
+import type {
+ AutoQueueStatusPayload,
+ AutoQueueKindStatus,
+ PlatformCooldownView,
+} from "../status";
import { PolicyTreeEditor } from "./PolicyTreeEditor";
// Live panel for the auto-queue runners: one section per kind (auto-transcribe /
@@ -162,6 +166,10 @@ function KindPanel({
</span>
</div>
+ {status.cooldowns.length > 0 && (
+ <CooldownBanner cooldowns={status.cooldowns} />
+ )}
+
<PolicyTreeEditor
kind={kind}
initialEnabled={status.policy.enabled}
@@ -255,3 +263,51 @@ function formatClock(ms: number): string {
minute: "2-digit",
});
}
+
+// Wall-clock that re-renders each second, starting null so SSR and the first
+// client render agree (no hydration mismatch from Date.now()).
+function useNow(): number | null {
+ const [now, setNow] = useState<number | null>(null);
+ useEffect(() => {
+ setNow(Date.now());
+ const id = setInterval(() => setNow(Date.now()), 1000);
+ return () => clearInterval(id);
+ }, []);
+ return now;
+}
+
+// Platforms paused by a rate-limit/network backoff. A manual Sync on one of
+// these is refused until it lapses; the runner skips it meanwhile.
+function CooldownBanner({ cooldowns }: { cooldowns: PlatformCooldownView[] }) {
+ const now = useNow();
+ return (
+ <div className="flex flex-col gap-1 rounded-md border border-amber-300 bg-amber-50 px-3 py-2 text-sm dark:border-amber-800 dark:bg-amber-950">
+ <span className="font-medium text-amber-800 dark:text-amber-200">
+ Rate-limit cooldown — auto-download and manual sync are paused on:
+ </span>
+ <ul className="flex flex-wrap gap-x-4 gap-y-1 text-amber-700 dark:text-amber-300">
+ {cooldowns.map((c) => {
+ const secs =
+ now === null
+ ? null
+ : Math.max(0, Math.ceil((c.untilMs - now) / 1000));
+ return (
+ <li key={c.platform} className="tabular-nums">
+ <span className="font-mono">{c.platform}</span>
+ {secs !== null && (
+ <> — {formatCooldown(secs)} left (attempt {c.fails})</>
+ )}
+ </li>
+ );
+ })}
+ </ul>
+ </div>
+ );
+}
+
+function formatCooldown(secs: number): string {
+ if (secs < 60) return `${secs}s`;
+ const m = Math.floor(secs / 60);
+ const s = secs % 60;
+ return s === 0 ? `${m}m` : `${m}m ${s}s`;
+}
diff --git a/editor/app/auto-queue/status.ts b/editor/app/auto-queue/status.ts
@@ -18,12 +18,23 @@ import type { AutoQueuePolicy } from "yt-dlp-transcript-common/jobs/autoQueuePol
// show "cornbreadman: 12 pending" and which leaf is currently being serviced.
// Backs a passive UI poll, like /api/scheduler/status.
+// A platform currently in a rate-limit (429) / network cooldown. `untilMs` is
+// the epoch ms the cooldown lapses; the client renders remaining time itself so
+// it counts down live. Only the download kind ever has these.
+export type PlatformCooldownView = {
+ platform: string;
+ untilMs: number;
+ fails: number;
+};
+
export type AutoQueueKindStatus = {
kind: AutoQueueKind;
policy: AutoQueuePolicy;
runner: AutoRunnerStatus;
pendingByLeaf: Record<string, number>;
picks: AutoQueuePick[];
+ // Platforms paused by a rate-limit/network backoff (download kind only).
+ cooldowns: PlatformCooldownView[];
};
export type AutoQueueStatusPayload = {
@@ -37,7 +48,24 @@ async function buildKind(kind: AutoQueueKind): Promise<AutoQueueKindStatus> {
const runner = getAutoRunnerStatus(kind);
const state = await readAutoQueueState(paths);
const pendingByLeaf = await computeLeafPendingCounts(kind, paths);
- return { kind, policy, runner, pendingByLeaf, picks: state[kind].picks };
+ const now = Date.now();
+ // Surface platforms still inside their cooldown window so the operator can see
+ // why an otherwise-pending platform isn't being serviced (and a manual sync on
+ // it would be refused). A manual sync/download 429 writes into this same map.
+ const cooldowns: PlatformCooldownView[] = Object.entries(
+ state[kind].platformBackoff,
+ )
+ .filter(([, e]) => e.until > now)
+ .map(([platform, e]) => ({ platform, untilMs: e.until, fails: e.fails }))
+ .sort((a, b) => b.untilMs - a.untilMs);
+ return {
+ kind,
+ policy,
+ runner,
+ pendingByLeaf,
+ picks: state[kind].picks,
+ cooldowns,
+ };
}
export async function buildAutoQueueStatusPayload(): Promise<AutoQueueStatusPayload> {
diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts
@@ -10,6 +10,11 @@ import {
downloadQueueKey,
resolveQueueKey,
} from "yt-dlp-transcript-common/lib/queueKeys";
+import { detectPlatform } from "yt-dlp-transcript-common/lib/platform";
+import {
+ platformCooldownRemainingMs,
+ recordDownloadBackoff,
+} from "yt-dlp-transcript-common/jobs/downloadBackoff";
import {
countNotYetDownloaded,
readChannelConfig,
@@ -76,6 +81,27 @@ async function runPipelineAction(
if (!channelConfig.url) {
return { ok: false, error: "Channel has no `url` configured" };
}
+ // Platform key shared with the auto-download runner (detectPlatform(url) ??
+ // "unknown") — drives the 429 cooldown both honor.
+ const platform = detectPlatform(channelConfig.url) ?? "unknown";
+ // A manually clicked Sync RESPECTS the per-platform rate-limit cooldown: if
+ // auto-download (or a prior sync) hit a 429, refuse rather than re-storm the
+ // source. Neutral (info) result so the UI shows it as a notice, not an error.
+ // Other download modes still run (and will write their own backoff on a 429).
+ if (mode === "sync") {
+ const remainingMs = await platformCooldownRemainingMs(platform, paths);
+ if (remainingMs > 0) {
+ const secs = Math.ceil(remainingMs / 1000);
+ return {
+ ok: false,
+ info: true,
+ error:
+ `${platform} is in a rate-limit cooldown (${secs}s remaining). ` +
+ `Auto-download or a prior sync hit HTTP 429 — sync will run once the ` +
+ `cooldown lapses.`,
+ };
+ }
+ }
// store-playlist only fetches the video list (no media written), so it is not
// gated; every other mode downloads audio/subtitles into transcriptsDir.
if (mode !== "store-playlist") {
@@ -136,6 +162,9 @@ async function runPipelineAction(
handlingOverride: options?.handlingOverride,
setProgress,
progressBaseline,
+ // On a 429/network failure, record the shared per-platform cooldown so
+ // the auto-download runner (and a later manual sync) back off too.
+ onPlatformBackoff: () => recordDownloadBackoff(platform, paths),
});
// The channel report (snapshot) is regenerated automatically after this
// job finishes, via the global debounced scheduler hooked into
diff --git a/editor/app/jobs/[id]/page.tsx b/editor/app/jobs/[id]/page.tsx
@@ -6,6 +6,7 @@ import { listAllJobs } from "yt-dlp-transcript-common/jobs/listJobs";
import { getPaths } from "yt-dlp-transcript-common/lib/paths";
import { JobLogTail } from "./components/JobLogTail";
import { BookmarkJobButton } from "../components/BookmarkJobButton";
+import { jobKindLabel } from "../jobKindLabels";
export const dynamic = "force-dynamic";
@@ -23,9 +24,8 @@ export async function generateMetadata({
? `${job.channelSlug}/${job.videoId}`
: job.channelSlug
: null;
- const subject = scope
- ? `${job.kind ?? "job"} (${scope})`
- : (job.kind ?? "job");
+ const kindLabel = job.kind ? jobKindLabel(job.kind) : "job";
+ const subject = scope ? `${kindLabel} (${scope})` : kindLabel;
return { title: `${subject} — Job` };
}
@@ -51,7 +51,9 @@ export default async function JobDetailPage({
<span className="font-mono">{id}</span>
</div>
<div className="flex items-center justify-between gap-3">
- <h1 className="text-2xl font-semibold">{job.kind ?? "job"}</h1>
+ <h1 className="text-2xl font-semibold" title={job.kind ?? undefined}>
+ {job.kind ? jobKindLabel(job.kind) : "job"}
+ </h1>
{job.bookmarkable && <BookmarkJobButton jobId={id} />}
</div>
{job.channelSlug && (
diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts
@@ -144,6 +144,7 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
drainable:
j.status === "running" && isDrainableKind(j.kind) && j.draining !== true,
bookmarkable: Boolean(j.spec),
+ background: j.background === true,
}));
const channels = channelSlugs.map((slug) => ({
diff --git a/editor/app/jobs/components/BookmarksList.tsx b/editor/app/jobs/components/BookmarksList.tsx
@@ -12,6 +12,7 @@ import {
renameBookmarkAction,
runBookmarkAction,
} from "../bookmarkActions";
+import { jobKindLabel } from "../jobKindLabels";
// The saved-jobs panel rendered on /jobs and /jobs/active. Each bookmark has a
// "Run again" button (streams the re-launched job's log), an inline rename, a
@@ -143,7 +144,9 @@ function BookmarkRow({
</button>
</>
)}
- <span className="font-mono text-xs text-zinc-500">{spec.kind}</span>
+ <span className="text-xs text-zinc-500" title={spec.kind}>
+ {jobKindLabel(spec.kind)}
+ </span>
{spec.bucket && (
<span className="font-mono text-xs text-zinc-500">
bucket {spec.bucket}
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 { jobKindLabel } from "../jobKindLabels";
import {
clearJobsFilters,
defaultJobsFilters,
@@ -127,8 +128,9 @@ export function JobsTable({ jobs }: { jobs: JobListEntry[] }) {
? "border-zinc-300 dark:border-zinc-700 text-zinc-400 dark:text-zinc-600 line-through"
: "border-zinc-400 dark:border-zinc-600 bg-zinc-100 dark:bg-zinc-800 text-zinc-800 dark:text-zinc-200"
}`}
+ title={kind}
>
- {kind}
+ {jobKindLabel(kind)}
</button>
);
})}
@@ -220,7 +222,9 @@ export function JobsTable({ jobs }: { jobs: JobListEntry[] }) {
{j.id}
</Link>
</td>
- <td className="px-3 py-2">{j.kind ?? "—"}</td>
+ <td className="px-3 py-2" title={j.kind ?? undefined}>
+ {j.kind ? jobKindLabel(j.kind) : "—"}
+ </td>
<td className="px-3 py-2 font-mono text-xs">
{j.channelSlug ? (
<Link
diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx
@@ -50,6 +50,10 @@ export type RunningJobsListItem = {
drainable?: boolean;
// True when the job carries a replay descriptor and can be bookmarked.
bookmarkable?: boolean;
+ // Background work (e.g. an auto-download unit) queues BEHIND a manual job on
+ // 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;
};
type Props = {
@@ -107,6 +111,14 @@ function JobRow({
<span className="text-xs font-medium" title={job.kind}>
{jobKindLabel(job.kind)}
</span>
+ {job.background && (
+ <span
+ className="text-[10px] uppercase tracking-wide px-1.5 py-0.5 rounded bg-amber-100 text-amber-800 dark:bg-amber-900/60 dark:text-amber-200"
+ title="Background work — yields to a manually-triggered sync/download on the same platform"
+ >
+ auto
+ </span>
+ )}
<Link
href={`/jobs/${job.id}`}
className="font-mono text-xs underline hover:text-zinc-900 dark:hover:text-zinc-100"
@@ -132,6 +144,7 @@ function JobRow({
{job.status === "queued" && (
<span className="text-xs text-zinc-500">
queue <code className="font-mono">{job.queueKey || "—"}</code>
+ {job.background && " · waiting behind a manual job"}
</span>
)}
<div className="ml-auto flex items-center gap-2">
diff --git a/editor/app/jobs/jobKindLabels.ts b/editor/app/jobs/jobKindLabels.ts
@@ -3,7 +3,10 @@
// 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",
- "auto-download": "Auto-download",
+ // 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",
diff --git a/editor/app/widget/components/MonitorWidget.tsx b/editor/app/widget/components/MonitorWidget.tsx
@@ -7,6 +7,7 @@ import type {
DiskStatusView,
} from "../../jobs/active/buildActiveJobs";
import type { RunningJobsListItem } from "../../jobs/components/RunningJobsList";
+import { jobKindLabel } from "../../jobs/jobKindLabels";
import type { WorkersPayload, WorkerView } from "../../workers/components/WorkersView";
import type { WidgetConfig } from "../lib/config";
@@ -308,7 +309,17 @@ function JobRow({
>
{job.status}
</span>
- <span className="font-mono">{job.kind}</span>
+ <span className="font-mono" title={job.kind}>
+ {jobKindLabel(job.kind)}
+ </span>
+ {job.background && (
+ <span
+ className="uppercase tracking-wide px-1 py-0.5 rounded text-[9px] bg-amber-100 text-amber-800 dark:bg-amber-900/60 dark:text-amber-200"
+ title="Background work — yields to a manual sync/download"
+ >
+ auto
+ </span>
+ )}
{job.channelSlug && (
<span className="font-mono text-zinc-500 truncate max-w-[12rem]">
{job.channelSlug}
diff --git a/editor/e2e/queues.spec.ts b/editor/e2e/queues.spec.ts
@@ -1,5 +1,6 @@
+import { mkdir, writeFile } from "node:fs/promises";
import { test, expect } from "@playwright/test";
-import { resetData } from "./helpers";
+import { resetData, resolvePath } from "./helpers";
async function setQueueViaCustom(
page: import("@playwright/test").Page,
@@ -30,6 +31,43 @@ test("default queue is platform:<platform>", async ({ page }) => {
await expect(firstRow).toContainText("platform:youtube");
});
+test("Sync is refused while its platform is in a rate-limit cooldown", async ({
+ page,
+}) => {
+ await resetData("slow-pipeline-channel");
+ // Seed a live download cooldown for youtube (the slow-channel platform) — as
+ // if auto-download or a prior sync had just hit a 429. A manual Sync must
+ // refuse rather than re-storm the source.
+ const until = Date.now() + 5 * 60_000;
+ const stateFile = resolvePath("test-transcripts/.auto-queue/state.json");
+ await mkdir(resolvePath("test-transcripts/.auto-queue"), { recursive: true });
+ await writeFile(
+ stateFile,
+ JSON.stringify({
+ transcription: { runtime: {}, picks: [], platformBackoff: {} },
+ download: {
+ runtime: {},
+ picks: [],
+ platformBackoff: { youtube: { until, fails: 1 } },
+ },
+ }),
+ );
+
+ await page.goto("/channels/slow-channel");
+ await page.getByRole("button", { name: "Sync" }).click();
+ // The refusal is a neutral notice (info), naming the platform + cooldown.
+ await expect(page.getByLabel("Sync notice")).toContainText(
+ /rate-limit cooldown/i,
+ { timeout: 10_000 },
+ );
+ // No job was created (refused before enqueue) — the /jobs table has no
+ // slow-channel row.
+ await page.goto("/jobs");
+ await expect(
+ page.getByRole("row").filter({ hasText: "slow-channel" }),
+ ).toHaveCount(0);
+});
+
test("unknown platform infers its queue from the URL domain", async ({
page,
}) => {