From fd22290a34800bbf803e8c7a9f6d44776458663c Mon Sep 17 00:00:00 2001 From: "MUH. IQRAM BAHRING" Date: Fri, 7 Aug 2026 03:51:19 +0800 Subject: [PATCH] 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. --- open-sse/executors/freebuff.js | 194 ++++++++- open-sse/executors/opencode.js | 19 + open-sse/handlers/chatCore.js | 89 +++- open-sse/services/proxyPoolFitness.js | 126 ++++++ open-sse/utils/error.js | 9 +- .../dashboard/providers/[id]/ConnectionRow.js | 8 +- .../dashboard/proxy-fitness/page.js | 384 ++++++++++++++++++ .../proxy-pools/[id]/fitness/clear/route.js | 25 ++ .../proxy-pools/fitness/clear-all/route.js | 22 + src/app/api/proxy-pools/fitness/route.js | 13 + src/app/api/proxy-pools/route.js | 3 + src/lib/network/connectionProxy.js | 38 +- src/shared/components/NoAuthProxyCard.js | 5 +- src/shared/components/Sidebar.js | 1 + src/sse/handlers/chat.js | 18 + src/sse/services/auth.js | 25 +- tests/unit/proxy-pool-fitness.test.js | 56 +++ 17 files changed, 1000 insertions(+), 35 deletions(-) create mode 100644 open-sse/services/proxyPoolFitness.js create mode 100644 src/app/(dashboard)/dashboard/proxy-fitness/page.js create mode 100644 src/app/api/proxy-pools/[id]/fitness/clear/route.js create mode 100644 src/app/api/proxy-pools/fitness/clear-all/route.js create mode 100644 src/app/api/proxy-pools/fitness/route.js create mode 100644 tests/unit/proxy-pool-fitness.test.js diff --git a/open-sse/executors/freebuff.js b/open-sse/executors/freebuff.js index 655060aa..4f23e630 100644 --- a/open-sse/executors/freebuff.js +++ b/open-sse/executors/freebuff.js @@ -8,6 +8,7 @@ import { DEFAULT_RETRY_CONFIG, resolveRetryEntry, } from "../config/runtimeConfig.js"; +import { markPoolUnfit, clearPoolUnfit } from "../services/proxyPoolFitness.js"; /** * Freebuff Executor — OpenAI-compatible chat completions on @@ -87,8 +88,110 @@ const FREE_ROOT_AGENT_BY_MODEL = { // don't share one session row). Re-claims are driven by the cache expiring or // by a 428 from chat — no early re-claim, so we never POST /session while our // own row is still active (which could come back as a spurious model_locked). -const sessionCache = new Map(); // `${token}::${model}` -> { instanceId, expiresAt } -const inflight = new Map(); // dedupe concurrent claims for the same key +// All state lives on globalThis so Next dev (Turbopack) bundles share ONE copy. +const FB_STATE_KEY = "__9routerFreebuffState__"; +const fbState = (globalThis[FB_STATE_KEY] ??= { + sessionCache: new Map(), // `${token}::${model}` -> { instanceId, expiresAt } + inflight: new Map(), // dedupe concurrent claims for the same key + modelLockCooldowns: new Map(), // `${token}::${model}` -> expiresAt (ms) + poolLimitCooldowns: new Map(), // `${proxyKey}::${model}` -> expiresAt (ms) +}); +const sessionCache = fbState.sessionCache; +const inflight = fbState.inflight; +const modelLockCooldowns = fbState.modelLockCooldowns; +const poolLimitCooldowns = fbState.poolLimitCooldowns; + +const MODEL_LOCK_COOLDOWN_MS = 10 * 60 * 1000; // session bound to another model (~1h) — re-check every 10 min +const POOL_LIMITED_COOLDOWN_MS = 5 * 60 * 1000; // IP tier refuses this model — try a different pool/relay + +// Cooldown maps need pruning: expired entries are cleared on write (sweep) and +// on read, so long-running servers don't accumulate one entry per (account,model) +// / (proxy,model) forever. +function setCooldown(map, key, until) { + const now = Date.now(); + for (const [k, v] of map) { + if (v <= now) map.delete(k); + } + map.set(key, until); +} + +function getCooldown(map, key) { + const until = map.get(key); + if (until == null) return null; + if (until <= Date.now()) { + map.delete(key); + return null; + } + return until; +} + +function proxyKeyOf(proxyOptions) { + return proxyOptions?.vercelRelayUrl || proxyOptions?.connectionProxyUrl || "direct"; +} + +function sessionGateFromText(text) { + let parsed = {}; + try { parsed = JSON.parse(String(text || "")); } catch { parsed = {}; } + return classifySessionGate(parsed.error || parsed.error_type || "", parsed.message || "", parsed.currentModel || null); +} + +// Parse a 409/428/410 body into { kind, currentModel }. `msg` may be a whole +// error string containing a JSON tail (requestSession errors embed the body). +function sessionGateFromError(error) { + const msg = String(error?.message || ""); + const start = msg.indexOf("{"); + if (start < 0) return null; + try { + const parsed = JSON.parse(msg.slice(start)); + return classifySessionGate(parsed.error || "", parsed.message || "", parsed.currentModel || null); + } catch { + return null; + } +} + +function classifySessionGate(code, message, currentModel) { + if (code === "session_superseded") return { kind: "superseded" }; + if (code === "model_locked") return { kind: "model_locked", currentModel }; + // session_model_mismatch with the limited-tier message is an IP-tier refusal; + // without it (or unknown) treat it as a model lock so we don't reclaim in a loop. + if (code === "session_model_mismatch") { + return /limited/i.test(String(message || "")) + ? { kind: "limited_ip" } + : { kind: "model_locked", currentModel }; + } + return { kind: "stale" }; // 428/410/unknown → reclaim +} + +// Applies cooldowns and throws for non-reclaimable gates. Never returns for them. +function throwSessionGateError(gate, { token, model, proxyKey, poolId, log }) { + if (gate.kind === "model_locked") { + const until = Date.now() + MODEL_LOCK_COOLDOWN_MS; + setCooldown(modelLockCooldowns, `${token}::${model}`, until); + const label = gate.currentModel ? `"${gate.currentModel}"` : "another model"; + const err = new Error( + `Freebuff session is locked to ${label} — it cannot serve ${model}. End the session on freebuff.com or wait for it to expire (~1h).`, + ); + err.status = 409; + err.resetsAtMs = until; + log?.warn?.("AUTH", `Freebuff model_locked (session=${label}, requested=${model}) — model cooldown ${MODEL_LOCK_COOLDOWN_MS / 60000}min`); + throw err; + } + if (gate.kind === "limited_ip") { + const until = Date.now() + POOL_LIMITED_COOLDOWN_MS; + setCooldown(poolLimitCooldowns, `${proxyKey}::${model}`, until); + const scope = `freebuff::${model}`; + if (poolId) markPoolUnfit(poolId, scope, until, "limited_ip"); + // Pool-scoped, not account-scoped: the caller retries via another pool + // instead of locking the account (resetsAtMs intentionally absent). + const err = new Error( + `Freebuff limited-mode IP rejected ${model} — this IP only allows DeepSeek V4 Flash / MiMo 2.5. Use a full-access proxy or a different model.`, + ); + err.status = 409; + err.poolScoped = { poolId, scope, reason: "limited_ip" }; + log?.warn?.("AUTH", `Freebuff limited-IP refused ${model} (proxy=${proxyKey.slice(0, 40)}…) — cooldown ${POOL_LIMITED_COOLDOWN_MS / 60000}min`); + throw err; + } +} function sessionOrigin() { return new URL(PROVIDERS.freebuff.baseUrl).origin; // https://www.codebuff.com @@ -185,7 +288,11 @@ async function requestSession(token, model, proxyOptions) { async function ensureSession(token, model, proxyOptions, force = false) { const key = sessionCacheKey(token, model); + // Lazy prune: drop stale rows so the cache never accumulates expired entries. const cached = sessionCache.get(key); + if (cached && cached.expiresAt <= Date.now()) { + sessionCache.delete(key); + } if (!force && cached && cached.expiresAt > Date.now()) { return { instanceId: cached.instanceId, status: "active" }; } @@ -262,6 +369,32 @@ export function resetSessionCache() { inflight.clear(); } +// Periodic sweeper: drop stale session rows + expired cooldowns so long-running +// servers never accumulate state for accounts/models no longer in use. +// Returns how many entries were removed. +export function pruneSessionState(now = Date.now()) { + let removed = 0; + for (const [key, entry] of sessionCache) { + if (entry?.expiresAt && entry.expiresAt <= now) { + sessionCache.delete(key); + removed += 1; + } + } + for (const [key, until] of modelLockCooldowns) { + if (until <= now) { + modelLockCooldowns.delete(key); + removed += 1; + } + } + for (const [key, until] of poolLimitCooldowns) { + if (until <= now) { + poolLimitCooldowns.delete(key); + removed += 1; + } + } + return removed; +} + export class FreebuffExecutor extends BaseExecutor { constructor() { super("freebuff", PROVIDERS.freebuff); @@ -282,6 +415,12 @@ export class FreebuffExecutor extends BaseExecutor { cost_mode: "free", }; body.provider = { allow_fallbacks: false }; + // Freebuff agents (base2-free-*) own reasoning: the backend applies the + // agent's reasoningOptions.effort server-side, so a client-sent + // reasoning_effort / reasoning.effort collides with that default → + // 400 "both provided with conflicting values". Mirror the CLI: send none. + delete body.reasoning_effort; + delete body.reasoning; // Free-tier gate: first system message must open with the CLI marker. return injectFreebuffMarker(body); } @@ -292,10 +431,32 @@ export class FreebuffExecutor extends BaseExecutor { throw new Error("Freebuff requires a connected Freebuff login (no access token found)"); } + // Fail fast while a known-dead (account,model) / (proxy,model) pair is in + // cooldown — no session claim, no run registration, no upstream spam. + const proxyKey = proxyKeyOf(proxyOptions); + const poolId = proxyOptions?.proxyPoolId || null; + const scope = `freebuff::${model}`; + const lockUntil = getCooldown(modelLockCooldowns, `${token}::${model}`); + if (lockUntil) { + const err = new Error(`Freebuff session locked to another model — retry after ${new Date(lockUntil).toLocaleTimeString()}`); + err.status = 409; + err.resetsAtMs = lockUntil; + throw err; + } + const poolUntil = getCooldown(poolLimitCooldowns, `${proxyKey}::${model}`); + if (poolUntil) { + const err = new Error(`Freebuff limited-mode IP rejected ${model} — retry with a full-access proxy after ${new Date(poolUntil).toLocaleTimeString()}`); + err.status = 409; + err.poolScoped = { poolId, scope, reason: "limited_ip" }; + throw err; + } + let session; try { session = await ensureSession(token, model, proxyOptions); } catch (error) { + const gate = sessionGateFromError(error); + if (gate) throwSessionGateError(gate, { token, model, proxyKey, poolId, log }); log?.error?.("AUTH", `Freebuff session failed: ${error.message}`); throw error; } @@ -396,9 +557,17 @@ export class FreebuffExecutor extends BaseExecutor { // 409 session_superseded — another instance took over the session // 409 session_model_mismatch — session bound to a different model // 410 session_expired — the active session's expires_at passed - // In every case: abandon the run, force a fresh session claim + a fresh - // run, then retry exactly once. + // model_locked / limited-tier mismatches are NOT reclaimable — the server + // keeps refusing until the session expires or the IP tier changes, so we + // set a cooldown and fail fast instead of force re-claiming in a loop. if (SESSION_STALE_CODES.has(response.status)) { + const text = await response.text().catch(() => ""); + const gate = sessionGateFromText(text); + if (gate.kind === "model_locked" || gate.kind === "limited_ip") { + markFinished("cancelled"); + throwSessionGateError(gate, { token, model, proxyKey, poolId, log }); + } + log?.debug?.("AUTH", `Freebuff ${response.status} session gate — re-claiming session`); markFinished("cancelled"); try { @@ -406,21 +575,34 @@ export class FreebuffExecutor extends BaseExecutor { runId = await startRun(token, model, proxyOptions); activeRunId = runId; } catch (error) { + const gate2 = sessionGateFromError(error); + if (gate2) throwSessionGateError(gate2, { token, model, proxyKey, poolId, log }); log?.error?.("AUTH", `Freebuff session re-claim failed: ${error.message}`); throw error; } ({ response, transformedBody } = await doChat()); if (SESSION_STALE_CODES.has(response.status)) { - const text = await response.text().catch(() => ""); + const text2 = await response.text().catch(() => ""); + const gate3 = sessionGateFromText(text2); + if (gate3.kind === "model_locked" || gate3.kind === "limited_ip") { + throwSessionGateError(gate3, { token, model, proxyKey, poolId, log }); + } const err = new Error( - `Freebuff session gate refused (${response.status}) — another freebuff instance may be holding the session. ${text.slice(0, 160)}`, + `Freebuff session gate refused (${response.status}) — another freebuff instance may be holding the session. ${text2.slice(0, 160)}`, ); err.status = response.status; throw err; } } + // A successful chat means the pair is healthy again — lift any cooldowns. + if (response.ok) { + modelLockCooldowns.delete(`${token}::${model}`); + poolLimitCooldowns.delete(`${proxyKey}::${model}`); + if (poolId) clearPoolUnfit(poolId, scope); + } + // The authToken has no refresh path — when it dies, the user re-logs in. // Drop the cached session for this token so a re-login starts clean. if (response.status === 401) { diff --git a/open-sse/executors/opencode.js b/open-sse/executors/opencode.js index f7aee211..f1af1ebd 100644 --- a/open-sse/executors/opencode.js +++ b/open-sse/executors/opencode.js @@ -5,6 +5,12 @@ import { injectReasoningContent } from "../utils/reasoningContentInjector.js"; // Models that use /zen/v1/messages (claude format) const MESSAGES_MODELS = new Set(); +// OpenCode free tier is limited per egress IP — a 429/403 with a limit-ish +// body means the POOL's IP is exhausted, not the account. Declare it +// pool-scoped so chatCore marks the pool unfit, retries via another pool, and +// it shows up (clearable) on the Proxy Fitness page. +const IP_LIMIT_BODY = /limit|rate|quota|exhausted|capacity|too many|retry/i; + export class OpenCodeExecutor extends BaseExecutor { constructor() { super("opencode", PROVIDERS.opencode); @@ -29,4 +35,17 @@ export class OpenCodeExecutor extends BaseExecutor { "Accept": "text/event-stream" }; } + + parseError(response, bodyText) { + const status = response?.status || 0; + const text = String(bodyText || ""); + if ((status === 429 || status === 403) && IP_LIMIT_BODY.test(text)) { + return { + status, + message: text.slice(0, 300) || `OpenCode free limit (${status})`, + poolScoped: { reason: "ip-limit" }, + }; + } + return null; // fall through to default parsing + } } diff --git a/open-sse/handlers/chatCore.js b/open-sse/handlers/chatCore.js index 9f17cf94..d21db965 100644 --- a/open-sse/handlers/chatCore.js +++ b/open-sse/handlers/chatCore.js @@ -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, diff --git a/open-sse/services/proxyPoolFitness.js b/open-sse/services/proxyPoolFitness.js new file mode 100644 index 00000000..72fc4644 --- /dev/null +++ b/open-sse/services/proxyPoolFitness.js @@ -0,0 +1,126 @@ +// Pool fitness registry — shared, in-memory, engine layer. +// +// Rotation strategies can opt into region/provider-aware pool selection: an +// executor that learns a pool's egress is unfit for a provider/model (e.g. +// Freebuff limited-mode IP) marks it here, and the pool picker skips it for +// that scope until the cooldown expires. +// +// Scope format: `provider::model` (e.g. "freebuff::openai/gpt-5.6-luna"). +// All functions are fail-open: unknown pool/scope ⇒ fit. + +// Scope format: `provider::model` (e.g. "freebuff::openai/gpt-5.6-luna"). +// All functions are fail-open: unknown pool/scope ⇒ fit. +// +// State lives on globalThis so Next dev (Turbopack) never splits one Map into +// several per-bundle copies — the executor/chatCore marks and the +// /api/proxy-pools/fitness reader must share the SAME registry. + +const FITNESS_STATE_KEY = "__9routerPoolFitness__"; +const fitness = (globalThis[FITNESS_STATE_KEY] ??= new Map()); // poolId -> Map + +export const POOL_UNFIT_MS = 5 * 60 * 1000; + +export function markPoolUnfit(poolId, scope, until = Date.now() + POOL_UNFIT_MS, reason = "") { + if (!poolId || !scope) return; + const byScope = fitness.get(poolId) || new Map(); + byScope.set(scope, { until, reason }); + fitness.set(poolId, byScope); +} + +export function clearPoolUnfit(poolId, scope) { + const byScope = fitness.get(poolId); + if (!byScope) return; + byScope.delete(scope); + if (byScope.size === 0) fitness.delete(poolId); +} + +// "provider::model" -> "provider::*" (null when the scope has no provider part) +function providerWildcardScope(scope) { + const sep = String(scope || "").indexOf("::"); + if (sep < 0) return null; + return `${scope.slice(0, sep)}::*`; +} + +export function isPoolFit(poolId, scope, now = Date.now()) { + if (!poolId) return true; + const byScope = fitness.get(poolId); + if (!byScope) return true; + // A provider-wide mark ("provider::*", e.g. manual blocks) also covers any + // model lookup for that provider. + const candidates = [scope, providerWildcardScope(scope)]; + for (const key of candidates) { + if (!key) continue; + const entry = byScope.get(key); + if (!entry) continue; + if (entry.until <= now) { + byScope.delete(key); + if (byScope.size === 0) fitness.delete(poolId); + continue; + } + return false; + } + return true; +} + +// Keep only pool ids that are not in cooldown for the scope. +export function fitPoolIds(poolIds, scope, now = Date.now()) { + return (poolIds || []).filter((id) => isPoolFit(id, scope, now)); +} + +// Clear every mark — or only scopes belonging to one provider (`provider::*`). +export function clearAllPoolUnfit(provider = null) { + if (provider) { + const prefix = `${provider}::`; + for (const [poolId, byScope] of fitness) { + for (const scope of [...byScope.keys()]) { + if (scope.startsWith(prefix)) byScope.delete(scope); + } + if (byScope.size === 0) fitness.delete(poolId); + } + return; + } + fitness.clear(); +} + +// Snapshot of live (non-expired) marks — expired entries are pruned here so +// consumers never see stale data and memory stays bounded. +// Test helper: drop all marks (module state is globalThis-backed). +export function resetPoolFitness() { + fitness.clear(); +} + +// Sweep all expired marks. Returns how many scope entries were removed. +export function pruneExpired(now = Date.now()) { + let removed = 0; + for (const [poolId, byScope] of fitness) { + for (const [scope, entry] of byScope) { + if (entry.until <= now) { + byScope.delete(scope); + removed += 1; + } + } + if (byScope.size === 0) fitness.delete(poolId); + } + return removed; +} + +export function poolFitnessSnapshot(now = Date.now()) { + const out = {}; + for (const [poolId, byScope] of fitness) { + let pruned = false; + for (const [scope, entry] of byScope) { + if (entry.until <= now) { + byScope.delete(scope); + pruned = true; + } + } + if (byScope.size === 0) { + fitness.delete(poolId); + continue; + } + if (pruned || byScope.size > 0) { + out[poolId] = Object.fromEntries(byScope); + } + } + return out; +} diff --git a/open-sse/utils/error.js b/open-sse/utils/error.js index 315723e3..d44a8344 100644 --- a/open-sse/utils/error.js +++ b/open-sse/utils/error.js @@ -69,7 +69,14 @@ export async function parseUpstreamError(response, executor = null) { const parsed = executor.parseError(response, bodyText); if (parsed && typeof parsed === "object") { const msg = parsed.message || DEFAULT_ERROR_MESSAGES[response.status] || `Upstream error: ${response.status}`; - return { statusCode: parsed.status || response.status, message: msg, resetsAtMs: parsed.resetsAtMs }; + return { + statusCode: parsed.status || response.status, + message: msg, + resetsAtMs: parsed.resetsAtMs, + // Executors declare IP/pool-scoped failures (e.g. per-IP rate limits) + // here; chatCore completes poolId/scope and retries via another pool. + poolScoped: parsed.poolScoped, + }; } } catch { /* fall through to default parsing */ } } diff --git a/src/app/(dashboard)/dashboard/providers/[id]/ConnectionRow.js b/src/app/(dashboard)/dashboard/providers/[id]/ConnectionRow.js index 51c8e2a5..2e4f92e4 100644 --- a/src/app/(dashboard)/dashboard/providers/[id]/ConnectionRow.js +++ b/src/app/(dashboard)/dashboard/providers/[id]/ConnectionRow.js @@ -43,7 +43,7 @@ export default function ConnectionRow({ connection, proxyPools, isOAuth, isFirst } if (selectedProxyIds.length > 1) { - const strategyLabel = rotationStrategy === "random" ? "Random" : rotationStrategy === "round-robin" ? "Round Robin" : rotationStrategy === "failover" ? "Failover" : "Multiple"; + const strategyLabel = rotationStrategy === "random" ? "Random" : rotationStrategy === "round-robin" ? "Round Robin" : rotationStrategy === "failover" ? "Failover" : rotationStrategy === "smart" ? "Smart" : "Multiple"; return `${selectedProxyIds.length} pools (${strategyLabel})`; } @@ -318,7 +318,7 @@ export default function ConnectionRow({ connection, proxyPools, isOAuth, isFirst Proxy {showProxyDropdown && ( -
+
{/* Rotation Strategy Selector */}
@@ -331,18 +331,20 @@ export default function ConnectionRow({ connection, proxyPools, isOAuth, isFirst + {rotationStrategy !== "none" && (

{rotationStrategy === "random" && "Randomly select proxy on each request"} {rotationStrategy === "round-robin" && "Rotate proxies in order across requests"} {rotationStrategy === "failover" && "Try next proxy on failure"} + {rotationStrategy === "smart" && "Skip pools whose egress IP is blocked for this provider/model (e.g. Freebuff limited-IP). See Proxy Fitness to clear/block."}

)}
{/* Proxy Pool Selection */} -
+
+ )} + +
+
+ +

+ Shows which provider is blocked on which proxy IP (from provider region gates, per-IP + limits, or manual marks). Smart rotation skips these pools while the block is active. +

+ + {/* Filters */} +
+
+ + + + {providerMenuOpen && ( +
+ +
+ {providerOptions.map((provider) => ( + + ))} +
+ )} +
+ +
+ + setSearch(e.target.value)} + placeholder="e.g. 104.28, vercel-relay…" + icon="search" + /> +
+
+ + {/* Records table */} + +
+ + + + + + + + + + + + + + {filteredRecords.length === 0 ? ( + + + + ) : ( + filteredRecords.map((rec) => ( + + + + + + + + + + )) + )} + +
ProviderModelIP / ProxyPoolReasonUntilActions
+ {records.length === 0 + ? "No active blocks. Blocks appear here when a provider region-gates a proxy IP." + : "No blocks match the current filters."} +
+ + block + {rec.provider} + + + {rec.model || all models} + +
+ {maskProxyUrl(rec.proxyUrl)} + {rec.egress?.ip && ( + + egress {rec.egress.ip}{rec.egress.country ? ` · ${rec.egress.country}` : ""} + {rec.egress.isUnstable && ( + = 2 ? `Egress IP changed ${rec.egress.ipCount}× — relay egress is not stable` : "Egress is unstable"} + > + unstable + + )} + + )} +
+
{rec.poolName}{rec.reason}{fmtTime(rec.until)} +
+ +
+
+
+
+ + )} + + setConfirmClearAll(false)} + onConfirm={handleClearAll} + title="Clear all blocks" + message={providerFilter !== "all" + ? `Clear all active blocks for provider "${providerFilter}"? This lets those pools be selected again immediately.` + : "Clear all active proxy blocks? This lets all pools be selected again immediately."} + confirmText={clearingAll ? "Clearing..." : "Clear All"} + cancelText="Cancel" + variant="danger" + /> +
+ ); +} \ No newline at end of file diff --git a/src/app/api/proxy-pools/[id]/fitness/clear/route.js b/src/app/api/proxy-pools/[id]/fitness/clear/route.js new file mode 100644 index 00000000..ce587750 --- /dev/null +++ b/src/app/api/proxy-pools/[id]/fitness/clear/route.js @@ -0,0 +1,25 @@ +import { NextResponse } from "next/server"; +import { clearPoolUnfit } from "open-sse/services/proxyPoolFitness.js"; + +// POST /api/proxy-pools/[id]/fitness/clear +// Body: { scope: "provider::model" } — clears the mark for this pool + scope. +export async function POST(request, { params }) { + try { + const { id } = await params; + let body = {}; + try { + body = await request.json(); + } catch { + body = {}; + } + const scope = typeof body?.scope === "string" ? body.scope.trim() : ""; + if (!id || !scope) { + return NextResponse.json({ error: "pool id and scope are required" }, { status: 400 }); + } + clearPoolUnfit(id, scope); + return NextResponse.json({ ok: true, poolId: id, scope }); + } catch (error) { + console.log("Error clearing pool fitness:", error); + return NextResponse.json({ error: "Failed to clear pool fitness" }, { status: 500 }); + } +} \ No newline at end of file diff --git a/src/app/api/proxy-pools/fitness/clear-all/route.js b/src/app/api/proxy-pools/fitness/clear-all/route.js new file mode 100644 index 00000000..26c29e3a --- /dev/null +++ b/src/app/api/proxy-pools/fitness/clear-all/route.js @@ -0,0 +1,22 @@ +import { NextResponse } from "next/server"; +import { clearAllPoolUnfit } from "open-sse/services/proxyPoolFitness.js"; + +// POST /api/proxy-pools/fitness/clear-all +// Body: { provider?: string } — clears every mark, or only marks scoped to the +// given provider ("provider::*") when provided. +export async function POST(request) { + try { + let body = {}; + try { + body = await request.json(); + } catch { + body = {}; + } + const provider = typeof body?.provider === "string" && body.provider.trim() ? body.provider.trim() : null; + clearAllPoolUnfit(provider); + return NextResponse.json({ ok: true, provider: provider || null }); + } catch (error) { + console.log("Error clearing proxy fitness:", error); + return NextResponse.json({ error: "Failed to clear proxy fitness" }, { status: 500 }); + } +} \ No newline at end of file diff --git a/src/app/api/proxy-pools/fitness/route.js b/src/app/api/proxy-pools/fitness/route.js new file mode 100644 index 00000000..6929b8d3 --- /dev/null +++ b/src/app/api/proxy-pools/fitness/route.js @@ -0,0 +1,13 @@ +import { NextResponse } from "next/server"; +import { poolFitnessSnapshot } from "open-sse/services/proxyPoolFitness.js"; + +// GET /api/proxy-pools/fitness — in-memory snapshot of pool fitness marks. +// Returns { pools: { [poolId]: { [scope]: { until, reason } } } }. +export async function GET() { + try { + return NextResponse.json({ pools: poolFitnessSnapshot() }); + } catch (error) { + console.log("Error reading proxy fitness:", error); + return NextResponse.json({ error: "Failed to read proxy fitness" }, { status: 500 }); + } +} \ No newline at end of file diff --git a/src/app/api/proxy-pools/route.js b/src/app/api/proxy-pools/route.js index 3dcfa06f..3d05aa3e 100644 --- a/src/app/api/proxy-pools/route.js +++ b/src/app/api/proxy-pools/route.js @@ -1,5 +1,6 @@ import { NextResponse } from "next/server"; import { createProxyPool, getProviderConnections, getProxyPools } from "@/models"; +import { getPoolGeo } from "open-sse/services/poolGeo.js"; function toBoolean(value) { if (value === "true") return true; @@ -65,6 +66,8 @@ export async function GET(request) { const enrichedProxyPools = proxyPools.map((pool) => ({ ...pool, boundConnectionCount: usageMap.get(pool.id) || 0, + // Egress geo from the background probe cache (null until first probe). + egress: getPoolGeo(pool.id) || null, })); return NextResponse.json({ proxyPools: enrichedProxyPools }); diff --git a/src/lib/network/connectionProxy.js b/src/lib/network/connectionProxy.js index 135859f7..af819be3 100644 --- a/src/lib/network/connectionProxy.js +++ b/src/lib/network/connectionProxy.js @@ -1,4 +1,5 @@ import { getProxyPoolById } from "@/models"; +import { fitPoolIds } from "open-sse/services/proxyPoolFitness.js"; // Safely normalize any value into a trimmed string. function normalizeString(value) { @@ -13,24 +14,41 @@ const rotateState = new Map(); // providerId → { index } * Pick one proxy pool ID from a list based on strategy. * round-robin: cycle sequentially (in-memory, resets on restart) * random: uniform random pick + * smart: region-aware — skip pools unfit for `scope`, round-robin on the fit subset * none/single: return first entry + * @param {string[]} poolIds + * @param {string} strategy + * @param {string} providerId + * @param {{ scope?: string, excludeIds?: string[] }} [opts] */ -export function pickProxyPoolId(poolIds, strategy, providerId) { +export function pickProxyPoolId(poolIds, strategy, providerId, opts = {}) { if (!poolIds || poolIds.length === 0) return null; - if (poolIds.length === 1) return poolIds[0]; + const { scope = null, excludeIds = [] } = opts || {}; - if (strategy === "round-robin") { + let eligible = poolIds.filter((id) => !(excludeIds || []).includes(id)); + // Region-aware filtering is opt-in via the "smart" strategy. + if (strategy === "smart" && scope) eligible = fitPoolIds(eligible, scope); + + if (eligible.length === 0) { + // Every candidate is unfit/excluded — fall back to the first non-excluded + // pool rather than deadlocking; its egress may have recovered. + eligible = poolIds.filter((id) => !(excludeIds || []).includes(id)); + if (eligible.length === 0) return null; + } + if (eligible.length === 1) return eligible[0]; + + if (strategy === "round-robin" || strategy === "smart") { const state = rotateState.get(providerId) || { index: -1 }; - state.index = (state.index + 1) % poolIds.length; + state.index = (state.index + 1) % eligible.length; rotateState.set(providerId, state); - return poolIds[state.index]; + return eligible[state.index]; } if (strategy === "random") { - return poolIds[Math.floor(Math.random() * poolIds.length)]; + return eligible[Math.floor(Math.random() * eligible.length)]; } - return poolIds[0]; // "none" or unknown + return eligible[0]; // "none" or unknown } /** @@ -66,7 +84,8 @@ function normalizeLegacyProxy(providerSpecificData = {}) { */ export async function resolveConnectionProxyConfig( providerSpecificData = {}, - connectionId = null + connectionId = null, + excludePoolIds = null ) { try { // Handle new multi-proxy format @@ -85,7 +104,8 @@ export async function resolveConnectionProxyConfig( * ----------------------------- */ if (proxyPoolIds.length > 0) { - const selectedPoolId = pickProxyPoolId(proxyPoolIds, proxyRotationStrategy, connectionId); + const scope = providerSpecificData?.proxyPoolScope || null; + const selectedPoolId = pickProxyPoolId(proxyPoolIds, proxyRotationStrategy, connectionId, { scope, excludeIds: excludePoolIds }); if (selectedPoolId) { const proxyPool = await getProxyPoolById(selectedPoolId); diff --git a/src/shared/components/NoAuthProxyCard.js b/src/shared/components/NoAuthProxyCard.js index 6229e60c..5a823cfc 100644 --- a/src/shared/components/NoAuthProxyCard.js +++ b/src/shared/components/NoAuthProxyCard.js @@ -11,6 +11,7 @@ const STRATEGIES = [ { value: "none", label: "None (single pool)" }, { value: "round-robin", label: "Round-robin" }, { value: "random", label: "Random" }, + { value: "smart", label: "Smart" }, ]; export default function NoAuthProxyCard({ providerId }) { @@ -121,7 +122,9 @@ export default function NoAuthProxyCard({ providerId }) { : isRotation ? rotateStrategy === "round-robin" ? `Rotating through all ${proxyPools.length} active pools in order. State is in-memory (resets on restart).` - : `Picking a random pool from ${proxyPools.length} active pools each request.` + : rotateStrategy === "smart" + ? `Skip pools whose egress IP is blocked for this provider (e.g. per-IP limits). See Proxy Fitness to clear/block.` + : `Picking a random pool from ${proxyPools.length} active pools each request.` : `Uses the selected pool above. Set to Round-robin or Random to rotate across all active pools.`}

diff --git a/src/shared/components/Sidebar.js b/src/shared/components/Sidebar.js index 812e319e..2f329da2 100644 --- a/src/shared/components/Sidebar.js +++ b/src/shared/components/Sidebar.js @@ -25,6 +25,7 @@ const navItems = [ { href: "/dashboard/usage", label: "Usage", icon: "bar_chart" }, { href: "/dashboard/quota", label: "Quota Tracker", icon: "data_usage" }, { href: "/dashboard/token-saver", label: "Token Saver", icon: "savings" }, + { href: "/dashboard/proxy-fitness", label: "Proxy Fitness", icon: "network_check" }, // { href: "/dashboard/pxpipe", label: "PXPIPE", icon: "image" }, { href: "/dashboard/cli-tools", label: "CLI Tools", icon: "terminal" }, ]; diff --git a/src/sse/handlers/chat.js b/src/sse/handlers/chat.js index 88353c50..e8231160 100644 --- a/src/sse/handlers/chat.js +++ b/src/sse/handlers/chat.js @@ -17,6 +17,7 @@ import { errorResponse, unavailableResponse } from "open-sse/utils/error.js"; import { handleComboChat, handleFusionChat, detectRequiredCapabilities } from "open-sse/services/combo.js"; import { augmentModelsWithCapacityAdapter, withCapacityAdapterStripping, getActiveAdapterStrategy } from "open-sse/services/capacityAdapter.js"; import { handleBypassRequest } from "open-sse/utils/bypassHandler.js"; +import { resolveConnectionProxyConfig } from "@/lib/network/connectionProxy"; import { HTTP_STATUS } from "open-sse/config/runtimeConfig.js"; import { detectFormatByEndpoint } from "open-sse/translator/formats.js"; import * as log from "../utils/logger.js"; @@ -284,6 +285,23 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re pxpipeTransform: chatSettings.pxpipeEnabled ? await getPxpipeTransform() : null, onPxpipeEvent: appendPxpipeEvent, providerThinking, + // Pool-scoped failure recovery: re-resolve proxy config excluding the + // failed pool so the request retries via another pool, not a dead end. + resolveProxyConfig: async (creds, excludePoolIds = []) => { + const psd = { ...(creds?.providerSpecificData || {}) }; + if (psd.proxyPoolIds?.length) psd.proxyPoolScope = `${provider}::${model}`; + const resolved = await resolveConnectionProxyConfig(psd, creds?.connectionId || creds?.id, excludePoolIds); + if (!resolved?.proxyPoolId) return null; + return { + connectionProxyEnabled: resolved.connectionProxyEnabled, + connectionProxyUrl: resolved.connectionProxyUrl, + connectionNoProxy: resolved.connectionNoProxy, + connectionProxyPoolId: resolved.proxyPoolId || null, + vercelRelayUrl: resolved.vercelRelayUrl || "", + proxyPoolId: resolved.proxyPoolId || null, + strictProxy: resolved.strictProxy === true, + }; + }, // Detect source format by endpoint + body sourceFormatOverride: request?.url ? detectFormatByEndpoint(new URL(request.url).pathname, body) : null, onCredentialsRefreshed: async (newCreds) => { diff --git a/src/sse/services/auth.js b/src/sse/services/auth.js index 8eb0545c..c2e80572 100644 --- a/src/sse/services/auth.js +++ b/src/sse/services/auth.js @@ -47,10 +47,16 @@ export async function getProviderCredentials(provider, excludeConnectionIds = nu const override = (settings.providerStrategies || {})[providerId] || {}; const strategy = override.rotateStrategy || "none"; let pickedId = override.proxyPoolId || null; + let poolIds = []; if (strategy !== "none") { const allPools = await getProxyPools({ isActive: true }); - const poolIds = allPools.filter(p => p.proxyUrl).map(p => p.id); - pickedId = pickProxyPoolId(poolIds, strategy, providerId); + poolIds = allPools.filter(p => p.proxyUrl).map(p => p.id); + // Scope region-aware ("smart") filtering to this provider/model so + // pools marked unfit here are skipped. + const scope = `${providerId}::${model || "*"}`; + pickedId = pickProxyPoolId(poolIds, strategy, providerId, { scope }); + } else if (override.proxyPoolId) { + poolIds = [override.proxyPoolId]; } const resolvedProxy = await resolveConnectionProxyConfig({ proxyPoolId: pickedId || "" }); return { @@ -64,6 +70,13 @@ export async function getProviderCredentials(provider, excludeConnectionIds = nu connectionNoProxy: resolvedProxy.connectionNoProxy, connectionProxyPoolId: resolvedProxy.proxyPoolId || null, vercelRelayUrl: resolvedProxy.vercelRelayUrl || "", + proxyPoolId: resolvedProxy.proxyPoolId || null, + strictProxy: resolvedProxy.strictProxy === true, + // Let chatCore's pool-scoped retry rotate across the same candidate + // pool set (excluding the failed pool) instead of reusing it — this + // is what makes per-IP limit retries work for no-auth providers. + proxyPoolIds: poolIds, + proxyRotationStrategy: strategy, }, }; } @@ -172,7 +185,11 @@ export async function getProviderCredentials(provider, excludeConnectionIds = nu connection = availableConnections[0]; } - const resolvedProxy = await resolveConnectionProxyConfig(connection.providerSpecificData || {}, connection.id); + // Scope the region-aware picker to this provider/model (e.g. freebuff::gpt-5.6-luna) + const psdForProxy = connection.providerSpecificData?.proxyPoolIds?.length + ? { ...connection.providerSpecificData, proxyPoolScope: `${providerId}::${model || ""}` } + : connection.providerSpecificData; + const resolvedProxy = await resolveConnectionProxyConfig(psdForProxy || {}, connection.id); return { authType: connection.authType, @@ -193,6 +210,8 @@ export async function getProviderCredentials(provider, excludeConnectionIds = nu connectionNoProxy: resolvedProxy.connectionNoProxy, connectionProxyPoolId: resolvedProxy.proxyPoolId || null, vercelRelayUrl: resolvedProxy.vercelRelayUrl || "", + proxyPoolId: resolvedProxy.proxyPoolId || null, + strictProxy: resolvedProxy.strictProxy === true, }, connectionId: connection.id, // Include current status for optimization check diff --git a/tests/unit/proxy-pool-fitness.test.js b/tests/unit/proxy-pool-fitness.test.js new file mode 100644 index 00000000..aae84a0d --- /dev/null +++ b/tests/unit/proxy-pool-fitness.test.js @@ -0,0 +1,56 @@ +import { describe, it, expect, beforeEach } from "vitest"; +import { + markPoolUnfit, + clearPoolUnfit, + clearAllPoolUnfit, + isPoolFit, + fitPoolIds, + poolFitnessSnapshot, + pruneExpired, + resetPoolFitness, +} from "open-sse/services/proxyPoolFitness.js"; + +describe("proxy pool fitness registry", () => { + beforeEach(() => resetPoolFitness()); + + it("marks a pool unfit for a scope and prunes on expiry", () => { + markPoolUnfit("p1", "freebuff::gpt-5.6-luna", Date.now() + 60_000, "limited_ip"); + expect(isPoolFit("p1", "freebuff::gpt-5.6-luna")).toBe(false); + expect(isPoolFit("p1", "freebuff::other-model")).toBe(true); + expect(isPoolFit("p2", "freebuff::gpt-5.6-luna")).toBe(true); + + markPoolUnfit("p1", "freebuff::gpt-5.6-luna", Date.now() - 1000); // expired + expect(isPoolFit("p1", "freebuff::gpt-5.6-luna")).toBe(true); // pruned on read + }); + + it("provider-wide mark (provider::*) covers any model lookup", () => { + markPoolUnfit("p1", "opencode::*", Date.now() + 60_000, "manual"); + expect(isPoolFit("p1", "opencode::sonnet-4.6")).toBe(false); + expect(isPoolFit("p1", "freebuff::gpt-5.6-luna")).toBe(true); + }); + + it("fitPoolIds filters unfit pools; snapshot drops expired marks", () => { + markPoolUnfit("p1", "fb::m1", Date.now() + 60_000); + markPoolUnfit("p1", "fb::m2", Date.now() - 1000); // expired + expect(fitPoolIds(["p1", "p2"], "fb::m1")).toEqual(["p2"]); + + const snap = poolFitnessSnapshot(); + expect(snap.p1["fb::m1"]).toBeDefined(); + expect(snap.p1["fb::m2"]).toBeUndefined(); + }); + + it("clear per scope, clear-all per provider, clear-all global, pruneExpired", () => { + markPoolUnfit("p1", "freebuff::m1", Date.now() + 60_000); + markPoolUnfit("p1", "kiro::m3", Date.now() + 60_000); + + clearPoolUnfit("p1", "freebuff::m1"); + expect(isPoolFit("p1", "freebuff::m1")).toBe(true); + + clearAllPoolUnfit("kiro"); + expect(poolFitnessSnapshot().p1).toBeUndefined(); + + markPoolUnfit("p2", "x::y", Date.now() - 1000); + expect(pruneExpired()).toBe(1); + expect(poolFitnessSnapshot()).toEqual({}); + }); +});