Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -198,7 +198,7 @@ Providers can expose a built-in shorthand, such as `agy` for `google-antigravity
| `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. |
| `proxy?` | `string \| null` | Per-provider egress route. Omit it to inherit the global proxy decision; use `"direct"` or `null` to force direct egress; or provide an absolute `http://`, `https://`, `socks5://`, or `socks5h://` proxy URL. An empty string is rejected. |
| `noProxy?` | `string \| string[]` | Destinations this provider reaches directly, using `NO_PROXY` host-pattern syntax. A match bypasses both this provider's own proxy and an inherited global proxy. |
| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, models? }` | Optional client-side outbound request-start pacing, separate from upstream usage, billing, and rate-limit indicators. RPM is converted to an even interval; `minIntervalMs` may impose a longer interval. Provider limits apply across all models, while `models` entries use exact upstream model IDs (for example `nvidia/llama-3.1-nemotron-ultra-253b-v1`) and can only add delay. Queue waits do not consume the upstream response-header timeout. HTTP, Responses WebSocket, and explicit adapter `fetchResponse`/`runTurn` dispatches are covered. |
| `requestPacing?` | `{ enabled, requestsPerMinute?, minIntervalMs?, maxConcurrentRequests?, models? }` | Optional client-side outbound request pacing, separate from upstream usage, billing, and rate-limit indicators. RPM is converted to an even interval; `minIntervalMs` may impose a longer interval. `maxConcurrentRequests` caps how many requests to this provider may be in flight at once, counted until the upstream response body closes; a `models` entry may set a tighter per-model cap. Provider limits apply across all models, while `models` entries use exact upstream model IDs (for example `nvidia/llama-3.1-nemotron-ultra-253b-v1`) and can only add delay or tighten the concurrency cap. Queue waits do not consume the upstream response-header timeout. HTTP, explicit adapter `fetchResponse` dispatches, and `runTurn` attempts are covered; a configured concurrency cap moves Responses WebSocket turns to HTTP/SSE so the lease can be returned. |
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
| `upstreamHttpVersion?` | `"auto" \| "http1.1" \| "h1" \| "http2" \| "h2"` | Pin the HTTP version used for upstream requests to this provider. Defaults to `auto`, which lets Bun negotiate. An explicit pin requires an HTTPS target and fails locally when it cannot be honored. Set `http1.1` when a provider's HTTP/2 SSE stream stalls instead of delivering events — the symptom is a long-running streaming request that produces nothing and eventually times out. For Cursor, `http1.1`/`h1` selects its `RunSSE` + `BidiAppend` compatibility transport for inference and also pins live model discovery. Management `POST`/`PATCH` accept `null` to clear it back to `auto`. |
| `responsesPath?` | `string` | Relative resource path for key-auth `openai-responses` requests. It must start with `/` and contain no scheme, query, or fragment. |
| `chatCompletionsPath?` | `string` | Relative resource path for `openai-chat` requests, the mirror of `responsesPath` and subject to the same shape rules. Needed when one upstream serves Chat Completions and Responses under different prefixes: a per-model wire override changes the adapter and leaves `baseUrl` alone, so without this an opted-in Chat request would be sent to the Responses base. Z.AI is the shipped example. |
Expand Down
11 changes: 8 additions & 3 deletions src/adapters/physical-send.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,10 @@ import type { AdapterFetchContext } from "./base";
import type { SendClass } from "../lib/request-execution-budget";
import type { AttemptRecoveryKind } from "../usage/log";
import { abortError, SendBudgetExhaustedError } from "../lib/upstream-retry";
import { trackProviderRequestSlotBody, type ProviderRequestSlot } from "../providers/request-pacing";

type PacedFetch = typeof globalThis.fetch & {
waitForPacing?: (signal?: AbortSignal) => Promise<void>;
waitForPacing?: (signal?: AbortSignal) => Promise<ProviderRequestSlot | undefined>;
unpacedFetch?: typeof globalThis.fetch;
};

Expand Down Expand Up @@ -37,12 +38,16 @@ export function createAdapterPhysicalSend(ctx: AdapterFetchContext = {}, fallbac
ctx.onPhysicalSend?.({ ordinal, ...(options.recovery ? { recovery: options.recovery } : {}) });
return (executor.unpacedFetch ?? executor)(input, init);
}) as typeof globalThis.fetch;
let pacingSlot: ProviderRequestSlot | undefined;
try {
await executor.waitForPacing?.(ctx.abortSignal);
pacingSlot = await executor.waitForPacing?.(ctx.abortSignal);
if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal);
await options.beforeDispatch?.();
if (ctx.abortSignal?.aborted) throw abortError(ctx.abortSignal);
return await options.dispatch(physicalExecutor);
return trackProviderRequestSlotBody(pacingSlot, await options.dispatch(physicalExecutor));
} catch (error) {
pacingSlot?.release();
throw error;
} finally {
permit?.release();
}
Expand Down
11 changes: 8 additions & 3 deletions src/config/schema/leaf-validators.ts
Original file line number Diff line number Diff line change
Expand Up @@ -97,18 +97,23 @@ 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",
maxConcurrentRequests: z.number().int().min(1).max(1_000).optional(),
}).strict().refine(value => value.requestsPerMinute !== undefined
|| value.minIntervalMs !== undefined
|| value.maxConcurrentRequests !== undefined, {
message: "request pacing rules need requestsPerMinute, minIntervalMs, or maxConcurrentRequests",
});

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(),
maxConcurrentRequests: z.number().int().min(1).max(1_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.maxConcurrentRequests !== undefined
|| (value.models !== undefined && Object.keys(value.models).length > 0), {
message: "enabled request pacing needs a provider rule or model override",
});
Expand All @@ -117,7 +122,7 @@ 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";
return "requestPacing must contain enabled and a valid requestsPerMinute/minIntervalMs/maxConcurrentRequests provider rule or model overrides";
}

/**
Expand Down
Loading
Loading