commit 3bea5b510b9e7fdde09b27bbdbedce0fe79c1fcf
parent 674c2aae499be6288afee320c7c66b6d11a8623b
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 20 Jul 2026 11:23:07 -0400
Ask chat: client-side rate limiting — pace requests, honour retry hints, self-tune
A free-tier sweep fired provider calls as fast as the loop could produce them,
blew through the per-minute cap, and stopped dead on the first 429 (and a
rate-limit mid-answer was a hard error — only the sweep path handled 429). Add
one client-side limiter behind both request paths.
- New export/app/lib/rateLimit.ts: per-model free-tier-safe RPM defaults,
a serialized abort-aware min-interval gate (acquire), retry-hint parsing
(Retry-After secs/HTTP-date + Gemini RetryInfo body), session step-down,
withRateLimitRetry, an event emitter, and a strict parse-or-reset override
store (ytdlp-tb:ai:rpm, with a "*" wildcard for e2e).
- Move PausableError into askProvider.ts (beside ContextTooLargeError) so
ensureOk can throw it without a cycle; shared.ts re-exports it. It gains an
optional retryAfterMs.
- Wire the two choke points: askStream (acquire + retry — covers answers and
scripted gather) and postJson (split into postJsonOnce; acquire + retry).
ensureOk now takes the model and turns a 429 into a retry-hinted PausableError.
- Sweep driver unchanged: its PausableError -> pause/checkpoint is now the last
resort, firing only after retries exhaust (right for a daily-quota cap).
- UI: a "Requests per minute" field in ProviderSettings, and a
"Rate limited — retrying in Ns" line in the sweep progress strip (gated behind
a post-hydration flag to avoid a hydration mismatch).
- Tests: new rateLimit.test.ts; ask-workspace.spec split into transient-429
auto-recover (no pause) and persistent-429 pause + Resume; shared e2e setup
raises RPM so pacing doesn't slow the suite. Full export e2e 121/121 green.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Diffstat:
11 files changed, 975 insertions(+), 46 deletions(-)
diff --git a/export/CHANGELOG.md b/export/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **"Ask AI" now paces itself to your key's rate limit instead of failing.** A whole-corpus sweep on a free-tier key used to fire provider calls as fast as the loop could produce them, blow straight through the per-minute request cap, and stop dead on the first HTTP 429 (and a rate-limit *mid-answer* surfaced as a hard error, because only the sweep path ever handled 429). Now a single client-side limiter sits behind **every** AI call: it **spaces requests** to a conservative, free-tier-safe **requests-per-minute** default chosen per model (e.g. Gemini `*-pro` → 5/min, `*flash` → 10/min; Claude/OpenAI higher), so a paid key runs fast and a free key just runs *slowly* rather than erroring. When a 429 does land, the limiter **honours the provider's own retry hint** (the `Retry-After` header, or Gemini's `RetryInfo` retry-delay) — or an exponential backoff — and **re-issues the request** (safe for streaming: the retry happens before any answer text is emitted), and it **self-lowers** the rate after a 429 so later calls pace slower. Only a 429 that outlasts the retries falls through to the existing **pause/checkpoint** (the right home for a daily-quota cap you resume tomorrow). Provider settings gain an editable **Requests per minute** field (with the effective spacing, e.g. *"≈12s between requests"*, and a reset-to-default), and a paced sweep shows a **"Rate limited — retrying in Ns"** line in its progress strip so it never looks frozen. See `export/app/lib/rateLimit.ts` (the whole mechanism), `export/app/lib/askProvider.ts` (`PausableError` + `acquire`/retry in `askStream`, 429 in `ensureOk`), `export/app/lib/nativeTools/shared.ts` (`acquire`/retry in `postJson`), `export/app/ask/{useAskChat,ProviderSettings,PinnedResultsPanel,AskChat}.tsx`, and `export/e2e/ask-workspace.spec.ts`.
- **"Ask AI" — an integrated grounding workspace: pick your AI target, sweep any size, save & resume chats.** Search and the chat used to behave like two tabs, grounding was all-or-nothing, a sweep had hard-coded caps and couldn't be paused, and a conversation lived in a single unnamed slot that "New chat" wiped. Now they're one surface:
- **Grounding palette (the signature control).** A first-class `Whole search ⇄ Selection` toggle in the chat always states *what the AI is looking at*, alongside an action cluster — **Ask**, **Sweep**, and preset directives (*Summary*, *Contradictions*, *Timeline*) that pre-fill the composer/sweep box. "Selection" grounds in a **hand-picked subset** of results shown as dismissible chips; an empty selection parks on Whole search, so nothing changes for people who never select.
- **Per-result selection + one-click Ask.** Every result card gets a select checkbox (an inclusion set mirroring the channel exclusion set) and an **Ask** button that grounds the chat in *that one video* and crosses to it. The selection survives `/` ⇄ `/ask` and reload, and self-cleans against a new search.
diff --git a/export/app/ask/AskChat.tsx b/export/app/ask/AskChat.tsx
@@ -110,6 +110,10 @@ export default function AskChat() {
setMaxAnswerTokens={s.setMaxAnswerTokens}
sweepBatch={s.sweepBatch}
setSweepBatch={s.setSweepBatch}
+ rpmOverride={s.rpmOverride}
+ rpmDefault={s.rpmDefault}
+ effectiveRpm={s.effectiveRpm}
+ setRpm={s.setRpm}
buildDebugJson={s.buildDebugJson}
/>
@@ -195,6 +199,7 @@ export default function AskChat() {
canSweep={s.canSweep}
sweeping={s.sweeping}
sweepProgress={s.sweepProgress}
+ rateLimitNotice={s.rateLimitNotice}
totalGroundingVideos={s.totalGroundingVideos}
sweepCapped={s.sweepCapped}
hasReport={s.report.trim() !== ""}
diff --git a/export/app/ask/PinnedResultsPanel.tsx b/export/app/ask/PinnedResultsPanel.tsx
@@ -32,6 +32,7 @@ export function PinnedResultsPanel({
canSweep,
sweeping,
sweepProgress,
+ rateLimitNotice,
totalGroundingVideos,
sweepCapped,
hasReport,
@@ -58,6 +59,10 @@ export function PinnedResultsPanel({
sweeping: boolean;
// Per-batch progress while sweeping.
sweepProgress: { done: number; total: number; sections: number };
+ // A rate-limiter pacing/backoff notice ("Rate limited — retrying in Ns"), or
+ // null when nothing is waiting. Shown in the sweep strip so a paced run's waits
+ // are visible instead of looking frozen.
+ rateLimitNotice: string | null;
// Full matched-set size (before caps) + whether the search itself was capped.
totalGroundingVideos: number;
sweepCapped: boolean;
@@ -259,6 +264,14 @@ export function PinnedResultsPanel({
? ` · ${sweepProgress.sections} section${sweepProgress.sections === 1 ? "" : "s"} written`
: ""}
</span>
+ {rateLimitNotice && (
+ <span
+ data-testid="rate-limit-notice"
+ className="text-xs font-medium text-warning"
+ >
+ {rateLimitNotice}
+ </span>
+ )}
</>
) : (
<>
diff --git a/export/app/ask/ProviderSettings.tsx b/export/app/ask/ProviderSettings.tsx
@@ -28,6 +28,11 @@ type Props = {
setMaxAnswerTokens: (n: number) => void;
sweepBatch: number;
setSweepBatch: (n: number) => void;
+ // Requests-per-minute pacing for the current provider+model.
+ rpmOverride: number | null;
+ rpmDefault: number;
+ effectiveRpm: number;
+ setRpm: (n: number | null) => void;
buildDebugJson: () => string;
};
@@ -50,9 +55,19 @@ export function ProviderSettings(props: Props) {
setMaxAnswerTokens,
sweepBatch,
setSweepBatch,
+ rpmOverride,
+ rpmDefault,
+ effectiveRpm,
+ setRpm,
buildDebugJson,
} = props;
const info = PROVIDERS[provider];
+ // Derived spacing between requests at the effective RPM (e.g. 5 RPM → ≈12s).
+ const spacingSec = effectiveRpm > 0 ? 60 / effectiveRpm : 0;
+ const spacingLabel =
+ spacingSec >= 1
+ ? `${Math.round(spacingSec)}s`
+ : `${Math.max(1, Math.round(spacingSec * 1000))}ms`;
// Seed the panel open when there's no key yet, then hand control to the user.
// Must NOT be recomputed from `apiKey` on every render — otherwise typing the
// first character of the key flips it closed and yanks the panel away mid-entry.
@@ -221,6 +236,44 @@ export function ProviderSettings(props: Props) {
</span>
</label>
+ {/* Requests per minute: paces every AI call so a free-tier key stays
+ under its per-minute cap instead of erroring. Defaults are conservative
+ per-model guesses (shown as the placeholder); set a higher value for a
+ paid key. Auto-lowers itself after a rate-limit (429). */}
+ <label className="flex flex-col gap-1 text-xs text-muted-foreground">
+ Requests per minute
+ <div className="flex items-center gap-2">
+ <input
+ type="number"
+ min={1}
+ max={100000}
+ step={1}
+ value={rpmOverride ?? ""}
+ placeholder={String(rpmDefault)}
+ onChange={(e) => {
+ const v = e.target.value.trim();
+ setRpm(v === "" ? null : Number(v));
+ }}
+ aria-label="Requests per minute"
+ className="w-32 rounded-md border border-border bg-background px-2 py-1.5 text-sm text-foreground"
+ />
+ {rpmOverride != null && (
+ <button
+ type="button"
+ onClick={() => setRpm(null)}
+ className="rounded-md border border-border px-2 py-1 text-xs text-muted-foreground transition-colors hover:text-foreground"
+ >
+ Reset to default
+ </button>
+ )}
+ </div>
+ <span className="text-muted-foreground/70">
+ Paces requests to ≈{spacingLabel} apart so a free-tier key stays under
+ its per-minute limit. Blank uses the per-model default ({rpmDefault}).
+ It auto-lowers itself after a rate-limit (429).
+ </span>
+ </label>
+
{/* Debug export: a redacted snapshot for troubleshooting. */}
<div className="flex flex-col gap-1 text-xs text-muted-foreground">
<span>Debug</span>
diff --git a/export/app/ask/useAskChat.ts b/export/app/ask/useAskChat.ts
@@ -21,6 +21,13 @@ import {
} from "../lib/searchAgent";
import type { RetrievedVideo } from "../lib/askRetrieval";
import { PausableError } from "../lib/nativeTools/shared";
+import {
+ defaultRpmFor,
+ effectiveRpm as computeEffectiveRpm,
+ getRpmOverride,
+ onRateLimitEvent,
+ setRpmOverride,
+} from "../lib/rateLimit";
import { buildDebugExport } from "../lib/askDebug";
import {
buildApiMessages,
@@ -206,6 +213,18 @@ export function useAskChat() {
// Ring buffer of the most recent provider calls (request/response/finishReason),
// for the debug export. Runtime-only (not persisted); trimmed to DEBUG_TRACE_CAP.
const debugTraceRef = useRef<DebugCall[]>([]);
+ // A short human notice while the rate limiter is pacing/backing off (e.g.
+ // "Rate limited — retrying in 12s"), driven by onRateLimitEvent. Cleared when
+ // the wait ends. Surfaced in the sweep progress strip so a paced run never
+ // looks frozen. Runtime-only.
+ const [rateLimitNotice, setRateLimitNotice] = useState<string | null>(null);
+ // Bumped after an RPM override is written so the derived effective/override/
+ // default values recompute (the store itself lives in localStorage).
+ const [rpmTick, setRpmTick] = useState(0);
+ // False until after hydration. The RPM values read localStorage, which is absent
+ // during SSR — gate on this so the first client render matches the server (empty)
+ // and only then reflect any stored override, avoiding a hydration mismatch.
+ const [hydrated, setHydrated] = useState(false);
// The effective pinned grounding = the ACTIVE-mode grounding (whole search or
// the hand-picked selection), overlaid with any per-video "Load context"
@@ -266,6 +285,23 @@ export function useAskChat() {
next.length > DEBUG_TRACE_CAP ? next.slice(-DEBUG_TRACE_CAP) : next;
}, []);
+ // Surface the rate limiter's pacing/backoff waits as a short notice (the
+ // progress strip shows it during a sweep). "resumed" clears it.
+ useEffect(() => {
+ return onRateLimitEvent((e) => {
+ if (e.type === "resumed") {
+ setRateLimitNotice(null);
+ return;
+ }
+ const secs = Math.max(1, Math.round(e.ms / 1000));
+ setRateLimitNotice(
+ e.reason === "retry"
+ ? `Rate limited — retrying in ${secs}s`
+ : `Pacing requests — next in ${secs}s`,
+ );
+ });
+ }, []);
+
// Restore saved preferences + (if remembered) the provider's key/model, plus a
// persisted conversation.
useEffect(() => {
@@ -324,6 +360,8 @@ export function useAskChat() {
} catch {
/* storage unavailable / corrupt */
}
+ // Now safe to read localStorage-backed RPM overrides in render.
+ setHydrated(true);
/* eslint-enable react-hooks/set-state-in-effect */
}, []);
@@ -480,6 +518,46 @@ export function useAskChat() {
}
}
+ // Requests-per-minute pacing for the CURRENT provider+model. `rpmOverride` is
+ // the user's pinned value (null ⇒ using the default/self-tuned rate); `rpmDefault`
+ // is the free-tier-safe default (shown as the field placeholder); `effectiveRpm`
+ // drives the "≈Ns between requests" hint. All recompute on provider/model change
+ // and after a write (rpmTick).
+ const rpmModel = model.trim() || PROVIDERS[provider].defaultModel;
+ // rpmOverride / effectiveRpm read imperative stores (localStorage + the session
+ // step-down map) React can't track, so rpmTick — bumped on every write — and
+ // rateLimitNotice — bumped on every pacing event, i.e. after a step-down — are
+ // listed as recompute triggers even though the bodies don't reference them.
+ const rpmOverride = useMemo(
+ () => (hydrated ? getRpmOverride(provider, rpmModel) : null),
+ // eslint-disable-next-line react-hooks/exhaustive-deps
+ [provider, rpmModel, rpmTick, hydrated],
+ );
+ const rpmDefault = useMemo(
+ () => defaultRpmFor(provider, rpmModel),
+ [provider, rpmModel],
+ );
+ const effectiveRpm = useMemo(
+ () =>
+ hydrated
+ ? computeEffectiveRpm(provider, rpmModel)
+ : defaultRpmFor(provider, rpmModel),
+ // eslint-disable-next-line react-hooks/exhaustive-deps
+ [provider, rpmModel, rpmTick, rateLimitNotice, hydrated],
+ );
+ // Set (n) or clear (null) the RPM override for the current provider+model.
+ const setRpm = useCallback(
+ (n: number | null) => {
+ const clean =
+ n == null || !Number.isFinite(n) || n <= 0
+ ? null
+ : Math.min(100000, Math.floor(n));
+ setRpmOverride(provider, rpmModel, clean);
+ setRpmTick((t) => t + 1);
+ },
+ [provider, rpmModel],
+ );
+
// Assemble a redacted debug snapshot (chat state + recent raw provider calls)
// for offline debugging, and hand it back as a pretty JSON string. The API key
// is never included.
@@ -891,6 +969,7 @@ export function useAskChat() {
setBusy(false);
setSweeping(false);
setReportUpdating(false);
+ setRateLimitNotice(null);
sweepPauseRef.current = false;
sendingRef.current = false;
abortRef.current = null;
@@ -1378,6 +1457,13 @@ export function useAskChat() {
// configurable batch size
sweepBatch,
setSweepBatch,
+ // requests-per-minute pacing (current provider+model)
+ rpmOverride,
+ rpmDefault,
+ effectiveRpm,
+ setRpm,
+ // rate-limiter pacing/backoff notice for the progress strip
+ rateLimitNotice,
// pausable sweep: a checkpoint (non-null ⇒ paused/resumable) + controls.
sweepCheckpoint,
sweepPaused: sweepCheckpoint != null && !sweeping,
diff --git a/export/app/lib/askProvider.ts b/export/app/lib/askProvider.ts
@@ -9,6 +9,13 @@
//
// Everything streams so the answer renders as it arrives.
+import {
+ acquire,
+ noteRateLimited,
+ parseRetryAfterMs,
+ withRateLimitRetry,
+} from "./rateLimit";
+
export type Provider = "anthropic" | "openai" | "gemini";
export type ChatRole = "user" | "assistant";
@@ -108,18 +115,29 @@ export type AskOptions = {
};
// Stream a completion from the chosen provider, invoking onDelta for each text
-// chunk and resolving with the full concatenated answer.
+// chunk and resolving with the full concatenated answer. Paced by the shared
+// rate limiter: `acquire` spaces this call to the model's RPM, and the dispatch
+// is wrapped in `withRateLimitRetry` so a 429 is honoured (retry the provider's
+// hint) rather than surfacing as a hard error mid-answer. One insertion covers
+// streamed answers AND the scripted gather loop (via askOnce).
export async function askStream(opts: AskOptions): Promise<string> {
- switch (opts.provider) {
- case "anthropic":
- return askAnthropic(opts);
- case "openai":
- return askOpenAI(opts);
- case "gemini":
- return askGemini(opts);
- default:
- throw new Error(`unknown provider: ${opts.provider}`);
- }
+ const model = opts.model || PROVIDERS[opts.provider]?.defaultModel || "";
+ await acquire(opts.provider, model, opts.signal);
+ return withRateLimitRetry(
+ () => {
+ switch (opts.provider) {
+ case "anthropic":
+ return askAnthropic(opts);
+ case "openai":
+ return askOpenAI(opts);
+ case "gemini":
+ return askGemini(opts);
+ default:
+ throw new Error(`unknown provider: ${opts.provider}`);
+ }
+ },
+ { provider: opts.provider, model, signal: opts.signal },
+ );
}
// Non-streaming single-shot completion, used by the scripted search loop to get
@@ -191,7 +209,33 @@ export function isContextLengthError(detail: string): boolean {
);
}
-async function ensureOk(res: Response, provider: string): Promise<void> {
+// Thrown when a request fails in a way that is transient and RESUMABLE rather
+// than fatal — a rate-limit / quota hit (HTTP 429). The retry wrapper catches it
+// to re-issue the request; when retries are exhausted the sweep driver catches it
+// and PAUSES (checkpoints the report + remaining batches) instead of aborting, so
+// the run can resume once the limit resets. `reason` is a short human string for
+// the paused banner; `retryAfterMs` carries the provider's own retry hint (when
+// present) so a retry waits exactly as long as asked.
+//
+// Hosted here (beside ContextTooLargeError) rather than in nativeTools/shared so
+// ensureOk can throw it without a shared↔askProvider import cycle; shared
+// re-exports it for its existing importers.
+export class PausableError extends Error {
+ reason: string;
+ retryAfterMs?: number;
+ constructor(message: string, reason: string, retryAfterMs?: number) {
+ super(message);
+ this.name = "PausableError";
+ this.reason = reason;
+ this.retryAfterMs = retryAfterMs;
+ }
+}
+
+async function ensureOk(
+ res: Response,
+ provider: string,
+ model: string,
+): Promise<void> {
if (res.ok) return;
let detail = "";
try {
@@ -203,6 +247,15 @@ async function ensureOk(res: Response, provider: string): Promise<void> {
if (res.status === 400 && isContextLengthError(detail)) {
throw new ContextTooLargeError(msg);
}
+ // Rate-limit / quota → resumable. Parse the provider's retry hint, step the
+ // session RPM estimate down, and throw PausableError so the retry wrapper can
+ // re-issue this request BEFORE any SSE body has been consumed (no duplicate
+ // text). Previously a streaming 429 was a generic hard error.
+ if (res.status === 429) {
+ const retryAfterMs = parseRetryAfterMs(res.headers, detail) ?? undefined;
+ noteRateLimited(provider.toLowerCase() as Provider, model);
+ throw new PausableError(msg, "usage limit", retryAfterMs);
+ }
throw new Error(msg);
}
@@ -234,7 +287,7 @@ async function askAnthropic(opts: AskOptions): Promise<string> {
stream: true,
}),
});
- await ensureOk(res, "Anthropic");
+ await ensureOk(res, "Anthropic", model);
const full: string[] = [];
let finishReason: FinishReason | undefined;
let rawFinish: string | undefined;
@@ -303,7 +356,7 @@ async function askOpenAI(opts: AskOptions): Promise<string> {
],
}),
});
- await ensureOk(res, "OpenAI");
+ await ensureOk(res, "OpenAI", model);
const full: string[] = [];
let finishReason: FinishReason | undefined;
let rawFinish: string | undefined;
@@ -368,7 +421,7 @@ async function askGemini(opts: AskOptions): Promise<string> {
: {}),
}),
});
- await ensureOk(res, "Gemini");
+ await ensureOk(res, "Gemini", model);
const full: string[] = [];
let finishReason: FinishReason | undefined;
let rawFinish: string | undefined;
diff --git a/export/app/lib/nativeTools/shared.ts b/export/app/lib/nativeTools/shared.ts
@@ -5,8 +5,23 @@
// search budget is reached). Only the request/response shapes differ per
// provider; this module holds the common types, tool text, and POST helper.
-import type { ChatMessage, DebugCall } from "../askProvider";
-import { ContextTooLargeError, isContextLengthError } from "../askProvider";
+import type { ChatMessage, DebugCall, Provider } from "../askProvider";
+import {
+ ContextTooLargeError,
+ isContextLengthError,
+ PausableError,
+} from "../askProvider";
+import {
+ acquire,
+ noteRateLimited,
+ parseRetryAfterMs,
+ withRateLimitRetry,
+} from "../rateLimit";
+
+// Re-exported for existing importers: PausableError now lives in askProvider (so
+// ensureOk can throw it without a cycle), but the native gather code + the sweep
+// driver still import it from here.
+export { PausableError } from "../askProvider";
export type NativeGatherContext = {
apiKey: string;
@@ -134,27 +149,16 @@ export type ParsedToolCall = {
// and retries the turn with the scripted transport.
export class ToolsUnavailableError extends Error {}
-// Thrown when a request fails in a way that is transient and RESUMABLE rather
-// than fatal — a rate-limit / quota hit (HTTP 429). Mid-sweep the driver catches
-// this and PAUSES (checkpoints the report + remaining batches) instead of
-// aborting, so the run can resume once the limit resets. `reason` is a short
-// human string for the paused banner.
-export class PausableError extends Error {
- reason: string;
- constructor(message: string, reason: string) {
- super(message);
- this.name = "PausableError";
- this.reason = reason;
- }
-}
-
export function abortError(): DOMException {
return new DOMException("Aborted", "AbortError");
}
// POST JSON and parse the response. 400/404 → ToolsUnavailableError (so the agent
// can fall back to scripted); other non-OK statuses throw a normal Error the chat
-// surfaces (e.g. 401 bad key).
+// surfaces (e.g. 401 bad key). Paced by the shared rate limiter: `acquire` spaces
+// this call to the model's RPM, and the fetch is wrapped in `withRateLimitRetry`
+// so a 429 retries the provider's hint before falling through to a PausableError
+// (which the sweep driver still turns into a resumable pause once retries run out).
export async function postJson(
url: string,
headers: Record<string, string>,
@@ -163,6 +167,32 @@ export async function postJson(
signal?: AbortSignal,
debug?: { onDebug?: (rec: DebugCall) => void; model?: string },
): Promise<unknown> {
+ const providerId = provider.toLowerCase() as Provider;
+ // Model comes from the request body (anthropic/openai) or, for gemini (model in
+ // the URL), the debug hint — either way it keys the per-model pacing gate.
+ const bodyModel = (body as { model?: unknown })?.model;
+ const model =
+ (typeof bodyModel === "string" ? bodyModel : undefined) ?? debug?.model ?? "";
+
+ if (signal?.aborted) throw abortError();
+ await acquire(providerId, model, signal);
+ return withRateLimitRetry(
+ () => postJsonOnce(url, headers, body, provider, providerId, model, signal, debug),
+ { provider: providerId, model, signal },
+ );
+}
+
+// One request attempt, re-issuable by withRateLimitRetry.
+async function postJsonOnce(
+ url: string,
+ headers: Record<string, string>,
+ body: unknown,
+ provider: string,
+ providerId: Provider,
+ model: string,
+ signal?: AbortSignal,
+ debug?: { onDebug?: (rec: DebugCall) => void; model?: string },
+): Promise<unknown> {
if (signal?.aborted) throw abortError();
const t0 = Date.now();
const res = await fetch(url, {
@@ -188,10 +218,13 @@ export async function postJson(
if (res.status === 400 || res.status === 404) {
throw new ToolsUnavailableError(msg);
}
- // Rate-limit / quota → resumable. Let the sweep driver pause + checkpoint
- // rather than surface a hard error and lose the run.
+ // Rate-limit / quota → resumable. Parse the provider's retry hint + step the
+ // session RPM down; the retry wrapper re-issues this request, and only after
+ // its attempts are exhausted does this bubble to the sweep driver's pause.
if (res.status === 429) {
- throw new PausableError(msg, "usage limit");
+ const retryAfterMs = parseRetryAfterMs(res.headers, detail) ?? undefined;
+ noteRateLimited(providerId, model);
+ throw new PausableError(msg, "usage limit", retryAfterMs);
}
throw new Error(msg);
}
diff --git a/export/app/lib/rateLimit.test.ts b/export/app/lib/rateLimit.test.ts
@@ -0,0 +1,202 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+
+// Minimal in-memory localStorage so the override store's window-guarded functions
+// run in node exactly as they do in the browser.
+class MemStorage {
+ private map = new Map<string, string>();
+ getItem(k: string): string | null {
+ return this.map.has(k) ? (this.map.get(k) as string) : null;
+ }
+ setItem(k: string, v: string): void {
+ this.map.set(k, String(v));
+ }
+ removeItem(k: string): void {
+ this.map.delete(k);
+ }
+ clear(): void {
+ this.map.clear();
+ }
+}
+const mem = new MemStorage();
+(globalThis as { window?: unknown }).window = { localStorage: mem };
+
+import {
+ __resetSessionRpm,
+ __setClock,
+ acquire,
+ defaultRpmFor,
+ effectiveRpm,
+ getRpmOverride,
+ noteRateLimited,
+ parseRetryAfterMs,
+ setRpmOverride,
+ type Clock,
+} from "./rateLimit";
+
+const RPM_KEY = "ytdlp-tb:ai:rpm";
+
+function resetAll() {
+ mem.clear();
+ __resetSessionRpm();
+}
+
+// ─── defaultRpmFor: model-pattern matching ───
+
+test("defaultRpmFor matches models by pattern per provider", () => {
+ assert.equal(defaultRpmFor("gemini", "gemini-2.5-pro"), 5);
+ assert.equal(defaultRpmFor("gemini", "gemini-1.5-pro-latest"), 5);
+ assert.equal(defaultRpmFor("gemini", "gemini-2.5-flash"), 10);
+ // flash-lite must win over the plain "flash" rule (its id contains "flash").
+ assert.equal(defaultRpmFor("gemini", "gemini-2.5-flash-lite"), 15);
+ assert.equal(defaultRpmFor("gemini", "something-else"), 10); // fallback
+ assert.equal(defaultRpmFor("anthropic", "claude-haiku-4-5"), 50);
+ assert.equal(defaultRpmFor("anthropic", "claude-opus-4-8"), 50);
+ assert.equal(defaultRpmFor("openai", "gpt-4o-mini"), 60);
+});
+
+// ─── parseRetryAfterMs: header seconds, HTTP-date, Gemini RetryInfo ───
+
+test("parseRetryAfterMs reads Retry-After header in seconds", () => {
+ const h = new Headers({ "retry-after": "2" });
+ assert.equal(parseRetryAfterMs(h), 2000);
+});
+
+test("parseRetryAfterMs reads Retry-After as an HTTP-date (against the clock)", () => {
+ const base = Date.parse("2030-01-01T00:00:00Z");
+ const restore = __setClock({ now: () => base, sleep: () => Promise.resolve() });
+ try {
+ const h = new Headers({ "retry-after": new Date(base + 5000).toUTCString() });
+ assert.equal(parseRetryAfterMs(h), 5000);
+ } finally {
+ restore();
+ }
+});
+
+test("parseRetryAfterMs reads Gemini's RetryInfo body form", () => {
+ const body = JSON.stringify({
+ error: {
+ code: 429,
+ details: [
+ { "@type": "type.googleapis.com/google.rpc.RetryInfo", retryDelay: "38s" },
+ ],
+ },
+ });
+ assert.equal(parseRetryAfterMs(null, body), 38000);
+});
+
+test("parseRetryAfterMs falls back to a loose scan, else null", () => {
+ assert.equal(parseRetryAfterMs(null, "please retry in 7s and try again"), 7000);
+ assert.equal(parseRetryAfterMs(new Headers(), "no hint here"), null);
+ assert.equal(parseRetryAfterMs(null), null);
+});
+
+// ─── limiter spacing + abort, via the injected clock ───
+
+test("acquire spaces same-model calls by 60000/rpm (injected clock, no real time)", async () => {
+ resetAll();
+ let t = 1000;
+ const sleeps: number[] = [];
+ const clock: Clock = {
+ now: () => t,
+ sleep: (ms) => {
+ sleeps.push(ms);
+ t += ms; // advance virtual time as if we slept
+ return Promise.resolve();
+ },
+ };
+ const restore = __setClock(clock);
+ try {
+ // anthropic default 50 RPM → 1200ms spacing.
+ await acquire("anthropic", "spacing-model");
+ await acquire("anthropic", "spacing-model");
+ await acquire("anthropic", "spacing-model");
+ assert.deepEqual(sleeps, [1200, 1200]); // first call is free, next two paced
+ } finally {
+ restore();
+ }
+});
+
+test("acquire aborts a pacing wait promptly", async () => {
+ resetAll();
+ const clock: Clock = {
+ now: () => 5000, // constant, so the 2nd acquire has a nonzero wait
+ sleep: (_ms, signal) =>
+ new Promise<void>((_resolve, reject) => {
+ if (signal?.aborted) {
+ reject(new DOMException("Aborted", "AbortError"));
+ return;
+ }
+ signal?.addEventListener(
+ "abort",
+ () => reject(new DOMException("Aborted", "AbortError")),
+ { once: true },
+ );
+ }),
+ };
+ const restore = __setClock(clock);
+ try {
+ await acquire("anthropic", "abort-model"); // sets last = 5000
+ const ac = new AbortController();
+ const pending = acquire("anthropic", "abort-model", ac.signal);
+ ac.abort();
+ await assert.rejects(pending, (e: Error) => e.name === "AbortError");
+ } finally {
+ restore();
+ }
+});
+
+// ─── step-down arithmetic ───
+
+test("noteRateLimited steps the session RPM down by 0.6, floored, min 1", () => {
+ resetAll();
+ assert.equal(noteRateLimited("anthropic", "claude-x"), 30); // floor(50 * 0.6)
+ assert.equal(noteRateLimited("anthropic", "claude-x"), 18); // floor(30 * 0.6)
+ assert.equal(noteRateLimited("anthropic", "claude-x"), 10); // floor(18 * 0.6)
+ // gemini pro starts at 5 → 3 → 1 → 1 (never below 1).
+ assert.equal(noteRateLimited("gemini", "gemini-2.5-pro"), 3);
+ assert.equal(noteRateLimited("gemini", "gemini-2.5-pro"), 1);
+ assert.equal(noteRateLimited("gemini", "gemini-2.5-pro"), 1);
+});
+
+test("effectiveRpm precedence: override ?? session step-down ?? default", () => {
+ resetAll();
+ // default
+ assert.equal(effectiveRpm("gemini", "gemini-2.5-pro"), 5);
+ // session step-down beats default
+ noteRateLimited("gemini", "gemini-2.5-pro"); // → 3
+ assert.equal(effectiveRpm("gemini", "gemini-2.5-pro"), 3);
+ // an explicit override beats the session step-down
+ setRpmOverride("gemini", "gemini-2.5-pro", 42);
+ assert.equal(effectiveRpm("gemini", "gemini-2.5-pro"), 42);
+});
+
+// ─── override store: parse-or-reset + get/set + wildcard ───
+
+test("override store round-trips and clears", () => {
+ resetAll();
+ assert.equal(getRpmOverride("openai", "gpt-4o"), null);
+ setRpmOverride("openai", "gpt-4o", 120);
+ assert.equal(getRpmOverride("openai", "gpt-4o"), 120);
+ setRpmOverride("openai", "gpt-4o", null);
+ assert.equal(getRpmOverride("openai", "gpt-4o"), null);
+});
+
+test("override store resets a malformed payload rather than trapping", () => {
+ resetAll();
+ mem.setItem(RPM_KEY, "}{ not json");
+ assert.equal(getRpmOverride("openai", "gpt-4o"), null);
+ // The malformed key was cleared.
+ assert.equal(mem.getItem(RPM_KEY), null);
+});
+
+test("a '*' wildcard override applies to every model", () => {
+ resetAll();
+ mem.setItem(RPM_KEY, JSON.stringify({ "*": 9999 }));
+ assert.equal(getRpmOverride("gemini", "gemini-2.5-pro"), 9999);
+ assert.equal(getRpmOverride("anthropic", "claude-haiku-4-5"), 9999);
+ // An exact key still wins over the wildcard.
+ setRpmOverride("anthropic", "claude-haiku-4-5", 5);
+ assert.equal(getRpmOverride("anthropic", "claude-haiku-4-5"), 5);
+ assert.equal(getRpmOverride("gemini", "gemini-2.5-pro"), 9999);
+});
diff --git a/export/app/lib/rateLimit.ts b/export/app/lib/rateLimit.ts
@@ -0,0 +1,394 @@
+// 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<Provider, { rules: RpmRule[]; fallback: number }> = {
+ 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 "<provider>:<model>" → 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<string, number> {
+ 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<string, number> = {};
+ for (const [k, v] of Object.entries(parsed as Record<string, unknown>)) {
+ 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 "<provider>:<model>" 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<string, number>();
+
+// 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<Listener>();
+
+// 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<void>;
+};
+
+function abortError(): DOMException {
+ return new DOMException("Aborted", "AbortError");
+}
+
+function realSleep(ms: number, signal?: AbortSignal): Promise<void> {
+ if (ms <= 0) return Promise.resolve();
+ return new Promise<void>((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<void>; last: number };
+const gates = new Map<string, Gate>();
+
+// 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<void> {
+ 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<void>((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<T>(
+ fn: () => Promise<T>,
+ opts: {
+ provider: Provider;
+ model: string;
+ signal?: AbortSignal;
+ attempts?: number;
+ },
+): Promise<T> {
+ 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");
+}
diff --git a/export/e2e/ask-workspace.spec.ts b/export/e2e/ask-workspace.spec.ts
@@ -227,14 +227,17 @@ test.describe("grounding workspace", () => {
await expect(page.getByText("what happened")).toBeVisible();
});
- test("a 429 mid-sweep PAUSES (not errors); Resume — even after a model switch — finishes it", async ({
+ // A transient 429 (with a Retry-After hint) is now RETRIED below the sweep
+ // driver — the run paces itself and carries on, no pause banner. Only a 429 that
+ // outlasts the retry cap falls through to the pause/checkpoint (next test).
+ test("a transient 429 (Retry-After) mid-sweep auto-recovers and completes, with no pause", async ({
page,
}) => {
await installRoutes(page);
// Force batch size 2 so the 3 fixtures span 2 batches.
await page.addInitScript(() => localStorage.setItem("ytdlp-tb:ai:sweepchunk", "2"));
- let batch2Attempts = 0;
+ let batch2FirstRoundHits = 0;
await page.route("https://api.anthropic.com/**", async (route) => {
if (route.request().method() === "OPTIONS") {
await route.fulfill({ status: 204, headers: CORS });
@@ -252,14 +255,14 @@ test.describe("grounding workspace", () => {
const text = typeof lastUser?.content === "string" ? lastUser.content : "";
if (Array.isArray(body.tools) && /running report/.test(system)) {
const batch = /Batch (\d+) of/.exec(text)?.[1] ?? "1";
- // Batch 2's FIRST attempt hits a rate limit (429) → pause. Its later
- // (resumed) attempt succeeds.
- if (batch === "2") {
- batch2Attempts += 1;
- if (batch2Attempts === 1) {
+ // Batch 2's first request 429s ONCE, with a 1s Retry-After hint. The
+ // limiter waits the hint and re-issues; the retry succeeds.
+ if (batch === "2" && !isToolResult) {
+ batch2FirstRoundHits += 1;
+ if (batch2FirstRoundHits === 1) {
await route.fulfill({
status: 429,
- headers: { ...CORS, "content-type": "application/json" },
+ headers: { ...CORS, "content-type": "application/json", "retry-after": "1" },
body: JSON.stringify({ type: "error", error: { message: "rate limited" } }),
});
return;
@@ -296,10 +299,89 @@ test.describe("grounding workspace", () => {
await page.getByRole("button", { name: /Build report from all 3 results/ }).click();
- // Batch 1 folds in; batch 2's 429 pauses the run (not a hard error).
+ // Both batches fold in — the 429 was absorbed by a retry, not a pause.
+ await expect(page.getByText("Finding recorded in batch 1.")).toBeVisible({ timeout: 15_000 });
+ await expect(page.getByText("Finding recorded in batch 2.")).toBeVisible({ timeout: 15_000 });
+ // No pause banner, and the sweep reported completion.
+ await expect(page.getByTestId("sweep-paused")).toHaveCount(0);
+ await expect(page.getByText(/Built a report from 3 results/)).toBeVisible();
+ });
+
+ // A 429 that persists past the retry cap is the genuine last resort: it falls
+ // through to the existing pause/checkpoint so the user can Resume later (e.g. a
+ // daily-quota cap that clears tomorrow). Guards that regression path.
+ test("a persistent 429 past the retry cap PAUSES; Resume — even after a model switch — finishes it", async ({
+ page,
+ }) => {
+ await installRoutes(page);
+ await page.addInitScript(() => localStorage.setItem("ytdlp-tb:ai:sweepchunk", "2"));
+
+ let batch2Hits = 0;
+ await page.route("https://api.anthropic.com/**", async (route) => {
+ if (route.request().method() === "OPTIONS") {
+ await route.fulfill({ status: 204, headers: CORS });
+ return;
+ }
+ const body = route.request().postDataJSON() as {
+ system?: string;
+ tools?: unknown[];
+ messages?: { role: string; content: unknown }[];
+ };
+ const system = body.system ?? "";
+ const msgs = body.messages ?? [];
+ const lastUser = [...msgs].reverse().find((m) => m.role === "user");
+ const isToolResult = Array.isArray(lastUser?.content);
+ const text = typeof lastUser?.content === "string" ? lastUser.content : "";
+ if (Array.isArray(body.tools) && /running report/.test(system)) {
+ const batch = /Batch (\d+) of/.exec(text)?.[1] ?? "1";
+ // Batch 2 429s on every attempt of the INITIAL run (enough to exhaust the
+ // retry cap → pause). The later resumed run's attempts succeed.
+ if (batch === "2" && !isToolResult) {
+ batch2Hits += 1;
+ if (batch2Hits <= 4) {
+ await route.fulfill({
+ status: 429,
+ headers: { ...CORS, "content-type": "application/json", "retry-after": "1" },
+ body: JSON.stringify({ type: "error", error: { message: "rate limited" } }),
+ });
+ return;
+ }
+ }
+ if (isToolResult) {
+ await route.fulfill({
+ status: 200,
+ headers: { ...CORS, "content-type": "application/json" },
+ body: toolUse("finish", {}),
+ });
+ return;
+ }
+ await route.fulfill({
+ status: 200,
+ headers: { ...CORS, "content-type": "application/json" },
+ body: toolUse("update_report", {
+ section: `Batch ${batch} finding`,
+ content: `Finding recorded in batch ${batch}.`,
+ }),
+ });
+ return;
+ }
+ await route.fulfill({
+ status: 200,
+ headers: { ...CORS, "content-type": "text/event-stream" },
+ body: sse("DONE"),
+ });
+ });
+
+ await searchThenChat(page, 3);
+ await keyIn(page);
+ await expect(page.getByText(/Grounded in 3 results/)).toBeVisible();
+
+ await page.getByRole("button", { name: /Build report from all 3 results/ }).click();
+
+ // Batch 1 folds in; batch 2's retries all 429 → the run pauses (not errors).
await expect(page.getByText("Finding recorded in batch 1.")).toBeVisible({ timeout: 15_000 });
const paused = page.getByTestId("sweep-paused");
- await expect(paused).toBeVisible();
+ await expect(paused).toBeVisible({ timeout: 15_000 });
await expect(paused).toContainText("usage limit");
// The partial report survived.
await expect(page.getByText("Finding recorded in batch 2.")).toHaveCount(0);
@@ -312,7 +394,6 @@ test.describe("grounding workspace", () => {
// The resumed batch completes and the report is now whole.
await expect(page.getByText("Finding recorded in batch 2.")).toBeVisible({ timeout: 15_000 });
await expect(page.getByTestId("sweep-paused")).toHaveCount(0);
- // Both batches are in the report.
await expect(page.getByText("Finding recorded in batch 1.")).toBeVisible();
});
});
diff --git a/export/e2e/helpers.ts b/export/e2e/helpers.ts
@@ -21,6 +21,14 @@ async function fulfillJson(route: Route, body: unknown) {
}
export async function installRoutes(page: Page) {
+ // Neutralise the client-side AI rate limiter for the whole suite: a "*" RPM
+ // override raises every model's ceiling so the per-request pacing gate never
+ // adds real waits (which would slow every AI test). Throttle-specific tests opt
+ // back in by driving 429s + Retry-After headers, which exercise the RETRY path
+ // regardless of this ceiling.
+ await page.addInitScript(() => {
+ localStorage.setItem("ytdlp-tb:ai:rpm", JSON.stringify({ "*": 100000 }));
+ });
await page.route("**/summaries/manifest.json", async (route) => {
await fulfillJson(route, summariesManifest());
});