// Client-side rate limiting for the bring-your-own-key AI chat. The export site // is fully static and hosts no inference — every provider call goes browser → // provider with the visitor's own key, so pacing to a free tier's // requests-per-minute (RPM) cap has to happen right here on the client. This // module is the whole mechanism: // // • DEFAULT_RPM / defaultRpmFor — conservative, free-tier-safe per-model // ceilings so a free key just works out of the box; // • acquire() — a serialized min-interval gate that spaces same-model // calls to the RPM (strict spacing, not a token bucket: // a bucket lets a burst instantly trip a 5/min cap); // • parseRetryAfterMs() — read a provider's own "retry in N" hint off a 429; // • noteRateLimited() — step the SESSION RPM estimate down when we guess high; // • withRateLimitRetry() — re-issue a whole request on a 429, honouring the hint; // • onRateLimitEvent() — a tiny emitter so the UI can show "pacing / retrying // in Ns" without threading a callback through four layers. // // TPM (tokens/minute) and RPD (requests/day) caps are NOT modelled — they surface // as 429s and fall through to the existing sweep pause/checkpoint. import { PausableError, type Provider } from "./askProvider"; // ─── free-tier-safe default RPM table ─── type RpmRule = { pattern: RegExp; rpm: number }; // APPROXIMATE, tier-dependent, and deliberately conservative ceilings for the // published FREE tiers. They are guesses that WILL drift as providers change // their limits — the adaptive step-down (noteRateLimited) is what makes exactness // unnecessary, so do NOT treat these as truth anywhere else. Rules are matched in // order (first hit wins); `fallback` covers anything unmatched. const DEFAULT_RPM: Record = { gemini: { rules: [ // gemini free tier: pro is the stingy one (~5 RPM); flash-lite is highest; // plain flash sits in between. flash-lite MUST be tested before flash // (its id also contains "flash"). { pattern: /pro/, rpm: 5 }, { pattern: /flash-lite/, rpm: 15 }, { pattern: /flash/, rpm: 10 }, ], fallback: 10, }, // Anthropic + OpenAI paid/dev tiers are far higher; these are conservative // floors that a real key easily clears (and the RPM field lets a user raise). anthropic: { rules: [], fallback: 50 }, // all claude-* openai: { rules: [], fallback: 60 }, }; // The free-tier-safe default RPM for a provider+model (before any override or // session step-down). Matched by model-id pattern; case-insensitive. export function defaultRpmFor(provider: Provider, model: string): number { const table = DEFAULT_RPM[provider]; if (!table) return 10; const m = (model || "").toLowerCase(); for (const r of table.rules) if (r.pattern.test(m)) return r.rpm; return table.fallback; } // ─── override store (localStorage) ─── // One key holds a map of ":" → RPM the user pinned (for a paid // key). Strict parse-or-reset + quota-swallowing write, matching the // askChatStorage house style. const RPM_KEY = "ytdlp-tb:ai:rpm"; function storeKey(provider: Provider, model: string): string { return `${provider}:${model}`; } function loadOverrides(): Record { if (typeof window === "undefined") return {}; let raw: string | null = null; try { raw = window.localStorage.getItem(RPM_KEY); } catch { return {}; } if (raw === null) return {}; try { const parsed: unknown = JSON.parse(raw); if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { throw new Error("bad shape"); } const out: Record = {}; for (const [k, v] of Object.entries(parsed as Record)) { if (typeof v === "number" && Number.isFinite(v) && v > 0) out[k] = v; } return out; } catch { // Malformed payload → reset the key rather than trap the UI in a broken state. try { window.localStorage.removeItem(RPM_KEY); } catch { /* ignore */ } return {}; } } // The user's pinned RPM for this provider+model, or null if none. A "*" catch-all // key raises the ceiling for EVERY model at once — used by the e2e suite (see the // shared test setup) so the throttle doesn't slow the whole suite; the settings // UI only ever writes exact ":" keys. export function getRpmOverride(provider: Provider, model: string): number | null { const map = loadOverrides(); const exact = map[storeKey(provider, model)]; if (typeof exact === "number") return exact; const wild = map["*"]; return typeof wild === "number" ? wild : null; } // Pin (or, with null, clear) the RPM override for a provider+model. export function setRpmOverride( provider: Provider, model: string, rpm: number | null, ): void { if (typeof window === "undefined") return; const map = loadOverrides(); const key = storeKey(provider, model); if (rpm == null) delete map[key]; else map[key] = rpm; try { window.localStorage.setItem(RPM_KEY, JSON.stringify(map)); } catch { /* quota / serialization errors are non-fatal */ } } // ─── session step-down (self-tuning) ─── // SESSION-ONLY RPM estimates, lowered by noteRateLimited when a 429 says we // guessed too high. NOT persisted — one blip shouldn't permanently hobble a paid // key; a reload starts fresh from the defaults. const sessionRpm = new Map(); // Step the session RPM estimate for a provider+model down after a rate-limit, so // later calls pace slower. Returns the new estimate. floor(rpm * 0.6), min 1. export function noteRateLimited(provider: Provider, model: string): number { const key = storeKey(provider, model); const current = sessionRpm.get(key) ?? defaultRpmFor(provider, model); const next = Math.max(1, Math.floor(current * 0.6)); sessionRpm.set(key, next); return next; } // The RPM actually used to pace a provider+model: override ?? session step-down // ?? default. A user's explicit override is authoritative (it wins over the // self-tuning) — the step-down only kicks in for the default-driven case. export function effectiveRpm(provider: Provider, model: string): number { const override = getRpmOverride(provider, model); if (override != null) return override; const session = sessionRpm.get(storeKey(provider, model)); if (session != null) return session; return defaultRpmFor(provider, model); } // Test seam: clear the session step-down map so unit tests start clean. export function __resetSessionRpm(): void { sessionRpm.clear(); } // ─── event emitter ─── // Emitted so the UI can surface pacing/retry waits without a callback threaded // through useAskChat → runReportChunk → *Gather → postJson. `waiting` fires when // a call is about to sleep (pacing between requests, or backing off a 429); // `resumed` fires when that sleep ends (or is aborted), so a listener can clear // its notice. export type RateLimitEvent = | { type: "waiting"; provider: Provider; model: string; ms: number; reason: "pacing" | "retry"; } | { type: "resumed"; provider: Provider; model: string }; type Listener = (e: RateLimitEvent) => void; const listeners = new Set(); // Subscribe to rate-limit events; returns an unsubscribe fn. export function onRateLimitEvent(fn: Listener): () => void { listeners.add(fn); return () => { listeners.delete(fn); }; } function emitEvent(e: RateLimitEvent): void { for (const fn of listeners) { try { fn(e); } catch { /* a listener throwing must not break pacing */ } } } // ─── injectable clock (testability) ─── export type Clock = { now: () => number; // Resolve after `ms`, or reject with an AbortError if `signal` aborts first. sleep: (ms: number, signal?: AbortSignal) => Promise; }; function abortError(): DOMException { return new DOMException("Aborted", "AbortError"); } function realSleep(ms: number, signal?: AbortSignal): Promise { if (ms <= 0) return Promise.resolve(); return new Promise((resolve, reject) => { if (signal?.aborted) { reject(abortError()); return; } const onAbort = () => { clearTimeout(t); reject(abortError()); }; const t = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); signal?.addEventListener("abort", onAbort, { once: true }); }); } const realClock: Clock = { now: () => Date.now(), sleep: realSleep }; let clock: Clock = realClock; // Test seam: swap in a fake clock so limiter spacing + abort-during-wait can be // unit-tested without real time. Returns a restore fn. export function __setClock(next: Clock): () => void { const prev = clock; clock = next; return () => { clock = prev; }; } // ─── the gate: acquire ─── // Per-key gate state: a promise chain that serializes acquirers for the same // provider+model, and the clock time the last acquirer was released (so the next // waits until last + spacing). type Gate = { chain: Promise; last: number }; const gates = new Map(); // Wait until it's this call's turn AND the min-interval since the previous // same-model call has elapsed. Serialized per provider:model and abort-aware, so // Stop cancels a long spacing wait promptly. Spacing = ceil(60000 / rpm) // (5 RPM → 12s). Strict spacing rather than a token bucket, so a burst can't // instantly trip a per-minute cap. export async function acquire( provider: Provider, model: string, signal?: AbortSignal, ): Promise { if (signal?.aborted) throw abortError(); const key = storeKey(provider, model); const gate = gates.get(key) ?? { chain: Promise.resolve(), last: 0 }; gates.set(key, gate); // Chain onto the prior acquirer so same-key calls are strictly serialized. const prior = gate.chain; let release!: () => void; gate.chain = new Promise((r) => { release = r; }); try { await prior; const rpm = effectiveRpm(provider, model); const spacing = rpm > 0 ? Math.ceil(60000 / rpm) : 0; const waitMs = gate.last === 0 ? 0 : Math.max(0, gate.last + spacing - clock.now()); if (waitMs > 0) { emitEvent({ type: "waiting", provider, model, ms: waitMs, reason: "pacing" }); try { await clock.sleep(waitMs, signal); } finally { emitEvent({ type: "resumed", provider, model }); } } gate.last = clock.now(); } finally { // Release the NEXT acquirer whether we paced, ran clean, or aborted. release(); } } // ─── retry-hint parsing ─── // Protobuf Duration string form, e.g. "38s", "1.5s", "0.200s" → ms. function parseDurationMs(s: string): number | null { const m = /^\s*(\d+(?:\.\d+)?)s\s*$/.exec(s); if (!m) return null; return Math.round(parseFloat(m[1]) * 1000); } type HeaderBag = { get(name: string): string | null } | null | undefined; // Read a provider's own "retry in N" hint off a 429. Covers the `Retry-After` // header in SECONDS or as an HTTP-date (Anthropic/OpenAI), PLUS Gemini's body // form — `error.details[]` with a RetryInfo carrying `retryDelay: "38s"`. Returns // null when no hint is present, so the caller falls back to exponential backoff. export function parseRetryAfterMs( headers: HeaderBag, bodyText?: string, ): number | null { const header = headers?.get?.("retry-after"); if (header) { const secs = Number(header); if (Number.isFinite(secs)) return Math.max(0, secs * 1000); const date = Date.parse(header); if (Number.isFinite(date)) return Math.max(0, date - clock.now()); } if (bodyText) { try { const parsed: unknown = JSON.parse(bodyText); const details = (parsed as { error?: { details?: unknown } })?.error ?.details; if (Array.isArray(details)) { for (const d of details) { const delay = (d as { retryDelay?: unknown })?.retryDelay; if (typeof delay === "string") { const ms = parseDurationMs(delay); if (ms != null) return ms; } } } } catch { /* not JSON — fall through to the loose scan */ } // Last resort: a bare "retryDelay": "38s" / "retry in 38s" mention. const m = /retry[^0-9]{0,16}?(\d+(?:\.\d+)?)\s*s\b/i.exec(bodyText); if (m) return Math.round(parseFloat(m[1]) * 1000); } return null; } // ─── retry wrapper ─── // Wrap a whole request attempt so a 429 can be retried by RE-ISSUING it (the fn // must build + send the request itself). Safe for streaming because ensureOk // throws BEFORE the SSE body is consumed — nothing has streamed yet, so a retry // can't duplicate text; a mid-stream failure is never a PausableError, so it // bubbles out on the first throw. On each PausableError we wait the provider's // hint (or an exponential backoff); after `attempts` failures we re-throw the // last one so the sweep driver's existing pause/checkpoint fires (the genuine // last resort — e.g. a daily-quota cap that won't clear for hours). export async function withRateLimitRetry( fn: () => Promise, opts: { provider: Provider; model: string; signal?: AbortSignal; attempts?: number; }, ): Promise { const attempts = opts.attempts ?? 4; let lastErr: PausableError | null = null; for (let attempt = 0; attempt < attempts; attempt++) { if (opts.signal?.aborted) throw abortError(); try { return await fn(); } catch (e) { if (!(e instanceof PausableError)) throw e; lastErr = e; // Out of attempts → let it bubble to the pause/checkpoint path. if (attempt === attempts - 1) break; const backoff = Math.min(60000, 1000 * 2 ** attempt); const waitMs = e.retryAfterMs != null ? e.retryAfterMs : backoff; emitEvent({ type: "waiting", provider: opts.provider, model: opts.model, ms: waitMs, reason: "retry", }); try { await clock.sleep(waitMs, opts.signal); } finally { emitEvent({ type: "resumed", provider: opts.provider, model: opts.model }); } } } throw lastErr ?? new Error("rate limit retry failed"); }