// Shared accessors for the per-platform download rate-limit cooldown that both // the auto-download runner and manual sync/download paths respect. // // The runner owns this state in memory — ONE object shared by every lane's // runner (`sharedAutoQueueState`) — and persists it WHOLE to // `.auto-queue/state.json` (the "download" kind's platformBackoff). These // helpers let code OUTSIDE the runner — a clicked Sync, a metadata scan, the // video page — observe and extend the same cooldown, so a 429 hit by either path // pauses the other. // // WRITE THROUGH THE LIVE OBJECT WHEN THERE IS ONE (release 8). A disk-only // read-modify-write while a runner is live is erased by the next persist of ANY // lane — its in-memory object never saw the cooldown — and the download lane's // own merge-from-disk never runs while that lane is paused, stopped or at // capacity, which is exactly when an operator is riding out a rate limit. With // no runner live there is no in-memory copy, and the disk round trip is right. // Per-platform one-download-at-a-time is enforced by the job queue (see // autoRunner.ts / registry.ts), so a lost cooldown could only ever cost // politeness, never concurrency. import { getPaths, type Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; import { type AutoQueueState, liveAutoQueueState, readAutoQueueState, writeAutoQueueState, } from "./autoQueueState"; import { type PacingSettings, type PlatformHoldEntry, type SubtitleDeferral, type SubtitleDeferralState, clearPlatformPacing, deferSubtitles, effectivePaceSeconds, escalatePlatform, heldPlatformSentence, mergeBackoffEntry, prunePlatformPacing, pruneSubtitleDeferrals, settlePlatformClean, } 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. export async function platformCooldownRemainingMs( platform: string, paths: Paths = getPaths(), ): Promise { const state = (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); const entry = state.download.platformBackoff[platform]; const now = Date.now(); return entry && entry.until > now ? entry.until - now : 0; } // 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( paths: Paths, mutate: (state: AutoQueueState) => T, ): Promise { const live = await liveAutoQueueState(paths); let state: AutoQueueState; if (live) { // Fold in whatever is on disk first (`until` and `fails` each to their // max, as the runner's own merge does), so a cooldown written while no // runner held the object is escalated from, not forgotten. const onDisk = (await readAutoQueueState(paths)).download.platformBackoff; for (const [pf, e] of Object.entries(onDisk)) { live.download.platformBackoff[pf] = mergeBackoffEntry( live.download.platformBackoff[pf], e, ); } state = live; } else { state = await readAutoQueueState(paths); } 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 { 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 { 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 { 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 { 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. Only // while it is held AND its probe is still in the future (review H2): once the // probe is due, a manual run IS the probe — refusing it too would leave a // platform with the lane off, or nothing pending on it, held for ever. `what` // names the refused thing ("Sync", "This download"). export async function heldPlatformRefusal( platform: string, what: string, paths: Paths = getPaths(), now: number = Date.now(), ): Promise { const state = (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); const hold = state.download.platformHolds[platform]; const entry = state.download.platformBackoff[platform]; if (!hold || !entry || entry.until <= now) return null; const probeMinutes = pacingSettings()?.holdProbeMinutes ?? Math.round((hold.probeAt - hold.since) / 60_000); return heldPlatformSentence(platform, hold, entry.fails, now, { probeMinutes, what, }); } // A MANUAL RUN THAT CAME BACK CLEAN (review H2) — a Sync, a download, a // metadata scan that the source answered without a rate limit or a network // failure — settles the platform exactly as a clean lane unit does: the // backoff and any hold clear, and it counts toward the pace's easing. Returns // the line for the job's log, or null when there was nothing to settle (and // then writes nothing). export async function recordPlatformClean( platform: string, paths: Paths = getPaths(), pacing: PacingSettings | undefined = pacingSettings(), ): Promise { const peek = (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); const d = peek.download; if (!d.platformBackoff[platform] && !d.platformHolds[platform] && !d.platformPace[platform]) { return null; } return mutateDownloadState(paths, (state) => { const now = Date.now(); prunePlatformPacing(state.download, now); const hadBackoff = Boolean(state.download.platformBackoff[platform]); const fx = settlePlatformClean(state.download, platform, { baseSeconds: staticSleepRequestsSeconds(platform), now, ...(pacing ? { pacing } : {}), }); if (fx.releasedHold) { return `${platform} answered cleanly: its hold and backoff are cleared (as a clean probe would).\n`; } if (hadBackoff) return `${platform} answered cleanly: its rate-limit backoff is cleared.\n`; if (fx.decayed) return `${platform} pace eased to ${fx.paceSeconds}s between requests.\n`; return null; }); } // "Clear hold" (review H2): the operator's word beats the machine, as with an // auto-pause. The hold, the backoff and the raised pace all go. Returns the // sentence for the job log, or null when the platform had none of them. export async function clearPlatformHold( platform: string, paths: Paths = getPaths(), ): Promise { return mutateDownloadState(paths, (state) => { const now = Date.now(); const pace = state.download.platformPace[platform]; const was = clearPlatformPacing(state.download, platform); if (!was.hold && was.fails === 0 && was.paceSeconds === null) return null; const parts = [ was.hold ? `the hold (since ${new Date(was.hold.since).toISOString().slice(0, 16).replace("T", " ")} UTC)` : null, was.fails > 0 ? `the backoff (${was.fails} failures)` : null, pace ? `the pace (${effectivePaceSeconds(pace, now)}s → base ${pace.baseSeconds}s)` : null, ].filter(Boolean); return `Cleared by hand for ${platform}: ${parts.join(", ")}. The next failure starts from the bottom of the backoff.\n`; }); }