Archilyzer · Source

archilyzer

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

commit f2a898a3dba12d5f3e3d1351e3275e15f3fb63de
parent a210dace521bb77bc17e488119c7275a44fbdef1
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Thu,  1 Oct 2026 23:06:48 -0400

common: a subtitle 429 never fails a download, and the pace adapts per platform

The subtitle fetch alone answering 429 is a new class, subs_rate_limit:
the youtube-handling primary runs with --ignore-errors (yt-dlp reports the
subtitle failure as a WARNING and carries on), the media pass fetches the
audio with every subtitle refused, and the record is a success carrying the
class. A primary that died on its subtitles alone is run once more without
them. The video's subtitles are deferred (subtitleDeferrals, 6 h, 7 days
from the third); the platform backoff and pace are not touched.

A platform-level rate limit doubles the platform's --sleep-requests up to
pacing.sleepRequestsCapSeconds; clean units ease it back; a backoff failing
at the 30-minute cap pacing.holdAfterFailsAtCap times holds the platform to
one probe per pacing.holdProbeMinutes, and a clean probe clears the hold.
channelExtraArgs reads the live pace, so every spawn uses it; the download
lane now honours sleepBetweenDownloadsSeconds plus the pace above its base,
as a batch does; YouTube's primary gets --sleep-subtitles. New idle reasons
held and paced. download-missing-subs skips deferred videos, records a
subtitle 429 and goes on, and clears a deferral on success.

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

Diffstat:
MSETTINGS.md | 27++++++++++++++++++++++++++-
Mcommon/controller/autoRunner.ts | 79++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Mcommon/jobs/autoQueueState.test.ts | 39+++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/autoQueueState.ts | 53++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/downloadBackoff.test.ts | 65+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/downloadBackoff.ts | 136++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Mcommon/jobs/platformBackoff.test.ts | 63+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/platformBackoff.ts | 349+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/unitOutcome.test.ts | 160++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/unitOutcome.ts | 74+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Mcommon/lib/availability.test.ts | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/availability.ts | 38++++++++++++++++++++++++++++++++++++++
Mcommon/lib/downloadOutcome.ts | 11+++++++++--
Mcommon/lib/settingsDocs.ts | 8++++++++
Mcommon/lib/settingsSchema.test.ts | 6+++++-
Mcommon/lib/settingsSchema.ts | 73++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/views/activeJobs.ts | 4++++
Mcommon/ytdlp/channelArgs.test.ts | 69++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/ytdlp/channelArgs.ts | 63++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcommon/ytdlp/downloadOneManaged.ts | 117+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcommon/ytdlp/managedDownloadsSleep.test.ts | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/ytdlp/platformArgs.mjs | 40++++++++++++++++++++++++++++++++++++++++
Mcommon/ytdlp/runYtdlp.ts | 120++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Acommon/ytdlp/subtitleRateLimit.test.ts | 147+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/operations/components/dispatch.ts | 10++++++++++
Msettings.json.example | 6++++++
26 files changed, 1810 insertions(+), 63 deletions(-)

diff --git a/SETTINGS.md b/SETTINGS.md @@ -21,6 +21,7 @@ A copied example PINS every default it spells — including each lane's `autoQue | [`cookieMode`](#cookiemode) | `"when-required"` | | [`social`](#social) | object — see below | | [`sleepBetweenDownloadsSeconds`](#sleepbetweendownloadsseconds) | `10` | +| [`pacing`](#pacing) | object — see below | | [`downloadFormat`](#downloadformat) | `"auto"` | | [`minFreeDiskGB`](#minfreediskgb) | `5` | | [`resumeMarginGB`](#resumemargingb) | `2` | @@ -184,10 +185,34 @@ Default: ## `sleepBetweenDownloadsSeconds` -Pause (seconds) inserted between per-video yt-dlp invocations in managed batch downloads. yt-dlp's own `-t sleep` only paces requests within one invocation, so without this the managed loop hammers the source IP back-to-back. 0 disables. Per-channel override available. +Pause (seconds) inserted between per-video yt-dlp invocations in managed batch downloads, and between two auto-download units on one platform (release 17; the lane ignored it before). yt-dlp's own `-t sleep` only paces requests within one invocation, so without this the managed loop hammers the source IP back-to-back. The adaptive pace above its base (see `pacing`) is added to it. 0 disables. Per-channel override available for batch downloads. Default: `10` +## `pacing` + +How the download pace adapts to rate limits, per platform (release 17). Every yt-dlp spawn against a platform paces its requests (`--sleep-requests`) at the platform's adaptive pace, and the download lane waits sleepBetweenDownloadsSeconds plus the pace above its fixed value between units. A rate limit doubles the pace; clean units ease it back; a rate limit that outlasts the cooldown cap holds the platform to one probe at a time. The live pace, cooldowns and holds are in `.auto-queue/state.json`, shown on /operations/download. See common/jobs/platformBackoff.ts. + +#### `pacing` + +| Key | Default | Description | +|---|---|---| +| `sleepRequestsCapSeconds` | `16` | Ceiling (seconds) on a platform's adaptive `--sleep-requests`. The pace starts at the platform's fixed value (1 s for YouTube and Rumble) and doubles on every platform-level rate limit until it reaches this. A subtitle-only 429 never moves it. Clamped 1–120; default 16. | +| `decayAfterCleanUnits` | `5` | How many clean auto-download units on a platform halve its pace one step back toward the fixed value. Clamped 1–1000; default 5. | +| `holdAfterFailsAtCap` | `3` | How many consecutive failures AT the 30-minute cooldown cap put a platform in a hold: the lane then runs one probe unit per holdProbeMinutes instead of one per cooldown, and a manual Sync or download on it is refused with the next probe's time. A clean probe clears the hold and the backoff. Clamped 1–100; default 3. | +| `holdProbeMinutes` | `60` | Minutes between probes while a platform is held. Clamped 1–1440; default 60. | + +Default: + +```json +{ + "sleepRequestsCapSeconds": 16, + "decayAfterCleanUnits": 5, + "holdAfterFailsAtCap": 3, + "holdProbeMinutes": 60 +} +``` + ## `downloadFormat` Default yt-dlp `-f` download format for every channel that doesn't set its own (ChannelConfig.downloadFormat). "auto" picks per-source: `original` for Odysee (whose HLS rungs are CDN-truncated), `bestaudio/worst` elsewhere. See common/ytdlp/downloadFormat.ts. diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -70,12 +70,17 @@ import { writeAutoQueueState, } from "../jobs/autoQueueState"; import { + currentPaceSeconds, + downloadGapMs, isCoolingDown, + isPlatformHeld, isVideoDeferred, mergeBackoffEntry, pruneDeferred, - pruneExpired, + prunePlatformPacing, + pruneSubtitleDeferrals, } from "../jobs/platformBackoff"; +import { staticSleepRequestsSeconds } from "../ytdlp/platformArgs.mjs"; import { applyUnitOutcome } from "../jobs/unitOutcome"; import { type DownloadFailureClass } from "../lib/availability"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; @@ -216,6 +221,14 @@ export type AutoRunnerIdleReason = // Every pending video left after the platform gates was rate-limited // recently and is deferred (videoDeferrals). Download only. | "deferred" + // Every pending platform is HELD: its rate limit outlasted the cooldown cap, + // and the lane runs one probe per `pacing.holdProbeMinutes` until one comes + // back clean (release 17, slice RL). Download only. + | "held" + // Every pending platform is inside the gap between two units: + // `sleepBetweenDownloadsSeconds` plus the adaptive pace above its base + // (release 17, slice RL). Download only. + | "paced" // No enabled, non-degraded worker slot exists. Transcription only. | "no-workers" // The worker pool is pause-all'd. Transcription only, and DISTINCT from @@ -1147,7 +1160,10 @@ async function runLoop( 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()); + // A HELD platform's entry is kept however old: a hold ends with a clean + // probe, not with time (release 17, slice RL). + prunePlatformPacing(kindState, Date.now()); + pruneSubtitleDeferrals(kindState.subtitleDeferrals, Date.now()); // And lapsed per-video deferrals (live ones survive a restart by design — a // restart must not re-hit the rate-limited video at fails+1). pruneDeferred(kindState.videoDeferrals, Date.now()); @@ -1228,6 +1244,13 @@ async function runLoop( // "unknown" for unrecognized hosts). Unused for transcription. const platformInFlight = new Map<string, number>(); const PER_PLATFORM_CAP = 1; + // THE GAP BETWEEN TWO UNITS ON ONE PLATFORM (release 17, slice RL): the epoch + // ms before which the platform takes no new unit. Set when a unit settles to + // `sleepBetweenDownloadsSeconds` plus the adaptive pace above its base + // (downloadGapMs) — the lane used to start the next unit the moment one + // settled, ignoring the setting every batch download honours. In memory: a + // restart is itself a gap. + const platformNextStartAt = new Map<string, number>(); // Channels this runner has already scanned (or tried to). See the pre-pick in // next() for why a scan is never retried inside one runner's lifetime. const scannedThisRun = new Set<string>(); @@ -1619,6 +1642,10 @@ async function runLoop( // Download only: some pending video was dropped because it is deferred // after a recent rate limit (see unitOutcome.ts). let anyDeferred = false; + // Download only: a platform skipped for its hold (one probe at a time) or + // for the gap between units — each its own answer to "why idle?". + let anyHeld = false; + let anyPaced = false; platformSkip.clear(); if (kind === "download") { // Merge in any cooldown a manual sync/import wrote to the shared state @@ -1646,8 +1673,16 @@ async function runLoop( platformSkip.add(pf); // A platform skipped for a cooldown is a different answer to "why is // it idle?" than one skipped for being busy, so the two are tracked - // apart rather than both reading as "capped". - anyCooling = true; + // apart rather than both reading as "capped". A held platform's + // cooldown is its next probe. + if (isPlatformHeld(kindState.platformHolds, pf)) anyHeld = true; + else anyCooling = true; + } + } + for (const [pf, at] of platformNextStartAt) { + if (at > now && !platformSkip.has(pf)) { + platformSkip.add(pf); + anyPaced = true; } } if (platformSkip.size > 0) { @@ -1756,9 +1791,13 @@ async function runLoop( else if (countPending(pending) === 0) { live.idleReason = anyCooling ? "cooldown" - : anyDeferred - ? "deferred" - : "capped"; + : anyHeld + ? "held" + : anyPaced + ? "paced" + : anyDeferred + ? "deferred" + : "capped"; } else live.idleReason = "capped"; return null; } @@ -1958,6 +1997,8 @@ async function runLoop( // and a rate-limited video is also deferred so the next pick after the // cooldown is a DIFFERENT video — see unitOutcome.ts. Any other outcome // retires the video for the session; a success clears the cooldown. + const settledAt = Date.now(); + const unitSettings = getSettings(); const effect = applyUnitOutcome( kindState, { @@ -1967,9 +2008,23 @@ async function runLoop( outcome: result.outcome, ...(result.failureClass ? { failureClass: result.failureClass } : {}), }, - Date.now(), + settledAt, + Math.random, + unitSettings.pacing, ); if (effect.line) onLog(effect.line); + // The gap before this platform's next unit, at the pace this outcome + // left it (release 17, slice RL). + if (unitPlatform) { + const base = staticSleepRequestsSeconds(unitPlatform); + const gap = downloadGapMs( + unitSettings.sleepBetweenDownloadsSeconds, + currentPaceSeconds(kindState.platformPace, unitPlatform, base), + base, + ); + if (gap > 0) platformNextStartAt.set(unitPlatform, settledAt + gap); + else platformNextStartAt.delete(unitPlatform); + } // An operation lane retires (operation, video) pairs inside // runOperationPick — it is the only thing that knows WHICH operation // ran — so a bare id here would be a key nothing ever reads. @@ -2265,7 +2320,13 @@ async function launchUnit(args: LaunchArgs): Promise<UnitResult> { // Complete-but-malformed source: terminal and kept on disk. Treat as skipped // (not failed) so it doesn't drive backoff and isn't re-picked for download. if (unitStatus === "corrupt-full-source") return { outcome: "skipped" }; - if (unitStatus && unitStatus.startsWith("ok")) return { outcome: "transcribed" }; + // A success whose SUBTITLE fetch alone answered 429 says so: the runner + // defers the video's subtitles and leaves the platform alone (release 17). + if (unitStatus && unitStatus.startsWith("ok")) { + return unitFailureClass === "subs_rate_limit" + ? { outcome: "transcribed", failureClass: "subs_rate_limit" } + : { outcome: "transcribed" }; + } // failureClass drives the runner's per-platform backoff (rate_limit/network). return { outcome: "failed", failureClass: unitFailureClass }; } diff --git a/common/jobs/autoQueueState.test.ts b/common/jobs/autoQueueState.test.ts @@ -102,6 +102,9 @@ test("every lane is keyed, and a state file written before a lane existed coerce picks: [], platformBackoff: {}, videoDeferrals: {}, + platformPace: {}, + platformHolds: {}, + subtitleDeferrals: {}, }); assert.deepEqual(back.backfill.picks, []); // And the lanes it does name are untouched. @@ -168,6 +171,42 @@ test("a state file written before videoDeferrals existed coerces it to {}; corru }); }); +test("the pacing memory round-trips; a file written before it coerces to {} (release 17, RL)", async () => { + await withPaths(async (paths) => { + const state = emptyAutoQueueState(); + state.download.platformPace = { + youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 2 }, + }; + state.download.platformHolds = { youtube: { since: 10, probeAt: 20 } }; + state.download.subtitleDeferrals = { + v1: { count: 2, lastAt: 5, until: 6, channelSlug: "alpha" }, + }; + await writeAutoQueueState(paths, state); + const back = await readAutoQueueState(paths); + assert.deepEqual(back.download.platformPace, state.download.platformPace); + assert.deepEqual(back.download.platformHolds, state.download.platformHolds); + assert.deepEqual(back.download.subtitleDeferrals, state.download.subtitleDeferrals); + assert.deepEqual(back.transcription.platformPace, {}); + + await writeFile( + paths.autoQueueStateFile, + JSON.stringify({ + download: { + platformPace: { youtube: { sleepRequestsSeconds: "x" }, rumble: { sleepRequestsSeconds: 2 } }, + platformHolds: { youtube: { since: 1 } }, + subtitleDeferrals: { v1: { count: 0, lastAt: 1, until: 2, channelSlug: "a" } }, + }, + }), + ); + const old = await readAutoQueueState(paths); + assert.deepEqual(old.download.platformPace, { + rumble: { sleepRequestsSeconds: 2, baseSeconds: 0, cleanUnits: 0 }, + }); + assert.deepEqual(old.download.platformHolds, {}); + assert.deepEqual(old.download.subtitleDeferrals, {}); + }); +}); + // --- write safety ------------------------------------------------------------ test("overlapping writes do not collide on the tmp file", async () => { diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts @@ -7,8 +7,14 @@ import { } from "./autoQueuePolicy"; import { type PlatformBackoffState, + type PlatformHoldState, + type PlatformPaceState, + type SubtitleDeferralState, type VideoDeferralState, coercePlatformBackoff, + coercePlatformHolds, + coercePlatformPace, + coerceSubtitleDeferrals, coerceVideoDeferrals, } from "./platformBackoff"; @@ -48,6 +54,17 @@ export type AutoQueueKindState = { // Persisted so a restart does not re-hit the same video at fails+1. A file // written before the key existed coerces to `{}`; an older build ignores it. videoDeferrals: VideoDeferralState; + // THE PACING MEMORY (release 17, slice RL) — download kind only, `{}` + // elsewhere; a file written before a key existed coerces it to `{}` and an + // older build drops it on its next write. See platformBackoff.ts. + // Per-platform `--sleep-requests` above the static value, while it is above. + platformPace: PlatformPaceState; + // Platforms whose rate limit outlasted the cooldown cap: one probe per + // `pacing.holdProbeMinutes` until a clean one. + platformHolds: PlatformHoldState; + // Videos whose SUBTITLE fetch answered 429 (their media did not): the count, + // and until when the batch subtitle fetch leaves them alone. + subtitleDeferrals: SubtitleDeferralState; }; // KEYED BY EVERY LANE, including the two with no executor yet. A lane whose @@ -65,6 +82,9 @@ export function emptyAutoQueueKindState(): AutoQueueKindState { picks: [], platformBackoff: {}, videoDeferrals: {}, + platformPace: {}, + platformHolds: {}, + subtitleDeferrals: {}, }; } @@ -122,6 +142,9 @@ function coerceKindState(value: unknown): AutoQueueKindState { picks, platformBackoff: coercePlatformBackoff(r.platformBackoff), videoDeferrals: coerceVideoDeferrals(r.videoDeferrals), + platformPace: coercePlatformPace(r.platformPace), + platformHolds: coercePlatformHolds(r.platformHolds), + subtitleDeferrals: coerceSubtitleDeferrals(r.subtitleDeferrals), }; } @@ -137,9 +160,11 @@ export async function readAutoQueueState(paths: Paths): Promise<AutoQueueState> } if (!raw || typeof raw !== "object") return emptyAutoQueueState(); const r = raw as Record<string, unknown>; - return Object.fromEntries( + const state = Object.fromEntries( LANES.map((lane) => [lane, coerceKindState(r[lane])]), ) as AutoQueueState; + notePace(state); + return state; } // Write the state atomically (tmp + rename), creating the .auto-queue dir on @@ -162,10 +187,14 @@ export async function writeAutoQueueState( picks: k.picks.slice(0, AUTO_QUEUE_PICK_LOG_LIMIT), platformBackoff: k.platformBackoff ?? {}, videoDeferrals: k.videoDeferrals ?? {}, + platformPace: k.platformPace ?? {}, + platformHolds: k.platformHolds ?? {}, + subtitleDeferrals: k.subtitleDeferrals ?? {}, }); const out = Object.fromEntries( LANES.map((lane) => [lane, trim(state[lane] ?? emptyAutoQueueKindState())]), ) as AutoQueueState; + notePace(out); await writeJsonAtomic(paths.autoQueueStateFile, out, { mkdir: true }); } @@ -181,6 +210,9 @@ type AutoQueueStateHolder = { // receives ONE object. loading: Promise<AutoQueueState> | null; loadingFile: string | null; + // The download lane's pace as last read or written, by ANY reader of the + // file — for `livePlatformPaceSeconds` when no runner holds the object. + lastPace: PlatformPaceState | null; }; declare global { @@ -195,11 +227,30 @@ function getHolder(): AutoQueueStateHolder { stateFile: null, loading: null, loadingFile: null, + lastPace: null, }; } return globalThis.__yttAutoQueueState__; } +function notePace(state: AutoQueueState): void { + getHolder().lastPace = state.download?.platformPace ?? {}; +} + +// THE PACE A SPAWN USES, SYNCHRONOUSLY (release 17, slice RL). Every yt-dlp +// argv in the repo is built by `channelExtraArgs`, which is synchronous and is +// called from a dozen places; threading an awaited state read through each +// would be a dozen chances to forget one. So this answers from what the process +// already holds: the shared object when a runner holds it (the one every lane +// and `recordDownloadBackoff` write through), else the pace the last read or +// write of the file saw (the status poll reads it every few seconds), else +// undefined — the caller then uses the static pace. Never reads the disk. +export function livePlatformPaceSeconds(platform: string): number | undefined { + const holder = getHolder(); + const pace = holder.state?.download?.platformPace ?? holder.lastPace; + return pace?.[platform]?.sleepRequestsSeconds; +} + // THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER, and it has to be. // // `writeAutoQueueState` serializes the WHOLE file — all four lanes — so two diff --git a/common/jobs/downloadBackoff.test.ts b/common/jobs/downloadBackoff.test.ts @@ -5,9 +5,14 @@ import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; import { + clearSubtitleDeferral, + heldPlatformRefusal, platformCooldownRemainingMs, + readSubtitleDeferrals, recordDownloadBackoff, + recordSubtitleDeferral, } from "./downloadBackoff"; +import { PACING_DEFAULTS } from "./platformBackoff"; import { BACKOFF_BASE_MS } from "./platformBackoff"; import { emptyAutoQueueState, @@ -151,3 +156,63 @@ test("with no runner live, recordDownloadBackoff stays a disk round trip and loa assert.equal(onDisk.download.platformBackoff.youtube?.fails, 1); }); }); + +// ── release 17, slice RL ───────────────────────────────────────────────────── + +test("a manual rate limit doubles the pace; a network failure does not", async () => { + await withTempPaths(async (paths) => { + await recordDownloadBackoff("youtube", paths, "rate_limit", PACING_DEFAULTS); + await recordDownloadBackoff("youtube", paths, "network", PACING_DEFAULTS); + const onDisk = await readAutoQueueState(paths); + assert.equal(onDisk.download.platformBackoff.youtube.fails, 2); + assert.deepEqual(onDisk.download.platformPace.youtube, { + sleepRequestsSeconds: 2, + baseSeconds: 1, + cleanUnits: 0, + }); + }); +}); + +test("a manual 429 at the cap holds the platform, and the refusal names the next probe", async () => { + await withTempPaths(async (paths) => { + const state = emptyAutoQueueState(); + state.download.platformBackoff.youtube = { until: Date.now() - 1, fails: 7 }; + await writeAutoQueueState(paths, state); + const pacing = { ...PACING_DEFAULTS, holdAfterFailsAtCap: 3, holdProbeMinutes: 60 }; + assert.equal(await heldPlatformRefusal("youtube", "Sync", paths), null); + await recordDownloadBackoff("youtube", paths, "rate_limit", pacing); + const onDisk = await readAutoQueueState(paths); + assert.ok(onDisk.download.platformHolds.youtube); + assert.equal( + onDisk.download.platformBackoff.youtube.until, + onDisk.download.platformHolds.youtube.probeAt, + ); + const text = await heldPlatformRefusal("youtube", "Sync", paths); + assert.match(text ?? "", /^youtube is held: .*\(8 failures in a row, held since \d\d:\d\d UTC\)\..* the next probe is in 60 min\. Sync will run once a probe comes back clean\.$/); + }); +}); + +test("subtitle deferrals count up through the file and clear on success", async () => { + await withTempPaths(async (paths) => { + const a = await recordSubtitleDeferral("v1", "alpha", paths, 1000); + const b = await recordSubtitleDeferral("v1", "alpha", paths, 2000); + assert.equal(a.count, 1); + assert.equal(b.count, 2); + const onDisk = await readAutoQueueState(paths); + assert.equal(onDisk.download.subtitleDeferrals.v1.count, 2); + // Nothing platform-wide moved. + assert.deepEqual(onDisk.download.platformBackoff, {}); + assert.deepEqual(onDisk.download.platformPace, {}); + await clearSubtitleDeferral("v1", paths); + assert.deepEqual(await readSubtitleDeferrals(paths), {}); + }); +}); + +test("with a runner live, a subtitle deferral is written through the shared object", async () => { + await withTempPaths(async (paths) => { + const live = await sharedAutoQueueState(paths); + await recordSubtitleDeferral("v9", "beta", paths); + assert.equal(live.download.subtitleDeferrals.v9.count, 1); + assert.equal((await readAutoQueueState(paths)).download.subtitleDeferrals.v9.count, 1); + }); +}); diff --git a/common/jobs/downloadBackoff.ts b/common/jobs/downloadBackoff.ts @@ -19,13 +19,26 @@ // politeness, never concurrency. import { getPaths, type Paths } from "../lib/paths"; +import { getSettings } from "../lib/settings"; import { type AutoQueueState, liveAutoQueueState, readAutoQueueState, writeAutoQueueState, } from "./autoQueueState"; -import { mergeBackoffEntry, nextBackoff, pruneExpired } from "./platformBackoff"; +import { + type PacingSettings, + type PlatformHoldEntry, + type SubtitleDeferral, + type SubtitleDeferralState, + deferSubtitles, + escalatePlatform, + heldPlatformSentence, + mergeBackoffEntry, + prunePlatformPacing, + pruneSubtitleDeferrals, +} from "./platformBackoff"; +import { staticSleepRequestsSeconds } from "../ytdlp/platformArgs.mjs"; // Milliseconds remaining in the platform's current cooldown window, or 0 if it // is not cooling down. Prefers the live shared object; else reads the file. @@ -40,13 +53,23 @@ export async function platformCooldownRemainingMs( return entry && entry.until > now ? entry.until - now : 0; } -// Record a rate-limit/network failure against a platform, escalating its -// exponential cooldown — through the live shared object when a runner holds -// one, else as a read-modify-write of the file. -export async function recordDownloadBackoff( - platform: string, - paths: Paths = getPaths(), -): Promise<void> { +// The `pacing` settings block, read defensively: a settings read that throws +// (a CLI with no corpus, a test's stub paths) paces on the defaults. +function pacingSettings(): PacingSettings | undefined { + try { + return getSettings().pacing; + } catch { + return undefined; + } +} + +// THE ONE WRITE-THROUGH: mutate the live shared object when a runner holds +// one (after folding in the backoff on disk, as the runner's own merge does), +// else read-modify-write the file. Every writer below goes through it. +async function mutateDownloadState<T>( + paths: Paths, + mutate: (state: AutoQueueState) => T, +): Promise<T> { const live = await liveAutoQueueState(paths); let state: AutoQueueState; if (live) { @@ -64,11 +87,96 @@ export async function recordDownloadBackoff( } else { state = await readAutoQueueState(paths); } - const now = Date.now(); - pruneExpired(state.download.platformBackoff, now); - state.download.platformBackoff[platform] = nextBackoff( - state.download.platformBackoff[platform], - now, - ); + const out = mutate(state); await writeAutoQueueState(paths, state); + return out; +} + +// Record a rate-limit/network failure against a platform, escalating its +// exponential cooldown — through the live shared object when a runner holds +// one, else as a read-modify-write of the file. A rate limit also doubles the +// platform's pace, and a backoff failing at the cap often enough holds it +// (release 17, slice RL) — the same escalation the lane applies to its units. +export async function recordDownloadBackoff( + platform: string, + paths: Paths = getPaths(), + failureClass: "rate_limit" | "network" = "rate_limit", + pacing: PacingSettings | undefined = pacingSettings(), +): Promise<void> { + await mutateDownloadState(paths, (state) => { + const now = Date.now(); + prunePlatformPacing(state.download, now); + escalatePlatform(state.download, platform, failureClass, now, { + baseSeconds: staticSleepRequestsSeconds(platform), + ...(pacing ? { pacing } : {}), + }); + }); +} + +// Record a SUBTITLE 429 against one video (release 17, slice RL): the count +// climbs and the batch subtitle fetch leaves it alone for 6 h, or for 7 days +// from the third. The platform's backoff, hold and pace are not touched. +export async function recordSubtitleDeferral( + videoId: string, + channelSlug: string, + paths: Paths = getPaths(), + now: number = Date.now(), +): Promise<SubtitleDeferral> { + return mutateDownloadState(paths, (state) => { + pruneSubtitleDeferrals(state.download.subtitleDeferrals, now); + return deferSubtitles(state.download.subtitleDeferrals, videoId, channelSlug, now); + }); +} + +// A video's subtitles came down: its deferral (and its count) is forgotten. +// A no-op — no write — when it had none. +export async function clearSubtitleDeferral( + videoId: string, + paths: Paths = getPaths(), +): Promise<void> { + const state = + (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); + if (!state.download.subtitleDeferrals[videoId]) return; + await mutateDownloadState(paths, (s) => { + delete s.download.subtitleDeferrals[videoId]; + }); +} + +// Every recorded subtitle deferral (live object first, else the file). +export async function readSubtitleDeferrals( + paths: Paths = getPaths(), +): Promise<SubtitleDeferralState> { + const state = + (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); + return state.download.subtitleDeferrals; +} + +// A platform's hold, when it is held, with the backoff's failure count. +export async function platformHold( + platform: string, + paths: Paths = getPaths(), +): Promise<{ hold: PlatformHoldEntry; fails: number } | null> { + const state = + (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); + const hold = state.download.platformHolds[platform]; + if (!hold) return null; + return { hold, fails: state.download.platformBackoff[platform]?.fails ?? 0 }; +} + +// The refusal for a manual Sync / download on a HELD platform, or null when it +// is not held. `what` names the refused thing ("Sync", "This download"). +export async function heldPlatformRefusal( + platform: string, + what: string, + paths: Paths = getPaths(), + now: number = Date.now(), +): Promise<string | null> { + const held = await platformHold(platform, paths); + if (!held) return null; + const probeMinutes = + pacingSettings()?.holdProbeMinutes ?? Math.round((held.hold.probeAt - held.hold.since) / 60_000); + return heldPlatformSentence(platform, held.hold, held.fails, now, { + probeMinutes, + what, + }); } diff --git a/common/jobs/platformBackoff.test.ts b/common/jobs/platformBackoff.test.ts @@ -16,7 +16,14 @@ import { deferVideo, isVideoDeferred, pruneDeferred, + PACING_DEFAULTS, + downloadGapMs, + failsAtCap, + FAILS_TO_REACH_CAP, + heldPlatformSentence, + prunePlatformPacing, } from "./platformBackoff"; +import { sanitizePacing } from "../lib/settingsSchema"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/platformBackoff.test.ts @@ -164,3 +171,59 @@ test("mergeBackoffEntry takes the max of until and of fails separately", () => { { until: 300, fails: 7 }, ); }); + +// ── release 17, slice RL ───────────────────────────────────────────────────── + +test("PACING_DEFAULTS are the settings schema's defaults", () => { + assert.deepEqual({ ...PACING_DEFAULTS }, sanitizePacing({})); + assert.deepEqual(sanitizePacing({ sleepRequestsCapSeconds: 999, holdProbeMinutes: 0 }), { + ...PACING_DEFAULTS, + sleepRequestsCapSeconds: 120, + holdProbeMinutes: 1, + }); +}); + +test("the cap is reached at the sixth failure; failsAtCap counts from there", () => { + assert.equal(FAILS_TO_REACH_CAP, 6); + assert.equal(BACKOFF_BASE_MS * 2 ** (FAILS_TO_REACH_CAP - 2) < BACKOFF_MAX_MS, true); + assert.deepEqual([5, 6, 7, 8].map(failsAtCap), [0, 1, 2, 3]); +}); + +test("the lane gap is the operator's sleep plus the pace above its base", () => { + assert.equal(downloadGapMs(0, 1, 1), 0); + assert.equal(downloadGapMs(30, 1, 1), 30_000); + assert.equal(downloadGapMs(30, 4, 1), 33_000); + assert.equal(downloadGapMs(0, 8, 1), 7_000); + assert.equal(downloadGapMs(10, 2, 0), 12_000); + assert.equal(downloadGapMs(-5, 0.5, 1), 0); +}); + +test("a held platform's backoff survives the prune; a hold with no backoff is dropped", () => { + const now = 10 * BACKOFF_MAX_MS; + const state = { + platformBackoff: { + youtube: { until: 0, fails: 9 }, + rumble: { until: 0, fails: 2 }, + } as PlatformBackoffState, + platformPace: {}, + platformHolds: { + youtube: { since: 0, probeAt: 0 }, + odysee: { since: 0, probeAt: 0 }, + }, + }; + prunePlatformPacing(state, now); + assert.deepEqual(Object.keys(state.platformBackoff), ["youtube"]); + assert.deepEqual(Object.keys(state.platformHolds), ["youtube"]); +}); + +test("the held sentence names the hold and the next probe", () => { + const since = Date.UTC(2026, 9, 1, 14, 2); + assert.equal( + heldPlatformSentence("youtube", { since, probeAt: since + 42 * 60_000 }, 8, since, { + probeMinutes: 60, + what: "Sync", + }), + "youtube is held: its rate limit outlasted the cooldown cap (8 failures in a row, held since 14:02 UTC). " + + "Auto-download probes it once every 60 min — the next probe is in 42 min. Sync will run once a probe comes back clean.", + ); +}); diff --git a/common/jobs/platformBackoff.ts b/common/jobs/platformBackoff.ts @@ -177,3 +177,352 @@ export function coerceVideoDeferrals(value: unknown): VideoDeferralState { } return out; } + +// --------------------------------------------------------------------------- +// The adaptive pace, the hold and the subtitle deferral (release 17, slice RL). +// +// Measured 2026-10-01: every one of the day's 78 HTTP 429s was YouTube's +// subtitle (timedtext) fetch, per video, after the media formats had already +// resolved — and the platform backoff read `fails: 80`, retrying one unit every +// 30 minutes into the same refusal. Three answers, each pure here and persisted +// beside `platformBackoff` in `.auto-queue/state.json` (download lane only): +// +// - A SUBTITLE 429 IS NOT A PLATFORM SIGNAL. `classifyDownloadFailure` says +// `subs_rate_limit` when the subtitle fetch is the only failure; the download +// goes on to the media, and the video's subtitles are deferred here — the +// platform backoff and the pace are not touched. Three deferrals and the +// subtitles are left alone for a week. +// - THE PACE ADAPTS. `--sleep-requests` for a platform starts at its static +// value (PLATFORM_ARGS), doubles on every platform-level `rate_limit` up to +// `pacing.sleepRequestsCapSeconds`, and halves back toward the base after +// every `pacing.decayAfterCleanUnits` clean units. +// - A BLOCK IS NOT A BURST. A platform whose backoff has failed at the +// 30-minute cap `pacing.holdAfterFailsAtCap` times in a row is HELD: its +// cooldown becomes one probe per `pacing.holdProbeMinutes`, and a clean probe +// clears the hold and the backoff together. + +// The `pacing` settings block (settingsSchema.ts), restated here so jobs/ need +// not import the schema module. +export type PacingSettings = { + sleepRequestsCapSeconds: number; + decayAfterCleanUnits: number; + holdAfterFailsAtCap: number; + holdProbeMinutes: number; +}; + +export const PACING_DEFAULTS: Readonly<PacingSettings> = Object.freeze({ + sleepRequestsCapSeconds: 16, + decayAfterCleanUnits: 5, + holdAfterFailsAtCap: 3, + holdProbeMinutes: 60, +}); + +// The current pace of one platform, when it differs from the static one. An +// entry back at its base is deleted, so `{}` means "every platform at its base". +export type PlatformPaceEntry = { + // Seconds between the HTTP requests one yt-dlp process makes + // (`--sleep-requests`), for every spawn against this platform. + sleepRequestsSeconds: number; + // The platform's static value the pace decays back to (PLATFORM_ARGS; 0 for + // a platform with no entry). Stored so a reader with no access to the args + // table (the status view) can say "base 1 s". + baseSeconds: number; + // Clean units since the last doubling or decay step. + cleanUnits: number; +}; + +export type PlatformPaceState = Record<string, PlatformPaceEntry>; + +// A platform in a HOLD. Its backoff entry's `until` is the next probe. +export type PlatformHoldEntry = { + // Epoch ms the hold began. + since: number; + // Epoch ms of the next probe (mirrors the backoff entry's `until`). + probeAt: number; +}; + +export type PlatformHoldState = Record<string, PlatformHoldEntry>; + +// A video whose subtitles answered 429 while its media did not. +export type SubtitleDeferral = { + // How many times its subtitle fetch was rate-limited. + count: number; + // Epoch ms of the most recent one. + lastAt: number; + // Epoch ms until which the batch subtitle fetch (download-missing-subs) + // leaves it alone: VIDEO_RATE_LIMIT_DEFER_MS after a strike, and + // SUBTITLE_HOLD_MS once `count` reaches SUBTITLE_HOLD_AFTER. + until: number; + channelSlug: string; +}; + +export type SubtitleDeferralState = Record<string, SubtitleDeferral>; + +export const SUBTITLE_HOLD_AFTER = 3; +export const SUBTITLE_HOLD_MS = 7 * 24 * 60 * 60_000; // 7 days + +// The smallest `fails` whose cooldown is the cap (60 s doubling → 32 min ≥ 30). +export const FAILS_TO_REACH_CAP = (() => { + let f = 1; + while (BACKOFF_BASE_MS * 2 ** (f - 1) < BACKOFF_MAX_MS) f++; + return f; +})(); + +// How many consecutive failures were spent AT the cap. +export function failsAtCap(fails: number): number { + return Math.max(0, fails - FAILS_TO_REACH_CAP + 1); +} + +export function isPlatformHeld(holds: PlatformHoldState, platform: string): boolean { + return holds[platform] !== undefined; +} + +// The pace a platform's spawns use now: its entry, else its base. +export function currentPaceSeconds( + pace: PlatformPaceState, + platform: string, + baseSeconds: number, +): number { + return pace[platform]?.sleepRequestsSeconds ?? baseSeconds; +} + +// What the lane waits between two units on one platform, and what a batch +// sleeps between two videos: the operator's `sleepBetweenDownloadsSeconds` +// plus the pace ABOVE ITS BASE. The base pace is already paid inside every +// spawn (`--sleep-requests`); what the gap adds is the part a rate limit added. +export function downloadGapMs( + sleepBetweenDownloadsSeconds: number, + paceSeconds: number, + baseSeconds: number, +): number { + const sleep = Math.max(0, sleepBetweenDownloadsSeconds); + const extra = Math.max(0, paceSeconds - baseSeconds); + return Math.round((sleep + extra) * 1000); +} + +export type PlatformPacingState = { + platformBackoff: PlatformBackoffState; + platformPace: PlatformPaceState; + platformHolds: PlatformHoldState; +}; + +export type EscalationEffect = { + entry: PlatformBackoffEntry; + // The pace after this failure (unchanged on a network failure). + paceSeconds: number; + // True when this failure put the platform into its hold. + enteredHold: boolean; + held: boolean; +}; + +// One platform-level failure (`rate_limit` or `network`): the backoff +// escalates exactly as before, a rate limit doubles the pace, and a backoff +// that has now failed at the cap `holdAfterFailsAtCap` times holds the +// platform — its next try becomes a probe `holdProbeMinutes` away. Mutates. +export function escalatePlatform( + state: PlatformPacingState, + platform: string, + failureClass: "rate_limit" | "network", + now: number, + opts: { pacing?: PacingSettings; baseSeconds: number; rand?: () => number }, +): EscalationEffect { + const pacing = opts.pacing ?? PACING_DEFAULTS; + const entry = nextBackoff(state.platformBackoff[platform], now, opts.rand); + if (failureClass === "rate_limit") { + const cur = currentPaceSeconds(state.platformPace, platform, opts.baseSeconds); + const cap = Math.max(pacing.sleepRequestsCapSeconds, opts.baseSeconds); + state.platformPace[platform] = { + sleepRequestsSeconds: Math.min(cap, Math.max(1, cur * 2)), + baseSeconds: opts.baseSeconds, + cleanUnits: 0, + }; + } + const wasHeld = isPlatformHeld(state.platformHolds, platform); + const held = wasHeld || failsAtCap(entry.fails) >= pacing.holdAfterFailsAtCap; + if (held) { + entry.until = now + pacing.holdProbeMinutes * 60_000; + state.platformHolds[platform] = { + since: state.platformHolds[platform]?.since ?? now, + probeAt: entry.until, + }; + } + state.platformBackoff[platform] = entry; + return { + entry, + paceSeconds: currentPaceSeconds(state.platformPace, platform, opts.baseSeconds), + enteredHold: held && !wasHeld, + held, + }; +} + +export type CleanEffect = { + // True when this clean unit was the probe that lifted a hold. + releasedHold: boolean; + // The pace after this unit, and whether it stepped down. + paceSeconds: number; + decayed: boolean; +}; + +// One clean unit on a platform: the backoff and any hold clear (a clean probe +// is what lifts a hold), and every `decayAfterCleanUnits` of them halves the +// pace back toward its base — an entry that reaches the base is deleted. +export function settlePlatformClean( + state: PlatformPacingState, + platform: string, + opts: { pacing?: PacingSettings; baseSeconds: number }, +): CleanEffect { + const pacing = opts.pacing ?? PACING_DEFAULTS; + const releasedHold = isPlatformHeld(state.platformHolds, platform); + clearBackoff(state.platformBackoff, platform); + delete state.platformHolds[platform]; + const pace = state.platformPace[platform]; + let decayed = false; + if (pace) { + pace.cleanUnits += 1; + if (pace.cleanUnits >= Math.max(1, pacing.decayAfterCleanUnits)) { + decayed = true; + const next = pace.sleepRequestsSeconds / 2; + if (next <= pace.baseSeconds) delete state.platformPace[platform]; + else state.platformPace[platform] = { ...pace, sleepRequestsSeconds: next, cleanUnits: 0 }; + } + } + return { + releasedHold, + paceSeconds: currentPaceSeconds(state.platformPace, platform, opts.baseSeconds), + decayed, + }; +} + +// Prune the backoff map without dropping a HELD platform's entry: a hold lasts +// until a clean probe, however long the lane was stopped, and its `fails` is +// the escalation memory that keeps it a hold. A hold whose backoff entry is +// gone (a hand-edited file) is dropped with it. Mutates. +export function prunePlatformPacing(state: PlatformPacingState, now: number): void { + const before = { ...state.platformBackoff }; + pruneExpired(state.platformBackoff, now); + for (const pf of Object.keys(state.platformHolds)) { + if (before[pf] && !state.platformBackoff[pf]) state.platformBackoff[pf] = before[pf]; + if (!state.platformBackoff[pf]) delete state.platformHolds[pf]; + } +} + +// Record one subtitle 429 against a video: the count climbs, and the batch +// subtitle fetch leaves it alone for 6 h — or for 7 days from the third. A +// manual fetch never reads this. Mutates; returns the entry. +export function deferSubtitles( + state: SubtitleDeferralState, + videoId: string, + channelSlug: string, + now: number, +): SubtitleDeferral { + const count = (state[videoId]?.count ?? 0) + 1; + const entry: SubtitleDeferral = { + count, + lastAt: now, + until: now + (count >= SUBTITLE_HOLD_AFTER ? SUBTITLE_HOLD_MS : VIDEO_RATE_LIMIT_DEFER_MS), + channelSlug, + }; + state[videoId] = entry; + return entry; +} + +export function isSubtitleDeferred( + state: SubtitleDeferralState, + videoId: string, + now: number, +): boolean { + const e = state[videoId]; + return e !== undefined && e.until > now; +} + +// Drop a deferral a week after its window lapsed: the count is what makes the +// third strike a week, so it outlives the 6 h window it opened. Mutates. +export function pruneSubtitleDeferrals(state: SubtitleDeferralState, now: number): void { + for (const [id, e] of Object.entries(state)) { + if (e.until + SUBTITLE_HOLD_MS <= now) delete state[id]; + } +} + +// ---- coercion (the persisted shapes; mirrors coercePlatformBackoff) ---- + +const finite = (v: unknown): v is number => typeof v === "number" && Number.isFinite(v); + +export function coercePlatformPace(value: unknown): PlatformPaceState { + const out: PlatformPaceState = {}; + if (!value || typeof value !== "object") return out; + for (const [pf, raw] of Object.entries(value as Record<string, unknown>)) { + if (!raw || typeof raw !== "object") continue; + const r = raw as Record<string, unknown>; + if (finite(r.sleepRequestsSeconds) && r.sleepRequestsSeconds > 0) { + out[pf] = { + sleepRequestsSeconds: r.sleepRequestsSeconds, + baseSeconds: finite(r.baseSeconds) && r.baseSeconds >= 0 ? r.baseSeconds : 0, + cleanUnits: finite(r.cleanUnits) && r.cleanUnits >= 0 ? Math.floor(r.cleanUnits) : 0, + }; + } + } + return out; +} + +export function coercePlatformHolds(value: unknown): PlatformHoldState { + const out: PlatformHoldState = {}; + if (!value || typeof value !== "object") return out; + for (const [pf, raw] of Object.entries(value as Record<string, unknown>)) { + if (!raw || typeof raw !== "object") continue; + const r = raw as Record<string, unknown>; + if (finite(r.since) && finite(r.probeAt)) out[pf] = { since: r.since, probeAt: r.probeAt }; + } + return out; +} + +export function coerceSubtitleDeferrals(value: unknown): SubtitleDeferralState { + const out: SubtitleDeferralState = {}; + if (!value || typeof value !== "object") return out; + for (const [id, raw] of Object.entries(value as Record<string, unknown>)) { + if (!raw || typeof raw !== "object") continue; + const r = raw as Record<string, unknown>; + if ( + finite(r.count) && + r.count >= 1 && + finite(r.lastAt) && + finite(r.until) && + typeof r.channelSlug === "string" + ) { + out[id] = { + count: Math.floor(r.count), + lastAt: r.lastAt, + until: r.until, + channelSlug: r.channelSlug, + }; + } + } + return out; +} + +// ---- words (shared by the refusals, the status view and the doctor) ---- + +function clockOf(ms: number): string { + return new Date(ms).toISOString().slice(11, 16) + " UTC"; +} + +function minutesText(ms: number): string { + const m = Math.max(1, Math.ceil(ms / 60_000)); + return m >= 120 ? `${Math.floor(m / 60)}h ${m % 60}m` : `${m} min`; +} + +// The sentence a manual Sync / download on a held platform is refused with. +export function heldPlatformSentence( + platform: string, + hold: PlatformHoldEntry, + fails: number, + now: number, + opts: { probeMinutes: number; what?: string }, +): string { + const what = opts.what ?? "This"; + return ( + `${platform} is held: its rate limit outlasted the cooldown cap ` + + `(${fails} failures in a row, held since ${clockOf(hold.since)}). ` + + `Auto-download probes it once every ${opts.probeMinutes} min — the next probe ` + + `is in ${minutesText(hold.probeAt - now)}. ${what} will run once a probe comes back clean.` + ); +} diff --git a/common/jobs/unitOutcome.test.ts b/common/jobs/unitOutcome.test.ts @@ -2,6 +2,9 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import { BACKOFF_BASE_MS, + FAILS_TO_REACH_CAP, + PACING_DEFAULTS, + SUBTITLE_HOLD_MS, VIDEO_RATE_LIMIT_DEFER_MS, isVideoDeferred, } from "./platformBackoff"; @@ -10,7 +13,13 @@ import { type UnitOutcomeState, applyUnitOutcome } from "./unitOutcome"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/unitOutcome.test.ts const noJitter = () => 0.5; -const empty = (): UnitOutcomeState => ({ platformBackoff: {}, videoDeferrals: {} }); +const empty = (): UnitOutcomeState => ({ + platformBackoff: {}, + videoDeferrals: {}, + platformPace: {}, + platformHolds: {}, + subtitleDeferrals: {}, +}); const unit = ( videoId: string, over: Partial<Parameters<typeof applyUnitOutcome>[1]> = {}, @@ -34,7 +43,7 @@ test("rate_limit backs the platform off AND defers the video 6h", () => { }); assert.equal( r.line, - "Auto-download: youtube rate_limit — backing off 60s (attempt 1). v1 deferred 6h; next video after cooldown.", + "Auto-download: youtube rate_limit — backing off 60s (attempt 1). v1 deferred 6h; next video after cooldown. Pace now 2s between requests.", ); }); @@ -103,3 +112,150 @@ test("fails climbs only across DISTINCT ids: the deferred video is not re-picked assert.equal(s.platformBackoff.youtube.fails, 3); for (const id of ids) assert.equal(isVideoDeferred(s.videoDeferrals, id, now), true); }); + +// ── release 17, slice RL: the subtitle deferral, the pace, the hold ───────── + +test("subs_rate_limit retires the video, defers its subtitles, and touches no platform state", () => { + const s = empty(); + s.platformBackoff.youtube = { until: 10, fails: 3 }; + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 2 }; + const before = structuredClone({ + platformBackoff: s.platformBackoff, + platformPace: s.platformPace, + platformHolds: s.platformHolds, + }); + const now = 1_000_000; + const r = applyUnitOutcome( + s, + unit("v1", { outcome: "transcribed", failureClass: "subs_rate_limit" }), + now, + noJitter, + ); + assert.equal(r.markCompleted, true); + assert.equal( + r.line, + "Auto-download: v1 downloaded; its subtitles were rate-limited (1×) — deferred 6h for download-missing-subs. youtube is not backed off.", + ); + assert.deepEqual(s.subtitleDeferrals.v1, { + count: 1, + lastAt: now, + until: now + VIDEO_RATE_LIMIT_DEFER_MS, + channelSlug: "alpha", + }); + // Not cleared like a success, not escalated like a failure, no video deferral. + assert.deepEqual( + { platformBackoff: s.platformBackoff, platformPace: s.platformPace, platformHolds: s.platformHolds }, + before, + ); + assert.deepEqual(s.videoDeferrals, {}); +}); + +test("the third subtitle deferral leaves the subtitles alone for 7 days", () => { + const s = empty(); + let now = 0; + for (let i = 1; i <= 3; i++) { + const r = applyUnitOutcome( + s, + unit("v1", { outcome: "transcribed", failureClass: "subs_rate_limit" }), + now, + noJitter, + ); + assert.equal(s.subtitleDeferrals.v1.count, i); + if (i < 3) { + assert.equal(s.subtitleDeferrals.v1.until, now + VIDEO_RATE_LIMIT_DEFER_MS); + now += VIDEO_RATE_LIMIT_DEFER_MS + 1; // the 6 h window lapses; the count stays + } else { + assert.equal(s.subtitleDeferrals.v1.until, now + SUBTITLE_HOLD_MS); + assert.match(r.line ?? "", /\(3×\) — left alone for 7 days/); + } + } + // A week after the hold lapsed the record is pruned, count and all. + applyUnitOutcome(s, unit("v2", { outcome: "skipped" }), now + 2 * SUBTITLE_HOLD_MS + 1, noJitter); + assert.deepEqual(s.subtitleDeferrals, {}); +}); + +test("the pace doubles on each rate limit, never on network, and stops at the cap", () => { + const s = empty(); + const pacing = { ...PACING_DEFAULTS, sleepRequestsCapSeconds: 8 }; + let now = 0; + const paces: number[] = []; + for (let i = 0; i < 5; i++) { + applyUnitOutcome(s, unit(`v${i}`, { failureClass: "rate_limit" }), now, noJitter, pacing); + paces.push(s.platformPace.youtube.sleepRequestsSeconds); + now = s.platformBackoff.youtube.until; + } + assert.deepEqual(paces, [2, 4, 8, 8, 8]); + applyUnitOutcome(s, unit("vn", { failureClass: "network" }), now, noJitter, pacing); + assert.equal(s.platformPace.youtube.sleepRequestsSeconds, 8); + // A platform with no static pace starts from nothing and goes to 1 s. + applyUnitOutcome(s, unit("o1", { platform: "odysee", failureClass: "rate_limit" }), now, noJitter, pacing); + assert.deepEqual(s.platformPace.odysee, { sleepRequestsSeconds: 1, baseSeconds: 0, cleanUnits: 0 }); +}); + +test("the pace eases one step per N clean units, and an entry back at its base is gone", () => { + const s = empty(); + const pacing = { ...PACING_DEFAULTS, decayAfterCleanUnits: 2 }; + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0 }; + const lines: (string | null)[] = []; + for (let i = 0; i < 4; i++) { + lines.push(applyUnitOutcome(s, unit(`c${i}`, { outcome: "transcribed" }), 0, noJitter, pacing).line); + } + assert.deepEqual(lines, [ + null, + "Auto-download: youtube pace eased to 2s between requests.", + null, + "Auto-download: youtube pace eased to 1s between requests.", + ]); + assert.deepEqual(s.platformPace, {}); + // A subtitle-429 unit is not clean: it does not count toward the decay. + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 1 }; + applyUnitOutcome(s, unit("x", { outcome: "transcribed", failureClass: "subs_rate_limit" }), 0, noJitter, pacing); + assert.equal(s.platformPace.youtube.cleanUnits, 1); +}); + +test("a backoff failing at the cap N times holds the platform: one probe per holdProbeMinutes", () => { + const s = empty(); + const pacing = { ...PACING_DEFAULTS, holdAfterFailsAtCap: 3, holdProbeMinutes: 60 }; + let now = 0; + const lines: string[] = []; + // Climb to the cap, then two failures AT it: no hold yet. + for (let i = 0; i < FAILS_TO_REACH_CAP + 1; i++) { + lines.push(applyUnitOutcome(s, unit(`v${i}`, { failureClass: "rate_limit" }), now, noJitter, pacing).line ?? ""); + now = s.platformBackoff.youtube.until; + } + assert.deepEqual(s.platformHolds, {}); + assert.equal(s.platformBackoff.youtube.fails, FAILS_TO_REACH_CAP + 1); + // The third at the cap: held, and the next try is a probe an hour away. + const r = applyUnitOutcome(s, unit("h1", { failureClass: "rate_limit" }), now, noJitter, pacing); + assert.deepEqual(s.platformHolds.youtube, { since: now, probeAt: now + 60 * 60_000 }); + assert.equal(s.platformBackoff.youtube.until, now + 60 * 60_000); + assert.match(r.line ?? "", /^Auto-download: youtube rate_limit — held after 8 failures in a row; next probe in 3600s\./); + // A failed probe keeps the hold (and its `since`) and re-arms the hour. + const since = now; + now = s.platformBackoff.youtube.until; + const r2 = applyUnitOutcome(s, unit("h2", { failureClass: "network" }), now, noJitter, pacing); + assert.deepEqual(s.platformHolds.youtube, { since, probeAt: now + 60 * 60_000 }); + assert.match(r2.line ?? "", /still held/); + // A clean probe clears the hold AND the backoff. + const r3 = applyUnitOutcome(s, unit("h3", { outcome: "transcribed" }), now + 60 * 60_000, noJitter, pacing); + assert.deepEqual(s.platformHolds, {}); + assert.equal(s.platformBackoff.youtube, undefined); + assert.equal( + r3.line, + "Auto-download: youtube probe came back clean — the hold and the backoff are cleared.", + ); +}); + +test("a subtitle 429 on a probe neither lifts nor deepens the hold", () => { + const s = empty(); + s.platformBackoff.youtube = { until: 100, fails: 9 }; + s.platformHolds.youtube = { since: 0, probeAt: 100 }; + applyUnitOutcome( + s, + unit("p", { outcome: "transcribed", failureClass: "subs_rate_limit" }), + 100, + noJitter, + ); + assert.deepEqual(s.platformHolds.youtube, { since: 0, probeAt: 100 }); + assert.deepEqual(s.platformBackoff.youtube, { until: 100, fails: 9 }); +}); diff --git a/common/jobs/unitOutcome.ts b/common/jobs/unitOutcome.ts @@ -17,20 +17,42 @@ // - anything else retires it; a success (`transcribed`) also clears the // platform's cooldown. // - every branch prunes lapsed deferrals, keeping the map bounded. +// +// Release 17, slice RL (see platformBackoff.ts, "The adaptive pace"): +// - `rate_limit` also doubles the platform's pace, and a backoff that has +// failed at the cap `pacing.holdAfterFailsAtCap` times HOLDS the platform: +// its next try is a probe `pacing.holdProbeMinutes` away. +// - `subs_rate_limit` (the subtitle fetch alone answered 429 and the media came +// down) retires the video like any success and defers its SUBTITLES — the +// platform's backoff, hold and pace are not touched, either way. +// - a clean unit (`transcribed` with no subtitle 429) clears the backoff and +// any hold — a clean probe is what lifts one — and counts toward the pace's +// decay. import type { DownloadFailureClass } from "../lib/availability"; import { + type PacingSettings, type PlatformBackoffState, + type PlatformHoldState, + type PlatformPaceState, + type SubtitleDeferralState, type VideoDeferralState, - clearBackoff, + SUBTITLE_HOLD_AFTER, + deferSubtitles, deferVideo, - nextBackoff, + escalatePlatform, pruneDeferred, + pruneSubtitleDeferrals, + settlePlatformClean, } from "./platformBackoff"; +import { staticSleepRequestsSeconds } from "../ytdlp/platformArgs.mjs"; export type UnitOutcomeState = { platformBackoff: PlatformBackoffState; videoDeferrals: VideoDeferralState; + platformPace: PlatformPaceState; + platformHolds: PlatformHoldState; + subtitleDeferrals: SubtitleDeferralState; }; export type FinishedUnit = { @@ -55,22 +77,32 @@ export function applyUnitOutcome( unit: FinishedUnit, now: number, rand: () => number = Math.random, + pacing?: PacingSettings, ): UnitOutcomeEffect { pruneDeferred(state.videoDeferrals, now); + pruneSubtitleDeferrals(state.subtitleDeferrals, now); const pf = unit.platform; if ( pf !== null && (unit.failureClass === "rate_limit" || unit.failureClass === "network") ) { - const entry = nextBackoff(state.platformBackoff[pf], now, rand); - state.platformBackoff[pf] = entry; + const fx = escalatePlatform(state, pf, unit.failureClass, now, { + baseSeconds: staticSleepRequestsSeconds(pf), + rand, + ...(pacing ? { pacing } : {}), + }); + const entry = fx.entry; const secs = Math.round((entry.until - now) / 1000); - const head = `Auto-download: ${pf} ${unit.failureClass} — backing off ${secs}s (attempt ${entry.fails}).`; + const head = fx.held + ? `Auto-download: ${pf} ${unit.failureClass} — ${fx.enteredHold ? "held" : "still held"} after ${entry.fails} failures in a row; next probe in ${secs}s.` + : `Auto-download: ${pf} ${unit.failureClass} — backing off ${secs}s (attempt ${entry.fails}).`; + const pace = + unit.failureClass === "rate_limit" ? ` Pace now ${fx.paceSeconds}s between requests.` : ""; if (unit.failureClass === "rate_limit") { deferVideo(state.videoDeferrals, unit.videoId, unit.channelSlug, now); return { markCompleted: false, - line: `${head} ${unit.videoId} deferred 6h; next video after cooldown.`, + line: `${head} ${unit.videoId} deferred 6h; next video after cooldown.${pace}`, }; } return { @@ -78,8 +110,36 @@ export function applyUnitOutcome( line: `${head} ${unit.videoId} will retry after cooldown.`, }; } + if (unit.failureClass === "subs_rate_limit") { + // The media came down; only the subtitle fetch was refused. Not a platform + // signal (the timedtext 429 is per video): nothing platform-wide moves. + const d = deferSubtitles(state.subtitleDeferrals, unit.videoId, unit.channelSlug, now); + const left = + d.count >= SUBTITLE_HOLD_AFTER + ? "left alone for 7 days" + : "deferred 6h for download-missing-subs"; + return { + markCompleted: true, + line: `Auto-download: ${unit.videoId} downloaded; its subtitles were rate-limited (${d.count}×) — ${left}. ${pf ?? "the platform"} is not backed off.`, + }; + } if (pf !== null && unit.outcome === "transcribed") { - clearBackoff(state.platformBackoff, pf); + const fx = settlePlatformClean(state, pf, { + baseSeconds: staticSleepRequestsSeconds(pf), + ...(pacing ? { pacing } : {}), + }); + if (fx.releasedHold) { + return { + markCompleted: true, + line: `Auto-download: ${pf} probe came back clean — the hold and the backoff are cleared.`, + }; + } + if (fx.decayed) { + return { + markCompleted: true, + line: `Auto-download: ${pf} pace eased to ${fx.paceSeconds}s between requests.`, + }; + } } return { markCompleted: true, line: null }; } diff --git a/common/lib/availability.test.ts b/common/lib/availability.test.ts @@ -2,7 +2,9 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import { classifyDownloadFailure, + hasSubtitleRateLimit, isSoftBlock, + isSubtitleRateLimitOnly, parseUnavailableFromStderr, } from "./availability"; @@ -126,3 +128,59 @@ test("a genuinely removed or unavailable video is still deleted / per_video", () assert.equal(parseUnavailableFromStderr(members), "members_only"); assert.equal(classifyDownloadFailure(members, "members_only"), "per_video"); }); + +// ── release 17, slice RL: a subtitle 429 is not a platform failure ────────── + +const SUB_429 = + "ERROR: [youtube] abc: Unable to download video subtitles for 'en': HTTP Error 429: Too Many Requests"; + +test("a subtitle 429 alone is subs_rate_limit — as an ERROR or a WARNING", () => { + assert.equal(classifyDownloadFailure(SUB_429, "error"), "subs_rate_limit"); + assert.equal( + classifyDownloadFailure( + "WARNING: Unable to download video subtitles for 'en-orig': HTTP Error 429: Too Many Requests", + undefined, + ), + "subs_rate_limit", + ); + // yt-dlp's own re-extraction after a --load-info-json subtitle failure is a + // WARNING; a second subtitle 429 from the URL is still subtitle-only. + assert.equal( + classifyDownloadFailure( + [ + SUB_429, + "WARNING: The info failed to download: ERROR: Unable to download video subtitles for 'en': HTTP Error 429: Too Many Requests; trying with URL https://www.youtube.com/watch?v=abc", + SUB_429, + ].join("\n"), + "error", + ), + "subs_rate_limit", + ); + assert.equal(isSubtitleRateLimitOnly(SUB_429), true); +}); + +test("a subtitle 429 beside any other ERROR is a platform rate limit", () => { + const webpage429 = + "ERROR: [youtube] abc: Unable to download webpage: HTTP Error 429: Too Many Requests"; + assert.equal(classifyDownloadFailure(`${SUB_429}\n${webpage429}`, "error"), "rate_limit"); + assert.equal(classifyDownloadFailure(webpage429, "error"), "rate_limit"); + // The soft block and the bot check are platform signals, whatever else failed. + assert.equal( + classifyDownloadFailure( + `${SUB_429}\nWARNING: This content isn't available, try again later.`, + "error", + ), + "rate_limit", + ); + assert.equal( + isSubtitleRateLimitOnly(`${SUB_429}\nSign in to confirm you’re not a bot`), + false, + ); +}); + +test("a subtitle failure that is not a rate limit is not subs_rate_limit", () => { + const sub404 = + "ERROR: [youtube] abc: Unable to download video subtitles for 'en': HTTP Error 404: Not Found"; + assert.equal(hasSubtitleRateLimit(sub404), false); + assert.equal(classifyDownloadFailure(sub404, "error"), "unknown"); +}); diff --git a/common/lib/availability.ts b/common/lib/availability.ts @@ -231,12 +231,47 @@ export function parseUnavailableFromStderr(stderr: string): Availability { // `per_video` — the rest of the batch can continue. Rate-limit and network // errors are batch-level signals: continuing would just hammer the source or // waste cycles. Everything else is `unknown` and treated as fatal to be safe. +// +// `subs_rate_limit` (release 17, slice RL) is the one class that is not a +// failure of the download: the SUBTITLE fetch answered 429 and nothing else +// failed. YouTube's timedtext endpoint refuses per video while the webpage, +// player and media requests in the same spawn succeed (measured 2026-09-25 and +// 2026-10-01), so it says nothing about the platform: the downloader goes on to +// the media, the record is a success carrying this class, and only the video's +// subtitles are deferred (jobs/platformBackoff.ts). export type DownloadFailureClass = | "per_video" | "rate_limit" + | "subs_rate_limit" | "network" | "unknown"; +// yt-dlp's subtitle-fetch failure, as an ERROR (the default) or as a WARNING +// (under --ignore-errors, which is how the managed downloader asks it to carry +// on to the media): `Unable to download video subtitles for 'en': HTTP Error +// 429: Too Many Requests`. +const SUBTITLE_RATE_LIMIT_RE = + /unable to download video subtitles for [^\n]*?(?:http error 429|too many requests|rate[- ]?limit)/i; + +export function hasSubtitleRateLimit(stderr: string): boolean { + return SUBTITLE_RATE_LIMIT_RE.test(stderr); +} + +// True when a subtitle 429 is the ONLY failure in the tail: every ERROR line +// is a subtitle-download line (yt-dlp's own "The info failed to download … +// trying with URL" retry is a WARNING and is allowed), and neither the soft +// block nor the bot check — both platform signals — is anywhere in it. +export function isSubtitleRateLimitOnly(stderr: string): boolean { + if (!hasSubtitleRateLimit(stderr)) return false; + if (isSoftBlock(stderr) || isBotCheck(stderr)) return false; + for (const line of stderr.split("\n")) { + if (/^\s*ERROR:/i.test(line) && !/unable to download video subtitles/i.test(line)) { + return false; + } + } + return true; +} + const PER_VIDEO_CLASSES: ReadonlyArray<Availability> = [ "private", "members_only", @@ -269,6 +304,9 @@ export function classifyDownloadFailure( // that also names a removed video) must still back off. Erring this way costs // one cooldown; erring the other way keeps requesting into the block. if (isSoftBlock(stderrTail)) return "rate_limit"; + // The subtitle fetch alone was refused: the media is unaffected, and the + // platform is not to back off for it (release 17, slice RL). + if (isSubtitleRateLimitOnly(stderrTail)) return "subs_rate_limit"; if ( availabilityClass !== undefined && PER_VIDEO_CLASSES.includes(availabilityClass) diff --git a/common/lib/downloadOutcome.ts b/common/lib/downloadOutcome.ts @@ -66,7 +66,12 @@ export type DownloadAttemptKind = // The live-chat pass a "chat-only" filter verdict runs INSTEAD of a download: // --skip-download --write-subs --sub-langs live_chat, no media, no archive // line. Recorded with n: 1, since it is the only real attempt there is. - | "live-chat-only"; + | "live-chat-only" + // The primary again, with every subtitle refused, after a primary that + // failed ONLY on its subtitle fetch (a 429): the media is what the download + // is for, and the subtitles are deferred (release 17, slice RL). Recorded + // with n: 1 beside the primary it re-runs. + | "primary-without-subs"; export type AudioCheckProbeVerdict = "clean" | "partial" | "malformed"; @@ -101,7 +106,9 @@ export type DownloadOutcomeRecord = { // 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. + // on success / skipped-filtered — except `subs_rate_limit`, which is set on a + // SUCCESS whose subtitle fetch alone answered 429: the media came down, and + // the video's subtitles are deferred (release 17, slice RL). failureClass?: DownloadFailureClass; fellBackToTranscribe?: boolean; // Set when status is "failed-short-audio": the measured shortfall, so the UI diff --git a/common/lib/settingsDocs.ts b/common/lib/settingsDocs.ts @@ -16,6 +16,7 @@ import { BACKFILL_SETTINGS_FIELD_DOCS, BUILD_PIPELINE_SETTINGS_FIELD_DOCS, DIARIZATION_SETTINGS_FIELD_DOCS, + PACING_SETTINGS_FIELD_DOCS, DIGEST_SETTINGS_FIELD_DOCS, SAVED_VIDEO_BACKUP_SETTINGS_FIELD_DOCS, SOCIAL_LINK_FIELD_DOCS, @@ -126,6 +127,13 @@ export function blockTables(d: SiteSettings): Partial<Record<keyof SiteSettings, { path: "workers[].remote", docs: REMOTE_WORKER_CONFIG_FIELD_DOCS }, { path: "workers[].llm", docs: LLM_WORKER_CONFIG_FIELD_DOCS }, ], + pacing: [ + { + path: "pacing", + docs: PACING_SETTINGS_FIELD_DOCS, + defaults: fromObject(d.pacing), + }, + ], archiveStorage: [ { path: "archiveStorage", diff --git a/common/lib/settingsSchema.test.ts b/common/lib/settingsSchema.test.ts @@ -32,6 +32,7 @@ import type { BuildPipelineSettings, DiarizationSettings, DigestSettings, + PacingSettingsBlock, ReportDebouncePreset, SavedVideoBackupSettings, SocialLink, @@ -66,6 +67,8 @@ type PreSchemaSiteSettings = { // Release 16 slice XL — the one key it adds. social: SocialSettings; sleepBetweenDownloadsSeconds: number; + // Release 17 slice RL — the adaptive pace and the hold. + pacing: PacingSettingsBlock; downloadFormat: DownloadFormatPreset; minFreeDiskGB: number; resumeMarginGB: number; @@ -95,7 +98,7 @@ type PreSchemaSiteSettings = { type Same<A, B> = [A] extends [B] ? ([B] extends [A] ? true : false) : false; const shapeUnchanged: Same<SiteSettings, PreSchemaSiteSettings> = true; -test("SiteSettings keeps its 32 fields, in file order", () => { +test("SiteSettings keeps its 33 fields, in file order", () => { assert.equal(shapeUnchanged, true); assert.deepEqual(Object.keys(siteSettingsSchema.shape), [ "adminTitle", @@ -107,6 +110,7 @@ test("SiteSettings keeps its 32 fields, in file order", () => { "cookieMode", "social", "sleepBetweenDownloadsSeconds", + "pacing", "downloadFormat", "minFreeDiskGB", "resumeMarginGB", diff --git a/common/lib/settingsSchema.ts b/common/lib/settingsSchema.ts @@ -636,6 +636,74 @@ export const ARCHIVE_STORAGE_SETTINGS_FIELD_DOCS: FieldDocs<ArchiveStorageSettin "happen.", }; +// Each field is documented in PACING_SETTINGS_FIELD_DOCS below (rendered into SETTINGS.md). +export type PacingSettingsBlock = { + sleepRequestsCapSeconds: number; + decayAfterCleanUnits: number; + holdAfterFailsAtCap: number; + holdProbeMinutes: number; +}; + +export const PACING_SETTINGS_FIELD_DOCS: FieldDocs<PacingSettingsBlock> = { + sleepRequestsCapSeconds: + "Ceiling (seconds) on a platform's adaptive `--sleep-requests`. The pace " + + "starts at the platform's fixed value (1 s for YouTube and Rumble) and " + + "doubles on every platform-level rate limit until it reaches this. A " + + "subtitle-only 429 never moves it. Clamped 1–120; default 16.", + decayAfterCleanUnits: + "How many clean auto-download units on a platform halve its pace one step " + + "back toward the fixed value. Clamped 1–1000; default 5.", + holdAfterFailsAtCap: + "How many consecutive failures AT the 30-minute cooldown cap put a " + + "platform in a hold: the lane then runs one probe unit per " + + "holdProbeMinutes instead of one per cooldown, and a manual Sync or " + + "download on it is refused with the next probe's time. A clean probe " + + "clears the hold and the backoff. Clamped 1–100; default 3.", + holdProbeMinutes: + "Minutes between probes while a platform is held. Clamped 1–1440; " + + "default 60.", +}; + +export const PACING_DEFAULT_SLEEP_REQUESTS_CAP_SECONDS = 16; +export const PACING_DEFAULT_DECAY_AFTER_CLEAN_UNITS = 5; +export const PACING_DEFAULT_HOLD_AFTER_FAILS_AT_CAP = 3; +export const PACING_DEFAULT_HOLD_PROBE_MINUTES = 60; + +function clampInt(v: unknown, min: number, max: number, dflt: number): number { + if (typeof v !== "number" || !Number.isFinite(v)) return dflt; + return Math.min(max, Math.max(min, Math.floor(v))); +} + +export function sanitizePacing(v: unknown): PacingSettingsBlock { + const r = (v && typeof v === "object" ? v : {}) as Record<string, unknown>; + return { + sleepRequestsCapSeconds: clampInt( + r.sleepRequestsCapSeconds, + 1, + 120, + PACING_DEFAULT_SLEEP_REQUESTS_CAP_SECONDS, + ), + decayAfterCleanUnits: clampInt( + r.decayAfterCleanUnits, + 1, + 1000, + PACING_DEFAULT_DECAY_AFTER_CLEAN_UNITS, + ), + holdAfterFailsAtCap: clampInt( + r.holdAfterFailsAtCap, + 1, + 100, + PACING_DEFAULT_HOLD_AFTER_FAILS_AT_CAP, + ), + holdProbeMinutes: clampInt( + r.holdProbeMinutes, + 1, + 1440, + PACING_DEFAULT_HOLD_PROBE_MINUTES, + ), + }; +} + export const SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS = 600; export const SLEEP_BETWEEN_DOWNLOADS_DEFAULT_SECONDS = 10; @@ -1480,7 +1548,10 @@ export const siteSettingsSchema = z.object({ "Per-platform settings of the social-post fetchers. Today one key: where the X fetchers' login comes from (`social.x.cookieSource`, chosen in the X account session section of /settings). See common/social/xCookieSource.ts.", ), sleepBetweenDownloadsSeconds: settingsField((v): number => clampSleepBetweenDownloadsSeconds(v)).describe( - "Pause (seconds) inserted between per-video yt-dlp invocations in managed batch downloads. yt-dlp's own `-t sleep` only paces requests within one invocation, so without this the managed loop hammers the source IP back-to-back. 0 disables. Per-channel override available.", + "Pause (seconds) inserted between per-video yt-dlp invocations in managed batch downloads, and between two auto-download units on one platform (release 17; the lane ignored it before). yt-dlp's own `-t sleep` only paces requests within one invocation, so without this the managed loop hammers the source IP back-to-back. The adaptive pace above its base (see `pacing`) is added to it. 0 disables. Per-channel override available for batch downloads.", + ), + pacing: settingsField((v): PacingSettingsBlock => sanitizePacing(v)).describe( + "How the download pace adapts to rate limits, per platform (release 17). Every yt-dlp spawn against a platform paces its requests (`--sleep-requests`) at the platform's adaptive pace, and the download lane waits sleepBetweenDownloadsSeconds plus the pace above its fixed value between units. A rate limit doubles the pace; clean units ease it back; a rate limit that outlasts the cooldown cap holds the platform to one probe at a time. The live pace, cooldowns and holds are in `.auto-queue/state.json`, shown on /operations/download. See common/jobs/platformBackoff.ts.", ), downloadFormat: settingsField((v): DownloadFormatPreset => isDownloadFormatPreset(v) ? v : "auto").describe( diff --git a/common/views/activeJobs.ts b/common/views/activeJobs.ts @@ -400,6 +400,10 @@ function autoIdleNote( return "every pending platform is in a rate-limit cooldown"; case "deferred": return "every pending video was rate-limited recently and is deferred"; + case "held": + return "every pending platform is held after repeated rate limits — one probe at a time"; + case "paced": + return "every pending platform is pausing between downloads"; case "no-workers": return "no enabled worker"; case "workers-paused": diff --git a/common/ytdlp/channelArgs.test.ts b/common/ytdlp/channelArgs.test.ts @@ -1,6 +1,14 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { channelExtraArgs, platformArgs, PLATFORM_ARGS } from "./channelArgs"; +import { + channelExtraArgs, + channelPaceSeconds, + pacedPlatformArgs, + platformArgs, + PLATFORM_ARGS, + staticSleepRequestsSeconds, + withSleepRequests, +} from "./channelArgs"; import type { ChannelConfig } from "../lib/channelConfig"; // Run with: @@ -80,3 +88,62 @@ test("platformArgs returns a copy, not the table's array", () => { a.push("--mutated"); assert.deepEqual(PLATFORM_ARGS.rumble, RUMBLE); }); + +// ── release 17, slice RL: the adaptive pace ───────────────────────────────── + +test("the pace raises --sleep-requests in place, and a channel's own still wins", () => { + assert.deepEqual( + channelExtraArgs(cfg({ url: "https://www.youtube.com/@x" }), undefined, 4), + ["--sleep-requests", "4"], + ); + assert.deepEqual( + channelExtraArgs(cfg({ url: "https://rumble.com/c/x" }), undefined, 8), + ["--impersonate", "chrome", "--sleep-requests", "8"], + ); + const own = channelExtraArgs( + cfg({ url: "https://www.youtube.com/@x", ytdlpExtraArgs: ["--sleep-requests", "3"] }), + undefined, + 8, + ); + assert.deepEqual(own, ["--sleep-requests", "8", "--sleep-requests", "3"]); +}); + +test("a platform with no static pace gets one appended; a pace never lowers one", () => { + assert.deepEqual(withSleepRequests([], 2), ["--sleep-requests", "2"]); + assert.deepEqual(withSleepRequests(["--sleep-requests", "1"], 0.5), ["--sleep-requests", "1"]); + assert.deepEqual(withSleepRequests(["--sleep-requests", "1"], undefined), ["--sleep-requests", "1"]); + assert.deepEqual(pacedPlatformArgs("odysee", 2), ["--sleep-requests", "2"]); + assert.equal(staticSleepRequestsSeconds("youtube"), 1); + assert.equal(staticSleepRequestsSeconds("odysee"), 0); + assert.equal(staticSleepRequestsSeconds(null), 0); +}); + +test("with no pace recorded, the static args are unchanged", () => { + globalThis.__yttAutoQueueState__ = undefined; + assert.equal(channelPaceSeconds(cfg({ url: "https://www.youtube.com/@x" })), 1); + assert.deepEqual( + channelExtraArgs(cfg({ url: "https://www.youtube.com/@x" })), + ["--sleep-requests", "1"], + ); +}); + +test("the live pace (the shared state) reaches every channelExtraArgs call", () => { + globalThis.__yttAutoQueueState__ = { + state: null, + stateFile: null, + loading: null, + loadingFile: null, + lastPace: { youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0 } }, + }; + try { + assert.deepEqual( + channelExtraArgs(cfg({ url: "https://www.youtube.com/@x" })), + ["--sleep-requests", "4"], + ); + assert.equal(channelPaceSeconds(cfg({ url: "https://www.youtube.com/@x" })), 4); + // Another platform is untouched. + assert.deepEqual(channelExtraArgs(cfg({ url: "https://rumble.com/c/x" })), RUMBLE); + } finally { + globalThis.__yttAutoQueueState__ = undefined; + } +}); diff --git a/common/ytdlp/channelArgs.ts b/common/ytdlp/channelArgs.ts @@ -1,6 +1,7 @@ // The per-channel argv every yt-dlp invocation in this repo appends: the cookie // source when the resolved policy calls for one, then the platform's fixed -// args (PLATFORM_ARGS), then the channel's own `ytdlpExtraArgs` verbatim. +// args (PLATFORM_ARGS) at the platform's CURRENT pace, then the channel's own +// `ytdlpExtraArgs` verbatim. // // Lifted out of ytdlp/downloadOneManaged.ts (where it was `channelConfigArgs`, // which still delegates here) because the clip-window fetch has to honour the @@ -10,18 +11,37 @@ // // This is the ONE builder. `configArgs` (runYtdlp.ts), the metadata scan, the // quick availability check and the new-channel probe all go through it (the -// probe has no config yet, so it calls `platformArgs` directly). Before +// probe has no config yet, so it calls `pacedPlatformArgs` directly). Before // release 5 there were four copies and the probe had none, so a Rumble // channel could not even be created. +// +// THE PACE ADAPTS (release 17, slice RL). `--sleep-requests` is the platform's +// static value until a platform-level rate limit doubles it; the current value +// lives in the download lane's persisted state (`platformPace`, see +// jobs/platformBackoff.ts) and is read here synchronously through +// `livePlatformPaceSeconds`, so EVERY spawn against the platform — listing, +// prefetch, primary, availability, metadata scan, clip — slows down together. +// The channel's own args still come last and win. import type { ChannelConfig } from "../lib/channelConfig"; import { cookieArgs } from "../lib/cookiePolicy"; import { detectPlatform, type Platform } from "../lib/platform"; +import { livePlatformPaceSeconds } from "../jobs/autoQueueState"; // The platform args table lives in `platformArgs.mjs` (plain JS so umtool's // `.mjs` scripts can import the same copy); re-exported here for TS callers. -import { platformArgs } from "./platformArgs.mjs"; -export { PLATFORM_ARGS, platformArgs, platformArgsForUrl } from "./platformArgs.mjs"; +import { + platformArgs, + staticSleepRequestsSeconds, + withSleepRequests, +} from "./platformArgs.mjs"; +export { + PLATFORM_ARGS, + platformArgs, + platformArgsForUrl, + staticSleepRequestsSeconds, + withSleepRequests, +} from "./platformArgs.mjs"; export function channelPlatform( config: Pick<ChannelConfig, "platform" | "url">, @@ -29,13 +49,46 @@ export function channelPlatform( return config.platform ?? detectPlatform(config.url); } +// The key the pacing state is kept under — the same one the rate-limit +// cooldown uses (`detectPlatform(url) ?? "unknown"`, autoRunner/runYtdlp), so a +// 429 recorded by any path paces every path. +export function pacingPlatformKey(config: Pick<ChannelConfig, "url">): string { + return detectPlatform(config.url) ?? "unknown"; +} + +// The platform's current `--sleep-requests`: the adaptive pace while one is +// recorded, else the static value. +export function channelPaceSeconds( + config: Pick<ChannelConfig, "platform" | "url">, +): number { + return ( + livePlatformPaceSeconds(pacingPlatformKey(config)) ?? + staticSleepRequestsSeconds(channelPlatform(config)) + ); +} + +// PLATFORM_ARGS for a platform at a given pace (default: its current one) — +// for the spawns that have a URL and no channel config. +export function pacedPlatformArgs( + platform: Platform | null | undefined, + paceSeconds: number | undefined = platform + ? livePlatformPaceSeconds(platform) + : undefined, +): string[] { + return withSleepRequests(platformArgs(platform), paceSeconds); +} + export function channelExtraArgs( config: ChannelConfig, cookies?: string, + // The pace to apply; default: the platform's current one. Tests pass it. + paceSeconds: number | undefined = livePlatformPaceSeconds( + pacingPlatformKey(config), + ), ): string[] { const args: string[] = []; args.push(...cookieArgs(cookies)); - args.push(...platformArgs(channelPlatform(config))); + args.push(...withSleepRequests(platformArgs(channelPlatform(config)), paceSeconds)); if (config.ytdlpExtraArgs?.length) args.push(...config.ytdlpExtraArgs); return args; } diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts @@ -5,6 +5,7 @@ import { execa } from "execa"; import { AUTH_RETRY_CLASSES, classifyDownloadFailure, + hasSubtitleRateLimit, parseUnavailableFromStderr, } from "../lib/availability"; import { @@ -63,7 +64,7 @@ import { import type { Paths } from "../lib/paths"; import { transcribeWithWorker } from "../controller/transcribeOne"; import { extractVideoId, outputArgsForUrl } from "./runYtdlp"; -import { channelExtraArgs } from "./channelArgs"; +import { channelExtraArgs, channelPaceSeconds } from "./channelArgs"; import { ARCHIVE_MARKER, runOneYtdlp as runOneYtdlpRaw, @@ -166,6 +167,16 @@ export type ManagedDownloadOpts = { // prefetch pass already wrote (fed in via sourceArgs as --load-info-json), so // we drop --write-info-json (the file is already on disk). When false, the // legacy single-call behavior: write the info json during this download. +// +// A SUBTITLE 429 IS A WARNING HERE, NOT THE END (release 17, slice RL). Under +// yt-dlp's default, a subtitle file it cannot fetch raises: the video fails, +// and — reading from --load-info-json — yt-dlp re-extracts from the URL once and +// asks the throttled endpoint again. `--ignore-errors` makes it report the +// failure as `WARNING: Unable to download video subtitles for …` and carry on +// (exit 0, the line still in the log for the classifier); the caller then sees +// "no transcript" and fetches the media (attempt 3). Any other error is still +// an ERROR and a non-zero exit. `--sleep-subtitles` is the pace before each +// subtitle request — the request YouTube throttles — never below `-t sleep`'s 5. function youtubeHandlingArgs( config: ChannelConfig, reuseInfoJson = false, @@ -177,10 +188,20 @@ function youtubeHandlingArgs( config.subLangs ?? "en.*,live_chat", ]; if (!reuseInfoJson) args.push("--write-info-json"); - args.push("--skip-download", "-t", "sleep"); + args.push( + "--skip-download", + "-t", + "sleep", + "--sleep-subtitles", + String(Math.max(YOUTUBE_SLEEP_SUBTITLES_FLOOR, channelPaceSeconds(config))), + "--ignore-errors", + ); return args; } +// `-t sleep`'s own `--sleep-subtitles`: the adaptive pace only ever raises it. +const YOUTUBE_SLEEP_SUBTITLES_FLOOR = 5; + // Format/extract args for a transcribe-handling download under a resolved // persistence plan. "ytdlp" mode is the legacy path: yt-dlp extracts the audio // itself (-x) and optionally keeps its bestaudio source via -k. "app" mode omits @@ -1187,6 +1208,72 @@ async function runManagedDownload( }); if (primaryRes.archiveLine) lastArchiveLine = primaryRes.archiveLine; + // ---------- A subtitle 429 is not a failed download (release 17, RL) ---------- + // YouTube's timedtext endpoint refuses per video while the media requests + // succeed. When the subtitle fetch is all that failed, the download goes on + // to the media and the record carries `subs_rate_limit` (a success): the + // caller defers the video's SUBTITLES, never the platform. + // - youtube handling runs the primary with --ignore-errors, so yt-dlp + // exits 0 with the failure as a WARNING; the media pass below fetches + // the audio when no transcript track arrived. + // - any other primary that died on its subtitles alone (a channel whose + // own args ask for subtitles) is run ONCE more with every subtitle + // refused: a second spawn, the media, no timedtext request. + let subsRateLimited = false; + if ( + !audioCheckEnabled && + !audioCheckCorruptSource && + !audioCheckCorruptFullSource + ) { + if ( + attemptSucceeded(primaryRes.exitCode) && + hasSubtitleRateLimit(primaryRes.stderrTail) + ) { + subsRateLimited = true; + } else if ( + !attemptSucceeded(primaryRes.exitCode) && + !opts.signal.aborted && + classifyDownloadFailure(primaryRes.stderrTail, primaryAvail) === + "subs_rate_limit" + ) { + subsRateLimited = true; + opts.onLog( + `The subtitle fetch for ${canonicalId ?? opts.videoUrl} was rate-limited (HTTP 429) and nothing else failed; ` + + `downloading without subtitles — they are deferred, the platform is not backed off.\n`, + ); + const noSubsRes = await runOneYtdlp(opts, channelDir, [ + "--ignore-config", + "--restrict-filenames", + ...FULL_LOG_PROGRESS_ARGS, + ...mediaArgs, + "--print", + `after_video:${ARCHIVE_MARKER} %(extractor)s %(id)s`, + ...channelConfigArgs(opts.channelConfig, primaryCookies), + // After the channel's own args: yt-dlp keeps the last occurrence. + "--no-write-subs", + "--no-write-auto-subs", + ...sourceArgs(opts.videoUrl, infoJsonPath), + ]); + lastFullTail = noSubsRes.stderrTail; + attempts.push({ + n: 1, + kind: "primary-without-subs", + handling: opts.channelConfig.handling, + usedCookies: Boolean(primaryCookies), + ytdlpExitCode: noSubsRes.exitCode, + availabilityClass: attemptSucceeded(noSubsRes.exitCode) + ? undefined + : parseUnavailableFromStderr(noSubsRes.stderrTail), + error: attemptSucceeded(noSubsRes.exitCode) + ? undefined + : trimError(noSubsRes.stderrTail), + }); + if (noSubsRes.archiveLine) lastArchiveLine = noSubsRes.archiveLine; + primaryRes = noSubsRes; + if (!attemptSucceeded(noSubsRes.exitCode)) subsRateLimited = false; + } + } + let lastSucceeded = !audioCheckCorruptSource && !audioCheckCorruptFullSource && @@ -1339,7 +1426,10 @@ async function runManagedDownload( ) { const hasTranscript = await hasAnyTranscriptOnDisk(videoDir); const noCaptions = await metadataReportsNoCaptions(videoDir); - const noSubsFallback = !hasTranscript && noCaptions; + // A subtitle 429 left no transcript: the media is fetched exactly as for a + // video with no captions, and the subtitles wait (release 17, slice RL). + const subsDeferredFallback = subsRateLimited && !hasTranscript; + const noSubsFallback = !hasTranscript && (noCaptions || subsDeferredFallback); // Forced only when the fallback would NOT have run on its own: a video with // no transcript and no captions takes today's path whoever asked. const forced = !noSubsFallback && opts.forceMedia === true; @@ -1368,6 +1458,11 @@ async function runManagedDownload( hasTranscript ? "a transcript is on disk" : "captions exist" }\n`, ); + } else if (subsDeferredFallback) { + opts.onLog( + `The subtitles for ${videoId} were rate-limited (HTTP 429); downloading the media anyway — ` + + `the subtitles are deferred for download-missing-subs, the platform is not backed off${plan.persist ? " (keeping source video)" : ""}.\n`, + ); } else { opts.onLog( `No subs available for ${videoId}; falling back to audio download + whisper${plan.persist ? " (keeping source video)" : ""}.\n`, @@ -1406,7 +1501,10 @@ async function runManagedDownload( ...channelConfigArgs(fallbackConfig, fallbackCookieOverride), // yt-dlp takes the LAST occurrence of an option, so the refusals go // after the channel's own args (the chat-only pass's rule). - ...(keepTranscript ? ["--no-write-subs", "--no-write-auto-subs"] : []), + // A deferred subtitle fetch is not asked again in the same minute. + ...(keepTranscript || subsDeferredFallback + ? ["--no-write-subs", "--no-write-auto-subs"] + : []), ...sourceArgs( opts.videoUrl, path.join(videoDir, "metadata.info.json"), @@ -1581,6 +1679,7 @@ async function runManagedDownload( lastFullTail, fellBackToTranscribe, shortAudio: shortAudioInfo, + subsRateLimited, }); } @@ -1599,17 +1698,23 @@ async function writeOutcome( lastFullTail: string; fellBackToTranscribe?: boolean; shortAudio?: DownloadOutcomeRecord["shortAudio"]; + // The subtitle fetch alone answered 429 and the download went on. + subsRateLimited?: boolean; }, ): Promise<DownloadOutcomeRecord> { const { status, attempts } = o; 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. + // is still caught. Skipped for successes and filter-skips — except that a + // success whose subtitles were rate-limited says so (`subs_rate_limit`), so + // its caller defers the subtitles (release 17, slice RL). const isFailure = status === "failed" || status === "failed-corrupt-source"; const failureClass = isFailure ? classifyDownloadFailure(o.lastFullTail, attempts.at(-1)?.availabilityClass) - : undefined; + : o.subsRateLimited && status.startsWith("ok") + ? ("subs_rate_limit" as const) + : undefined; const record: DownloadOutcomeRecord = { videoId: o.videoId, webpageUrl: opts.videoUrl, diff --git a/common/ytdlp/managedDownloadsSleep.test.ts b/common/ytdlp/managedDownloadsSleep.test.ts @@ -246,3 +246,61 @@ test("declinedWithoutMediaFetch: a chat-only pass is a fetch", () => { false, ); }); + +// ── release 17, slice RL ───────────────────────────────────────────────────── + +// Runs two fetched videos with an injected pace and returns the slept ms. +async function gapFor( + paceSeconds: number, + outcomes: DownloadOutcomeRecord[] = [FETCHED, FETCHED], +): Promise<{ slept: number[]; deferred: string[] }> { + const slept: number[] = []; + const deferred: string[] = []; + let i = 0; + const urls = outcomes.map( + (_, n) => `https://www.youtube.com/watch?v=vid${String(n).padStart(8, "0")}`, + ); + const channelConfig = { + handling: "youtube", + url: "https://www.youtube.com/@c/videos", + } as never; + await runManagedDownloads( + { + channelSlug: "c", + mode: "download-missing" as never, + channelConfig, + paths: getPaths(), + onLog: () => {}, + signal: new AbortController().signal, + abortOnError: false, + }, + urls, + channelConfig, + undefined, + { + downloadOne: async () => outcomes[i++], + sleep: async (ms) => { + slept.push(ms); + }, + paceSeconds: () => paceSeconds, + recordSubtitleDeferral: async (id, slug) => { + deferred.push(`${slug}/${id}`); + }, + }, + ); + return { slept, deferred }; +} + +test("the gap is sleepBetweenDownloadsSeconds plus the pace above its base", async () => { + // youtube's base pace is 1 s: at the base the gap is the setting alone. + assert.deepEqual((await gapFor(1)).slept, [30_000]); + // After two rate limits (pace 4 s) the batch waits 3 s more per video. + assert.deepEqual((await gapFor(4)).slept, [33_000]); +}); + +test("a subtitle 429 on a success is recorded as a subtitle deferral, and the batch goes on", async () => { + const subs = { ...outcome("ok", [PREFETCH, PRIMARY], "subs_rate_limit"), videoId: "s1" }; + const r = await gapFor(1, [subs, FETCHED]); + assert.deepEqual(r.deferred, ["c/s1"]); + assert.deepEqual(r.slept, [30_000]); +}); diff --git a/common/ytdlp/platformArgs.mjs b/common/ytdlp/platformArgs.mjs @@ -58,3 +58,43 @@ export function platformArgs(platform) { export function platformArgsForUrl(url) { return platformArgs(detectPlatform(url)); } + +// THE STATIC PACE a platform's spawns carry (`--sleep-requests` in its +// PLATFORM_ARGS entry), or 0 for a platform with none. The adaptive pace +// (common/jobs/platformBackoff.ts, release 17 slice RL) starts here, doubles on +// a rate limit and decays back to it. +/** + * @param {string | null | undefined} platform + * @returns {number} + */ +export function staticSleepRequestsSeconds(platform) { + if (!platform) return 0; + /** @type {Record<string, readonly string[] | undefined>} */ + const table = PLATFORM_ARGS; + const args = table[platform]; + if (!args) return 0; + const i = args.indexOf("--sleep-requests"); + const v = i >= 0 ? Number(args[i + 1]) : 0; + return Number.isFinite(v) && v > 0 ? v : 0; +} + +// `args` with its `--sleep-requests` raised to `seconds` — replaced when the +// args carry a lower one, appended when they carry none. Never LOWERS a pace: +// a value below what the args already say is ignored. Returns a copy. +/** + * @param {readonly string[]} args + * @param {number | undefined | null} seconds + * @returns {string[]} + */ +export function withSleepRequests(args, seconds) { + const out = [...args]; + if (seconds == null || !Number.isFinite(seconds) || seconds <= 0) return out; + const i = out.indexOf("--sleep-requests"); + if (i < 0) { + out.push("--sleep-requests", String(seconds)); + return out; + } + const cur = Number(out[i + 1]); + if (!Number.isFinite(cur) || seconds > cur) out[i + 1] = String(seconds); + return out; +} diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -13,7 +13,13 @@ import { getSettings } from "../lib/settings"; import { patchChannelConfig } from "../controller/channels"; import { diskGate } from "../lib/diskSpace"; import { detectPlatform } from "../lib/platform"; -import { channelExtraArgs, platformArgs } from "./channelArgs"; +import { + channelExtraArgs, + channelPaceSeconds, + channelPlatform, + pacedPlatformArgs, + staticSleepRequestsSeconds, +} from "./channelArgs"; import { isRealAudioFile } from "../lib/videoStatus"; import { readVttProvenance } from "../lib/subtitleProvenance"; import { extractVideoId } from "../lib/videoId"; @@ -55,9 +61,21 @@ import { isFullSweepDue, resolveFullSweepIntervalMinutes, } from "../jobs/deepSync"; -import { platformCooldownRemainingMs } from "../jobs/downloadBackoff"; +import { + clearSubtitleDeferral, + platformCooldownRemainingMs, + readSubtitleDeferrals, + recordSubtitleDeferral, +} from "../jobs/downloadBackoff"; +import { + SUBTITLE_HOLD_AFTER, + downloadGapMs, + isSubtitleDeferred, +} from "../jobs/platformBackoff"; +import { hasSubtitleRateLimit } from "../lib/availability"; import { computeKeepWindow, type KeepWindow } from "../controller/keptVideos"; import { downloadOneManaged } from "./downloadOneManaged"; +import { STDERR_TAIL_BYTES } from "./runOneYtdlp"; import { resolveDownloadFormatPreset, resolveDownloadFormatSelector, @@ -491,8 +509,8 @@ export async function probeChannelMeta(opts: { "--print", "%(channel,uploader,playlist_title,playlist,uploader_id,title)s", // No channel config exists yet; the platform's fixed args still apply - // (a Rumble probe 403s without `--impersonate`). - ...platformArgs(detectPlatform(url)), + // (a Rumble probe 403s without `--impersonate`), at its current pace. + ...pacedPlatformArgs(detectPlatform(url)), url, ]; log(`$ ${paths.ytdlpBin} ${args.join(" ")}\n`); @@ -901,6 +919,10 @@ async function downloadPlaylistManaged( export type ManagedDownloadsDeps = { downloadOne?: typeof downloadOneManaged; sleep?: (ms: number, signal: AbortSignal) => Promise<void>; + // Where a video's subtitle 429 is recorded (release 17, slice RL). + recordSubtitleDeferral?: (videoId: string, channelSlug: string) => Promise<unknown>; + // The platform's current pace (default: channelPaceSeconds). + paceSeconds?: () => number; }; // True when the download filter declined the video before any media request: @@ -972,6 +994,14 @@ export async function runManagedDownloads( const sleepSeconds = opts.channelConfig.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds; + // THE PACE ABOVE ITS BASE joins the gap (release 17, slice RL): after a + // rate limit doubled the platform's pace, a batch spaces its videos further + // apart too. Read per gap, so a 429 in this batch slows the rest of it. + const basePace = staticSleepRequestsSeconds(channelPlatform(effectiveChannelConfig)); + const paceNow = deps.paceSeconds ?? (() => channelPaceSeconds(effectiveChannelConfig)); + const recordSubs = + deps.recordSubtitleDeferral ?? + ((id: string, slug: string) => recordSubtitleDeferral(id, slug, opts.paths)); // Title-filter rejections from this batch, written to the channel-level // metadata-scan store ONCE at the end. Per-video writes would rewrite the // whole file for every non-matching video — the exact shape of work a @@ -1050,6 +1080,16 @@ export async function runManagedDownloads( } finally { task?.end(); } + // A subtitle 429 on a download that otherwise succeeded: the video's + // subtitles are deferred, the platform is not backed off, the batch + // goes on (release 17, slice RL). Best-effort, like the backoff. + if (outcome.failureClass === "subs_rate_limit") { + try { + await recordSubs(outcome.videoId, opts.channelSlug); + } catch { + /* shared-state write is best-effort */ + } + } if (outcome.status === "skipped-filtered") { skippedCount++; } else if (outcome.status === "chat-only") { @@ -1103,15 +1143,16 @@ export async function runManagedDownloads( // declined slept 30 s after nothing but its metadata prefetch. Every // other outcome still sleeps — a real fetch, success or failure, and // every failure, per-video ones included (see declinedWithoutMediaFetch). + const gapMs = downloadGapMs(sleepSeconds, paceNow(), basePace); if ( - sleepSeconds > 0 && + gapMs > 0 && !isLast && !willAbortLoop && !opts.signal.aborted && !declinedWithoutMediaFetch(outcome) ) { - opts.onLog(`Sleeping ${sleepSeconds}s before next download...\n`); - await sleep(sleepSeconds * 1000, opts.signal); + opts.onLog(`Sleeping ${gapMs / 1000}s before next download...\n`); + await sleep(gapMs, opts.signal); } }), ), @@ -1189,12 +1230,15 @@ function trackOnDisk(entries: string[], track: string): boolean { // Fetch only the missing subtitle tracks for one already-downloaded video. // Uses the literal canonical output path (outputArgsForUrl) so the .vtt files // land in the same data/<canonicalId>/ dir as the rest of the video. +// THE BATCH SUBTITLE FETCH, per video. Returns the stderr it printed so the +// caller can tell a subtitle 429 (deferred, release 17 slice RL) from any +// other failure; throws on any non-zero exit as before. async function downloadSubsForUrl( opts: RunYtdlpOpts, root: string, url: string, subLangs: string, -): Promise<void> { +): Promise<string> { const args: string[] = [ "--ignore-config", "--restrict-filenames", @@ -1213,7 +1257,7 @@ async function downloadSubsForUrl( "--", url, ]; - await runChildAndStream(opts, root, args); + return runChildAndStream(opts, root, args); } async function downloadMissingSubs(opts: RunYtdlpOpts): Promise<void> { @@ -1242,6 +1286,13 @@ async function downloadMissingSubs(opts: RunYtdlpOpts): Promise<void> { let unidentifiable = 0; let noMetadata = 0; let noExpectedTracks = 0; + // Videos whose subtitles answered 429 recently (release 17, slice RL): left + // alone until their deferral lapses — 6 h after a strike, 7 days from the + // third. The video page's own download still fetches them. + const subDeferrals = await readSubtitleDeferrals(opts.paths).catch( + () => ({}) as Awaited<ReturnType<typeof readSubtitleDeferrals>>, + ); + const deferredNow: string[] = []; for (const url of urls) { // Post-reconcile a video lives in data/<canonicalId>/, so the URL's @@ -1251,6 +1302,10 @@ async function downloadMissingSubs(opts: RunYtdlpOpts): Promise<void> { unidentifiable++; continue; } + if (isSubtitleDeferred(subDeferrals, dirId, Date.now())) { + deferredNow.push(dirId); + continue; + } const videoDir = path.join(dataDir, dirId); let metaRaw: string; try { @@ -1300,6 +1355,14 @@ async function downloadMissingSubs(opts: RunYtdlpOpts): Promise<void> { : "") + (noMetadata ? `, ${noMetadata} without metadata (skipped)` : "") + (unidentifiable ? `, ${unidentifiable} unidentifiable URLs` : "") + + (deferredNow.length + ? `, ${deferredNow.length} deferred after a subtitle rate limit (${deferredNow + .map((id) => { + const d = (subDeferrals as Record<string, { count: number; until: number }>)[id]; + return `${id} ${d.count}×${d.count >= SUBTITLE_HOLD_AFTER ? " — left alone" : ""} until ${new Date(d.until).toISOString().slice(0, 16).replace("T", " ")} UTC`; + }) + .join(", ")})` + : "") + `.\n`, ); @@ -1324,16 +1387,38 @@ async function downloadMissingSubs(opts: RunYtdlpOpts): Promise<void> { label: extractVideoId(url) ?? url, kind: "download", }); + const id = extractVideoId(url); try { - await downloadSubsForUrl( + const stderr = await downloadSubsForUrl( task ? { ...opts, onLog: task.onLog } : opts, root, url, subLangs, ); + // yt-dlp exited clean. A subtitle 429 it reported anyway (a WARNING) + // is still a deferral; otherwise the subtitles came down and the + // video's deferral, if it had one, is forgotten. + if (id) { + if (hasSubtitleRateLimit(stderr)) { + await recordSubtitleDeferral(id, opts.channelSlug, opts.paths).catch(() => {}); + } else { + await clearSubtitleDeferral(id, opts.paths).catch(() => {}); + } + } } catch (err) { failed++; - if (abortOnError && !firstFailure) firstFailure = err as Error; + // A SUBTITLE 429 IS PER VIDEO (release 17, slice RL): the video is + // deferred and the batch goes on to the next one, abortOnError or + // not. Any other failure keeps today's abort. + const tail = (err as { stderrTail?: string }).stderrTail ?? ""; + if (id && classifyDownloadFailure(tail, undefined) === "subs_rate_limit") { + const d = await recordSubtitleDeferral(id, opts.channelSlug, opts.paths).catch(() => null); + (task?.onLog ?? opts.onLog)( + `Subtitles for ${id} were rate-limited (HTTP 429)${d ? ` — ${d.count}×, deferred` : ""}; continuing with the next video.\n`, + ); + } else if (abortOnError && !firstFailure) { + firstFailure = err as Error; + } } finally { task?.end(); } @@ -1939,11 +2024,13 @@ async function safeBackfillAvailability(opts: RunYtdlpOpts): Promise<void> { } } +// Returns the last STDERR_TAIL_BYTES of stderr (and attaches it to the thrown +// error as `stderrTail`), so a caller can classify what yt-dlp said. async function runChildAndStream( opts: RunYtdlpOpts, cwd: string, args: string[], -): Promise<void> { +): Promise<string> { opts.onLog(`$ ${opts.paths.ytdlpBin} ${args.join(" ")}\n`); const child = execa(opts.paths.ytdlpBin, args, { cwd, @@ -1953,6 +2040,10 @@ async function runChildAndStream( reject: false, }); child.all?.on("data", (c: Buffer) => opts.onLog(c.toString("utf8"))); + let stderrTail = ""; + child.stderr?.on("data", (c: Buffer) => { + stderrTail = (stderrTail + c.toString("utf8")).slice(-STDERR_TAIL_BYTES); + }); const result = await child; // yt-dlp exit code convention: 101 = "break-on-existing" / "max-downloads" // (clean stop, not an error). Treat it the same as 0. @@ -1961,11 +2052,14 @@ async function runChildAndStream( result.exitCode !== 101 && !opts.signal.aborted ) { - throw new Error(`yt-dlp exited with code ${result.exitCode}`); + throw Object.assign(new Error(`yt-dlp exited with code ${result.exitCode}`), { + stderrTail, + }); } if (result.exitCode === 101) { opts.onLog(`yt-dlp stopped on existing entry (exit 101).\n`); } + return stderrTail; } // The three sync-state stamps (CHANNEL_SYNC_STATE_KEYS). Each re-reads the diff --git a/common/ytdlp/subtitleRateLimit.test.ts b/common/ytdlp/subtitleRateLimit.test.ts @@ -0,0 +1,147 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { chmod, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { existsSync } from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import type { ChannelConfig } from "../lib/channelConfig"; +import { downloadOneManaged } from "./downloadOneManaged"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test ytdlp/subtitleRateLimit.test.ts +// +// A SUBTITLE 429 NEVER FAILS A DOWNLOAD (release 17, slice RL). The yt-dlp here +// is a temp node script that plays the 2026-10-01 shape: the prefetch and the +// media requests succeed, the subtitle fetch answers 429 — as a WARNING when +// asked to --ignore-errors (yt-dlp's own semantics), as an ERROR otherwise. It +// records every argv so the spawns can be counted and read. + +const ID = "H64QQZuw-aA"; +const VIDEO = `https://www.youtube.com/watch?v=${ID}`; +const SUB_429 = `Unable to download video subtitles for 'en': HTTP Error 429: Too Many Requests`; + +const FAKE = `#!/usr/bin/env node +const fs = require("node:fs"); +const path = require("node:path"); +const argv = process.argv.slice(2); +fs.appendFileSync(process.env.FAKE_LOG, JSON.stringify(argv) + "\\n"); +const has = (f) => argv.includes(f); +const dir = path.join("data", ${JSON.stringify(ID)}); +fs.mkdirSync(dir, { recursive: true }); +if (has("--no-write-auto-subs") && has("--write-info-json") && !has("-x") && !has("-f")) { + // the metadata prefetch: captions are listed + fs.writeFileSync(path.join(dir, "metadata.info.json"), JSON.stringify({ + id: ${JSON.stringify(ID)}, title: "t", automatic_captions: { en: [{}] }, subtitles: {}, + })); + process.exit(0); +} +const subsOn = (has("--write-subs") || has("--write-auto-subs")) && + !(argv.lastIndexOf("--no-write-subs") > argv.lastIndexOf("--write-subs")); +if (subsOn) { + if (has("--ignore-errors")) { + process.stderr.write("WARNING: [youtube] ${ID}: ${SUB_429}\\n"); + if (has("--skip-download")) { process.stdout.write("DLOM_ARCHIVE youtube ${ID}\\n"); process.exit(0); } + } else { + process.stderr.write("ERROR: [youtube] ${ID}: ${SUB_429}\\n"); + process.exit(1); + } +} +// the media +fs.writeFileSync(path.join(dir, "audio.mp3"), "audio"); +process.stdout.write("DLOM_ARCHIVE youtube ${ID}\\n"); +process.exit(0); +`; + +async function run(config: Partial<ChannelConfig>): Promise<{ + record: Awaited<ReturnType<typeof downloadOneManaged>>; + spawns: string[][]; + audio: boolean; + transcript: boolean; + log: string; +}> { + const root = await mkdtemp(path.join(tmpdir(), "subs-rl-")); + try { + const bin = path.join(root, "fake-ytdlp.cjs"); + await writeFile(bin, FAKE); + await chmod(bin, 0o755); + const logFile = path.join(root, "argv.log"); + process.env.FAKE_LOG = logFile; + const paths = { + channelsDir: path.join(root, "channels"), + ytdlpBin: bin, + } as Paths; + let log = ""; + const record = await downloadOneManaged({ + channelSlug: "c", + channelConfig: { + handling: "youtube", + url: "https://www.youtube.com/@c/videos", + ...config, + } as ChannelConfig, + paths, + videoUrl: VIDEO, + onLog: (s) => { + log += s; + }, + signal: new AbortController().signal, + }); + const spawns = (await readFile(logFile, "utf8")) + .split("\n") + .filter(Boolean) + .map((l) => JSON.parse(l) as string[]); + const dir = path.join(paths.channelsDir, "c", "data", ID); + return { + record, + spawns, + audio: existsSync(path.join(dir, "audio.mp3")), + transcript: existsSync(path.join(dir, "transcript.en.vtt")), + log, + }; + } finally { + delete process.env.FAKE_LOG; + await rm(root, { recursive: true, force: true }); + } +} + +test("youtube handling: the subtitle 429 is a warning, the media comes down, the record says subs_rate_limit", async () => { + const r = await run({}); + assert.equal(r.record.status, "ok"); + assert.equal(r.record.failureClass, "subs_rate_limit"); + assert.equal(r.record.fellBackToTranscribe, true); + assert.equal(r.audio, true); + assert.equal(r.transcript, false); + assert.deepEqual( + r.record.attempts.map((a) => a.kind), + ["metadata-prefetch", "primary", "no-subs-fallback"], + ); + // The primary asked yt-dlp to carry on past a subtitle failure, and paced + // the subtitle request at no less than `-t sleep`'s 5 s. + const primary = r.spawns[1]; + assert.ok(primary.includes("--ignore-errors")); + const ss = primary.indexOf("--sleep-subtitles"); + assert.ok(ss > primary.indexOf("sleep"), "--sleep-subtitles comes after -t sleep"); + assert.equal(primary[ss + 1], "5"); + // The media pass does not ask the throttled endpoint again. + const media = r.spawns[2]; + assert.ok(media.lastIndexOf("--no-write-subs") >= 0); + assert.ok(media.lastIndexOf("--no-write-auto-subs") >= 0); + assert.match(r.log, /subtitles for H64QQZuw-aA were rate-limited \(HTTP 429\); downloading the media anyway/); + // Three spawns: prefetch, subtitles, media — the same as a video with no captions. + assert.equal(r.spawns.length, 3); +}); + +test("a primary that died on its subtitles alone is run once more without them", async () => { + // A transcribe-handling channel whose own args ask for subtitles. + const r = await run({ handling: "transcribe", ytdlpExtraArgs: ["--write-subs"] }); + assert.equal(r.record.status, "ok"); + assert.equal(r.record.failureClass, "subs_rate_limit"); + assert.deepEqual( + r.record.attempts.map((a) => a.kind), + ["metadata-prefetch", "primary", "primary-without-subs"], + ); + assert.equal(r.audio, true); + const again = r.spawns[2]; + // The refusals come after the channel's own --write-subs (last flag wins). + assert.ok(again.lastIndexOf("--no-write-subs") > again.lastIndexOf("--write-subs")); +}); diff --git a/editor/app/operations/components/dispatch.ts b/editor/app/operations/components/dispatch.ts @@ -153,6 +153,10 @@ export function idleReasonText( return "every pending platform is in a rate-limit cooldown"; case "deferred": return "every pending video was rate-limited recently and is deferred"; + case "held": + return "every pending platform is held after repeated rate limits — one probe at a time"; + case "paced": + return "every pending platform is pausing between downloads"; case "no-workers": return "no enabled worker to run it"; case "workers-paused": @@ -206,6 +210,12 @@ export function formatRecency( // are noise and are dropped. export function formatCooldown(secs: number): string { if (secs < 60) return `${secs}s`; + // The days arm is for a subtitle hold (7 days), release 17. + if (secs >= 86400) { + const d = Math.floor(secs / 86400); + const h = Math.floor((secs % 86400) / 3600); + return h === 0 ? `${d}d` : `${d}d ${h}h`; + } if (secs >= 3600) { const h = Math.floor(secs / 3600); const m = Math.floor((secs % 3600) / 60); diff --git a/settings.json.example b/settings.json.example @@ -9,6 +9,12 @@ "x": {} }, "sleepBetweenDownloadsSeconds": 10, + "pacing": { + "sleepRequestsCapSeconds": 16, + "decayAfterCleanUnits": 5, + "holdAfterFailsAtCap": 3, + "holdProbeMinutes": 60 + }, "downloadFormat": "auto", "minFreeDiskGB": 5, "resumeMarginGB": 2,