diff --git a/assets/request-pacing-dashboard.jpg b/assets/request-pacing-dashboard.jpg new file mode 100644 index 0000000000..8a4f6008cf Binary files /dev/null and b/assets/request-pacing-dashboard.jpg differ diff --git a/docs-site/src/content/docs/ja/reference/configuration/providers.md b/docs-site/src/content/docs/ja/reference/configuration/providers.md index a990b52263..c680faf8ed 100644 --- a/docs-site/src/content/docs/ja/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ja/reference/configuration/providers.md @@ -56,6 +56,7 @@ account を削除しても mapping は保持され、同じ id を再追加す | --- | --- | --- | | `adapter` | `string` | `openai-chat`、`openai-responses`、`anthropic`、`google`、`kiro`、`cursor`、`azure-openai` (または別名 `azure`) のいずれか。 | | `baseUrl` | `string` |アップストリーム API のベース URL。ほとんどの組み込み固定エンドポイントは不一致を無視します。衝突安全キー プリセットは、古い同じ名前のカスタム宛先を保持します。 | +| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | オプションの送信開始間隔調整。プロバイダー制限は全モデルに適用され、モデル別設定は遅延を増やす場合のみ有効です。キュー待機は応答ヘッダーのタイムアウトを消費しません。 | | `responsesPath?` | `string` |キー認証 `openai-responses` リクエストの相対リソース パス。 `/` で始まり、スキーム、クエリ、またはフラグメントが含まれていない必要があります。 | | `supportsServiceTier?` | `boolean` | `service_tier` ケイパビリティの 3 状態です。`true`: fast モードが注入でき、呼び出し元の値も保持されます。`false`: フィールドは削除され、注入もされません (非対応と文書化されたアップストリームには送りません)。未設定: 未分類 — 呼び出し元の値はそのまま保持され、fast モードは注入しません。レジストリは正規 OpenAI (`true`)、DeepSeek、Volcengine Ark (`false`) を分類します。実際にティアをサポートするカスタム ゲートウェイにのみ明示的に設定してください。 | | `preserveResponsesReasoningContent?` | `boolean` | リプレイされる Responses reasoning アイテムの平文 reasoning コンテンツを消去せずに保持します (消去は ChatGPT バックエンドのルールです)。DeepSeek のように reasoning リプレイを受け入れるアップストリームで有効にしてください。プロキシ生成の `ocxr1` エンベロープは常に削除されます。 | diff --git a/docs-site/src/content/docs/ko/reference/configuration/providers.md b/docs-site/src/content/docs/ko/reference/configuration/providers.md index 7aed7c7771..eaea77f373 100644 --- a/docs-site/src/content/docs/ko/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ko/reference/configuration/providers.md @@ -56,6 +56,7 @@ managed map을 활성화하면 privacy-safe selector를 만들고, 이후 계정 | --- | --- | --- | | `adapter` | `string` | `openai-chat`, `openai-responses`, `anthropic`, `google`, `kiro`, `cursor`, `azure-openai` 중 하나이며, `azure`는 별칭입니다. | | `baseUrl` | `string` | 상위 API 기본 URL입니다. 대부분의 내장 고정 엔드포인트는 불일치를 무시합니다. 충돌 안전 키 프리셋은 같은 이름의 이전 사용자 지정 목적지를 보존합니다. | +| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | 선택적 아웃바운드 요청 시작 속도 조절입니다. Provider 제한은 모든 모델에 적용되고 정확한 모델 override는 지연을 더 늘릴 때만 적용됩니다. 큐 대기는 응답 헤더 타임아웃을 소모하지 않습니다. | | `responsesPath?` | `string` | 키 인증 `openai-responses` 요청의 상대 리소스 경로입니다. 반드시 `/`로 시작해야 하며 스킴, query, fragment를 포함하면 안 됩니다. | | `supportsServiceTier?` | `boolean` | `service_tier` 케이퍼빌리티 3상태입니다. `true`: fast 모드가 주입할 수 있고 호출자 값도 보존합니다. `false`: 필드를 제거하고 절대 주입하지 않습니다(미지원으로 문서화된 업스트림에는 볼 수 없습니다). 미설정: 미분류 — 호출자가 준 값은 그대로 보존하고 fast 모드는 주입하지 않습니다. 레지스트리는 정식 OpenAI(`true`), DeepSeek, Volcengine Ark(`false`)를 분류하며, 실제로 티어를 지원하는 커스텀 게이트웨이에만 명시적으로 설정하세요. | | `preserveResponsesReasoningContent?` | `boolean` | 리플레이되는 Responses reasoning 항목의 평문 reasoning 내용을 지우지 않고 유지합니다(지우는 것은 ChatGPT 백엔드 규칙입니다). DeepSeek처럼 reasoning 리플레이를 허용하는 업스트림에 켜세요. 프록시가 만든 `ocxr1` 봉투는 항상 제거됩니다. | diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index 62203e285d..05f288d12e 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -66,6 +66,7 @@ differing backup and rewrites known legacy namespaced selected ids to bare ids. | --- | --- | --- | | `adapter` | `string` | One of `openai-chat`, `openai-responses`, `anthropic`, `google`, `kiro`, `cursor`, `azure-openai` (or alias `azure`). | | `baseUrl` | `string` | Upstream API base URL. Most built-in fixed endpoints ignore a mismatch; collision-safe key presets preserve an older same-named custom destination. | +| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | Optional outbound request-start pacing. RPM is converted to an even interval; `minIntervalMs` may impose a longer interval. Provider limits apply across all models, while exact model overrides can only add delay. Queue waits do not consume the upstream response-header timeout. HTTP, Responses WebSocket, and adapter `fetchResponse` transports are covered; custom `runTurn` transports are not. | | `responsesPath?` | `string` | Relative resource path for key-auth `openai-responses` requests. It must start with `/` and contain no scheme, query, or fragment. | | `supportsServiceTier?` | `boolean` | Tri-state `service_tier` capability. `true`: fast mode may inject and caller values are preserved. `false`: the field is stripped and never injected (the upstream documented as not supporting it must not receive it). Absent: the provider is unclassified — caller-supplied values are preserved untouched and fast mode never injects. The registry classifies canonical OpenAI (`true`), DeepSeek, and Volcengine Ark (`false`); set it explicitly only for custom gateways that genuinely support tiers. | | `preserveResponsesReasoningContent?` | `boolean` | Keep plaintext reasoning content on replayed Responses reasoning items instead of blanking it (blanking is the ChatGPT backend's rule). Enable for upstreams whose contract accepts reasoning replay, such as DeepSeek. Proxy-minted `ocxr1` envelopes are always stripped. | diff --git a/docs-site/src/content/docs/ru/reference/configuration/providers.md b/docs-site/src/content/docs/ru/reference/configuration/providers.md index b1bb26c78c..f0de52a3a9 100644 --- a/docs-site/src/content/docs/ru/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ru/reference/configuration/providers.md @@ -69,6 +69,7 @@ cross-route credential fallback не существует. Строки API GPT- | --- | --- | --- | | `adapter` | `string` | Один из `openai-chat`, `openai-responses`, `anthropic`, `google`, `kiro`, `cursor`, `azure-openai` (или alias `azure`). | | `baseUrl` | `string` | Базовый URL API upstream'а. Большинство built-in fixed-endpoint'ов игнорируют несовпадение; collision-safe key-preset'ы сохраняют старый custom destination с тем же именем. | +| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | Опциональная равномерная задержка начала исходящих запросов. Лимит провайдера действует на все модели, а правила моделей могут только увеличить задержку. Ожидание очереди не расходует таймаут заголовков ответа. | | `responsesPath?` | `string` | Relative resource path для key-auth запросов `openai-responses`. Должен начинаться с `/` и не может содержать scheme, query или fragment. | | `supportsServiceTier?` | `boolean` | Три состояния поддержки `service_tier`. `true`: fast mode может подставлять поле, значения вызывающего сохраняются. `false`: поле удаляется и никогда не подставляется (апстрим, для которого задокументировано отсутствие поддержки, не должен его получать). Не задано: провайдер не классифицирован — значения вызывающего сохраняются без изменений, fast mode не подставляет. Registry классифицирует canonical OpenAI (`true`), DeepSeek и Volcengine Ark (`false`); задавайте явно только для custom gateway'ев, реально поддерживающих tier'ы. | | `preserveResponsesReasoningContent?` | `boolean` | Сохранять plaintext reasoning content в replay'нутых Responses reasoning item'ах вместо очистки (очистка — правило ChatGPT backend'а). Включайте для upstream'ов, чей контракт принимает reasoning replay, например DeepSeek. Proxy-minted `ocxr1` envelope'ы удаляются всегда. | diff --git a/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md b/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md index 1d2a8c10c9..2ac6f65e93 100644 --- a/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md +++ b/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md @@ -56,6 +56,7 @@ selector,而不是分配一个新名称。 | --- | --- | --- | | `adapter` | `string` | `openai-chat`、`openai-responses`、`anthropic`、`google`、`kiro`、`cursor`、`azure-openai`(或别名 `azure`)之一。 | | `baseUrl` | `string` | 上游 API 基础 URL。大多数内置固定端点会忽略不匹配的值;具备冲突安全键的预设会保留一个更早、同名的自定义目标。 | +| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | 可选的出站请求启动节流。提供商限制适用于所有模型,模型规则只能增加延迟。排队等待不计入响应头超时。 | | `responsesPath?` | `string` | 用于 key-auth `openai-responses` 请求的相对资源路径。必须以 `/` 开头,且不能包含 scheme、query 或 fragment。 | | `supportsServiceTier?` | `boolean` | `service_tier` 能力的三态。`true`:fast 模式可以注入,调用方提供的值也会被保留。`false`:剥离该字段且绝不注入(已明确不支持的上游不会收到它)。未设置:未分类——调用方提供的值原样保留,fast 模式绝不注入。注册表已对官方 OpenAI(`true`)、DeepSeek 和 Volcengine Ark(`false`)分类;仅对真正支持分层的自定义网关显式设置。 | | `preserveResponsesReasoningContent?` | `boolean` | 在重放的 Responses reasoning 项中保留明文 reasoning 内容,而不是清空(清空是 ChatGPT 后端的规则)。对接受 reasoning 重放的上游(如 DeepSeek)启用。代理生成的 `ocxr1` 信封始终会被剥离。 | diff --git a/docs-site/src/content/docs/zh-tw/reference/configuration/providers.md b/docs-site/src/content/docs/zh-tw/reference/configuration/providers.md index 2f729dc355..c4f5419a0f 100644 --- a/docs-site/src/content/docs/zh-tw/reference/configuration/providers.md +++ b/docs-site/src/content/docs/zh-tw/reference/configuration/providers.md @@ -40,6 +40,7 @@ description: 供應商項目、認證、端點、模型目錄、配額、context | --- | --- | --- | | `adapter` | `string` | `openai-chat`、`openai-responses`、`anthropic`、`google`、`kiro`、`cursor`、`azure-openai`(或別名 `azure`)之一。 | | `baseUrl` | `string` | 上游 API base URL。多數內建固定端點忽略不符;碰撞安全的金鑰預設保留較舊的同名自訂目的地。 | +| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | 選用的出站請求啟動節流。供應商限制適用於所有模型,模型規則只能增加延遲。排隊等待不計入回應標頭逾時。 | | `responsesPath?` | `string` | Key-auth `openai-responses` 請求的相對資源路徑。必須以 `/` 開頭且不含 scheme、query 或 fragment。 | | `disabled?` | `boolean` | 將供應商保留在磁碟上但排除於路由與模型/目錄清單。 | | `apiKey?` | `string` | API 金鑰,或在請求時解析的 `${ENV_VAR}` / `$ENV_VAR` 參考。 | diff --git a/gui/src/components/provider-workspace/ProviderSettings.tsx b/gui/src/components/provider-workspace/ProviderSettings.tsx index 26ca85fa35..48151ad290 100644 --- a/gui/src/components/provider-workspace/ProviderSettings.tsx +++ b/gui/src/components/provider-workspace/ProviderSettings.tsx @@ -22,6 +22,31 @@ const ADAPTERS = ["openai-responses", "openai-chat", "anthropic", "google", "azu const EMPTY_MODELS: string[] = []; type ChoicesStatus = "idle" | "loading" | "ready" | "error"; +type PacingRule = { requestsPerMinute?: number; minIntervalMs?: number }; +type PacingStatus = { enabled: boolean; queued: number; nextSlotInMs: number; lastStartedAt?: number; lastModelId?: string }; + +function numberDraft(value: number | undefined): string { return value === undefined ? "" : String(value); } +function positiveRpm(value: string): number | undefined { + if (!value.trim()) return undefined; + const parsed = Number(value); + return Number.isFinite(parsed) && parsed >= 1 / 60 ? parsed : undefined; +} +function positiveInteger(value: string): number | undefined { + if (!value.trim()) return undefined; + const parsed = Number(value); + return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : undefined; +} +function pacingSignature(value: WorkspaceItem["requestPacing"] | undefined): string { + const models = Object.entries(value?.models ?? {}) + .sort(([left], [right]) => left.localeCompare(right)) + .map(([model, rule]) => [model, rule.requestsPerMinute ?? null, rule.minIntervalMs ?? null]); + return JSON.stringify([ + value?.enabled === true, + value?.requestsPerMinute ?? null, + value?.minIntervalMs ?? null, + models, + ]); +} export default function ProviderSettings({ item, availableModels = EMPTY_MODELS, apiBase, onUpdateProvider, onDirtyChange, onRegisterSave, @@ -53,6 +78,14 @@ export default function ProviderSettings({ const [baseUrlChoices, setBaseUrlChoices] = useState(); const [choicesStatus, setChoicesStatus] = useState(apiBase ? "loading" : "idle"); const [endpointChoice, setEndpointChoice] = useState(() => "custom"); + const [pacingEnabled, setPacingEnabled] = useState(item.requestPacing?.enabled === true); + const [pacingRpm, setPacingRpm] = useState(() => numberDraft(item.requestPacing?.requestsPerMinute)); + const [pacingDelay, setPacingDelay] = useState(() => numberDraft(item.requestPacing?.minIntervalMs)); + const [pacingModels, setPacingModels] = useState>(() => ({ ...(item.requestPacing?.models ?? {}) })); + const [pacingModelId, setPacingModelId] = useState(""); + const [pacingModelRpm, setPacingModelRpm] = useState(""); + const [pacingModelDelay, setPacingModelDelay] = useState(""); + const [pacingStatus, setPacingStatus] = useState(null); /* eslint-disable react-hooks/set-state-in-effect -- intentional form reset when saved provider fields change */ useEffect(() => { @@ -64,10 +97,14 @@ export default function ProviderSettings({ setNote(item.note ?? ""); setAllowPrivateNetwork(item.allowPrivateNetwork ?? false); setLiveModels(item.liveModels !== false); + setPacingEnabled(item.requestPacing?.enabled === true); + setPacingRpm(numberDraft(item.requestPacing?.requestsPerMinute)); + setPacingDelay(numberDraft(item.requestPacing?.minIntervalMs)); + setPacingModels({ ...(item.requestPacing?.models ?? {}) }); setMsg(null); setModeMsg(null); queueMicrotask(() => setEndpointChoice(matchChoiceId(baseUrlChoices, item.baseUrl))); - }, [item.adapter, item.baseUrl, item.defaultModel, item.authMode, item.apiKeyTransport, item.keyOptional, item.note, item.allowPrivateNetwork, item.liveModels, baseUrlChoices]); + }, [item.adapter, item.baseUrl, item.defaultModel, item.authMode, item.apiKeyTransport, item.keyOptional, item.note, item.allowPrivateNetwork, item.liveModels, item.requestPacing, baseUrlChoices]); /* eslint-enable react-hooks/set-state-in-effect */ // Account mode syncs on its own: a mode PATCH refresh must not reset an in-progress @@ -109,6 +146,27 @@ export default function ProviderSettings({ // eslint-disable-next-line react-hooks/exhaustive-deps -- item.baseUrl sync is handled by the form-reset effect }, [apiBase, item.name]); + useEffect(() => { + if (!apiBase) return; + let active = true; + const load = () => { + fetch(`${apiBase}/api/provider-request-pacing?name=${encodeURIComponent(item.name)}`) + .then(r => readJsonIfOk(r)) + .then(status => { if (active && status) setPacingStatus(status); }) + .catch(() => undefined); + }; + load(); + const timer = window.setInterval(load, 2_000); + return () => { active = false; window.clearInterval(timer); }; + }, [apiBase, item.name]); + + const pacingDraft = useMemo(() => ({ + enabled: pacingEnabled, + ...(positiveRpm(pacingRpm) !== undefined ? { requestsPerMinute: positiveRpm(pacingRpm) } : {}), + ...(positiveInteger(pacingDelay) !== undefined ? { minIntervalMs: positiveInteger(pacingDelay) } : {}), + ...(Object.keys(pacingModels).length > 0 ? { models: pacingModels } : {}), + }), [pacingDelay, pacingEnabled, pacingModels, pacingRpm]); + const dirty = adapter.trim() !== item.adapter || baseUrl.trim() !== item.baseUrl || defaultModel.trim() !== (item.defaultModel ?? "") @@ -117,8 +175,10 @@ export default function ProviderSettings({ || note.trim() !== (item.note ?? "") || allowPrivateNetwork !== (item.allowPrivateNetwork ?? false) || liveModels !== (item.liveModels !== false); + const pacingDirty = pacingSignature(pacingDraft) !== pacingSignature(item.requestPacing); + const formDirty = dirty || pacingDirty; - useEffect(() => { onDirtyChange?.(dirty); return () => onDirtyChange?.(false); }, [dirty, onDirtyChange]); + useEffect(() => { onDirtyChange?.(formDirty); return () => onDirtyChange?.(false); }, [formDirty, onDirtyChange]); const modelOptions = useMemo(() => { const set = new Set(availableModels); @@ -152,12 +212,28 @@ export default function ProviderSettings({ setSaving(true); setMsg(null); try { - const patch: ProviderUpdatePatch = { adapter: adapter.trim(), baseUrl: nextBaseUrl, defaultModel: defaultModel.trim(), authMode, note: note.trim(), allowPrivateNetwork }; - // Keep omitted legacy values omitted unless the user actually changes this toggle. - // Otherwise an unrelated settings save manufactures `liveModels: true` provenance. - if (liveModels !== (item.liveModels !== false)) patch.liveModels = liveModels; - if (supportsApiKeyTransport) patch.apiKeyTransport = apiKeyTransport; - else if (item.apiKeyTransport !== undefined) patch.apiKeyTransport = ""; + if (pacingEnabled && !pacingDraft.requestsPerMinute && !pacingDraft.minIntervalMs && !pacingDraft.models) { + setMsg({ ok: false, text: t("pws.pacingRuleRequired") }); return false; + } + const pacingOnly = pacingDirty && !dirty; + const patch: ProviderUpdatePatch = pacingOnly + ? { requestPacing: pacingDraft } + : { + adapter: adapter.trim(), + baseUrl: nextBaseUrl, + defaultModel: defaultModel.trim(), + authMode, + note: note.trim(), + allowPrivateNetwork, + ...(pacingDirty ? { requestPacing: pacingDraft } : {}), + }; + if (!pacingOnly) { + // Keep omitted legacy values omitted unless the user actually changes this toggle. + // Otherwise an unrelated settings save manufactures `liveModels: true` provenance. + if (liveModels !== (item.liveModels !== false)) patch.liveModels = liveModels; + if (supportsApiKeyTransport) patch.apiKeyTransport = apiKeyTransport; + else if (item.apiKeyTransport !== undefined) patch.apiKeyTransport = ""; + } const res = await onUpdateProvider(item.name, patch); setMsg(res.ok ? { ok: true, text: t("pws.settingsSaved") } : { ok: false, text: res.error || t("prov.saveFailed") }); return res.ok; @@ -201,6 +277,8 @@ export default function ProviderSettings({ setDefaultModel(item.defaultModel ?? ""); setAuthMode(initialAuth); setApiKeyTransport(item.apiKeyTransport ?? "x-api-key"); setNote(item.note ?? ""); setAllowPrivateNetwork(item.allowPrivateNetwork ?? false); setLiveModels(item.liveModels !== false); setMsg(null); + setPacingEnabled(item.requestPacing?.enabled === true); setPacingRpm(numberDraft(item.requestPacing?.requestsPerMinute)); + setPacingDelay(numberDraft(item.requestPacing?.minIntervalMs)); setPacingModels({ ...(item.requestPacing?.models ?? {}) }); setEndpointChoice(matchChoiceId(baseUrlChoices, item.baseUrl)); }; @@ -213,6 +291,15 @@ export default function ProviderSettings({ } }; + const addPacingModel = () => { + const modelId = pacingModelId.trim(); + const rpm = positiveRpm(pacingModelRpm); + const delay = positiveInteger(pacingModelDelay); + if (!modelId || (rpm === undefined && delay === undefined)) return; + setPacingModels(current => ({ ...current, [modelId]: { ...(rpm !== undefined ? { requestsPerMinute: rpm } : {}), ...(delay !== undefined ? { minIntervalMs: delay } : {}) } })); + setPacingModelId(""); setPacingModelRpm(""); setPacingModelDelay(""); + }; + return (
- {dirty && ( +
+
+

{t("pws.pacingTitle")}

{t("pws.pacingDesc")}

+ +
+
+ + +
+

{t("pws.pacingSlowerWins")}

+
+ {pacingStatus?.queued ?? 0} {t("pws.pacingQueued")} + {pacingStatus?.nextSlotInMs ?? 0} ms {t("pws.pacingNextSlot")} + {pacingStatus?.lastModelId ?? t("pws.pacingNone")} {t("pws.pacingLastModel")} +
+

{t("pws.pacingModelOverrides")}

+
+ + + + +
+ {Object.entries(pacingModels).length > 0 &&
{Object.entries(pacingModels).map(([model, rule]) =>
{model}{rule.requestsPerMinute !== undefined ? `${rule.requestsPerMinute} ${t("pws.pacingRpmUnit")}` : ""}{rule.requestsPerMinute !== undefined && rule.minIntervalMs !== undefined ? " · " : ""}{rule.minIntervalMs !== undefined ? `${rule.minIntervalMs} ms` : ""}
)}
} +
+ {formDirty && (
{t("pws.settingsUnsavedBar")}
diff --git a/gui/src/components/provider-workspace/types.ts b/gui/src/components/provider-workspace/types.ts index 97e6f1969d..b5757d06eb 100644 --- a/gui/src/components/provider-workspace/types.ts +++ b/gui/src/components/provider-workspace/types.ts @@ -98,6 +98,7 @@ export type ProviderUpdatePatch = { disabled?: boolean; allowPrivateNetwork?: boolean; liveModels?: boolean; + requestPacing?: WorkspaceItem["requestPacing"] | null; /** Dedicated field: the API PATCHes it alone for the canonical `openai` provider. */ codexAccountMode?: "direct" | "pool"; }; diff --git a/gui/src/i18n/de.ts b/gui/src/i18n/de.ts index 1894b3deb0..ee7a577330 100644 --- a/gui/src/i18n/de.ts +++ b/gui/src/i18n/de.ts @@ -1665,6 +1665,23 @@ export const de: Record = { "pws.removeDefaultConfirmBody": "Standardanbieter \"{name}\" entfernen? \"{defaultProvider}\" wird zum Standardanbieter. Dies kann nicht rückgängig gemacht werden.", "pws.removeConfirmTitle": "Anbieter entfernen", "pws.saveSettings": "Speichern", + "pws.pacingTitle": "Anfragetaktung", + "pws.pacingDesc": "Verteilt ausgehende Anfragestarts für diesen Anbieter gleichmäßig. Streaming-Antworten dürfen sich überlappen.", + "pws.pacingEnabled": "Aktiviert", + "pws.pacingRpm": "Anfragen pro Minute", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "Mindestintervall (ms)", + "pws.pacingSlowerWins": "Das langsamere Anbieterlimit gilt. Modellregeln können nur stärker verzögern.", + "pws.pacingQueued": "in Warteschlange", + "pws.pacingNextSlot": "bis zum nächsten Slot", + "pws.pacingLastModel": "letztes Modell", + "pws.pacingNone": "Keine", + "pws.pacingModelOverrides": "Modellregeln", + "pws.pacingModel": "Modell", + "pws.pacingAdd": "Regel hinzufügen", + "pws.pacingRemove": "Entfernen", + "pws.pacingRemoveModel": "Anfragetaktung für {model} entfernen", + "pws.pacingRuleRequired": "Legen Sie zuerst ein Anbieterlimit oder eine Modellregel fest.", "pws.saving": "Wird gespeichert…", "pws.settingsSaved": "Einstellungen gespeichert.", "pws.accountModeSaved": "Kontomodus gespeichert.", diff --git a/gui/src/i18n/en.ts b/gui/src/i18n/en.ts index 04b154d0be..2183305e84 100644 --- a/gui/src/i18n/en.ts +++ b/gui/src/i18n/en.ts @@ -1178,6 +1178,23 @@ export const en = { "pws.removeDefaultConfirmBody": "Remove default provider \"{name}\"? \"{defaultProvider}\" will become the default provider. This cannot be undone.", "pws.removeConfirmTitle": "Remove provider", "pws.saveSettings": "Save", + "pws.pacingTitle": "Request pacing", + "pws.pacingDesc": "Evenly delay outbound request starts for this provider. Streaming responses may overlap.", + "pws.pacingEnabled": "Enabled", + "pws.pacingRpm": "Requests per minute", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "Minimum interval (ms)", + "pws.pacingSlowerWins": "The slower provider limit wins. Model overrides can only add more delay.", + "pws.pacingQueued": "queued", + "pws.pacingNextSlot": "until next slot", + "pws.pacingLastModel": "last model", + "pws.pacingNone": "None", + "pws.pacingModelOverrides": "Model overrides", + "pws.pacingModel": "Model", + "pws.pacingAdd": "Add override", + "pws.pacingRemove": "Remove", + "pws.pacingRemoveModel": "Remove request pacing override for {model}", + "pws.pacingRuleRequired": "Enable request pacing only after setting a provider limit or a model override.", "pws.saving": "Saving…", "pws.settingsSaved": "Settings saved.", "pws.accountModeSaved": "Account mode saved.", diff --git a/gui/src/i18n/ja.ts b/gui/src/i18n/ja.ts index 690534344f..d9f3a17cce 100644 --- a/gui/src/i18n/ja.ts +++ b/gui/src/i18n/ja.ts @@ -1120,6 +1120,23 @@ export const ja: Record = { "pws.removeDefaultConfirmBody": "既定のプロバイダー \"{name}\" を削除しますか? \"{defaultProvider}\" が既定のプロバイダーになります。この操作は元に戻せません。", "pws.removeConfirmTitle": "プロバイダーを削除", "pws.saveSettings": "保存", + "pws.pacingTitle": "リクエスト間隔調整", + "pws.pacingDesc": "このプロバイダーへの送信開始を均等に遅延します。ストリーミング応答は重複できます。", + "pws.pacingEnabled": "有効", + "pws.pacingRpm": "1分あたりのリクエスト数", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "最小間隔 (ms)", + "pws.pacingSlowerWins": "より遅いプロバイダー制限が優先され、モデル設定は遅延を増やす場合のみ適用されます。", + "pws.pacingQueued": "待機中", + "pws.pacingNextSlot": "次のスロットまで", + "pws.pacingLastModel": "最後のモデル", + "pws.pacingNone": "なし", + "pws.pacingModelOverrides": "モデル別設定", + "pws.pacingModel": "モデル", + "pws.pacingAdd": "設定を追加", + "pws.pacingRemove": "削除", + "pws.pacingRemoveModel": "{model} のリクエスト間隔設定を削除", + "pws.pacingRuleRequired": "プロバイダー制限またはモデル別設定を追加してから有効にしてください。", "pws.saving": "保存中…", "pws.settingsSaved": "設定を保存しました。", "pws.accountModeSaved": "アカウントモードを保存しました。", diff --git a/gui/src/i18n/ko.ts b/gui/src/i18n/ko.ts index 440319d2ed..fd5a1d6e66 100644 --- a/gui/src/i18n/ko.ts +++ b/gui/src/i18n/ko.ts @@ -1692,6 +1692,23 @@ export const ko: Record = { "pws.removeDefaultConfirmBody": "기본 프로바이더 \"{name}\"을(를) 제거하시겠습니까? \"{defaultProvider}\"이(가) 기본 프로바이더가 됩니다. 이 작업은 되돌릴 수 없습니다.", "pws.removeConfirmTitle": "프로바이더 제거", "pws.saveSettings": "저장", + "pws.pacingTitle": "요청 속도 조절", + "pws.pacingDesc": "이 프로바이더로 나가는 요청 시작을 일정 간격으로 지연합니다. 스트리밍 응답은 서로 겹칠 수 있습니다.", + "pws.pacingEnabled": "사용", + "pws.pacingRpm": "분당 요청 수", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "최소 시작 간격(ms)", + "pws.pacingSlowerWins": "프로바이더의 더 느린 제한이 우선합니다. 모델별 설정은 지연을 더 늘릴 때만 적용됩니다.", + "pws.pacingQueued": "대기 중", + "pws.pacingNextSlot": "다음 슬롯까지", + "pws.pacingLastModel": "마지막 모델", + "pws.pacingNone": "없음", + "pws.pacingModelOverrides": "모델별 설정", + "pws.pacingModel": "모델", + "pws.pacingAdd": "모델 제한 추가", + "pws.pacingRemove": "제거", + "pws.pacingRemoveModel": "{model} 모델의 요청 속도 설정 제거", + "pws.pacingRuleRequired": "프로바이더 제한이나 모델별 제한을 하나 이상 설정한 뒤 요청 속도 조절을 켜세요.", "pws.saving": "저장 중…", "pws.settingsSaved": "설정이 저장되었습니다.", "pws.accountModeSaved": "계정 모드가 저장되었습니다.", diff --git a/gui/src/i18n/ru.ts b/gui/src/i18n/ru.ts index 9b36937968..7914b7512d 100644 --- a/gui/src/i18n/ru.ts +++ b/gui/src/i18n/ru.ts @@ -1162,6 +1162,23 @@ export const ru: Record = { "pws.removeDefaultConfirmBody": "Удалить провайдера по умолчанию \"{name}\"? \"{defaultProvider}\" станет провайдером по умолчанию. Это действие нельзя отменить.", "pws.removeConfirmTitle": "Удалить провайдера", "pws.saveSettings": "Сохранить", + "pws.pacingTitle": "Интервал запросов", + "pws.pacingDesc": "Равномерно задерживает начало исходящих запросов к провайдеру. Потоковые ответы могут пересекаться.", + "pws.pacingEnabled": "Включено", + "pws.pacingRpm": "Запросов в минуту", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "Минимальный интервал (мс)", + "pws.pacingSlowerWins": "Действует более медленный лимит провайдера. Правила моделей могут только увеличить задержку.", + "pws.pacingQueued": "в очереди", + "pws.pacingNextSlot": "до следующего слота", + "pws.pacingLastModel": "последняя модель", + "pws.pacingNone": "Нет", + "pws.pacingModelOverrides": "Правила моделей", + "pws.pacingModel": "Модель", + "pws.pacingAdd": "Добавить правило", + "pws.pacingRemove": "Удалить", + "pws.pacingRemoveModel": "Удалить интервал запросов для {model}", + "pws.pacingRuleRequired": "Сначала задайте лимит провайдера или правило модели.", "pws.saving": "Сохранение…", "pws.settingsSaved": "Настройки сохранены.", "pws.accountModeSaved": "Режим аккаунта сохранён.", diff --git a/gui/src/i18n/tr.ts b/gui/src/i18n/tr.ts index 183b536af7..b1724cb3d7 100644 --- a/gui/src/i18n/tr.ts +++ b/gui/src/i18n/tr.ts @@ -1169,6 +1169,23 @@ export const tr: Record = { "pws.removeDefaultConfirmBody": "Varsayılan sağlayıcı \"{name}\" kaldırılsın mı? \"{defaultProvider}\" varsayılan sağlayıcı olacaktır. Bu işlem geri alınamaz.", "pws.removeConfirmTitle": "Sağlayıcıyı kaldır", "pws.saveSettings": "Kaydet", + "pws.pacingTitle": "İstek aralığı", + "pws.pacingDesc": "Bu sağlayıcıya giden istek başlangıçlarını eşit aralıklarla geciktirir. Akış yanıtları çakışabilir.", + "pws.pacingEnabled": "Etkin", + "pws.pacingRpm": "Dakikadaki istek", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "En kısa aralık (ms)", + "pws.pacingSlowerWins": "Daha yavaş sağlayıcı sınırı geçerlidir. Model kuralları yalnızca ek gecikme getirir.", + "pws.pacingQueued": "kuyrukta", + "pws.pacingNextSlot": "sonraki aralığa", + "pws.pacingLastModel": "son model", + "pws.pacingNone": "Yok", + "pws.pacingModelOverrides": "Model kuralları", + "pws.pacingModel": "Model", + "pws.pacingAdd": "Kural ekle", + "pws.pacingRemove": "Kaldır", + "pws.pacingRemoveModel": "{model} için istek aralığı kuralını kaldır", + "pws.pacingRuleRequired": "Önce bir sağlayıcı sınırı veya model kuralı belirleyin.", "pws.saving": "Kaydediliyor…", "pws.settingsSaved": "Ayarlar kaydedildi.", "pws.accountModeSaved": "Hesap modu kaydedildi.", diff --git a/gui/src/i18n/zh-TW.ts b/gui/src/i18n/zh-TW.ts index 340c9e18ec..127da4c080 100644 --- a/gui/src/i18n/zh-TW.ts +++ b/gui/src/i18n/zh-TW.ts @@ -959,6 +959,23 @@ export const zhTW: Record = { "pws.removeDefaultConfirmBody": "移除預設供應商「{name}」?「{defaultProvider}」將成為預設供應商。此操作無法撤消。", "pws.removeConfirmTitle": "移除供應商", "pws.saveSettings": "儲存", + "pws.pacingTitle": "請求節流", + "pws.pacingDesc": "均勻延遲送往此供應商的請求啟動。串流回應可以重疊。", + "pws.pacingEnabled": "啟用", + "pws.pacingRpm": "每分鐘請求數", + "pws.pacingRpmUnit": "次/分鐘", + "pws.pacingDelay": "最小間隔(毫秒)", + "pws.pacingSlowerWins": "以較慢的供應商限制為準,模型規則只能增加延遲。", + "pws.pacingQueued": "排隊中", + "pws.pacingNextSlot": "距下個時段", + "pws.pacingLastModel": "上個模型", + "pws.pacingNone": "無", + "pws.pacingModelOverrides": "模型規則", + "pws.pacingModel": "模型", + "pws.pacingAdd": "新增規則", + "pws.pacingRemove": "移除", + "pws.pacingRemoveModel": "移除 {model} 的請求節流規則", + "pws.pacingRuleRequired": "請先設定供應商限制或模型規則,再啟用請求節流。", "pws.saving": "儲存中…", "pws.settingsSaved": "設定已儲存。", "pws.settingsUnsavedBar": "有未儲存的更改。", diff --git a/gui/src/i18n/zh.ts b/gui/src/i18n/zh.ts index 0385282d91..05915b1325 100644 --- a/gui/src/i18n/zh.ts +++ b/gui/src/i18n/zh.ts @@ -1685,6 +1685,23 @@ export const zh: Record = { "pws.removeDefaultConfirmBody": "移除默认提供方「{name}」?「{defaultProvider}」将成为默认提供方。此操作无法撤销。", "pws.removeConfirmTitle": "移除提供商", "pws.saveSettings": "保存", + "pws.pacingTitle": "请求节流", + "pws.pacingDesc": "均匀延迟发往此提供商的请求启动。流式响应可以重叠。", + "pws.pacingEnabled": "启用", + "pws.pacingRpm": "每分钟请求数", + "pws.pacingRpmUnit": "RPM", + "pws.pacingDelay": "最小间隔(毫秒)", + "pws.pacingSlowerWins": "以较慢的提供商限制为准,模型规则只能增加延迟。", + "pws.pacingQueued": "排队中", + "pws.pacingNextSlot": "距下个时隙", + "pws.pacingLastModel": "上个模型", + "pws.pacingNone": "无", + "pws.pacingModelOverrides": "模型规则", + "pws.pacingModel": "模型", + "pws.pacingAdd": "添加规则", + "pws.pacingRemove": "移除", + "pws.pacingRemoveModel": "移除 {model} 的请求节流规则", + "pws.pacingRuleRequired": "请先设置提供商限制或模型规则,再启用请求节流。", "pws.saving": "保存中…", "pws.settingsSaved": "设置已保存。", "pws.accountModeSaved": "账户模式已保存。", diff --git a/gui/src/provider-workspace/catalog.ts b/gui/src/provider-workspace/catalog.ts index 847d0ae326..90ed999548 100644 --- a/gui/src/provider-workspace/catalog.ts +++ b/gui/src/provider-workspace/catalog.ts @@ -45,6 +45,12 @@ export interface WorkspaceProvider { disabled?: boolean; note?: string; allowPrivateNetwork?: boolean; + requestPacing?: { + enabled?: boolean; + requestsPerMinute?: number; + minIntervalMs?: number; + models?: Record; + }; /** Codex account routing mode for the canonical `openai` forward provider. */ codexAccountMode?: "direct" | "pool"; } diff --git a/gui/src/styles/provider-workspace-settings.css b/gui/src/styles/provider-workspace-settings.css index 4fc0d52dce..3429106e7f 100644 --- a/gui/src/styles/provider-workspace-settings.css +++ b/gui/src/styles/provider-workspace-settings.css @@ -117,6 +117,33 @@ .pwi-settings-dirty { font-size: var(--text-caption); color: var(--amber); } +.pwi-pacing-card { + display: flex; flex-direction: column; gap: 12px; padding: 14px; + border: 1px solid var(--border); border-radius: var(--radius-sm); + background: color-mix(in srgb, var(--surface) 88%, var(--accent) 12%); +} +.pwi-pacing-head { display: flex; align-items: flex-start; justify-content: space-between; gap: 14px; } +.pwi-pacing-head h3, .pwi-pacing-card h4 { margin: 0; color: var(--text); font-size: var(--text-control); } +.pwi-pacing-head p { margin: 4px 0 0; color: var(--muted); font-size: var(--text-caption); line-height: 1.45; } +.pwi-pacing-toggle { display: flex; align-items: center; gap: 6px; white-space: nowrap; font-size: var(--text-label); } +.pwi-pacing-grid { display: grid; grid-template-columns: repeat(2, minmax(0, 1fr)); gap: 10px; align-items: end; } +.pwi-pacing-grid--model { grid-template-columns: minmax(160px, 2fr) repeat(2, minmax(110px, 1fr)) auto; } +.pwi-pacing-status { display: grid; grid-template-columns: repeat(3, minmax(0, 1fr)); gap: 8px; } +.pwi-pacing-status span { padding: 8px 10px; border: 1px solid var(--border-soft); border-radius: var(--radius-xs); color: var(--muted); font-size: var(--text-caption); } +.pwi-pacing-status strong { display: block; overflow: hidden; color: var(--text); font-size: var(--text-label); text-overflow: ellipsis; white-space: nowrap; } +.pwi-pacing-overrides { display: flex; flex-direction: column; gap: 6px; } +.pwi-pacing-row { display: grid; grid-template-columns: minmax(0, 1fr) auto auto; align-items: center; gap: 10px; padding: 7px 8px; border-radius: var(--radius-xs); background: var(--bg); } +.pwi-pacing-row code { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.pwi-pacing-row span { color: var(--muted); font-size: var(--text-caption); } + +@media (max-width: 760px) { + .pwi-pacing-head { flex-direction: column; } + .pwi-pacing-grid, .pwi-pacing-grid--model, .pwi-pacing-status { grid-template-columns: 1fr; } + .pwi-pacing-grid--model .btn { width: 100%; } + .pwi-pacing-row { grid-template-columns: minmax(0, 1fr) auto; } + .pwi-pacing-row span { grid-column: 1 / -1; grid-row: 2; } +} + /* ── ProviderJsonEditor ────────────────────────────────────── */ .pwi-json-panel { display: flex; flex-direction: column; gap: 8px; } diff --git a/gui/tests/provider-settings-live-models-provenance.test.tsx b/gui/tests/provider-settings-live-models-provenance.test.tsx index 097780177e..11d7cc6ab7 100644 --- a/gui/tests/provider-settings-live-models-provenance.test.tsx +++ b/gui/tests/provider-settings-live-models-provenance.test.tsx @@ -91,6 +91,7 @@ test("an unrelated settings save does not materialize an omitted liveModels valu expect(patches).toHaveLength(1); expect(Object.hasOwn(patches[0]!, "liveModels")).toBe(false); + expect(Object.hasOwn(patches[0]!, "requestPacing")).toBe(false); await act(async () => { root.unmount(); }); }); diff --git a/gui/tests/provider-settings-request-pacing.test.tsx b/gui/tests/provider-settings-request-pacing.test.tsx new file mode 100644 index 0000000000..1d12a1d07d --- /dev/null +++ b/gui/tests/provider-settings-request-pacing.test.tsx @@ -0,0 +1,79 @@ +import { afterEach, beforeEach, expect, test } from "bun:test"; +import { Window } from "happy-dom"; +import { act } from "react"; +import type { Root } from "react-dom/client"; +import ProviderSettings from "../src/components/provider-workspace/ProviderSettings"; +import type { ProviderUpdatePatch } from "../src/components/provider-workspace/types"; +import { LanguageProvider } from "../src/i18n/provider"; +import type { WorkspaceItem } from "../src/provider-workspace/catalog"; + +const globals = ["document", "window", "navigator", "localStorage", "IS_REACT_ACT_ENVIRONMENT"] as const; +let previousGlobals: Record<(typeof globals)[number], unknown>; +let testWindow: Window; + +beforeEach(() => { + previousGlobals = Object.fromEntries(globals.map(key => [key, Reflect.get(globalThis, key)])) as typeof previousGlobals; + testWindow = new Window({ url: "http://localhost/#providers/workspace" }); + Object.defineProperty(testWindow.navigator, "language", { configurable: true, value: "en-US" }); + Object.defineProperties(globalThis, { + document: { configurable: true, value: testWindow.document }, + window: { configurable: true, value: testWindow }, + navigator: { configurable: true, value: testWindow.navigator }, + localStorage: { configurable: true, value: testWindow.localStorage }, + }); + (globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true; +}); + +afterEach(() => { + testWindow.close(); + for (const key of globals) Object.defineProperty(globalThis, key, { configurable: true, value: previousGlobals[key] }); +}); + +async function setInput(input: HTMLInputElement, value: string): Promise { + await act(async () => { + Object.getOwnPropertyDescriptor(testWindow.HTMLInputElement.prototype, "value")!.set!.call(input, value); + input.dispatchEvent(new testWindow.Event("input", { bubbles: true })); + }); +} + +test("settings saves provider pacing and a slower model override", async () => { + const item: WorkspaceItem = { + name: "nvidia", + adapter: "openai-chat", + baseUrl: "https://integrate.api.nvidia.com/v1", + authMode: "key", + }; + const patches: ProviderUpdatePatch[] = []; + const container = document.createElement("div"); + document.body.append(container); + const { createRoot } = await import("react-dom/client"); + let root!: Root; + await act(async () => { + root = createRoot(container); + root.render( { patches.push(patch); return { ok: true }; }} + />); + }); + + await act(async () => { container.querySelector(".pwi-pacing-toggle input")!.click(); }); + const numbers = container.querySelectorAll('.pwi-pacing-card input[type="number"]'); + await setInput(numbers[0]!, "38"); + const modelInput = container.querySelector('.pwi-pacing-grid--model input[list]')!; + await setInput(modelInput, "deepseek-ai/deepseek-v4-flash-0731"); + await setInput(numbers[2]!, "10"); + await act(async () => { container.querySelector(".pwi-pacing-grid--model button")!.click(); }); + const save = container.querySelector(".pwi-settings-sticky-bar .btn-primary")!; + await act(async () => { save.click(); await Promise.resolve(); }); + + expect(patches).toEqual([{ + requestPacing: { + enabled: true, + requestsPerMinute: 38, + models: { "deepseek-ai/deepseek-v4-flash-0731": { requestsPerMinute: 10 } }, + }, + }]); + expect(container.textContent).toContain("deepseek-ai/deepseek-v4-flash-0731"); + await act(async () => { root.unmount(); }); +}); diff --git a/src/config.ts b/src/config.ts index ef3580e1a2..5a7ee67680 100644 --- a/src/config.ts +++ b/src/config.ts @@ -613,6 +613,33 @@ const retryOn429PolicySchema = z.object({ respectRetryAfter: z.boolean().optional(), }).strict(); +const requestPacingRuleSchema = z.object({ + // Keep the RPM-derived timer within the same one-hour bound as minIntervalMs. + requestsPerMinute: z.number().min(1 / 60).max(60_000).optional(), + minIntervalMs: z.number().int().min(1).max(3_600_000).optional(), +}).strict().refine(value => value.requestsPerMinute !== undefined || value.minIntervalMs !== undefined, { + message: "request pacing rules need requestsPerMinute or minIntervalMs", +}); + +const requestPacingSchema = z.object({ + enabled: z.boolean(), + requestsPerMinute: z.number().min(1 / 60).max(60_000).optional(), + minIntervalMs: z.number().int().min(1).max(3_600_000).optional(), + models: z.record(z.string().trim().min(1), requestPacingRuleSchema).optional(), +}).strict().refine(value => value.enabled === false + || value.requestsPerMinute !== undefined + || value.minIntervalMs !== undefined + || (value.models !== undefined && Object.keys(value.models).length > 0), { + message: "enabled request pacing needs a provider rule or model override", +}); + +export function requestPacingConfigError(value: unknown): string | null { + if (value === undefined) return null; + const parsed = requestPacingSchema.safeParse(value); + if (parsed.success) return null; + return "requestPacing must contain enabled and a valid requestsPerMinute/minIntervalMs provider rule or model overrides"; +} + /** * Zod schema for one provider entry: known fields are validated strictly while unknown * fields pass through (preserved for runtime extensions). @@ -620,6 +647,7 @@ const retryOn429PolicySchema = z.object({ const providerConfigSchema = z.object({ adapter: z.string().min(1), baseUrl: z.string().min(1), + requestPacing: requestPacingSchema.optional().catch(undefined), mcpMaxTools: z.number().int().positive().optional(), mcpMaxSchemaBytes: z.number().int().positive().optional(), mcpMaxResultBytes: z.number().int().positive().optional(), diff --git a/src/lib/upstream-reachability.ts b/src/lib/upstream-reachability.ts index 5147f1f695..8a9222d58f 100644 --- a/src/lib/upstream-reachability.ts +++ b/src/lib/upstream-reachability.ts @@ -23,6 +23,7 @@ * MUST stay a leaf module: imports nothing from server.ts or adapters. */ +import { RequestPacingQueueOverloadError } from "../providers/request-pacing"; import { UpstreamRetryEvidenceError } from "./upstream-retry"; export const PRE_CONNECT_REACHABILITY_CODES = new Set([ @@ -68,6 +69,9 @@ export type TransportFailureKind = "timeout" | "connect_neutral" | "connect_erro * account-attributed behavior. */ export function classifyTransportFailureKind(err: unknown): TransportFailureKind { + // A local pacing admission failure happened before transport classification and must + // remain a client-visible overload, not account or host health evidence. + if (err instanceof RequestPacingQueueOverloadError) throw err; const evidence = err instanceof UpstreamRetryEvidenceError ? err : undefined; const rejection = evidence ? evidence.cause : err; if (rejection instanceof Error && rejection.name === "TimeoutError") return "timeout"; diff --git a/src/providers/request-pacing.ts b/src/providers/request-pacing.ts new file mode 100644 index 0000000000..d76f12bdde --- /dev/null +++ b/src/providers/request-pacing.ts @@ -0,0 +1,270 @@ +import type { OcxProviderConfig, RequestPacingRule } from "../types"; + +export const REQUEST_PACING_MAX_QUEUE_DEPTH = 256; +export const REQUEST_PACING_MAX_QUEUE_AGE_MS = 60_000; + +let maxQueueDepth = REQUEST_PACING_MAX_QUEUE_DEPTH; +let maxQueueAgeMs = REQUEST_PACING_MAX_QUEUE_AGE_MS; + +export type RequestPacingQueueOverloadReason = "queue_full" | "queue_expired"; + +export class RequestPacingQueueOverloadError extends Error { + constructor( + public readonly providerName: string, + public readonly reason: RequestPacingQueueOverloadReason, + public readonly retryAfterSeconds: number, + ) { + super(reason === "queue_full" + ? `request pacing queue for provider '${providerName}' is full` + : `request pacing queue for provider '${providerName}' exceeded the maximum queued age`); + this.name = "RequestPacingQueueOverloadError"; + } +} + +interface Waiter { + modelId?: string; + providerIntervalMs: number; + modelIntervalMs: number; + queuedAt: number; + signal?: AbortSignal; + resolve: () => void; + reject: (reason: unknown) => void; + abort?: () => void; +} + +interface ProviderPacer { + queue: Waiter[]; + providerNextStartAt: number; + modelNextStartAt: Map; + timer?: ReturnType; + active: boolean; + lastStartedAt?: number; + lastModelId?: string; +} + +export interface ProviderRequestPacingStatus { + provider: string; + enabled: boolean; + queued: number; + nextSlotInMs: number; + lastStartedAt?: number; + lastModelId?: string; +} + +const pacers = new Map(); + +function abortReason(signal: AbortSignal): unknown { + return signal.reason ?? new DOMException("The operation was aborted", "AbortError"); +} + +function normalizedInterval(rule: RequestPacingRule | undefined): number { + if (!rule) return 0; + const rpmInterval = typeof rule.requestsPerMinute === "number" && rule.requestsPerMinute > 0 + ? 60_000 / rule.requestsPerMinute + : 0; + const fixedInterval = typeof rule.minIntervalMs === "number" && rule.minIntervalMs > 0 + ? rule.minIntervalMs + : 0; + return Math.max(rpmInterval, fixedInterval); +} + +export function requestPacingIntervalMs(provider: OcxProviderConfig, modelId?: string): number { + const policy = provider.requestPacing; + if (!policy?.enabled) return 0; + const override = modelId ? policy.models?.[modelId] : undefined; + return Math.max(normalizedInterval(policy), normalizedInterval(override)); +} + +function requestPacingIntervals(provider: OcxProviderConfig, modelId?: string): { + providerIntervalMs: number; + modelIntervalMs: number; +} { + const policy = provider.requestPacing; + if (!policy?.enabled) return { providerIntervalMs: 0, modelIntervalMs: 0 }; + return { + providerIntervalMs: normalizedInterval(policy), + modelIntervalMs: modelId ? normalizedInterval(policy.models?.[modelId]) : 0, + }; +} + +function waiterReadyAt(state: ProviderPacer, modelId: string | undefined): number { + return Math.max( + state.providerNextStartAt, + modelId ? (state.modelNextStartAt.get(modelId) ?? 0) : 0, + ); +} + +function pacingRetryAfterSeconds(state: ProviderPacer, modelId: string | undefined, now: number): number { + return Math.max(1, Math.ceil(Math.max(0, waiterReadyAt(state, modelId) - now) / 1000)); +} + +function rejectExpiredWaiters(providerName: string, state: ProviderPacer, now: number): void { + for (let index = state.queue.length - 1; index >= 0; index -= 1) { + const waiter = state.queue[index]!; + if (now - waiter.queuedAt < maxQueueAgeMs) continue; + state.queue.splice(index, 1); + if (waiter.abort) waiter.signal?.removeEventListener("abort", waiter.abort); + waiter.reject(new RequestPacingQueueOverloadError( + providerName, + "queue_expired", + pacingRetryAfterSeconds(state, waiter.modelId, now), + )); + } +} + +function runQueue(providerName: string, state: ProviderPacer): void { + if (state.active) return; + if (state.queue.length === 0) { + if (state.timer) clearTimeout(state.timer); + state.timer = undefined; + return; + } + if (state.timer) return; + const now = Date.now(); + for (const [modelId, readyAt] of state.modelNextStartAt) { + if (readyAt <= now) state.modelNextStartAt.delete(modelId); + } + rejectExpiredWaiters(providerName, state, now); + if (state.queue.length === 0) return; + + const providerReadyAt = Math.max(now, state.providerNextStartAt); + const waiterIndex = state.queue.findIndex(waiter => { + const modelReadyAt = waiter.modelId ? (state.modelNextStartAt.get(waiter.modelId) ?? 0) : 0; + return Math.max(providerReadyAt, modelReadyAt) <= now; + }); + if (waiterIndex < 0) { + let earliestAt = Number.POSITIVE_INFINITY; + for (const waiter of state.queue) { + const modelReadyAt = waiter.modelId ? (state.modelNextStartAt.get(waiter.modelId) ?? 0) : 0; + const readyAt = Math.max(providerReadyAt, modelReadyAt); + const expiresAt = waiter.queuedAt + maxQueueAgeMs; + earliestAt = Math.min(earliestAt, readyAt, expiresAt); + } + const delayMs = Math.max(0, earliestAt - now); + state.timer = setTimeout(() => { + state.timer = undefined; + runQueue(providerName, state); + }, delayMs); + return; + } + const waiter = state.queue[waiterIndex]!; + const delayMs = Math.max(0, state.providerNextStartAt - now); + if (delayMs > 0) { + state.timer = setTimeout(() => { + state.timer = undefined; + runQueue(providerName, state); + }, delayMs); + return; + } + + state.active = true; + state.queue.splice(waiterIndex, 1); + if (waiter.abort) waiter.signal?.removeEventListener("abort", waiter.abort); + const startedAt = Date.now(); + state.lastStartedAt = startedAt; + state.lastModelId = waiter.modelId; + state.providerNextStartAt = startedAt + waiter.providerIntervalMs; + if (waiter.modelId && waiter.modelIntervalMs > 0) { + state.modelNextStartAt.set(waiter.modelId, startedAt + waiter.modelIntervalMs); + } + state.active = false; + waiter.resolve(); + queueMicrotask(() => runQueue(providerName, state)); +} + +export async function waitForProviderRequestSlot( + providerName: string, + provider: OcxProviderConfig, + modelId?: string, + signal?: AbortSignal, +): Promise { + const intervals = requestPacingIntervals(provider, modelId); + if (Math.max(intervals.providerIntervalMs, intervals.modelIntervalMs) <= 0) return; + if (signal?.aborted) throw abortReason(signal); + + const state = pacers.get(providerName) ?? { + queue: [], providerNextStartAt: 0, modelNextStartAt: new Map(), active: false, + }; + pacers.set(providerName, state); + + // Give already-eligible or expired waiters a chance to leave before applying the + // admission bound to the newest request. This preserves FIFO-ish fairness while + // keeping the retained queue strictly bounded under burst load. + if (state.timer) { + clearTimeout(state.timer); + state.timer = undefined; + } + runQueue(providerName, state); + if (state.queue.length >= maxQueueDepth) { + throw new RequestPacingQueueOverloadError( + providerName, + "queue_full", + pacingRetryAfterSeconds(state, modelId, Date.now()), + ); + } + + await new Promise((resolve, reject) => { + const waiter: Waiter = { modelId, ...intervals, queuedAt: Date.now(), signal, resolve, reject }; + waiter.abort = () => { + const index = state.queue.indexOf(waiter); + if (index >= 0) state.queue.splice(index, 1); + if (state.timer) { + clearTimeout(state.timer); + state.timer = undefined; + } + reject(abortReason(signal!)); + runQueue(providerName, state); + }; + signal?.addEventListener("abort", waiter.abort, { once: true }); + state.queue.push(waiter); + // Abort may race between the eager check above and listener registration. + if (signal?.aborted) { + waiter.abort(); + return; + } + if (state.timer) { + clearTimeout(state.timer); + state.timer = undefined; + } + runQueue(providerName, state); + }); +} + +export function providerRequestPacingStatus( + providerName: string, + provider: OcxProviderConfig, + now = Date.now(), +): ProviderRequestPacingStatus { + const state = pacers.get(providerName); + let nextSlotAt = state?.providerNextStartAt ?? 0; + if (state && state.queue.length > 0) { + let earliestQueuedSlotAt = Number.POSITIVE_INFINITY; + for (const waiter of state.queue) { + earliestQueuedSlotAt = Math.min(earliestQueuedSlotAt, waiterReadyAt(state, waiter.modelId)); + } + if (Number.isFinite(earliestQueuedSlotAt)) nextSlotAt = earliestQueuedSlotAt; + } + return { + provider: providerName, + enabled: provider.requestPacing?.enabled === true, + queued: state?.queue.length ?? 0, + nextSlotInMs: Math.max(0, Math.ceil(nextSlotAt - now)), + ...(state?.lastStartedAt !== undefined ? { lastStartedAt: state.lastStartedAt } : {}), + ...(state?.lastModelId ? { lastModelId: state.lastModelId } : {}), + }; +} + +export function setProviderRequestPacingLimitsForTest(limits: { + maxQueueDepth?: number; + maxQueueAgeMs?: number; +}): void { + if (limits.maxQueueDepth !== undefined) maxQueueDepth = limits.maxQueueDepth; + if (limits.maxQueueAgeMs !== undefined) maxQueueAgeMs = limits.maxQueueAgeMs; +} + +export function resetProviderRequestPacingForTest(): void { + for (const state of pacers.values()) if (state.timer) clearTimeout(state.timer); + pacers.clear(); + maxQueueDepth = REQUEST_PACING_MAX_QUEUE_DEPTH; + maxQueueAgeMs = REQUEST_PACING_MAX_QUEUE_AGE_MS; +} diff --git a/src/server/auth-cors.ts b/src/server/auth-cors.ts index dd7742cc02..aeeafa0bf6 100644 --- a/src/server/auth-cors.ts +++ b/src/server/auth-cors.ts @@ -14,6 +14,7 @@ import { providerModelCostsConfigError, reasoningSummaryDeliveryRecordConfigError, retryOn429PolicyConfigError, + requestPacingConfigError, sanitizeModelCostsForDisplay, } from "../config"; import { providerDestinationConfigError } from "../lib/destination-policy"; @@ -453,6 +454,8 @@ export function providerManagementConfigError(name: unknown, provider: unknown): // modelCosts is a user-owned display overlay, not part of the canonical // forward seed; it is validated separately below (providerModelCostsConfigError). delete canonicalCandidate.modelCosts; + // requestPacing is a user-owned transport overlay, not part of the canonical seed. + delete canonicalCandidate.requestPacing; const canonical = seed && sameCanonicalProviderSeed(canonicalCandidate, seed); if (!canonical) { return `provider ${name} must equal the canonical built-in provider seed`; @@ -477,6 +480,10 @@ export function providerManagementConfigError(name: unknown, provider: unknown): // it before it reaches the management API response. return `provider ${JSON.stringify(redactSecretString(name))} ${retryOn429Error}`; } + const requestPacingError = requestPacingConfigError(raw.requestPacing); + if (requestPacingError) { + return `provider ${JSON.stringify(redactSecretString(name))} ${requestPacingError}`; + } const modelCostsError = providerModelCostsConfigError(raw.modelCosts); if (modelCostsError) { // The provider name is caller-controlled and can be token-shaped; redact and JSON-escape @@ -580,6 +587,7 @@ export function safeConfigDTO(config: OcxConfig): unknown { "keyOptional", "freeTier", "liveModels", + "requestPacing", "models", "contextWindow", "modelContextWindows", diff --git a/src/server/management/provider-routes.ts b/src/server/management/provider-routes.ts index 203cc6a584..c9da16cb8c 100644 --- a/src/server/management/provider-routes.ts +++ b/src/server/management/provider-routes.ts @@ -14,6 +14,7 @@ import { normalizeNonBlankStringArray, providerBaseUrlConfigError, providerHeadersConfigError, + requestPacingConfigError, readConfigAdmissionSnapshot, saveConfigPreservingClaudeCode, withConfigMutationLockSync, @@ -44,6 +45,7 @@ import { import { routedSlug, slugEquals } from "../../providers/slug-codec"; import { clearAccountQuotaCache, clearProviderQuotaCache, fetchProviderQuotaReports } from "../../providers/quota"; import { clearKeyCooldowns } from "../../providers/key-failover"; +import { providerRequestPacingStatus } from "../../providers/request-pacing"; import { CODEX_FORWARD_BASE_URL, isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; import { codexAccountNamespaceProviderCollisionError } from "../../codex/account-namespace-match"; import { clearThreadAccountMap } from "../../codex/routing"; @@ -177,6 +179,20 @@ function applyProviderPatchFields( next.liveModels = rawBody.liveModels; touched = true; } + if (Object.hasOwn(rawBody, "requestPacing")) { + const value = rawBody.requestPacing; + if (value === null) { + delete next.requestPacing; + } else { + if (!isPlainRecord(value)) return { error: "requestPacing must be a plain object or null" }; + const pacingError = requestPacingConfigError(value); + if (pacingError) return { error: pacingError }; + // `requestPacingConfigError` is the runtime narrowing boundary above; keep the + // assertion explicit because a generic plain record cannot express `enabled`. + next.requestPacing = structuredClone(value) as unknown as OcxProviderConfig["requestPacing"]; + } + touched = true; + } // The Models page edits the catalog hints in place; keep them on the existing // provider mutation path so validation, cache invalidation, and convergence stay unified (#1073). if (Object.hasOwn(rawBody, "contextWindow")) { @@ -304,6 +320,22 @@ export async function handleProviderRoutes(ctx: ManagementContext): Promise [ + providerName, + providerRequestPacingStatus(providerName, provider), + ]), + )); + } + if (url.pathname === "/api/providers" && req.method === "GET") { return jsonResponse(Object.entries(config.providers).map(([name, p]) => ({ name, adapter: p.adapter, baseUrl: publicProviderBaseUrl(p.baseUrl), defaultModel: p.defaultModel, @@ -312,6 +344,7 @@ export async function handleProviderRoutes(ctx: ManagementContext): Promise 0, allowPrivateNetwork: p.allowPrivateNetwork === true, liveModels: p.liveModels !== false, + requestPacing: p.requestPacing, models: p.models ?? [], contextWindow: p.contextWindow, modelContextWindows: p.modelContextWindows, @@ -437,6 +470,7 @@ export async function handleProviderRoutes(ctx: ManagementContext): Promise key === "requestPacing"); + if (applied.editorTouched && !pacingOnly) { const providerError = providerManagementConfigError(name, next); if (providerError) return jsonResponse({ error: providerError }, 400); const resolvedError = await providerDestinationResolvedError(name, next); @@ -575,7 +613,7 @@ export async function handleProviderRoutes(ctx: ManagementContext): Promise +): Promise { + try { + return await handleResponsesCompactImpl(...args); + } catch (error) { + const overload = requestPacingOverloadResponse(error); + if (overload) return overload; + throw error; + } +} diff --git a/src/server/responses/compact.ts b/src/server/responses/compact.ts index 49649d62e2..0cfee64e7e 100644 --- a/src/server/responses/compact.ts +++ b/src/server/responses/compact.ts @@ -490,7 +490,7 @@ export async function handleResponsesCompact( req.signal, connectMs, false, - providerFetch(sendProvider), + providerFetch(sendProvider, route.providerName, route.modelId), // Every credential-bearing forward send gets manual redirects, not only // pool sends: direct mode carries the caller's credential too (#914). sendProvider.authMode === "forward", diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 7fd2b87b03..bc106cfd24 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -114,6 +114,7 @@ import { ForwardAdmissionCredentialError, validateForwardAdmissionCredential } f import { createTranslatorBudget, isTranslatorBudgetExceededError, type TranslatorBudget } from "../../lib/translator-budget"; import { listOpenAiForwardSidecarCandidates, resolveFirstUsableOpenAiSidecar, type ResolvedOpenAiForwardSidecar } from "../../providers/openai-sidecar"; import { isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; +import { waitForProviderRequestSlot } from "../../providers/request-pacing"; import { slugsEquivalent } from "../../providers/slug-codec"; import { applyOpenAiVirtualModel, resolveOpenAiCompactModel } from "../../providers/openai-virtual-models"; import { isUsageDebugEnabled } from "../../usage/debug"; @@ -589,7 +590,7 @@ async function retryCodexPoolOnAlternateAccount( upstream.signal, connectMs, stream, - providerFetch(route.provider), + providerFetch(route.provider, route.providerName, route.modelId), // Credential-bearing forward send: never follow a redirect into a // dead-host rejection after the credential was seen (#914). route.provider.authMode === "forward", @@ -2257,7 +2258,7 @@ async function handleResponsesInner( method: request.method, headers: request.headers, body: request.body, - }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider), + }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, route.providerName, route.modelId), route.provider.authMode === "forward") // Every real attempt response — including an intermediate 5xx the // retry wrapper replaces — proves the host was reached (#914 review). @@ -2318,7 +2319,7 @@ async function handleResponsesInner( method: request.method, headers: request.headers, body: request.body, - }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider), + }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, route.providerName, route.modelId), route.provider.authMode === "forward") .then(res => { settleObservedHostResponse(); @@ -2917,7 +2918,7 @@ async function handleResponsesInner( : clampImageMaxRounds(config.images?.videoMaxRounds ?? 2), connectTimeoutMs: config.connectTimeoutMs ?? 200_000, stallTimeoutSec: config.stallTimeoutSec, - fetchImpl: providerFetch(route.provider), + fetchImpl: providerFetch(route.provider, route.providerName, route.modelId), onRequestBuilt: request => recordAdapterReasoning(logCtx, request), ...(vidPlan?.timeoutMs ? { videoTimeoutMs: vidPlan.timeoutMs } : {}), onUsage: usage => { @@ -3224,6 +3225,7 @@ async function handleResponsesInner( try { if (activeAdapter.fetchResponse) { noteAttemptSend(logCtx.activeAttempt, inputTokenEstimate); + await waitForProviderRequestSlot(route.providerName, route.provider, route.modelId, upstream.signal); upstreamResponse = await activeAdapter.fetchResponse(builtInitialRequest, { abortSignal: upstream.signal, timeoutMs: connectMs, @@ -3237,7 +3239,7 @@ async function handleResponsesInner( method: builtInitialRequest.method, headers: builtInitialRequest.headers, body: builtInitialRequest.body, - }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider)); + }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, route.providerName, route.modelId)); }, { abortSignal: upstream.signal, label: safeHostLabel(builtInitialRequest.url) }, ); @@ -3312,11 +3314,13 @@ async function handleResponsesInner( noteAttemptSend(logCtx.activeAttempt, retryEstimate, recovery); try { try { - return activeAdapter.fetchResponse - ? await activeAdapter.fetchResponse(retryRequest, { abortSignal: upstream.signal, timeoutMs: connectMs, stream: parsed.stream }) - : await fetchWithHeaderTimeout(retryRequest.url, { - method: retryRequest.method, headers: retryRequest.headers, body: retryRequest.body, - }, upstream.signal, connectMs, parsed.stream, providerFetch(route.provider)); + if (activeAdapter.fetchResponse) { + await waitForProviderRequestSlot(route.providerName, route.provider, route.modelId, upstream.signal); + return await activeAdapter.fetchResponse(retryRequest, { abortSignal: upstream.signal, timeoutMs: connectMs, stream: parsed.stream }); + } + return await fetchWithHeaderTimeout(retryRequest.url, { + method: retryRequest.method, headers: retryRequest.headers, body: retryRequest.body, + }, upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, route.providerName, route.modelId)); } finally { retryRequest.releaseBodyObservation?.(); } @@ -3621,6 +3625,7 @@ async function handleResponsesInner( try { if (activeAdapter.fetchResponse) { noteAttemptSend(logCtx.activeAttempt, continuationEstimate, replayKind); + await waitForProviderRequestSlot(route.providerName, route.provider, nextParsed.modelId, upstream.signal); return await activeAdapter.fetchResponse(builtContinuationRequest, { abortSignal: upstream.signal, timeoutMs: connectMs, @@ -3640,7 +3645,7 @@ async function handleResponsesInner( upstream.signal, connectMs, nextParsed.stream, - providerFetch(route.provider), + providerFetch(route.provider, route.providerName, nextParsed.modelId), ); }, { abortSignal: upstream.signal, label: safeHostLabel(builtContinuationRequest.url) }, diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index b4affdc0ab..b2d36bab94 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -97,6 +97,7 @@ import { } from "../relay"; import { hasResponsesItemIdRepair, relaySseWithResponsesItemIdRepair } from "../responses-item-id-repair"; import type { EffectiveSubagentRoster, SpawnAgentSurface } from "../../codex/catalog"; +import { waitForProviderRequestSlot } from "../../providers/request-pacing"; export function disableResponsesRequestTimeout(req: Request, server: Pick, "timeout"> | undefined): boolean { @@ -131,18 +132,33 @@ export function safeOriginLabel(url: string): string { -export function providerFetch(provider: OcxProviderConfig): typeof globalThis.fetch { +interface PaceAwareFetch { + waitForPacing?: (signal?: AbortSignal) => Promise; + unpacedFetch?: typeof globalThis.fetch; +} + +export function providerFetch(provider: OcxProviderConfig, providerName?: string, modelId?: string): typeof globalThis.fetch { const base = (provider as OcxProviderConfig & { fetch?: typeof globalThis.fetch }).fetch ?? globalThis.fetch; // ChatGPT Codex backend: streaming turns ride the responses_websockets // transport (measured ~3s faster TTFT than the SSE POST queue); everything // else keeps the provider's HTTP fetch. See ws-upstream.ts for the details. - const wrapped = (input: Parameters[0], init?: RequestInit) => { + const unpaced = async (input: Parameters[0], init?: RequestInit) => { if (typeof input === "string" && init && shouldUseCodexWsUpstream(input, init)) { return codexWsUpstreamFetch(input, init, base); } return base(input, init); }; - return wrapped as typeof globalThis.fetch; + const waitForPacing = (signal?: AbortSignal) => providerName + ? waitForProviderRequestSlot(providerName, provider, modelId, signal) + : Promise.resolve(); + const wrapped = async (input: Parameters[0], init?: RequestInit) => { + await waitForPacing(init?.signal ?? undefined); + return unpaced(input, init); + }; + const preconnect = (...args: Parameters): void => { + base.preconnect?.(...args); + }; + return Object.assign(wrapped, { preconnect, waitForPacing, unpacedFetch: Object.assign(unpaced, { preconnect }) }); } @@ -156,6 +172,9 @@ export async function fetchWithHeaderTimeout( executor: typeof globalThis.fetch = globalThis.fetch, manualRedirect = false, ): Promise { + const pacing = executor as typeof globalThis.fetch & PaceAwareFetch; + await pacing.waitForPacing?.(abortSignal); + const fetchExecutor = pacing.unpacedFetch ?? executor; const timeout = new AbortController(); const timer = setTimeout(() => { if (!timeout.signal.aborted) timeout.abort(new DOMException("Timeout elapsed", "TimeoutError")); @@ -167,7 +186,7 @@ export async function fetchWithHeaderTimeout( headers.set("accept-encoding", "identity"); } try { - return await executor(url, { + return await fetchExecutor(url, { ...init, headers, // Credential-bearing sends opt into manual redirects so a 3xx is relayed diff --git a/src/server/responses/pacing-overload.ts b/src/server/responses/pacing-overload.ts new file mode 100644 index 0000000000..6ed30b85a7 --- /dev/null +++ b/src/server/responses/pacing-overload.ts @@ -0,0 +1,13 @@ +import { formatErrorResponse } from "../../bridge"; +import { RequestPacingQueueOverloadError } from "../../providers/request-pacing"; + +/** Convert local request-pacing admission failures into a retryable HTTP response. */ +export function requestPacingOverloadResponse(error: unknown): Response | undefined { + if (!(error instanceof RequestPacingQueueOverloadError)) return undefined; + return formatErrorResponse( + 429, + "rate_limit_error", + error.message, + { retryAfter: String(error.retryAfterSeconds) }, + ); +} diff --git a/src/server/responses/policy-fallback.ts b/src/server/responses/policy-fallback.ts index 7e9cbd3c74..ff4b176d61 100644 --- a/src/server/responses/policy-fallback.ts +++ b/src/server/responses/policy-fallback.ts @@ -5,6 +5,7 @@ import { finishRequestAttempt, type RequestLogContext } from "../request-log"; import type { OcxConfig } from "../../types"; import type { RouteCandidateTrace, RouteDecisionTraceV1 } from "../../routing/trace"; import { handleResponses as handleResponsesCore } from "./core"; +import { requestPacingOverloadResponse } from "./pacing-overload"; type CoreHandler = typeof handleResponsesCore; type CoreOptions = Parameters[3]; @@ -122,7 +123,14 @@ export async function handleResponsesWithPolicyFallback( // Core owns the client-facing parse/decompression error. } - let response = await runCore(req, config, logCtx, options); + let response: Response; + try { + response = await runCore(req, config, logCtx, options); + } catch (error) { + const overload = requestPacingOverloadResponse(error); + if (overload) return overload; + throw error; + } const initialTrace = logCtx.routeDecision; const initialRequestedModel = logCtx.requestedModel; if (!rawBody || !isPolicyDecision(initialTrace)) return response; @@ -140,7 +148,13 @@ export async function handleResponsesWithPolicyFallback( finishFailedPolicyAttempt(logCtx, response.status); const retryRequest = requestWithCandidate(req, rawBody, next); try { - response = await runCore(retryRequest, config, logCtx, options); + try { + response = await runCore(retryRequest, config, logCtx, options); + } catch (error) { + const overload = requestPacingOverloadResponse(error); + if (overload) return overload; + throw error; + } } finally { logCtx.requestedModel = initialRequestedModel; logCtx.routeDecision = initialTrace; diff --git a/src/server/responses/upstream-error.ts b/src/server/responses/upstream-error.ts index 7de21398fb..c7f33fc9e5 100644 --- a/src/server/responses/upstream-error.ts +++ b/src/server/responses/upstream-error.ts @@ -5,7 +5,12 @@ * #553 looking for an adapter URL bug that does not exist. Name the likely cause and the * command that settles it. */ +import { RequestPacingQueueOverloadError } from "../../providers/request-pacing"; + export function describeUpstreamConnectFailure(err: unknown, connectMs: number): string { + // Local pacing admission is not a transport failure. Let the outer response boundary + // preserve its retryable 429 identity instead of laundering it into a 502. + if (err instanceof RequestPacingQueueOverloadError) throw err; if (err instanceof Error && err.name === "TimeoutError") { return `Provider connect timeout after ${connectMs}ms`; } diff --git a/src/types.ts b/src/types.ts index d24811f430..e79cec69f3 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1226,12 +1226,28 @@ export interface ProviderCostOverlay { cacheWrite: number; } +export interface RequestPacingRule { + /** Evenly spread request starts to this many requests per minute. */ + requestsPerMinute?: number; + /** Minimum delay between request starts. The slower configured value wins. */ + minIntervalMs?: number; +} + +export interface ProviderRequestPacingConfig extends RequestPacingRule { + /** False preserves legacy behavior with no client-side waiting. */ + enabled: boolean; + /** Exact upstream model-id overrides; other models inherit the provider rule. */ + models?: Record; +} + /** * One configured provider entry. `authMode` (default `"key"`) decides whether same-target 429 * retries are allowed; OAuth/forward credentials and local runtimes are never replayed. */ export interface OcxProviderConfig { adapter: string; + /** Optional outbound request-start pacing shared by this provider and its model overrides. */ + requestPacing?: ProviderRequestPacingConfig; /** Cursor MCP compatibility bounds; positive integers when configured. */ mcpMaxTools?: number; mcpMaxSchemaBytes?: number; diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index 1b09e8fdde..03aee9a3be 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -335,6 +335,15 @@ once the server observes the client disconnect (Bun propagates it asynchronously cancelled with 499 before any replay; because the propagation is async, a replay may precede the cancel if the interval elapses first (bounded by the same `attempts` budget). +Provider-level `requestPacing` is the proactive companion to `retryOn429`. It reserves outbound +request-start slots before transport work begins, so a known RPM ceiling does not have to fail once +before the proxy reacts. One provider-wide lane enforces the aggregate ceiling. Exact model lanes +may add a slower interval without lowering the provider-wide interval or blocking an otherwise +eligible sibling model. Queue wait is abort-aware and happens before the response-header timeout is +armed. The shared fetch boundary covers HTTP and Responses WebSocket sends; explicit adapter +`fetchResponse` calls reserve the same lane at their call sites. Custom `runTurn` transports stay +outside this policy because they own their complete request lifecycle. + [Decision Log] - 목적과 의도: Prevent Kiro progress from becoming a false final answer, reject invalid empty completion retries, and stop concurrent transient 429s from consuming independent retry budgets. - 기존 구현 및 제약 조건: Kiro text has no trustworthy phase; stop metadata arrives only at stream end; the private completion tool is adapter-owned; normal parallel tool traffic must remain parallel; client cancellation must interrupt all waits. diff --git a/tests/management-provider-validation.test.ts b/tests/management-provider-validation.test.ts index 7e1ee8c553..c5d32be225 100644 --- a/tests/management-provider-validation.test.ts +++ b/tests/management-provider-validation.test.ts @@ -439,6 +439,69 @@ describe("provider management validation", () => { expect(secretNameError).toContain("[REDACTED]"); }); + test("provider request pacing PATCH persists provider and model limits without catalog churn", async () => { + if (existsSync(TEST_DIR)) rmSync(TEST_DIR, { recursive: true }); + mkdirSync(TEST_DIR, { recursive: true }); + process.env.OPENCODEX_HOME = TEST_DIR; + const liveConfig: OcxConfig = { + port: 0, + hostname: "127.0.0.1", + defaultProvider: "nvidia", + providers: { + nvidia: { + adapter: "openai-chat", + baseUrl: "https://integrate.api.nvidia.com/v1", + apiKey: "sk-nvidia", + }, + }, + }; + saveConfig(liveConfig); + let catalogRefreshes = 0; + const request = async (path: string, init?: RequestInit) => { + const req = new Request(`http://127.0.0.1${path}`, init); + return handleManagementAPI(req, new URL(req.url), liveConfig, { + createManagementConvergeCodex: catalogConvergenceFactory(() => { catalogRefreshes += 1; }), + }); + }; + const policy = { + enabled: true, + requestsPerMinute: 38, + minIntervalMs: 1_600, + models: { "deepseek-ai/deepseek-v4-flash-0731": { requestsPerMinute: 10 } }, + }; + + const saved = await request("/api/providers?name=nvidia", { + method: "PATCH", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ requestPacing: policy }), + }); + expect(saved?.status).toBe(200); + expect(liveConfig.providers.nvidia?.requestPacing).toEqual(policy); + expect(loadConfig().providers.nvidia?.requestPacing).toEqual(policy); + expect(catalogRefreshes).toBe(0); + + const providers = await request("/api/providers"); + expect((await providers?.json()).find((row: { name: string }) => row.name === "nvidia").requestPacing).toEqual(policy); + const status = await request("/api/provider-request-pacing?name=nvidia"); + expect(await status?.json()).toMatchObject({ provider: "nvidia", enabled: true, queued: 0, nextSlotInMs: 0 }); + + const invalid = await request("/api/providers?name=nvidia", { + method: "PATCH", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ requestPacing: { enabled: true, requestsPerMinute: -1 } }), + }); + expect(invalid?.status).toBe(400); + expect(liveConfig.providers.nvidia?.requestPacing).toEqual(policy); + + const timerOverflow = await request("/api/providers?name=nvidia", { + method: "PATCH", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ requestPacing: { enabled: true, requestsPerMinute: 0.001 } }), + }); + expect(timerOverflow?.status).toBe(400); + expect(liveConfig.providers.nvidia?.requestPacing).toEqual(policy); + }); + test("provider discovery status is additive and omitted before an attempt", async () => { markProviderDiscoveryFailed("auth-broken", { reason: "http", httpStatus: 401 }); try { diff --git a/tests/request-pacing.test.ts b/tests/request-pacing.test.ts new file mode 100644 index 0000000000..7debbc372c --- /dev/null +++ b/tests/request-pacing.test.ts @@ -0,0 +1,188 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { + providerRequestPacingStatus, + RequestPacingQueueOverloadError, + requestPacingIntervalMs, + resetProviderRequestPacingForTest, + setProviderRequestPacingLimitsForTest, + waitForProviderRequestSlot, +} from "../src/providers/request-pacing"; +import { providerFetch } from "../src/server/responses/fetch-helpers"; +import { fetchWithHeaderTimeout } from "../src/server/responses/fetch-helpers"; +import { requestPacingOverloadResponse } from "../src/server/responses/pacing-overload"; +import type { OcxProviderConfig } from "../src/types"; + +afterEach(() => resetProviderRequestPacingForTest()); + +function provider(requestPacing: OcxProviderConfig["requestPacing"]): OcxProviderConfig { + return { adapter: "openai-chat", baseUrl: "https://example.test/v1", requestPacing }; +} + +describe("requestPacingIntervalMs", () => { + test("uses the slower of provider RPM, provider delay, and model override", () => { + const configured = provider({ + enabled: true, + requestsPerMinute: 120, + minIntervalMs: 700, + models: { + slow: { requestsPerMinute: 30 }, + attemptedFast: { requestsPerMinute: 600 }, + }, + }); + expect(requestPacingIntervalMs(configured, "ordinary")).toBe(700); + expect(requestPacingIntervalMs(configured, "slow")).toBe(2_000); + expect(requestPacingIntervalMs(configured, "attemptedFast")).toBe(700); + }); + + test("supports model-only pacing while unrelated models remain unpaced", () => { + const configured = provider({ enabled: true, models: { slow: { minIntervalMs: 900 } } }); + expect(requestPacingIntervalMs(configured, "slow")).toBe(900); + expect(requestPacingIntervalMs(configured, "other")).toBe(0); + }); +}); + +describe("provider request pacing queue", () => { + test("spaces concurrent starts in one provider FIFO and exposes queue state", async () => { + const starts: number[] = []; + const fetchImpl = Object.assign(async () => { + starts.push(Date.now()); + return new Response("ok"); + }, { preconnect() {} }) as typeof globalThis.fetch; + const configured = { + ...provider({ enabled: true, requestsPerMinute: 600 }), + fetch: fetchImpl, + } as OcxProviderConfig & { fetch: typeof globalThis.fetch }; + const send = providerFetch(configured, "demo", "model-a"); + const pending = [send("https://example.test/v1/chat/completions"), send("https://example.test/v1/chat/completions"), send("https://example.test/v1/chat/completions")]; + await Bun.sleep(10); + expect(providerRequestPacingStatus("demo", configured).queued).toBe(2); + await Promise.all(pending); + expect(starts).toHaveLength(3); + expect(starts[1] - starts[0]).toBeGreaterThanOrEqual(85); + expect(starts[2] - starts[1]).toBeGreaterThanOrEqual(85); + const status = providerRequestPacingStatus("demo", configured); + expect(status.queued).toBe(0); + expect(status.lastModelId).toBe("model-a"); + }); + + test("aborted queued requests leave immediately and never consume a start", async () => { + const configured = provider({ enabled: true, minIntervalMs: 1_000 }); + await waitForProviderRequestSlot("demo", configured, "first"); + const controller = new AbortController(); + const queued = waitForProviderRequestSlot("demo", configured, "cancelled", controller.signal); + expect(providerRequestPacingStatus("demo", configured).queued).toBe(1); + controller.abort(); + await expect(queued).rejects.toHaveProperty("name", "AbortError"); + expect(providerRequestPacingStatus("demo", configured).queued).toBe(0); + }); + + test("rejects newest admission when the provider queue is full", async () => { + setProviderRequestPacingLimitsForTest({ maxQueueDepth: 2, maxQueueAgeMs: 5_000 }); + const configured = provider({ enabled: true, minIntervalMs: 1_000 }); + await waitForProviderRequestSlot("demo", configured, "first"); + const controller = new AbortController(); + const queued = [ + waitForProviderRequestSlot("demo", configured, "second", controller.signal), + waitForProviderRequestSlot("demo", configured, "third", controller.signal), + ]; + expect(providerRequestPacingStatus("demo", configured).queued).toBe(2); + await expect(waitForProviderRequestSlot("demo", configured, "newest")).rejects.toMatchObject({ + name: "RequestPacingQueueOverloadError", + reason: "queue_full", + providerName: "demo", + }); + expect(providerRequestPacingStatus("demo", configured).queued).toBe(2); + controller.abort(); + await Promise.allSettled(queued); + }); + + test("expires a queued request at the bounded queued-age deadline", async () => { + setProviderRequestPacingLimitsForTest({ maxQueueAgeMs: 25 }); + const configured = provider({ enabled: true, minIntervalMs: 1_000 }); + await waitForProviderRequestSlot("demo", configured, "first"); + const queued = waitForProviderRequestSlot("demo", configured, "stale"); + expect(providerRequestPacingStatus("demo", configured).queued).toBe(1); + await expect(queued).rejects.toMatchObject({ + name: "RequestPacingQueueOverloadError", + reason: "queue_expired", + providerName: "demo", + }); + expect(providerRequestPacingStatus("demo", configured).queued).toBe(0); + }); + + test("maps pacing admission overload to 429 with Retry-After", async () => { + const response = requestPacingOverloadResponse(new RequestPacingQueueOverloadError("demo", "queue_full", 3)); + expect(response?.status).toBe(429); + expect(response?.headers.get("Retry-After")).toBe("3"); + expect(await response?.json()).toMatchObject({ error: { type: "rate_limit_error" } }); + }); + + test("a slow model override does not slow other models beyond the provider interval", async () => { + const starts: Array<{ model: string; at: number }> = []; + const configured = provider({ + enabled: true, + minIntervalMs: 80, + models: { slow: { minIntervalMs: 400 } }, + }); + await waitForProviderRequestSlot("demo", configured, "slow"); + starts.push({ model: "slow", at: Date.now() }); + await waitForProviderRequestSlot("demo", configured, "fast"); + starts.push({ model: "fast", at: Date.now() }); + expect(starts[1].at - starts[0].at).toBeGreaterThanOrEqual(65); + expect(starts[1].at - starts[0].at).toBeLessThan(250); + }); + + test("the same model still observes its slower model override", async () => { + const configured = provider({ + enabled: true, + minIntervalMs: 50, + models: { slow: { minIntervalMs: 180 } }, + }); + const first = Date.now(); + await waitForProviderRequestSlot("demo", configured, "slow"); + await waitForProviderRequestSlot("demo", configured, "slow"); + expect(Date.now() - first).toBeGreaterThanOrEqual(160); + }); + + test("a model waiting on its override does not block another eligible model", async () => { + const configured = provider({ + enabled: true, + minIntervalMs: 60, + models: { slow: { minIntervalMs: 350 } }, + }); + await waitForProviderRequestSlot("demo", configured, "slow"); + const started = Date.now(); + const secondSlow = waitForProviderRequestSlot("demo", configured, "slow"); + const fast = waitForProviderRequestSlot("demo", configured, "fast"); + await fast; + expect(Date.now() - started).toBeGreaterThanOrEqual(45); + expect(Date.now() - started).toBeLessThan(220); + await secondSlow; + }); + + test("disabled policies preserve the unpaced legacy path", async () => { + const configured = provider({ enabled: false, requestsPerMinute: 1 }); + const started = Date.now(); + await Promise.all([ + waitForProviderRequestSlot("demo", configured, "a"), + waitForProviderRequestSlot("demo", configured, "b"), + ]); + expect(Date.now() - started).toBeLessThan(50); + expect(providerRequestPacingStatus("demo", configured).enabled).toBe(false); + }); + + test("queue waiting does not consume the response-header timeout budget", async () => { + const fetchImpl = Object.assign(async () => { + await Bun.sleep(20); + return new Response("ok"); + }, { preconnect() {} }) as typeof globalThis.fetch; + const configured = { + ...provider({ enabled: true, minIntervalMs: 120 }), + fetch: fetchImpl, + } as OcxProviderConfig & { fetch: typeof globalThis.fetch }; + const executor = providerFetch(configured, "demo", "model-a"); + await fetchWithHeaderTimeout("https://example.test/v1/chat/completions", {}, new AbortController().signal, 50, false, executor); + const second = await fetchWithHeaderTimeout("https://example.test/v1/chat/completions", {}, new AbortController().signal, 50, false, executor); + expect(second.status).toBe(200); + }); +}); diff --git a/tests/routing-policy-fallback.test.ts b/tests/routing-policy-fallback.test.ts index 8ba1193e1b..dc812c794b 100644 --- a/tests/routing-policy-fallback.test.ts +++ b/tests/routing-policy-fallback.test.ts @@ -1,5 +1,6 @@ import { describe, expect, test } from "bun:test"; +import { RequestPacingQueueOverloadError } from "../src/providers/request-pacing"; import type { OcxConfig } from "../src/types"; import { beginRequestAttempt, type RequestLogContext } from "../src/server/request-log"; import type { RouteDecisionTraceV1 } from "../src/routing/trace"; @@ -92,6 +93,21 @@ describe("policy candidate fallback", () => { expect(logCtx.activeAttempt).toBe(logCtx.attempts?.[1]); }); + test("returns local pacing overload without switching policy candidates", async () => { + const trace = policyTrace(); + const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext; + let calls = 0; + const response = await handleResponsesWithPolicyFallback(request(), {} as OcxConfig, logCtx, {}, { + runCore: async () => { + calls += 1; + throw new RequestPacingQueueOverloadError("provider-a", "queue_full", 2); + }, + }); + expect(response.status).toBe(429); + expect(response.headers.get("Retry-After")).toBe("2"); + expect(calls).toBe(1); + }); + test("does not switch candidates for terminal client/input failures", async () => { const trace = policyTrace(); const logCtx = { requestedModel: "policy/daily", routeDecision: trace, attempts: [] } as unknown as RequestLogContext;