commit 5a9cebd38ea4acd082d4b55d5307bbfb3623a515
parent 935569938878d3ecc5ae8d3c90286f239ec367cf
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 23 Jun 2026 14:36:23 -0400
Merge feat/auto-queue-platform-backoff: auto-download shares the per-platform Sync queue
Diffstat:
28 files changed, 1193 insertions(+), 93 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -25,6 +25,15 @@ import {
recordPick,
writeAutoQueueState,
} from "../jobs/autoQueueState";
+import {
+ clearBackoff,
+ isCoolingDown,
+ nextBackoff,
+ 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";
@@ -46,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;
@@ -219,6 +234,9 @@ async function runLoop(
const state = await readAutoQueueState(paths);
const kindState = state[kind];
const runtime = kindState.runtime;
+ // Drop long-lapsed platform cooldowns on boot; entries still in (or recently
+ // out of) their window are kept so an Odysee 429 cooldown survives a restart.
+ pruneExpired(kindState.platformBackoff, Date.now());
const completed = new Set<string>();
const completedOrder: string[] = [];
let counters = { transcribed: 0, failed: 0, skipped: 0 };
@@ -284,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;
@@ -343,17 +366,35 @@ async function runLoop(
const exclude = new Set<string>([...live.inFlight.keys(), ...completed]);
removeIds(pending, exclude);
// Download only: drop videos whose platform already has an in-flight
- // download, so a busy platform yields to the next-priority free one.
+ // 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") {
- const busy = new Set<string>();
+ // 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) busy.add(pf);
+ if (n >= PER_PLATFORM_CAP) skip.add(pf);
+ }
+ for (const pf of Object.keys(kindState.platformBackoff)) {
+ if (isCoolingDown(kindState.platformBackoff, pf, now)) skip.add(pf);
}
- if (busy.size > 0) {
+ if (skip.size > 0) {
for (const leafId of Object.keys(pending)) {
pending[leafId] = pending[leafId].filter((id) => {
const slug = owner.get(id);
- return !(slug && busy.has(platformKey(slug, slugToPlatform)));
+ return !(slug && skip.has(platformKey(slug, slugToPlatform)));
});
}
}
@@ -373,26 +414,70 @@ async function runLoop(
continue;
}
- launchUnit({ kind, paths, pick, channelSlug, tracker, onLog, signal })
- .then((outcome) => {
- if (outcome === "transcribed") counters.transcribed++;
- else if (outcome === "failed") counters.failed++;
- else counters.skipped++;
+ const unitPlatform =
+ kind === "download" ? platformKey(channelSlug, slugToPlatform) : null;
+ let result: UnitResult = { outcome: "failed" };
+ launchUnit({
+ kind,
+ paths,
+ pick,
+ channelSlug,
+ tracker,
+ onLog,
+ signal,
+ onChildJob: (jid) => childJobIds.set(pick.videoId, jid),
+ })
+ .then((r) => {
+ result = r;
})
.catch((err) => {
- counters.failed++;
+ result = { outcome: "failed" };
onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`);
})
.finally(() => {
+ childJobIds.delete(pick.videoId);
live.inFlight.delete(pick.videoId);
for (const id of pick.path) {
live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1);
}
- if (kind === "download") {
- const pf = platformKey(channelSlug, slugToPlatform);
- platformInFlight.set(pf, Math.max(0, (platformInFlight.get(pf) ?? 0) - 1));
+ if (unitPlatform) {
+ platformInFlight.set(
+ unitPlatform,
+ Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1),
+ );
}
- markCompleted(pick.videoId);
+
+ if (result.outcome === "transcribed") counters.transcribed++;
+ else if (result.outcome === "failed") counters.failed++;
+ else counters.skipped++;
+
+ // Per-platform backoff: a rate-limit / network failure pauses the whole
+ // platform with an exponential cooldown, and the video is deliberately
+ // NOT marked completed so it's retried once the cooldown lapses. Any
+ // other outcome marks the video done-for-session; a success also clears
+ // the platform's cooldown.
+ const backoffHit =
+ unitPlatform !== null &&
+ (result.failureClass === "rate_limit" ||
+ result.failureClass === "network");
+ if (backoffHit && unitPlatform) {
+ const entry = nextBackoff(
+ kindState.platformBackoff[unitPlatform],
+ Date.now(),
+ );
+ kindState.platformBackoff[unitPlatform] = entry;
+ const secs = Math.round((entry.until - Date.now()) / 1000);
+ onLog(
+ `Auto-download: ${unitPlatform} ${result.failureClass} — backing off ${secs}s (attempt ${entry.fails}). ${pick.videoId} will retry after cooldown.`,
+ );
+ } else {
+ if (unitPlatform && result.outcome === "transcribed") {
+ clearBackoff(kindState.platformBackoff, unitPlatform);
+ }
+ markCompleted(pick.videoId);
+ }
+ // Persist so the cooldown (and reset) survive a restart.
+ persist();
// The runner job has no channelSlug, so the registry's own
// snapshot-on-finish hook doesn't fire — request a regen here so the
// channel's pending counts refresh and a restart sees fresh buckets.
@@ -409,9 +494,11 @@ async function runLoop(
startedAt: Date.now(),
});
for (const id of pick.path) live.active[id] = (live.active[id] ?? 0) + 1;
- if (kind === "download") {
- const pf = platformKey(channelSlug, slugToPlatform);
- platformInFlight.set(pf, (platformInFlight.get(pf) ?? 0) + 1);
+ if (unitPlatform) {
+ platformInFlight.set(
+ unitPlatform,
+ (platformInFlight.get(unitPlatform) ?? 0) + 1,
+ );
}
recordPick(kindState, {
at: Date.now(),
@@ -423,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(
@@ -439,14 +533,24 @@ 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
+// (the runner uses rate_limit/network to drive per-platform backoff); it's
+// derived from the returned outcome record via classifyDownloadFailure.
+type UnitResult = {
+ outcome: "transcribed" | "skipped" | "failed";
+ failureClass?: DownloadFailureClass;
};
-// Run a single unit of work. Returns the outcome for counters. Hard cancel
-// (signal) aborts an in-flight unit; drain is handled by the loop (it stops
-// launching new units), so the drain signal is intentionally NOT forwarded here.
-async function launchUnit(args: LaunchArgs): Promise<
- "transcribed" | "skipped" | "failed"
-> {
+// Run a single unit of work. Returns the outcome for counters/backoff. Hard
+// cancel (signal) aborts an in-flight unit; drain is handled by the loop (it
+// stops launching new units), so the drain signal is intentionally NOT
+// forwarded here.
+async function launchUnit(args: LaunchArgs): Promise<UnitResult> {
if (args.kind === "transcription") {
const res = await transcribeOneFromQueue({
paths: args.paths,
@@ -461,15 +565,19 @@ async function launchUnit(args: LaunchArgs): Promise<
onLog: args.onLog,
signal: args.signal,
});
- return res.outcome;
+ 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 "skipped";
+ if (!config) return { outcome: "skipped" };
const url = await findVideoSourceUrl(
args.paths,
args.channelSlug,
@@ -478,37 +586,77 @@ async function launchUnit(args: LaunchArgs): Promise<
);
if (!url) {
args.onLog(`Auto-download: no source URL for ${args.pick.videoId} — skipped`);
- return "skipped";
+ 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 {
- 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,
- });
- return "transcribed"; // generic "success" — feeds the done counter
- } catch (err) {
- if (args.signal.aborted || (err as Error)?.name === "AbortError") {
- return "skipped";
- }
- args.onLog(`Auto-download failed for ${args.pick.videoId}: ${String(err)}`);
- return "failed";
- } finally {
- task.end();
+ if (!res.ok) {
+ args.onLog(
+ `Auto-download: could not enqueue ${args.pick.videoId}: ${res.error}`,
+ );
+ return { outcome: "failed" };
}
+ 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/autoQueueState.test.ts b/common/jobs/autoQueueState.test.ts
@@ -0,0 +1,71 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { mkdtemp, rm, writeFile } from "node:fs/promises";
+import { tmpdir } from "node:os";
+import path from "node:path";
+import type { Paths } from "../lib/paths";
+import {
+ emptyAutoQueueState,
+ readAutoQueueState,
+ writeAutoQueueState,
+} from "./autoQueueState";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/autoQueueState.test.ts
+
+// Only autoQueueStateFile is touched by these functions, so a minimal Paths is fine.
+async function withPaths(
+ fn: (paths: Paths) => Promise<void>,
+): Promise<void> {
+ const dir = await mkdtemp(path.join(tmpdir(), "ttb-aqstate-"));
+ 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("platformBackoff round-trips through write/read", async () => {
+ await withPaths(async (paths) => {
+ const state = emptyAutoQueueState();
+ state.download.platformBackoff = {
+ odysee: { until: 123456, fails: 3 },
+ youtube: { until: 999, fails: 1 },
+ };
+ await writeAutoQueueState(paths, state);
+
+ const back = await readAutoQueueState(paths);
+ assert.deepEqual(back.download.platformBackoff, {
+ odysee: { until: 123456, fails: 3 },
+ youtube: { until: 999, fails: 1 },
+ });
+ assert.deepEqual(back.transcription.platformBackoff, {});
+ });
+});
+
+test("missing file yields empty state with empty platformBackoff", async () => {
+ await withPaths(async (paths) => {
+ const back = await readAutoQueueState(paths);
+ assert.deepEqual(back.download.platformBackoff, {});
+ assert.deepEqual(back.transcription.platformBackoff, {});
+ });
+});
+
+test("corrupt platformBackoff field coerces to empty without throwing", async () => {
+ await withPaths(async (paths) => {
+ // Write a valid state first so the .auto-queue dir exists, then overwrite
+ // the file with a hand-corrupted shape.
+ await writeAutoQueueState(paths, emptyAutoQueueState());
+ await writeFile(
+ paths.autoQueueStateFile,
+ JSON.stringify({
+ download: { platformBackoff: { odysee: { until: "nope" }, junk: 5 } },
+ }),
+ );
+
+ const back = await readAutoQueueState(paths);
+ assert.deepEqual(back.download.platformBackoff, {});
+ });
+});
diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts
@@ -5,6 +5,10 @@ import {
type AutoQueueRuntime,
emptyAutoQueueRuntime,
} from "./autoQueuePolicy";
+import {
+ type PlatformBackoffState,
+ coercePlatformBackoff,
+} from "./platformBackoff";
// Persistent fairness state for the auto-queue runners. Unlike in-flight worker
// counts (which are zero by definition after a restart and live only on the
@@ -30,6 +34,10 @@ export type AutoQueueKindState = {
runtime: AutoQueueRuntime;
// Newest-first, bounded to AUTO_QUEUE_PICK_LOG_LIMIT entries.
picks: AutoQueuePick[];
+ // Per-platform download cooldowns (rate-limit/network backoff). Only used by
+ // the "download" kind; empty for transcription. Persisted so an Odysee 429
+ // cooldown survives a server restart instead of re-storming on boot.
+ platformBackoff: PlatformBackoffState;
};
export type AutoQueueState = {
@@ -40,7 +48,7 @@ export type AutoQueueState = {
export const AUTO_QUEUE_PICK_LOG_LIMIT = 50;
export function emptyAutoQueueKindState(): AutoQueueKindState {
- return { runtime: emptyAutoQueueRuntime(), picks: [] };
+ return { runtime: emptyAutoQueueRuntime(), picks: [], platformBackoff: {} };
}
export function emptyAutoQueueState(): AutoQueueState {
@@ -93,7 +101,11 @@ function coerceKindState(value: unknown): AutoQueueKindState {
.filter((p): p is AutoQueuePick => p !== null)
.slice(0, AUTO_QUEUE_PICK_LOG_LIMIT)
: [];
- return { runtime: coerceRuntime(r.runtime), picks };
+ return {
+ runtime: coerceRuntime(r.runtime),
+ picks,
+ platformBackoff: coercePlatformBackoff(r.platformBackoff),
+ };
}
// Read the state file, tolerating a missing/corrupt file by returning empty
@@ -123,6 +135,7 @@ export async function writeAutoQueueState(
const trim = (k: AutoQueueKindState): AutoQueueKindState => ({
runtime: k.runtime,
picks: k.picks.slice(0, AUTO_QUEUE_PICK_LOG_LIMIT),
+ platformBackoff: k.platformBackoff ?? {},
});
const out: AutoQueueState = {
transcription: trim(state.transcription),
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/platformBackoff.test.ts b/common/jobs/platformBackoff.test.ts
@@ -0,0 +1,93 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import {
+ BACKOFF_BASE_MS,
+ BACKOFF_MAX_MS,
+ type PlatformBackoffState,
+ clearBackoff,
+ coercePlatformBackoff,
+ isCoolingDown,
+ nextBackoff,
+ pruneExpired,
+} from "./platformBackoff";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/platformBackoff.test.ts
+
+// Deterministic jitter: no jitter (returns 0.5 -> factor 1.0).
+const noJitter = () => 0.5;
+
+test("nextBackoff doubles the cooldown across consecutive failures", () => {
+ const now = 1_000_000;
+ const first = nextBackoff(undefined, now, noJitter);
+ assert.equal(first.fails, 1);
+ assert.equal(first.until - now, BACKOFF_BASE_MS); // 60s
+
+ const second = nextBackoff(first, now, noJitter);
+ assert.equal(second.fails, 2);
+ assert.equal(second.until - now, BACKOFF_BASE_MS * 2); // 120s
+
+ const third = nextBackoff(second, now, noJitter);
+ assert.equal(third.fails, 3);
+ assert.equal(third.until - now, BACKOFF_BASE_MS * 4); // 240s
+});
+
+test("nextBackoff caps the cooldown at BACKOFF_MAX_MS", () => {
+ const now = 0;
+ let entry = nextBackoff(undefined, now, noJitter);
+ for (let i = 0; i < 20; i++) entry = nextBackoff(entry, now, noJitter);
+ assert.equal(entry.until - now, BACKOFF_MAX_MS);
+ assert.ok(entry.fails > 1);
+});
+
+test("nextBackoff applies bounded jitter (±10%)", () => {
+ const now = 0;
+ const low = nextBackoff(undefined, now, () => 0); // factor 0.9
+ const high = nextBackoff(undefined, now, () => 1); // factor 1.1
+ assert.equal(low.until, Math.round(BACKOFF_BASE_MS * 0.9));
+ assert.equal(high.until, Math.round(BACKOFF_BASE_MS * 1.1));
+});
+
+test("isCoolingDown reflects the until window", () => {
+ const state: PlatformBackoffState = { odysee: { until: 500, fails: 1 } };
+ assert.equal(isCoolingDown(state, "odysee", 499), true);
+ assert.equal(isCoolingDown(state, "odysee", 500), false); // boundary: until is exclusive
+ assert.equal(isCoolingDown(state, "odysee", 501), false);
+ assert.equal(isCoolingDown(state, "youtube", 0), false); // absent platform
+});
+
+test("clearBackoff removes a platform (success resets escalation)", () => {
+ const state: PlatformBackoffState = { odysee: { until: 999, fails: 3 } };
+ clearBackoff(state, "odysee");
+ assert.deepEqual(state, {});
+ // Next failure starts over at attempt 1.
+ const fresh = nextBackoff(state["odysee"], 0, noJitter);
+ assert.equal(fresh.fails, 1);
+});
+
+test("pruneExpired drops only long-lapsed entries, keeps active/recent ones", () => {
+ const now = 10_000_000;
+ const state: PlatformBackoffState = {
+ active: { until: now + 1000, fails: 1 }, // still cooling down
+ recent: { until: now - 1000, fails: 2 }, // lapsed but within retain window
+ old: { until: now - BACKOFF_MAX_MS - 1, fails: 5 }, // long lapsed
+ };
+ pruneExpired(state, now);
+ assert.ok(state.active, "active kept");
+ assert.ok(state.recent, "recently-lapsed kept for escalation memory");
+ assert.equal(state.old, undefined, "long-lapsed dropped");
+});
+
+test("coercePlatformBackoff tolerates corrupt/missing shapes", () => {
+ assert.deepEqual(coercePlatformBackoff(undefined), {});
+ assert.deepEqual(coercePlatformBackoff("nope"), {});
+ assert.deepEqual(coercePlatformBackoff({ odysee: { until: 5 } }), {}); // missing fails
+ assert.deepEqual(coercePlatformBackoff({ odysee: { until: "x", fails: 1 } }), {});
+ assert.deepEqual(
+ coercePlatformBackoff({
+ odysee: { until: 5, fails: 2 },
+ junk: 42,
+ youtube: { until: 7, fails: 1, extra: "ignored" },
+ }),
+ { odysee: { until: 5, fails: 2 }, youtube: { until: 7, fails: 1 } },
+ );
+});
diff --git a/common/jobs/platformBackoff.ts b/common/jobs/platformBackoff.ts
@@ -0,0 +1,94 @@
+// Per-platform exponential backoff for the auto-download runner. Pure logic +
+// types, no I/O — so it's unit-testable and can be persisted by
+// autoQueueState.ts and driven by autoRunner.ts.
+//
+// When a download fails with a rate-limit (HTTP 429 / throttle) or network
+// error, the runner must stop launching downloads for THAT platform for a
+// while instead of immediately re-picking the next video and hammering the
+// source again (the original Odysee 429 storm). Each consecutive failure
+// doubles the cooldown up to a cap; a successful download resets it.
+
+export type PlatformBackoffEntry = {
+ // Epoch ms: downloads for this platform are paused until this instant.
+ until: number;
+ // Consecutive backoff-triggering failures, drives the exponential delay.
+ fails: number;
+};
+
+// Keyed by platform name (e.g. "odysee", "youtube", or "unknown").
+export type PlatformBackoffState = Record<string, PlatformBackoffEntry>;
+
+// First cooldown after a single failure; doubles from here.
+export const BACKOFF_BASE_MS = 60_000; // 1 minute
+// Ceiling on the cooldown — long rate limits shouldn't park a platform forever.
+export const BACKOFF_MAX_MS = 30 * 60_000; // 30 minutes
+// ±10% jitter so multiple platforms recovering at once don't sync up.
+const JITTER = 0.1;
+
+// Compute the next backoff entry for a platform given its previous one (or
+// undefined for the first failure). `now` and `rand` are injected for
+// deterministic tests; `rand` defaults to Math.random and returns [0, 1).
+export function nextBackoff(
+ prev: PlatformBackoffEntry | undefined,
+ now: number,
+ rand: () => number = Math.random,
+): PlatformBackoffEntry {
+ const fails = (prev?.fails ?? 0) + 1;
+ const base = Math.min(BACKOFF_BASE_MS * 2 ** (fails - 1), BACKOFF_MAX_MS);
+ const jitter = 1 + (rand() * 2 - 1) * JITTER;
+ return { until: now + Math.round(base * jitter), fails };
+}
+
+// True when the platform is currently in a cooldown window.
+export function isCoolingDown(
+ state: PlatformBackoffState,
+ platform: string,
+ now: number,
+): boolean {
+ const entry = state[platform];
+ return entry !== undefined && entry.until > now;
+}
+
+// Clear a platform's backoff (called on a successful download). Mutates in place.
+export function clearBackoff(
+ state: PlatformBackoffState,
+ platform: string,
+): void {
+ delete state[platform];
+}
+
+// Drop entries whose cooldown lapsed long enough ago that keeping their `fails`
+// count for escalation memory no longer matters. Mutates in place, keeping the
+// map bounded across a long-running server. We retain a lapsed entry for one
+// full BACKOFF_MAX_MS window so a platform that recovers and immediately
+// re-fails keeps escalating rather than restarting at the base delay.
+export function pruneExpired(
+ state: PlatformBackoffState,
+ now: number,
+ retainMs: number = BACKOFF_MAX_MS,
+): void {
+ for (const [platform, entry] of Object.entries(state)) {
+ if (entry.until + retainMs < now) delete state[platform];
+ }
+}
+
+// Defensive coercion for the persisted shape — any stored value that doesn't
+// match is dropped, so a hand-edited/corrupt state file can't crash the runner.
+// Mirrors coerceRuntime/coercePick in autoQueueState.ts.
+export function coercePlatformBackoff(value: unknown): PlatformBackoffState {
+ const out: PlatformBackoffState = {};
+ if (!value || typeof value !== "object") return out;
+ for (const [platform, raw] of Object.entries(value as Record<string, unknown>)) {
+ if (!raw || typeof raw !== "object") continue;
+ const r = raw as Record<string, unknown>;
+ if (
+ typeof r.until === "number" &&
+ Number.isFinite(r.until) &&
+ typeof r.fails === "number" &&
+ Number.isFinite(r.fails)
+ ) {
+ out[platform] = { until: r.until, fails: r.fails };
+ }
+ }
+ return out;
+}
diff --git a/common/jobs/progressParsers.test.ts b/common/jobs/progressParsers.test.ts
@@ -0,0 +1,61 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { createDownloadProgressParser } from "./progressParsers";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/progressParsers.test.ts
+
+// Subtitle-only (youtube --skip-download) progress lines, as seen in real job
+// logs: total=NA while downloading, a real total only on the finished line.
+const DL_NA = "DLOM_PROGRESS status=downloading downloaded=7168 total=NA frag=NA/NA speed=3442662 eta=NA";
+const DL_FIN = "DLOM_PROGRESS status=finished downloaded=177273 total=177273 frag=NA/NA speed=745217 eta=NA";
+
+test("subtitle-only download advances one step per completed track", () => {
+ const p = createDownloadProgressParser();
+ // Without a total, byte fraction can't move — but once the track count is
+ // known the bar steps per finished track.
+ assert.equal(p.feed("[info] vid: Downloading subtitles: en-orig, en"), null);
+
+ let u = p.feed(DL_NA); // first track, downloading
+ assert.equal(u?.fraction, 0);
+ assert.match(u?.detail ?? "", /subs 0\/2/);
+
+ u = p.feed(DL_FIN); // first track done
+ assert.equal(u?.fraction, 0.5);
+ assert.match(u?.detail ?? "", /subs 1\/2/);
+
+ u = p.feed(DL_NA); // second track, downloading
+ assert.equal(u?.fraction, 0.5);
+
+ u = p.feed(DL_FIN); // second track done
+ assert.equal(u?.fraction, 1);
+ assert.match(u?.detail ?? "", /subs 2\/2/);
+});
+
+test("without the subtitle announcement, behaviour is unchanged (jumps at finish)", () => {
+ const p = createDownloadProgressParser();
+ const downloading = p.feed(DL_NA);
+ // total=NA, no track count → no fraction, just whatever detail it can build.
+ assert.equal(downloading?.fraction, undefined);
+ const finished = p.feed(DL_FIN);
+ assert.equal(finished?.fraction, 1); // finished line still reports 100%
+});
+
+test("a real media download is byte-based, not stepped, even after a subs announce", () => {
+ const p = createDownloadProgressParser();
+ // youtube no-subs fallback: subs announced, then it switches to an audio
+ // download that DOES report a byte total — subtitle-step mode must release.
+ p.feed("[info] vid: Downloading subtitles: en-orig, en");
+ const u = p.feed(
+ "DLOM_PROGRESS status=downloading downloaded=500 total=1000 frag=NA/NA speed=1000 eta=1",
+ );
+ assert.equal(u?.fraction, 0.5); // 500/1000, not a subtitle step
+ assert.doesNotMatch(u?.detail ?? "", /subs /);
+});
+
+test("media download fraction is byte-based from the start", () => {
+ const p = createDownloadProgressParser();
+ const u = p.feed(
+ "DLOM_PROGRESS status=downloading downloaded=250 total=1000 frag=NA/NA speed=1000 eta=3",
+ );
+ assert.equal(u?.fraction, 0.25);
+});
diff --git a/common/jobs/progressParsers.ts b/common/jobs/progressParsers.ts
@@ -190,6 +190,16 @@ export function createDownloadProgressParser(): {
let lastProbeMs: number | undefined;
let diffMs = 0; // smoothed per-probe duration increase
let intervalSec: number | undefined;
+ // Subtitle-only downloads (youtube handling runs yt-dlp with --skip-download)
+ // fetch one or more tiny .vtt files. yt-dlp can't report a byte total for
+ // them, so DLOM_PROGRESS carries total=NA while downloading and the per-byte
+ // fraction never moves. Instead, count the announced subtitle tracks and
+ // advance the bar one step per completed track. `subsTotal > 0` means we're in
+ // this mode; it's cleared the moment a real media download starts (a
+ // "downloading" line with a concrete byte total — e.g. a youtube no-subs
+ // fallback that switches to an audio download).
+ let subsTotal = 0;
+ let subsDone = 0;
const adjustEtaSeconds = (etaSeconds: number): number => {
if (lastProbeMs === undefined || !intervalSec || intervalSec <= 0) {
@@ -242,6 +252,50 @@ export function createDownloadProgressParser(): {
}
return null;
}
+
+ // yt-dlp announces the subtitle tracks it's about to fetch, e.g.
+ // [info] <id>: Downloading subtitles: en-orig, en
+ // Use the count as the denominator for the step-progress bar below.
+ const subsAnnounce = line.match(/Downloading subtitles:\s*(.+?)\s*$/);
+ if (subsAnnounce) {
+ const tracks = subsAnnounce[1]
+ .split(",")
+ .map((t) => t.trim())
+ .filter(Boolean);
+ if (tracks.length > 0) {
+ subsTotal = tracks.length;
+ subsDone = 0;
+ }
+ return null;
+ }
+
+ if (line.includes(DOWNLOAD_PROGRESS_MARKER)) {
+ const field = (key: string): string | undefined =>
+ line.match(new RegExp(`\\b${key}=(\\S+)`))?.[1];
+ const status = field("status");
+ const total = num(field("total"));
+ // A "downloading" line with a real byte total means actual media is
+ // streaming (not a subtitle file) — leave subtitle-step mode.
+ if (status === "downloading" && total !== undefined && total > 0) {
+ subsTotal = 0;
+ }
+ if (subsTotal > 0) {
+ if (status === "finished" && subsDone < subsTotal) subsDone += 1;
+ const structured = parseStructuredDownloadProgress(
+ line,
+ adjustEtaSeconds,
+ );
+ const out: ProgressUpdate = {
+ fraction: clamp01(subsDone / subsTotal),
+ };
+ const parts: string[] = [];
+ if (structured?.detail) parts.push(structured.detail);
+ parts.push(`subs ${subsDone}/${subsTotal}`);
+ out.detail = parts.join(" · ");
+ return out;
+ }
+ }
+
const structured = parseStructuredDownloadProgress(line, adjustEtaSeconds);
if (structured) return structured;
// Legacy "[download] …%" lines (non-structured) carry no numeric ETA we
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/lib/downloadOutcome.ts b/common/lib/downloadOutcome.ts
@@ -1,7 +1,7 @@
// Client-safe types and constants for the per-video download outcome sidecar.
// Mirrors availability.ts: server-only I/O lives in downloadOutcome-server.ts.
-import type { Availability } from "./availability";
+import type { Availability, DownloadFailureClass } from "./availability";
import type { ChannelHandling } from "./channelConfig";
export type DownloadOutcomeStatus =
@@ -63,6 +63,13 @@ export type DownloadOutcomeRecord = {
startedAt: string;
finishedAt: string;
attempts: DownloadAttempt[];
+ // Set on a failed download: the failure classified against the FULL stderr
+ // tail of the last attempt (not the truncated per-attempt `error`), so a
+ // rate-limit that yt-dlp logs as a WARNING before failing with a different
+ // final line (e.g. Odysee "HTTP Error 429" → "No video formats found") is
+ // still recognised. Drives the auto-runner's per-platform backoff. Undefined
+ // on success / skipped-filtered.
+ failureClass?: DownloadFailureClass;
fellBackToTranscribe?: boolean;
// Set when status is "skipped-filtered": which app-level filter declined the
// download and why. Recorded so the UI/log can explain the skip.
diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts
@@ -3,6 +3,7 @@ import { appendFile, mkdir, readdir, readFile } from "node:fs/promises";
import { createWriteStream, type WriteStream } from "node:fs";
import { execa } from "execa";
import {
+ classifyDownloadFailure,
parseUnavailableFromStderr,
type Availability,
} from "../lib/availability";
@@ -317,6 +318,10 @@ async function runManagedDownload(
let status: DownloadOutcomeStatus = "failed";
let fellBackToTranscribe = false;
let lastArchiveLine: string | null = null;
+ // Full (untruncated) stderr tail of the most recent attempt, so a failed
+ // download can be classified (rate_limit/network) against everything yt-dlp
+ // printed — not just the last 3 lines stored on the attempt record.
+ let lastFullTail = "";
// ---------- Attempt 0: metadata prefetch + app-level filters ----------
// Split the per-video download into a cheap metadata-only pass followed by
@@ -345,6 +350,7 @@ async function runManagedDownload(
opts.videoUrl,
];
const prefetchRes = await runOneYtdlp(opts, channelDir, prefetchArgs);
+ lastFullTail = prefetchRes.stderrTail;
const prefetchAvail = attemptSucceeded(prefetchRes.exitCode)
? undefined
: parseUnavailableFromStderr(prefetchRes.stderrTail);
@@ -472,6 +478,7 @@ async function runManagedDownload(
primaryRes = await runOneYtdlp(opts, channelDir, primaryArgs);
}
+ lastFullTail = primaryRes.stderrTail;
const primaryAvail = attemptSucceeded(primaryRes.exitCode)
? undefined
: parseUnavailableFromStderr(primaryRes.stderrTail);
@@ -520,6 +527,7 @@ async function runManagedDownload(
...sourceArgs(opts.videoUrl, infoJsonPath),
];
const retryRes = await runOneYtdlp(opts, channelDir, retryArgs);
+ lastFullTail = retryRes.stderrTail;
const retryAvail = attemptSucceeded(retryRes.exitCode)
? undefined
: parseUnavailableFromStderr(retryRes.stderrTail);
@@ -587,6 +595,7 @@ async function runManagedDownload(
),
];
const fallbackRes = await runOneYtdlp(opts, channelDir, fallbackArgs);
+ lastFullTail = fallbackRes.stderrTail;
const fallbackAvail = attemptSucceeded(fallbackRes.exitCode)
? undefined
: parseUnavailableFromStderr(fallbackRes.stderrTail);
@@ -654,6 +663,13 @@ async function runManagedDownload(
// ---------- Sidecar ----------
const finishedAt = new Date().toISOString();
+ // Classify a failed download against the FULL stderr tail of the last attempt
+ // so a rate-limit logged as a WARNING (then masked by a different final error)
+ // is still caught. Skipped for successes and filter-skips.
+ const isFailure = status === "failed" || status === "failed-corrupt-source";
+ const failureClass = isFailure
+ ? classifyDownloadFailure(lastFullTail, attempts.at(-1)?.availabilityClass)
+ : undefined;
const record: DownloadOutcomeRecord = {
videoId,
webpageUrl: opts.videoUrl,
@@ -661,6 +677,7 @@ async function runManagedDownload(
startedAt,
finishedAt,
attempts,
+ ...(failureClass ? { failureClass } : {}),
...(fellBackToTranscribe ? { fellBackToTranscribe: true } : {}),
};
// Only write the sidecar if we know which dir to put it in. If the very first
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,21 +601,34 @@ async function runManagedDownloads(
skippedCount++;
} else if (outcome.status === "failed") {
failedCount++;
- if (abortOnError && !firstFailure) {
- const lastAttempt = outcome.attempts[outcome.attempts.length - 1];
- const failureClass = classifyDownloadFailure(
+ // 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}`,
- );
+ // 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,9 @@
# 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/*`.
- **"Cut release" can now create the release commit for you.** After turning `## [Unreleased]` into a dated semver heading, cutting a release used to leave the changelog edit sitting in your working tree to `git commit` by hand. A **Commit changelog** checkbox now sits next to the **Cut release** button (on `/changelog` for the editor and `/deploy` for the export), checked by default — leave it on and the cut is followed by a path-limited `git commit` of just that one CHANGELOG.md, with the message `Release <workspace> <version>` (e.g. `Release export 0.4.1`). To keep the release commit clean it commits *only* the changelog: if the working tree has any *other* uncommitted change, the cut is refused up front (nothing is written) with an error telling you to commit or stash those first — a dirty changelog itself is fine, so uncommitted `[Unreleased]` bullets get folded into the release commit. Uncheck the box to cut without committing, exactly as before. See `editor/app/deploy/cutReleaseAction.ts`, `editor/app/deploy/components/CutReleaseForm.tsx`, and the new `common/lib/git.ts`.
- **"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`.
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,
}) => {