Skip to content
Merged
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
181 changes: 101 additions & 80 deletions open-sse/handlers/chatCore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -704,8 +704,9 @@ export async function handleChatCore({
const recordKeyHealthStatus = (
status: number,
creds: Record<string, unknown> | null | undefined,
transport?: string
): void => recordKeyHealthStatusFor(status, creds, log, transport);
transport?: string,
failureDetail?: string
): void => recordKeyHealthStatusFor(status, creds, log, transport, failureDetail);
// ── Phase 9.2: Idempotency check ──
// Resolve the idempotency key once here and reuse it at the Phase 9.2 save site below,
// rather than re-deriving it. (#3821-review LEDGER-6)
Expand Down Expand Up @@ -3174,11 +3175,20 @@ export async function handleChatCore({
});

if (
res.response.status === 401 &&
stream &&
(res.response.ok ||
res.response.status === HTTP_STATUS.UNAUTHORIZED ||
res.response.status === HTTP_STATUS.FORBIDDEN) &&
executionConnectionId &&
!(await shouldIsolateProbeFailures())
) {
recordKeyHealthStatus(401, execCreds);
const failureDetail = res.response.ok
? ""
: await res.response
.clone()
.text()
.catch(() => "");
recordKeyHealthStatus(res.response.status, execCreds, res.transport, failureDetail);
}

if (isModelScope() && res.response.status === 429 && attempts < maxAttempts - 1) {
Expand Down Expand Up @@ -3529,13 +3539,6 @@ export async function handleChatCore({
// Non-stream: release semaphore immediately after reading full response body.
const status = rawResult.response.status;

// Use execution credentials captured during request processing
if (
rawResult._executionCredentials?.connectionId &&
rawResult._executionCredentials?.apiKey
) {
recordKeyHealthStatus(status, rawResult._executionCredentials, rawResult.transport);
}
releaseRawResultAccountSemaphore =
typeof rawResult._accountSemaphoreRelease === "function"
? rawResult._accountSemaphoreRelease
Expand All @@ -3561,6 +3564,19 @@ export async function handleChatCore({
contentType,
upstreamStream
);
// Use the exact execution credential selected for this request. Model capability
// failures stay in routing telemetry; authoritative success only recovers this key.
if (
rawResult._executionCredentials?.connectionId &&
(rawResult._executionCredentials.apiKey || rawResult._executionCredentials.accessToken)
) {
recordKeyHealthStatus(
status,
rawResult._executionCredentials,
rawResult.transport,
status >= 400 ? payload : ""
);
}
releaseRawResultAccountSemaphore();
releaseRawResultAccountSemaphore = () => {};

Expand Down Expand Up @@ -4241,79 +4257,84 @@ export async function handleChatCore({
`[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active`
);
} else {
// Kimi's 403 says "billing cycle" for both an exhausted subscription and a
// temporary request window. Read its official usage endpoint before making
// the connection terminal: a non-zero Weekly quota plus an empty Ratelimit
// window must recover automatically at the reported reset time.
let kimiRateLimitResetAt: string | null = null;
if (provider === "kimi-coding") {
try {
const { fetchAndPersistProviderLimits } =
await import("@/lib/usage/providerLimits");
const { usage } = await fetchAndPersistProviderLimits(errorConnectionId, "manual");
kimiRateLimitResetAt = getKimiTemporaryRateLimitResetAt(usage);
} catch {
// Preserve the existing quota handling when Kimi's usage endpoint is unavailable.
// Kimi's 403 says "billing cycle" for both an exhausted subscription and a
// temporary request window. Read its official usage endpoint before making
// the connection terminal: a non-zero Weekly quota plus an empty Ratelimit
// window must recover automatically at the reported reset time.
let kimiRateLimitResetAt: string | null = null;
if (provider === "kimi-coding") {
try {
const { fetchAndPersistProviderLimits } =
await import("@/lib/usage/providerLimits");
const { usage } = await fetchAndPersistProviderLimits(
errorConnectionId,
"manual"
);
kimiRateLimitResetAt = getKimiTemporaryRateLimitResetAt(usage);
} catch {
// Preserve the existing quota handling when Kimi's usage endpoint is unavailable.
}
}
}

// Providers with per-model quotas — lock the model only, not the connection
const quotaCooldownMs = kimiRateLimitResetAt
? Math.max(new Date(kimiRateLimitResetAt).getTime() - Date.now(), 0)
: retryAfterMs || COOLDOWN_MS.rateLimit;
const accountSemaphoreKey = resolveAccountSemaphoreKey({
provider,
model: currentModel,
connectionId: errorConnectionId,
credentials,
});
if (accountSemaphoreKey) {
markAccountSemaphoreBlocked(accountSemaphoreKey, quotaCooldownMs);
}
if (kimiRateLimitResetAt) {
await updateProviderConnection(errorConnectionId, {
testStatus: "unavailable",
rateLimitedUntil: kimiRateLimitResetAt,
backoffLevel: 0,
lastErrorType: PROVIDER_ERROR_TYPES.RATE_LIMITED,
lastError: message,
errorCode: statusCode,
});
console.warn(
`[provider] Node ${errorConnectionId} Kimi request window exhausted (${statusCode}) — retrying after ${kimiRateLimitResetAt}`
);
} else if (isModelScope() && errorConnectionId) {
const lockFn = provider === "antigravity" ? lockExactModel : lockModel;
lockFn(provider, errorConnectionId, model, "quota_exhausted", quotaCooldownMs);
console.warn(
`[provider] Node ${errorConnectionId} ModelScope model quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (connection stays active)`
);
} else if (
lockModelIfPerModelQuota(
// Providers with per-model quotas — lock the model only, not the connection
const quotaCooldownMs = kimiRateLimitResetAt
? Math.max(new Date(kimiRateLimitResetAt).getTime() - Date.now(), 0)
: retryAfterMs || COOLDOWN_MS.rateLimit;
const accountSemaphoreKey = resolveAccountSemaphoreKey({
provider,
errorConnectionId,
model,
"quota_exhausted",
quotaCooldownMs
)
) {
const quotaScope = getQuotaScopeLabelForProvider(provider, model);
console.warn(
`[provider] Node ${errorConnectionId} ${quotaScope}-only quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (cooldown_scope=${quotaScope}, ttl_source=${retryAfterMs ? "upstream" : "inferred"}, connection stays active)`
);
} else {
await writeTerminalStatus(
errorConnectionId,
{
testStatus: "credits_exhausted",
model: currentModel,
connectionId: errorConnectionId,
credentials,
});
if (accountSemaphoreKey) {
markAccountSemaphoreBlocked(accountSemaphoreKey, quotaCooldownMs);
}
if (kimiRateLimitResetAt) {
await updateProviderConnection(errorConnectionId, {
testStatus: "unavailable",
rateLimitedUntil: kimiRateLimitResetAt,
backoffLevel: 0,
lastErrorType: PROVIDER_ERROR_TYPES.RATE_LIMITED,
lastError: message,
lastErrorType: errorType,
errorCode: String(statusCode),
},
"production"
);
console.warn(`[provider] Node ${errorConnectionId} exhausted quota (${statusCode})`);
}
errorCode: statusCode,
});
console.warn(
`[provider] Node ${errorConnectionId} Kimi request window exhausted (${statusCode}) — retrying after ${kimiRateLimitResetAt}`
);
} else if (isModelScope() && errorConnectionId) {
const lockFn = provider === "antigravity" ? lockExactModel : lockModel;
lockFn(provider, errorConnectionId, model, "quota_exhausted", quotaCooldownMs);
console.warn(
`[provider] Node ${errorConnectionId} ModelScope model quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (connection stays active)`
);
} else if (
lockModelIfPerModelQuota(
provider,
errorConnectionId,
model,
"quota_exhausted",
quotaCooldownMs
)
) {
const quotaScope = getQuotaScopeLabelForProvider(provider, model);
console.warn(
`[provider] Node ${errorConnectionId} ${quotaScope}-only quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (cooldown_scope=${quotaScope}, ttl_source=${retryAfterMs ? "upstream" : "inferred"}, connection stays active)`
);
} else {
await writeTerminalStatus(
errorConnectionId,
{
testStatus: "credits_exhausted",
lastError: message,
lastErrorType: errorType,
errorCode: String(statusCode),
},
"production"
);
console.warn(
`[provider] Node ${errorConnectionId} exhausted quota (${statusCode})`
);
}
} // close probeIsolated3 else
}
} else if (errorType === PROVIDER_ERROR_TYPES.UNAUTHORIZED) {
Expand Down
41 changes: 35 additions & 6 deletions open-sse/handlers/chatCore/keyHealth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,13 @@
* handleChatCore. Translates an upstream HTTP status into the in-memory key-health state
* (apiKeyRotator) for the connection's currently-selected key, and persists the change to the
* provider connection so it survives process restarts:
* - 401 → record a failure (warning, then invalid at the threshold), always persisted.
* - genuine 401/403 credential rejection → record a failure (warning, then invalid at the
* threshold), always persisted.
* - 402 → terminal (insufficient balance); mark the current key invalid immediately (#5239),
* persisted on the active→invalid transition.
* - 2xx → record a success, persisted only when recovering from a warning/invalid state.
* Any other status only refreshes the tracked extra-key set. The handler binds its `log` once and
* delegates here, keeping the existing call sites unchanged.
* Model availability failures remain model/routing telemetry even when an upstream reports them
* with 401/403. Any other status only refreshes the tracked extra-key set.
*/

import {
Expand All @@ -21,18 +22,46 @@ import {
trackConnectionExtraKeys,
type KeyHealth,
} from "../../services/apiKeyRotator.ts";
import { isModelUnavailableError } from "../../services/modelFamilyFallback.ts";
import { updateProviderConnection } from "@/lib/db/providers";

type KeyHealthLog = {
warn?: (tag: string, message: string) => void;
error?: (tag: string, message: string) => void;
} | null;

const CREDENTIAL_FAILURE_PATTERNS = [
/\b(?:invalid|incorrect|expired|missing|revoked)\s+api[\s_-]?key\b/i,
/\bapi[\s_-]?key\s+(?:is\s+)?(?:invalid|incorrect|expired|missing|revoked|not\s+valid)\b/i,
/\bauthentication[\s_-]+(?:failed|error|required)\b/i,
/\b(?:invalid|expired|missing|revoked)\s+(?:token|credentials?|bearer)\b/i,
/\bunauthorized\b/i,
/\bnot\s+authenticated\b/i,
/\bforbidden\b/i,
/\baccess\s+denied\b/i,
];

function isModelCapabilityFailure(status: number, failureDetail: string): boolean {
if (!failureDetail) return false;
const normalizedDetail = failureDetail.replace(/[_-]+/g, " ");
// Model-family fallback already owns these phrases. Use a model-capable status for
// classification because some aggregators misreport the same model rejection as 401.
return isModelUnavailableError(status === 401 ? 403 : status, normalizedDetail);
}

function isCredentialFailure(status: number, failureDetail: string): boolean {
if (status !== 401 && status !== 403) return false;
if (isModelCapabilityFailure(status, failureDetail)) return false;
if (status === 401) return true;
return CREDENTIAL_FAILURE_PATTERNS.some((pattern) => pattern.test(failureDetail));
}

export function recordKeyHealthStatus(
status: number,
creds: Record<string, unknown> | null | undefined,
log?: KeyHealthLog,
transport?: string
transport?: string,
failureDetail = ""
): void {
// CLIProxyAPI owns a shared external credential pool. Its auth failures cannot be
// attributed to the native OmniRoute connection selected before proxy dispatch.
Expand All @@ -55,11 +84,11 @@ export function recordKeyHealthStatus(

trackConnectionExtraKeys(connId, extraKeys);

if (status === 401) {
if (isCredentialFailure(status, failureDetail)) {
const updatedHealth = recordKeyFailure(connId, currentKeyId);
log?.warn?.(
"AUTH",
`401 on connection ${connId.slice(0, 8)} - key marked as failed (failure #${updatedHealth.failures})`
`${status} on connection ${connId.slice(0, 8)} - key marked as failed (failure #${updatedHealth.failures})`
);

// Persist health status to DB on every failure (not just invalid transitions)
Expand Down
28 changes: 27 additions & 1 deletion open-sse/services/apiKeyRotator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,10 @@ const MAX_CONNECTION_EXTRA_KEYS = 500;
*/
export function trackConnectionExtraKeys(connectionId: string, extraKeys: string[]): void {
const validExtras = extraKeys.filter((k) => typeof k === "string" && k.trim().length > 0);
if (!_connectionExtraKeys.has(connectionId) && _connectionExtraKeys.size >= MAX_CONNECTION_EXTRA_KEYS) {
if (
!_connectionExtraKeys.has(connectionId) &&
_connectionExtraKeys.size >= MAX_CONNECTION_EXTRA_KEYS
) {
const oldest = _connectionExtraKeys.keys().next().value;
if (oldest !== undefined) _connectionExtraKeys.delete(oldest);
}
Expand Down Expand Up @@ -308,6 +311,29 @@ export function syncHealthFromDB(connectionId: string, health?: Record<string, K
}
}

/** Recover one authoritatively validated key without changing sibling-key health. */
export function recoverKeyHealth(
connectionId: string,
keyId: string,
providerSpecificData: unknown
): Record<string, unknown> | undefined {
const data =
providerSpecificData && typeof providerSpecificData === "object"
? (providerSpecificData as Record<string, unknown>)
: {};
const health = data.apiKeyHealth as Record<string, KeyHealth> | undefined;
const currentHealth = health?.[keyId];
if (!currentHealth || (currentHealth.status === "active" && currentHealth.failures === 0)) {
return undefined;
}

syncHealthFromDB(connectionId, health);
return {
...data,
apiKeyHealth: { ...health, [keyId]: recordKeySuccess(connectionId, keyId) },
};
}

/**
* Reset the rotation index for a connection.
* Call this when a key fails (401/403) to skip the bad key next time.
Expand Down
21 changes: 10 additions & 11 deletions src/app/api/providers/[id]/test/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ import {
import { providerAllowsOptionalApiKey } from "@/shared/constants/providers";
import { shouldUseApiKeyConnectionTest } from "./webSessionTestDispatch";
import { testCodexAppServerConnection, makeDiagnosis } from "./codexAppServerHealth";
import { removeConnectionHealth } from "@omniroute/open-sse/services/apiKeyRotator.ts";
import { recoverKeyHealth } from "@omniroute/open-sse/services/apiKeyRotator.ts";
import { shouldClearErrorStateOnValidProbe } from "@/lib/usage/providerLimits";
import { isConnectionUnavailableToAuxiliaryActivity } from "@/lib/exclusiveLeaseIsolation";
import { classifyAmbiguousOrAuthError, type ClassifyFailureArgs } from "./mistralAmbiguousAuth";
Expand Down Expand Up @@ -1125,7 +1125,11 @@ export async function testSingleConnection(connectionId: string, validationModel
lastError: clearErrorState ? null : result.valid ? connection.lastError : result.error,
lastErrorAt: clearErrorState ? null : result.valid ? connection.lastErrorAt : now,
lastTested: now,
lastErrorType: clearErrorState ? null : result.valid ? connection.lastErrorType : diagnosis.type,
lastErrorType: clearErrorState
? null
: result.valid
? connection.lastErrorType
: diagnosis.type,
lastErrorSource: clearErrorState
? null
: result.valid
Expand All @@ -1147,16 +1151,11 @@ export async function testSingleConnection(connectionId: string, validationModel

if (clearErrorState) {
updateData.backoffLevel = 0;
}

const psd = connection?.providerSpecificData as Record<string, unknown> | undefined;
updateData.providerSpecificData = {
...(psd || {}),
apiKeyHealth: {},
};

try {
removeConnectionHealth(connectionId);
} catch {}
if (result.valid && (connection.apiKey || connection.accessToken)) {
const recovered = recoverKeyHealth(connectionId, "primary", connection.providerSpecificData);
if (recovered) updateData.providerSpecificData = recovered;
}

// If token was refreshed, update tokens in DB
Expand Down
Loading
Loading