From 7cbf6fd81722d0be503f5301fbb2665304d0fc54 Mon Sep 17 00:00:00 2001 From: MythEclipse Date: Wed, 17 Jun 2026 04:59:16 +0700 Subject: [PATCH] feat: add per-model proxy cooldown for rate-limit handling MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ProxyPool: - Add cooldowns map (host:port::model -> expiry) and cooldownDuration (60s default) - Add markRateLimited(model) — puts current proxy in cooldown for a model - Add isProxyInCooldown, isIndexInCooldown, isCurrentInCooldown checks - Modify rotate(model?) — skip proxies in cooldown for the given model - Add setCooldownDuration(ms) for configuration SessionProxyPool: - Add markRateLimited(sessionId, model) — delegates to underlying pool - Modify rotateNow(sessionId, model?) — skip cooldown proxies - Modify acquire(sessionId, model?) — skip cooldown when picking least-loaded - Modify pickLeastUsedIndex(model?) — skip cooldown proxies in scan fetch-utils.ts: - Add extractModel helper — extracts model from context string - On HTTP 429: call markRateLimited(model) before markFailed - Pass model to rotate() / rotateNow() / acquire() throughout Co-Authored-By: Claude --- src/lib/fetch-utils.ts | 31 +++++++++-- src/lib/proxy-pool.ts | 118 +++++++++++++++++++++++++++++++++++------ 2 files changed, 129 insertions(+), 20 deletions(-) diff --git a/src/lib/fetch-utils.ts b/src/lib/fetch-utils.ts index 543dd83..560be2c 100644 --- a/src/lib/fetch-utils.ts +++ b/src/lib/fetch-utils.ts @@ -89,6 +89,17 @@ function logProxy(method: string, msg: string, extra?: Record): console.log(parts.join(" ")); } +// ─── Model extraction from context ─────────────────────────────────── + +/** + * Extract model identifier from the context string. + * Context format: "provider:modelname" (e.g. "openai:deepseek-v4-flash-free"). + * Returns undefined when no context given (backward-compatible). + */ +function extractModel(context?: string): string | undefined { + return context || undefined; +} + // ─── Error sanitization (prevent leaking upstream details) ────────────── /** @@ -148,7 +159,7 @@ export async function fetchWithRetry( init.proxy = proxyPool.getProxyUrl()!; usedProxy = true; } else { - const next = proxyPool.rotate(); + const next = proxyPool.rotate(extractModel(context)); if (!next) break; init.proxy = proxyPool.getProxyUrl()!; usedProxy = true; @@ -176,14 +187,19 @@ export async function fetchWithRetry( lastError = new Error(`Upstream returned ${response.status}`); logProxy("fetchWithRetry", `non-2xx attempt=${attempt + 1} status=${response.status}`, { context }); if (usedProxy && proxyPool && proxyPool.size > 0 && init.proxy) { + const model = extractModel(context); + // If rate-limited, put proxy in per-model cooldown + if (response.status === 429 && model) { + proxyPool.markRateLimited(model); + } proxyPool.markFailed(); - proxyPool.rotate(); + proxyPool.rotate(model); } } catch (err) { lastError = err; if (usedProxy && proxyPool && proxyPool.size > 0 && init.proxy) { proxyPool.markFailed(); - proxyPool.rotate(); + proxyPool.rotate(extractModel(context)); } // Log the actual error so operators can diagnose @@ -324,7 +340,12 @@ export async function fetchWithSessionRetry( // Non-2xx — rotate proxy immediately for the next attempt lastError = new Error(`Upstream returned ${response.status}`); lastResponse = response; - const rotated = sessionPool.rotateNow(sessionId); + const model = extractModel(context); + // If rate-limited, put proxy in per-model cooldown + if (response.status === 429 && model) { + sessionPool.markRateLimited(sessionId, model); + } + const rotated = sessionPool.rotateNow(sessionId, model); logProxy("fetchWithSessionRetry", `non-2xx attempt=${attempt + 1} status=${response.status} rotated=${rotated}`, { context, sessionId: sessionId.slice(0, 8), @@ -338,7 +359,7 @@ export async function fetchWithSessionRetry( ); } catch (err) { lastError = err; - const rotated = sessionPool.rotateNow(sessionId); + const rotated = sessionPool.rotateNow(sessionId, extractModel(context)); const errMsg = err instanceof Error ? err.message : String(err); logProxy("fetchWithSessionRetry", `failed attempt=${attempt + 1} err=${errMsg} rotated=${rotated}`, { context, diff --git a/src/lib/proxy-pool.ts b/src/lib/proxy-pool.ts index 0d8df6c..42168ee 100644 --- a/src/lib/proxy-pool.ts +++ b/src/lib/proxy-pool.ts @@ -38,6 +38,9 @@ export class ProxyPool { private failureThreshold = 3; /** host:port -> consecutive failure count */ private failures = new Map(); + /** host:port::model -> expiry epoch ms */ + private cooldowns = new Map(); + private cooldownDuration = 60000; // default 60s // -- Load -------------------------------------------------------------------- @@ -138,29 +141,30 @@ export class ProxyPool { /** * Advance to the next proxy (round-robin, wraps around). - * Skips proxies that have exceeded the failure threshold. + * Skips proxies that have exceeded the failure threshold or are in + * cooldown for the given model. * Returns the new current proxy or `null` if the pool is empty or * all proxies are failed. */ - rotate(): ProxyEntry | null { + rotate(model?: string): ProxyEntry | null { if (this.proxies.length === 0) return null; const oldIndex = this.currentIndex; const startIndex = this.currentIndex; - // Keep advancing until we find a non-failed proxy or loop back + // Keep advancing until we find a non-failed, non-cooldown proxy or loop back let checked = 0; do { this.currentIndex = (this.currentIndex + 1) % this.proxies.length; checked++; - if (!this.isFailed()) { + if (!this.isFailed() && !this.isCurrentInCooldown(model)) { const entry = this.proxies[this.currentIndex] ?? null; - logPool(`rotate ${oldIndex} -> ${this.currentIndex} (skipped ${checked - 1} failed)`); + logPool(`rotate ${oldIndex} -> ${this.currentIndex} (skipped ${checked - 1} failed/in-cooldown)`); return entry; } } while (this.currentIndex !== startIndex && checked <= this.proxies.length); - // All proxies failed — stay on current but log it - logPool(`rotate ${oldIndex} -> ${this.currentIndex} (all proxies failed)`); + // All proxies failed or in cooldown — stay on current but log it + logPool(`rotate ${oldIndex} -> ${this.currentIndex} (all proxies failed or in cooldown)`); return this.proxies[this.currentIndex] ?? null; } @@ -211,6 +215,69 @@ export class ProxyPool { setFailureThreshold(n: number): void { this.failureThreshold = n; } + + // -- Cooldown (per-model rate-limit) ----------------------------------------- + + /** + * Set the cooldown duration in milliseconds (default 60000). + * When a proxy gets rate-limited (429) for a specific model, it will + * be skipped for that model for this duration. + */ + setCooldownDuration(ms: number): void { + this.cooldownDuration = ms; + } + + /** Internal cooldown key format: host:port::model */ + private cooldownKey(host: string, port: number, model: string): string { + return `${host}:${port}::${model}`; + } + + /** + * Mark the **current** proxy as rate-limited for a specific model. + * The proxy enters a cooldown period during which it will be skipped + * for this model but remains available for other models. + */ + markRateLimited(model: string): void { + const entry = this.getCurrent(); + if (!entry) return; + const key = this.cooldownKey(entry.host, entry.port, model); + const expiry = Date.now() + this.cooldownDuration; + this.cooldowns.set(key, expiry); + logPool(`markRateLimited key=${key} expiry=${expiry} duration=${this.cooldownDuration}ms`); + } + + /** + * Check if a specific proxy (host:port) is in cooldown for a model. + * Returns true if the proxy is cooling down for that model. + */ + isProxyInCooldown(host: string, port: number, model: string): boolean { + const key = this.cooldownKey(host, port, model); + const expiry = this.cooldowns.get(key); + if (!expiry) return false; + if (Date.now() > expiry) { + this.cooldowns.delete(key); // lazy cleanup + return false; + } + return true; + } + + /** + * Check if a proxy at a given pool index is in cooldown for a model. + * Returns true if the proxy doesn't exist or is in cooldown. + */ + isIndexInCooldown(index: number, model: string): boolean { + const entry = this.proxies[index]; + if (!entry) return true; + return this.isProxyInCooldown(entry.host, entry.port, model); + } + + /** Check if the current proxy is in cooldown for the given model. */ + private isCurrentInCooldown(model?: string): boolean { + if (!model) return false; + const entry = this.getCurrent(); + if (!entry) return true; + return this.isProxyInCooldown(entry.host, entry.port, model); + } } // --- SessionProxyPool --------------------------------------------------------- @@ -269,7 +336,7 @@ export class SessionProxyPool { * If the session already has a proxy, returns the same one (resume). * Otherwise picks the least-loaded proxy. */ - acquire(sessionId: string): string | null { + acquire(sessionId: string, model?: string): string | null { if (this.pool.size === 0) return null; const existing = this.sessions.get(sessionId); @@ -278,7 +345,7 @@ export class SessionProxyPool { return this.formatProxyUrlAtIndex(existing.proxyIndex); } - const index = this.pickLeastUsedIndex(); + const index = this.pickLeastUsedIndex(model); if (index === -1) return null; this.sessions.set(sessionId, { proxyIndex: index, failures: 0 }); @@ -382,7 +449,7 @@ export class SessionProxyPool { * * @returns true if the session was moved to a different proxy. */ - rotateNow(sessionId: string): boolean { + rotateNow(sessionId: string, model?: string): boolean { const info = this.sessions.get(sessionId); if (!info) return false; @@ -396,13 +463,14 @@ export class SessionProxyPool { if (usedBy.size === 0) this.proxyUsage.delete(oldIndex); } - // Find the least-loaded proxy that is NOT the current one. - // Start scanning from (oldIndex + 1) so we don't immediately - // bounce back to index 0 when all usage counts are equal. + // Find the least-loaded proxy that is NOT the current one and NOT + // in cooldown for this model. Start scanning from (oldIndex + 1) + // so we don't immediately bounce back to index 0. let bestIndex = -1; let bestCount = Infinity; for (let step = 1; step <= this.pool.size; step++) { const i = (oldIndex + step) % this.pool.size; + if (model && (this.pool as any).isIndexInCooldown(i, model)) continue; const count = this.proxyUsage.get(i)?.size ?? 0; if (count < bestCount) { bestCount = count; @@ -438,6 +506,20 @@ export class SessionProxyPool { logPool(`markSuccess session=${sessionId.slice(0, 8)} proxyIndex=${info.proxyIndex}`); } + /** + * Mark this session's proxy as rate-limited for a specific model. + * Delegates to the underlying ProxyPool's cooldown so the proxy is + * skipped for this model on subsequent requests (different sessions). + */ + markRateLimited(sessionId: string, model: string): void { + const info = this.sessions.get(sessionId); + if (!info) return; + const savedIdx = (this.pool as any).currentIndex as number; + (this.pool as any).currentIndex = info.proxyIndex; + this.pool.markRateLimited(model); + (this.pool as any).currentIndex = savedIdx; + } + // -- Internals ---------------------------------------------------------------- /** Get the ProxyEntry at a given index. Forward reference to local type. */ @@ -456,13 +538,14 @@ export class SessionProxyPool { } /** Return the index of the proxy with the fewest active sessions, or -1. */ - private pickLeastUsedIndex(): number { + private pickLeastUsedIndex(model?: string): number { if (this.pool.size === 0) return -1; - let bestIndex = 0; + let bestIndex = -1; let bestCount = Infinity; for (let i = 0; i < this.pool.size; i++) { + if (model && (this.pool as any).isIndexInCooldown(i, model)) continue; const count = this.proxyUsage.get(i)?.size ?? 0; logPool(`pickLeastUsed proxy[${i}] count=${count}`); if (count < bestCount) { @@ -471,6 +554,11 @@ export class SessionProxyPool { } } + if (bestIndex === -1) { + logPool(`pickLeastUsed all proxies in cooldown`); + return 0; // fallback to first + } + logPool(`pickLeastUsed selected index=${bestIndex} count=${bestCount}`); return bestIndex; }