feat(proxy-fitness): pool fitness registry + Smart rotation + pool-scoped retry
Pool/IP fitness registry (globalThis-backed, provider::model scopes, 5-min cooldown, provider::* wildcard) fed by pool-scoped failures: freebuff limited-IP / model-locked gates and opencode free per-IP limits. chatCore retries a failed pool via another pool without locking the account; new Smart rotation strategy (per-connection + no-auth providers) skips unfit pools. Proxy Fitness dashboard page: active-block table with provider/IP filters, per-record Clear and provider-scoped Clear All; egress column via pool geo enrichment. Unit tests for the registry.
This commit is contained in:
@@ -29,6 +29,12 @@ import { getCapabilitiesForModel } from "../providers/capabilities.js";
|
||||
import { stripUnsupportedModalities } from "../translator/concerns/modality.js";
|
||||
import { prefetchRemoteImages } from "../translator/concerns/prefetch.js";
|
||||
import { resolveSessionId } from "../utils/sessionManager.js";
|
||||
import { markPoolUnfit, clearPoolUnfit } from "../services/proxyPoolFitness.js";
|
||||
|
||||
// Pool-scoped failure retry: when an executor tags an error as belonging to a
|
||||
// proxy pool (region gate, dead proxy, …), re-resolve the proxy config
|
||||
// excluding that pool and retry instead of failing the whole account.
|
||||
const MAX_POOL_RETRIES = 2;
|
||||
|
||||
/**
|
||||
* Core chat handler - shared between SSE and Worker
|
||||
@@ -57,7 +63,7 @@ export function stripContinuityFields(body) {
|
||||
return body;
|
||||
}
|
||||
|
||||
export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, headroomEnabled, headroomUrl, headroomCompressUserMessages, cavemanEnabled, cavemanLevel, ponytailEnabled, ponytailLevel, pxpipeEnabled, pxpipeMinChars, pxpipeTimeoutMs, pxpipeTransform, onPxpipeEvent, sourceFormatOverride, providerThinking }) {
|
||||
export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, headroomEnabled, headroomUrl, headroomCompressUserMessages, cavemanEnabled, cavemanLevel, ponytailEnabled, ponytailLevel, pxpipeEnabled, pxpipeMinChars, pxpipeTimeoutMs, pxpipeTransform, onPxpipeEvent, sourceFormatOverride, providerThinking, resolveProxyConfig }) {
|
||||
const { provider, model } = modelInfo;
|
||||
const requestStartTime = Date.now();
|
||||
// Stable per-session color so all lines of one CLI conversation share a tag
|
||||
@@ -291,16 +297,23 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
log, provider, model, reqTag
|
||||
});
|
||||
|
||||
const proxyOptions = {
|
||||
connectionProxyEnabled: credentials?.providerSpecificData?.connectionProxyEnabled === true,
|
||||
connectionProxyUrl: credentials?.providerSpecificData?.connectionProxyUrl || "",
|
||||
connectionNoProxy: credentials?.providerSpecificData?.connectionNoProxy || "",
|
||||
vercelRelayUrl: credentials?.providerSpecificData?.vercelRelayUrl || "",
|
||||
};
|
||||
// Build proxy options from the resolved provider-specific data. `strictProxy`
|
||||
// is forced for freebuff so a dead/limited pool can never leak the request
|
||||
// to the caller's real IP (the freebuff session tier is per-egress-IP).
|
||||
const buildProxyOptions = (psd = {}) => ({
|
||||
connectionProxyEnabled: psd?.connectionProxyEnabled === true,
|
||||
connectionProxyUrl: psd?.connectionProxyUrl || "",
|
||||
connectionNoProxy: psd?.connectionNoProxy || "",
|
||||
vercelRelayUrl: psd?.vercelRelayUrl || "",
|
||||
strictProxy: psd?.strictProxy === true || provider === "freebuff",
|
||||
proxyPoolId: psd?.proxyPoolId || psd?.connectionProxyPoolId || null,
|
||||
});
|
||||
|
||||
let proxyOptions = buildProxyOptions(credentials?.providerSpecificData || {});
|
||||
|
||||
if (proxyOptions.vercelRelayUrl) {
|
||||
const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown";
|
||||
const poolId = credentials?.providerSpecificData?.connectionProxyPoolId || "none";
|
||||
const poolId = proxyOptions.proxyPoolId || "none";
|
||||
log?.info?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | vercel-relay=${proxyOptions.vercelRelayUrl}`);
|
||||
} else if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionProxyUrl) {
|
||||
let maskedProxyUrl = proxyOptions.connectionProxyUrl;
|
||||
@@ -314,7 +327,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
// Keep raw if URL parsing fails
|
||||
}
|
||||
|
||||
const poolId = credentials?.providerSpecificData?.connectionProxyPoolId || "none";
|
||||
const poolId = proxyOptions.proxyPoolId || "none";
|
||||
const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown";
|
||||
log?.info?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | url=${maskedProxyUrl}`);
|
||||
}
|
||||
@@ -324,13 +337,65 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
log?.debug?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | no_proxy=${proxyOptions.connectionNoProxy}`);
|
||||
}
|
||||
|
||||
// Execute request — with pool-scoped retry: a failed pool (region gate,
|
||||
// per-IP limit, dead proxy) is marked unfit and the request retried via
|
||||
// another pool instead of failing the account. Covers both thrown errors
|
||||
// (executor.execute) and non-ok responses declared poolScoped via
|
||||
// parseError — poolId/scope are completed here from proxyOptions.
|
||||
const proxyScope = `${provider}::${model}`;
|
||||
let parsedNonOk = null;
|
||||
|
||||
const tryNextPool = async (poolScoped, reasonMsg) => {
|
||||
const failed = {
|
||||
poolId: poolScoped?.poolId || proxyOptions.proxyPoolId || null,
|
||||
scope: poolScoped?.scope || proxyScope,
|
||||
reason: poolScoped?.reason || "pool-scoped",
|
||||
};
|
||||
markPoolUnfit(failed.poolId, failed.scope, undefined, failed.reason);
|
||||
log?.warn?.("PROXY", `${provider.toUpperCase()} | pool ${failed.poolId || "?"} unfit for ${failed.scope} (${failed.reason}) — retry with another pool. ${reasonMsg || ""}`);
|
||||
try {
|
||||
const resolved = await resolveProxyConfig(credentials, [failed.poolId]);
|
||||
if (resolved?.proxyPoolId) {
|
||||
credentials.providerSpecificData = { ...(credentials.providerSpecificData || {}), ...resolved };
|
||||
proxyOptions = buildProxyOptions(credentials.providerSpecificData);
|
||||
return true;
|
||||
}
|
||||
} catch (resolverError) {
|
||||
// A resolver failure must not mask the original pool error.
|
||||
log?.warn?.("PROXY", `${provider.toUpperCase()} | pool re-resolve failed: ${resolverError.message}`);
|
||||
}
|
||||
return false;
|
||||
};
|
||||
|
||||
const executeWithPoolFallback = async (attempt = 0) => {
|
||||
let result;
|
||||
try {
|
||||
result = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions });
|
||||
} catch (error) {
|
||||
if (error?.poolScoped && typeof resolveProxyConfig === "function" && attempt < MAX_POOL_RETRIES) {
|
||||
if (await tryNextPool(error.poolScoped, error.message)) return executeWithPoolFallback(attempt + 1);
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
// Non-ok response that the executor declared pool/IP-scoped (e.g. opencode
|
||||
// free per-IP limit) — parse once, retry via another pool when possible.
|
||||
if (!result.response.ok) {
|
||||
const parsed = await parseUpstreamError(result.response, executor);
|
||||
if (parsed.poolScoped && typeof resolveProxyConfig === "function" && attempt < MAX_POOL_RETRIES) {
|
||||
if (await tryNextPool(parsed.poolScoped, parsed.message)) return executeWithPoolFallback(attempt + 1);
|
||||
}
|
||||
parsedNonOk = parsed;
|
||||
}
|
||||
return result;
|
||||
};
|
||||
|
||||
// Execute request
|
||||
let providerResponse, providerUrl, providerHeaders, finalBody;
|
||||
// Most executors return their registry format. Cursor AgentService is an
|
||||
// exception: it is decoded by the executor into OpenAI-compatible output.
|
||||
let providerResponseFormat = targetFormat;
|
||||
try {
|
||||
const result = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions });
|
||||
const result = await executeWithPoolFallback();
|
||||
providerResponse = result.response;
|
||||
providerUrl = result.url;
|
||||
providerHeaders = result.headers;
|
||||
@@ -359,7 +424,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
if (log?.errorLine) {
|
||||
log.errorLine(reqTag, "✗", `ERROR 502 · ${provider}/${model} · ${Date.now() - requestStartTime}ms\n ${errMsg}${error.stack ? `\n ${error.stack}` : ""}`);
|
||||
}
|
||||
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg);
|
||||
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg, error?.resetsAtMs || undefined);
|
||||
}
|
||||
|
||||
// Handle 401/403 - try token refresh (skip for noAuth providers)
|
||||
@@ -402,7 +467,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
// Provider returned error
|
||||
if (!providerResponse.ok) {
|
||||
trackPendingRequest(model, provider, connectionId, false, true);
|
||||
const { statusCode, message, resetsAtMs } = await parseUpstreamError(providerResponse, executor);
|
||||
const { statusCode, message, resetsAtMs } = parsedNonOk || await parseUpstreamError(providerResponse, executor);
|
||||
appendRequestLog({ model, provider, connectionId, status: `FAILED ${statusCode}` }).catch(() => { });
|
||||
saveRequestDetail(buildRequestDetail({
|
||||
provider, model, connectionId,
|
||||
|
||||
Reference in New Issue
Block a user