// Per-platform exponential backoff for the auto-download runner. Pure logic + // types, no I/O — so it's unit-testable and can be persisted by // autoQueueState.ts and driven by autoRunner.ts. // // When a download fails with a rate-limit (HTTP 429 / throttle) or network // error, the runner must stop launching downloads for THAT platform for a // while instead of immediately re-picking the next video and hammering the // source again (the original Odysee 429 storm). Each consecutive failure // doubles the cooldown up to a cap; a successful download resets it. export type PlatformBackoffEntry = { // Epoch ms: downloads for this platform are paused until this instant. until: number; // Consecutive backoff-triggering failures, drives the exponential delay. fails: number; }; // Keyed by platform name (e.g. "odysee", "youtube", or "unknown"). export type PlatformBackoffState = Record; // First cooldown after a single failure; doubles from here. export const BACKOFF_BASE_MS = 60_000; // 1 minute // Ceiling on the cooldown — long rate limits shouldn't park a platform forever. export const BACKOFF_MAX_MS = 30 * 60_000; // 30 minutes // ±10% jitter so multiple platforms recovering at once don't sync up. const JITTER = 0.1; // TEST-ONLY: E2E_BACKOFF_BASE_MS shortens the FIRST cooldown (and so every // doubling from it), so the e2e suite can watch one lapse in seconds instead // of a minute (pacing.spec). The cap and the hold arithmetic (FAILS_TO_REACH_CAP // below) stay on the real constants. Unset everywhere but the e2e test server. function backoffBaseMs(): number { const raw = Number(process.env.E2E_BACKOFF_BASE_MS); return Number.isFinite(raw) && raw > 0 ? raw : BACKOFF_BASE_MS; } // Compute the next backoff entry for a platform given its previous one (or // undefined for the first failure). `now` and `rand` are injected for // deterministic tests; `rand` defaults to Math.random and returns [0, 1). export function nextBackoff( prev: PlatformBackoffEntry | undefined, now: number, rand: () => number = Math.random, ): PlatformBackoffEntry { const fails = (prev?.fails ?? 0) + 1; const base = Math.min(backoffBaseMs() * 2 ** (fails - 1), BACKOFF_MAX_MS); const jitter = 1 + (rand() * 2 - 1) * JITTER; return { until: now + Math.round(base * jitter), fails }; } // Fold two records of the same platform's backoff — the live one and the one // on disk — into one. `until` and `fails` merge SEPARATELY, each to its max: // taking one record whole (the later `until`) could pair a short manual // cooldown with a lost escalation count, so the next 429 restarted at the // base delay (release 8 review, S). export function mergeBackoffEntry( a: PlatformBackoffEntry | undefined, b: PlatformBackoffEntry, ): PlatformBackoffEntry { if (!a) return b; return { until: Math.max(a.until, b.until), fails: Math.max(a.fails, b.fails) }; } // True when the platform is currently in a cooldown window. export function isCoolingDown( state: PlatformBackoffState, platform: string, now: number, ): boolean { const entry = state[platform]; return entry !== undefined && entry.until > now; } // Clear a platform's backoff (called on a successful download). Mutates in place. export function clearBackoff( state: PlatformBackoffState, platform: string, ): void { delete state[platform]; } // Drop entries whose cooldown lapsed long enough ago that keeping their `fails` // count for escalation memory no longer matters. Mutates in place, keeping the // map bounded across a long-running server. We retain a lapsed entry for one // full BACKOFF_MAX_MS window so a platform that recovers and immediately // re-fails keeps escalating rather than restarting at the base delay. export function pruneExpired( state: PlatformBackoffState, now: number, retainMs: number = BACKOFF_MAX_MS, ): void { for (const [platform, entry] of Object.entries(state)) { if (entry.until + retainMs < now) delete state[platform]; } } // Defensive coercion for the persisted shape — any stored value that doesn't // match is dropped, so a hand-edited/corrupt state file can't crash the runner. // Mirrors coerceRuntime/coercePick in autoQueueState.ts. export function coercePlatformBackoff(value: unknown): PlatformBackoffState { const out: PlatformBackoffState = {}; if (!value || typeof value !== "object") return out; for (const [platform, raw] of Object.entries(value as Record)) { if (!raw || typeof raw !== "object") continue; const r = raw as Record; if ( typeof r.until === "number" && Number.isFinite(r.until) && typeof r.fails === "number" && Number.isFinite(r.fails) ) { out[platform] = { until: r.until, fails: r.fails }; } } return out; } // --------------------------------------------------------------------------- // Per-video deferral (release 7, YouTube lane pacing). // // A platform cooldown alone does not pace a lane whose queue is ordered: with // `order: "listed"` the runner re-picks the SAME video the moment the cooldown // lapses, so one video whose subtitle fetch keeps answering 429 climbs `fails` // to the 30-minute cap and holds the whole platform (2026-09-24: one Short // retried 12x). A rate-limited video is therefore also DEFERRED for // VIDEO_RATE_LIMIT_DEFER_MS: the auto-download pick skips it, so after the // cooldown the runner moves on to the next video. Persisted beside // `platformBackoff` (a restart must not re-hit the same video at fails+1). // Manual Sync / download-missing do not consult it. export type VideoDeferral = { // Epoch ms: the auto-download pick skips this video until this instant. until: number; // The channel the video was picked from — for the status view's link. channelSlug: string; }; // Keyed by video id. export type VideoDeferralState = Record; export const VIDEO_RATE_LIMIT_DEFER_MS = 6 * 60 * 60_000; // 6 hours // Defer a video from `now` for `ms`. Mutates in place. export function deferVideo( state: VideoDeferralState, videoId: string, channelSlug: string, now: number, ms: number = VIDEO_RATE_LIMIT_DEFER_MS, ): void { state[videoId] = { until: now + ms, channelSlug }; } // True while the video's deferral window is open. export function isVideoDeferred( state: VideoDeferralState, videoId: string, now: number, ): boolean { const entry = state[videoId]; return entry !== undefined && entry.until > now; } // Drop deferrals whose window has lapsed. Mutates in place — keeps the map // bounded (unlike platform backoff there is no escalation memory to retain). export function pruneDeferred(state: VideoDeferralState, now: number): void { for (const [videoId, entry] of Object.entries(state)) { if (entry.until <= now) delete state[videoId]; } } // Defensive coercion for the persisted shape; mirrors coercePlatformBackoff. export function coerceVideoDeferrals(value: unknown): VideoDeferralState { const out: VideoDeferralState = {}; if (!value || typeof value !== "object") return out; for (const [videoId, raw] of Object.entries(value as Record)) { if (!raw || typeof raw !== "object") continue; const r = raw as Record; if ( typeof r.until === "number" && Number.isFinite(r.until) && typeof r.channelSlug === "string" ) { out[videoId] = { until: r.until, channelSlug: r.channelSlug }; } } 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 = 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; // Epoch ms of the last doubling or decay step. THE ONE TIME-BASED RULE // (review L2): the pace also eases one step per PACE_TIME_DECAY_MS with no // rate limit, so a raised pace cannot outlive its cause while the lane is // off or has nothing pending on the platform. 0 for an entry written before // the field existed (it eases at once). steppedAt: number; }; // One easing step per hour with no rate limit (see `steppedAt`). export const PACE_TIME_DECAY_MS = 60 * 60_000; export type PlatformPaceState = Record; // 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; // Whether a RATE LIMIT is among the failures that held it — false for a // hold reached through network failures alone, which is worded "failing", // not "rate-limited" (review L5). rateLimited: boolean; }; export type PlatformHoldState = Record; // 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; 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. // // A PLATFORM FLOOR (`floor.minSeconds`, the platform's PLATFORM_MIN_GAP_SECONDS // in ytdlp/platformArgs.mjs — BitChute's 60 s): the operator's sleep is raised // to at least the floor and then stretched by up to half again at random // (`floor.random`, in [0, 1)), so a platform that rate-limits is never asked on // a fixed beat. With no floor the gap is exactly what it always was. export function downloadGapMs( sleepBetweenDownloadsSeconds: number, paceSeconds: number, baseSeconds: number, floor: { minSeconds?: number; random?: number } = {}, ): number { let sleep = Math.max(0, sleepBetweenDownloadsSeconds); const min = floor.minSeconds ?? 0; if (min > 0) { const r = Math.min(1, Math.max(0, floor.random ?? Math.random())); sleep = Math.max(min, sleep) * (1 + 0.5 * r); } 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, steppedAt: now, }; } 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, rateLimited: (state.platformHolds[platform]?.rateLimited ?? false) || failureClass === "rate_limit", }; } 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; // False for a probe whose media came down but whose subtitles 429'd: it // lifts a hold, and does not count toward the pace's easing (review L1). countForDecay?: boolean; now?: 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 && opts.countForDecay !== false) { 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, steppedAt: opts.now ?? pace.steppedAt, }; } } 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 { decayPaceByTime(state.platformPace, now); 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]; } } // The pace an entry stands at NOW once the hourly easing is applied: one // halving per whole PACE_TIME_DECAY_MS since `steppedAt`, never below its base. // Pure — the readers (the args builder, the view, the doctor) use it so a // stale entry on disk never paces a spawn above what the rule allows. export function effectivePaceSeconds(entry: PlatformPaceEntry, now: number): number { const steps = Math.max(0, Math.floor((now - entry.steppedAt) / PACE_TIME_DECAY_MS)); let v = entry.sleepRequestsSeconds; for (let i = 0; i < steps && v > entry.baseSeconds; i++) v = v / 2; return Math.max(v, entry.baseSeconds); } // Apply the hourly easing to the stored map (an entry at its base is // deleted). Mutates. Called where the state is pruned and on every unit. export function decayPaceByTime(pace: PlatformPaceState, now: number): void { for (const [pf, e] of Object.entries(pace)) { const steps = Math.max(0, Math.floor((now - e.steppedAt) / PACE_TIME_DECAY_MS)); if (steps === 0) continue; const v = effectivePaceSeconds(e, now); if (v <= e.baseSeconds) delete pace[pf]; else pace[pf] = { ...e, sleepRequestsSeconds: v, steppedAt: e.steppedAt + steps * PACE_TIME_DECAY_MS }; } } // THE OPERATOR'S WORD (review H2): "Clear hold" on the lane page. The hold, // the backoff and the raised pace all go, as if the platform had never been // refused; the next failure starts the escalation from the bottom. Mutates; // returns what there was, for the job log. export function clearPlatformPacing( state: PlatformPacingState, platform: string, ): { hold: PlatformHoldEntry | null; fails: number; paceSeconds: number | null } { const out = { hold: state.platformHolds[platform] ?? null, fails: state.platformBackoff[platform]?.fails ?? 0, paceSeconds: state.platformPace[platform]?.sleepRequestsSeconds ?? null, }; delete state.platformHolds[platform]; clearBackoff(state.platformBackoff, platform); delete state.platformPace[platform]; return out; } // 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)) { if (!raw || typeof raw !== "object") continue; const r = raw as Record; 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, steppedAt: finite(r.steppedAt) ? r.steppedAt : 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)) { if (!raw || typeof raw !== "object") continue; const r = raw as Record; if (finite(r.since) && finite(r.probeAt)) { out[pf] = { since: r.since, probeAt: r.probeAt, rateLimited: typeof r.rateLimited === "boolean" ? r.rateLimited : true, }; } } 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)) { if (!raw || typeof raw !== "object") continue; const r = raw as Record; 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"; const why = hold.rateLimited ? "its rate limit outlasted the cooldown cap" : "it kept failing (network errors) past the cooldown cap"; return ( `${platform} is held: ${why} ` + `(${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.` ); }