Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
1 change: 1 addition & 0 deletions changelog.d/features/13700-per-model-concurrency.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **feat(sse):** opt-in per-model concurrency caps per connection — configure exact `modelConcurrency` ceilings inside `rateLimitOverrides` (dashboard editor or `PATCH /api/providers/[id]`), enforced as an atomic fourth gate alongside the global/provider/account semaphores with local queueing; unconfigured connections behave exactly as before. Also fixes `rateLimitOverrides.executionMaxWaitMs`: the schema accepted it but the DB allowlist rejected it, so saving it failed; a dashboard save now also preserves an API-set value ([#13700](https://github.com/diegosouzapw/OmniRoute/pull/13700))
46 changes: 46 additions & 0 deletions docs/architecture/RESILIENCE_GUIDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,52 @@ Each provider connection can declare a `max_concurrent` ceiling
Leave it empty for no limit. This is the single knob that drives the serialization
layer below — set it to the account's real concurrency (e.g. GLM ~1, MiniMax ~2).

### Per-model concurrency caps (`modelConcurrency`)

A connection can additionally declare exact per-model concurrency ceilings
inside its `rateLimitOverrides` map:

```json
{
"rateLimitOverrides": {
"maxConcurrent": 4,
"modelConcurrency": { "glm-5": 1, "glm-4.7": 3 }
}
}
```

Set it in the connection modal (**Rate limit overrides → Per-model
concurrency caps**, one `model=cap` per line) or via
`PATCH /api/providers/[id]` with the same JSON shape. Key semantics:

- **Connection-wide vs model-specific:** `maxConcurrent` remains the shared
connection-wide ceiling. When both apply, both gates are acquired
atomically in the same composite gate
(`global → provider → account → model`); the effective behavior is the
stricter applicable limit.
- **Exact model-key match:** the key is the model string passed to the
executor after routing resolution — normally the bare upstream model id
(`glm-5`), not a client-side `provider/model` alias (`zai/glm-5` does not
match `glm-5`). Values are positive-integer concurrent-request ceilings.
- **Local queueing, no discovery:** excess requests queue locally with the
existing queue/timeout semantics (typed `SEMAPHORE_TIMEOUT` /
`SEMAPHORE_QUEUE_FULL` admission errors). OmniRoute does not discover or
infer upstream policy — it enforces the exact ceilings the operator
configured. A saturated model gate never disables the provider and never
creates a permanent model lockout; upstream 429/cooldown/fallback behavior
remains the error backstop.
- **Per-connection, per-process scope:** caps are per database connection
and held in-memory, so two connections reusing the same upstream API key
do not coordinate with each other.
- **Unconfigured means unchanged:** omitting the map (or leaving the
dashboard field blank) adds no model gate. Example configuration without
asserting any universal provider limit:

```text
glm-5=1
glm-4.7=3
```

### Quota-share request serialization

When a quota-share dispatch targets a connection that declares a positive
Expand Down
6 changes: 6 additions & 0 deletions open-sse/executors/base/validationDispatch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,12 @@ export type ProviderCredentials = {
expiresAt?: string;
connectionId?: string; // T07: used for API key rotation index
maxConcurrent?: number | null;
/**
* Optional per-model concurrency ceilings for this connection (see
* ProviderCredentials in open-sse/types.d.ts). Normalized at credential
* selection; the chat core resolves the exact-model cap fail-open.
*/
modelConcurrency?: Record<string, number> | null;
rateLimitMaxConcurrent?: number | null;
providerSpecificData?: Record<string, unknown>;
requestEndpointPath?: string;
Expand Down
14 changes: 14 additions & 0 deletions open-sse/handlers/chatCore/executeProviderRequest.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import { getExecutionConnectionId } from "./executionCredentials.ts";
import {
resolveAccountSemaphoreKey,
resolveAccountSemaphoreMaxConcurrency,
resolveModelSemaphore,
} from "./executorHelpers.ts";
import {
materializeDeduplicatedExecutionResult,
Expand Down Expand Up @@ -255,6 +256,13 @@ export async function executeProviderRequest(
connectionId: attemptConnectionId,
credentials: execCreds,
});
// Opt-in per-model ceiling; joins the composite gate below.
const modelGate = resolveModelSemaphore({
provider,
model: modelToCall,
connectionId: attemptConnectionId,
credentials: execCreds,
});
const canonicalProviderKey = resolveProviderId(String(provider).trim().toLowerCase());
const providerConcurrency =
resilienceSettings.providerQuotaOverrides[canonicalProviderKey]?.providerConcurrency ??
Expand All @@ -263,6 +271,8 @@ export async function executeProviderRequest(
trace("pre_semaphore", {
semaphoreKey: accountSemaphoreKey,
max: accountSemaphoreMaxConcurrency,
modelSemaphoreKey: modelGate.key,
modelMax: modelGate.maxConcurrency,
});
if (accountSemaphoreKey && accountSemaphoreMaxConcurrency != null) {
updatePendingScope(pendingScope, {
Expand All @@ -289,6 +299,10 @@ export async function executeProviderRequest(
key: accountSemaphoreKey || "",
maxConcurrency: accountSemaphoreKey ? accountSemaphoreMaxConcurrency : null,
},
{
key: modelGate.key || "",
maxConcurrency: modelGate.key ? modelGate.maxConcurrency : null,
},
],
{
timeoutMs: maxWaitMs,
Expand Down
62 changes: 61 additions & 1 deletion open-sse/handlers/chatCore/executorHelpers.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
import { FORMATS } from "../../translator/formats.ts";
import { buildAccountSemaphoreKey } from "../../services/accountSemaphore.ts";
import {
buildAccountSemaphoreKey,
buildModelSemaphoreKey,
} from "../../services/accountSemaphore.ts";
import { getHeaderValueCaseInsensitive } from "./headers.ts";

function toFiniteNumberOrNull(value: unknown): number | null {
Expand Down Expand Up @@ -49,6 +52,27 @@ export function resolveAccountSemaphoreMaxConcurrency(
return toFiniteNumberOrNull(credentials?.maxConcurrent);
}

/**
* Resolve the per-model concurrency cap for one execution attempt.
*
* Exact-match on the model string passed to the executor after routing
* resolution (normally the bare upstream model id, e.g. "glm-5" — not a
* client-side `provider/model` alias). Missing credentials, a missing or
* malformed map, or a non-positive cap all resolve to null ("no model
* gate") so an unconfigured or corrupt map never blocks a request.
*/
export function resolveModelSemaphoreMaxConcurrency(
credentials: Record<string, unknown> | null | undefined,
model: string | null | undefined
): number | null {
if (!model || typeof model !== "string" || model.trim().length === 0) return null;
const map = credentials?.modelConcurrency;
if (!map || typeof map !== "object" || Array.isArray(map)) return null;
const cap = toFiniteNumberOrNull((map as Record<string, unknown>)[model]);
if (cap == null || !Number.isInteger(cap) || cap < 1) return null;
return cap;
}

export function resolveAccountSemaphoreKey({
provider,
model,
Expand All @@ -65,6 +89,42 @@ export function resolveAccountSemaphoreKey({
return buildAccountSemaphoreKey({ provider, accountKey });
}

/**
* Build the per-connection, per-model semaphore key for one execution
* attempt. Returns null when no positive cap is configured for
* `provider + connection + model`, in which case the caller must not add a
* model requirement to the composite gate (behavior unchanged).
*/
export function resolveModelSemaphoreKey({
provider,
model,
connectionId,
credentials,
}: {
provider: string | null | undefined;
model: string;
connectionId: string | null | undefined;
credentials: Record<string, unknown> | null | undefined;
}): string | null {
const accountKey = resolveAccountSemaphoreAccountKey(connectionId, credentials);
if (!accountKey || !provider) return null;
if (resolveModelSemaphoreMaxConcurrency(credentials, model) == null) return null;
return buildModelSemaphoreKey({ provider, accountKey, model });
}

/** Per-model gate (key + cap) for the composite semaphore; key null = no gate. */
export function resolveModelSemaphore(args: {
provider: string | null | undefined;
model: string;
connectionId: string | null | undefined;
credentials: Record<string, unknown> | null | undefined;
}): { key: string | null; maxConcurrency: number | null } {
return {
key: resolveModelSemaphoreKey(args),
maxConcurrency: resolveModelSemaphoreMaxConcurrency(args.credentials, args.model),
};
}

export function buildClaudePromptCacheLogMeta(
targetFormat: string,
finalBody: Record<string, unknown> | null | undefined,
Expand Down
17 changes: 17 additions & 0 deletions open-sse/services/accountSemaphore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,23 @@ export function buildAccountSemaphoreKey({
return `${String(provider)}:${String(accountKey)}`;
}

export interface ModelSemaphoreKeyParts extends AccountSemaphoreKeyParts {
model: string;
}

/**
* Collision-safe key for the per-connection, per-model concurrency gate.
* The `model:` infix keeps model gates disjoint from account gates even
* when a model id itself contains `:` characters.
*/
export function buildModelSemaphoreKey({
provider,
accountKey,
model,
}: ModelSemaphoreKeyParts): string {
return `${String(provider)}:${String(accountKey)}:model:${String(model)}`;
}

function isBypassed(maxConcurrency?: number | null): boolean {
return maxConcurrency == null || !Number.isFinite(maxConcurrency) || maxConcurrency <= 0;
}
Expand Down
7 changes: 4 additions & 3 deletions open-sse/services/rateLimitManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ import {
import { LimiterWedgeWatchdog, WATCHDOG_INTERVAL_MS } from "./rateLimitManager/wedgeWatchdog";
import { createCancellableJob } from "./rateLimitManager/queuedJobCancel";
import { toNumber } from "@/shared/utils/numeric";
import type { ConnectionRateLimitOverrides } from "@/lib/db/providers/columns";
import {
getExecutorTimeoutMs,
resolveConnectionTimeoutMs,
Expand Down Expand Up @@ -103,7 +104,7 @@ const enabledConnections = new Set<string>();

// Store per-connection rate limit overrides (RPM, TPM, TPD, minTime, maxConcurrent)
// Populated from provider_connections.rateLimitOverrides on startup and refresh.
const connectionRateLimitOverrides = new Map<string, Record<string, number>>();
const connectionRateLimitOverrides = new Map<string, ConnectionRateLimitOverrides>();

// Store learned limits for persistence (debounced)
// One learned entry per limiter key (provider:connection[:model]). The previous
Expand Down Expand Up @@ -266,7 +267,7 @@ export function resolveRequestQueueMaxWaitMs(
*/
export function resolveExecutionMaxWaitMs(connectionId?: string): number {
const override = connectionId
? (connectionRateLimitOverrides.get(connectionId) as Record<string, number> | undefined)
? (connectionRateLimitOverrides.get(connectionId) as ConnectionRateLimitOverrides | undefined)
?.executionMaxWaitMs
: undefined;
return resolveOverride(override, currentRequestQueueSettings.executionMaxWaitMs);
Expand Down Expand Up @@ -536,7 +537,7 @@ export function isRateLimitEnabled(connectionId) {
* connection so the next request gets a fresh limiter with the new settings.
*
* @param {string} connectionId
* @param {Record<string, number> | null} overrides - New overrides (null/undefined clears)
* @param {ConnectionRateLimitOverrides | null} overrides - New overrides (null/undefined clears)
*/
export function refreshConnectionRateLimits(connectionId, overrides) {
if (overrides === null || overrides === undefined) {
Expand Down
8 changes: 5 additions & 3 deletions open-sse/services/rateLimitManager/overrideUpdates.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,10 @@
* @module services/rateLimitManager/overrideUpdates
*/

import type { ConnectionRateLimitOverrides } from "@/lib/db/providers/columns";

export function buildOverrideUpdates(
overrides: Record<string, number>
overrides: ConnectionRateLimitOverrides
): Record<string, number> {
const updates: Record<string, number> = {};
if (typeof overrides.maxConcurrent === "number" && overrides.maxConcurrent > 0) {
Expand All @@ -29,13 +31,13 @@ export function buildOverrideUpdates(
}

export function loadOverrideMap(
target: Map<string, Record<string, number>>,
target: Map<string, ConnectionRateLimitOverrides>,
connections: Array<Record<string, unknown>>
): void {
target.clear();
for (const conn of connections) {
const overrides = conn.rateLimitOverrides;
if (overrides && typeof overrides === "object" && !Array.isArray(overrides))
target.set(String(conn.id), overrides as Record<string, number>);
target.set(String(conn.id), overrides as ConnectionRateLimitOverrides);
}
}
6 changes: 6 additions & 0 deletions open-sse/types.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,12 @@ export interface ProviderCredentials {
connectionId: string;
/** Optional per-account concurrency cap */
maxConcurrent?: number | null;
/**
* Optional per-model concurrency ceilings for this connection, keyed by
* the exact model string passed to the executor after routing resolution
* (normally the bare upstream model id). Absent/null means no model gate.
*/
modelConcurrency?: Record<string, number> | null;
/** User email associated with the connection */
email?: string;
/** API key (for apikey auth type) */
Expand Down
1 change: 1 addition & 0 deletions scripts/i18n/untranslatable-keys.json
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,7 @@
"providers.proxy",
"providers.proxyConfiguredBySource",
"providers.quotaWindowSpark",
"providers.rateLimitOverridesModelConcurrencyPlaceholder",
"providers.responses",
"providers.responsesApi",
"providers.responsesPath",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,12 @@ import { maskEmail } from "@/shared/utils/maskEmail";
import useEmailPrivacyStore from "@/store/emailPrivacyStore";
import { useNotificationStore } from "@/store/notificationStore";
import { type CodexServiceTier } from "@/lib/providers/requestDefaults";
import type { ConnectionRateLimitOverrides } from "@/lib/db/providers/columns";
import ModelConcurrencyField from "./ModelConcurrencyField";
import {
buildRateLimitOverridesFromForm,
modelConcurrencyFormValue,
} from "./rateLimitOverridesFromForm";
import { resolveDashboardProviderInfo } from "../../../providerPageUtils";
import {
isBaseUrlConfigurableProvider,
Expand Down Expand Up @@ -82,7 +88,7 @@ export interface EditConnectionModalConnection {
email?: string;
priority?: number;
maxConcurrent?: number | null;
rateLimitOverrides?: Record<string, number> | null;
rateLimitOverrides?: ConnectionRateLimitOverrides | null;
authType?: string;
provider?: string;
apiKey?: string;
Expand Down Expand Up @@ -128,6 +134,7 @@ export default function EditConnectionModal({
minTime: "",
maxWaitMs: "",
rateLimitMaxConcurrent: "",
modelConcurrency: "",
apiKey: "",
healthCheckInterval: "" as number | "",
baseUrl: "",
Expand Down Expand Up @@ -341,6 +348,7 @@ export default function EditConnectionModal({
connection.rateLimitOverrides?.maxConcurrent != null
? String(connection.rateLimitOverrides.maxConcurrent)
: "",
modelConcurrency: modelConcurrencyFormValue(connection.rateLimitOverrides),
apiKey: "",
// Unset per-connection override means "follow the global default" —
// surface that as an empty field (0 renders as an explicit opt-out).
Expand Down Expand Up @@ -546,16 +554,10 @@ export default function EditConnectionModal({
healthCheckInterval:
formData.healthCheckInterval === "" ? undefined : formData.healthCheckInterval,
};
const overrides: Record<string, number> = {};
if (formData.rpm.trim()) overrides.rpm = Number(formData.rpm);
if (formData.rpd.trim()) overrides.rpd = Number(formData.rpd);
if (formData.tpm.trim()) overrides.tpm = Number(formData.tpm);
if (formData.tpd.trim()) overrides.tpd = Number(formData.tpd);
if (formData.minTime.trim()) overrides.minTime = Number(formData.minTime);
if (formData.maxWaitMs.trim()) overrides.maxWaitMs = Number(formData.maxWaitMs);
if (formData.rateLimitMaxConcurrent.trim())
overrides.maxConcurrent = Number(formData.rateLimitMaxConcurrent);
updates.rateLimitOverrides = Object.keys(overrides).length > 0 ? overrides : null;
const rateLimit = buildRateLimitOverridesFromForm(formData, connection?.rateLimitOverrides);
const invalid = rateLimit.invalidModelConcurrency;
if (invalid) return setSaveError(t("rateLimitOverridesModelConcurrencyInvalid", invalid));
updates.rateLimitOverrides = rateLimit.overrides;
if (isAntigravityFamily) {
updates.projectId = trimmedCloudCodeProjectId || null;
}
Expand Down Expand Up @@ -1311,6 +1313,12 @@ export default function EditConnectionModal({
placeholder={t("inherit")}
hint={t("rateLimitOverridesMaxConcurrentHint")}
/>
<ModelConcurrencyField
value={formData.modelConcurrency}
onChange={(modelConcurrency) =>
setFormData({ ...formData, modelConcurrency })
}
/>
</div>
</div>
</div>
Expand Down
Loading