diff --git a/open-sse/services/poolGeo.js b/open-sse/services/poolGeo.js index 319934d3..862a013f 100644 --- a/open-sse/services/poolGeo.js +++ b/open-sse/services/poolGeo.js @@ -21,19 +21,6 @@ export const POOL_GEO_TTL_MS = 60 * 60 * 1000; // How many past egress IPs to remember for flapping detection. export const POOL_GEO_IP_HISTORY_MAX = 8; -// Debounce probe-failure logs: one line per failing pool per hour — the probe -// runs every 30 min, so a broken relay must not spam the log every pass. -const probeFailLogs = new Map(); // proxyUrl -> lastLoggedAt -const PROBE_FAIL_LOG_INTERVAL_MS = 60 * 60 * 1000; - -function logProbeFailure(kind, proxyUrl, detail) { - const key = String(proxyUrl || ""); - const now = Date.now(); - if (probeFailLogs.has(key) && now - probeFailLogs.get(key) < PROBE_FAIL_LOG_INTERVAL_MS) return; - probeFailLogs.set(key, now); - console.log(`[GeoProbe] ${kind} ${key.slice(0, 60)}${detail ? ` | ${detail.slice(0, 120)}` : ""}`); -} - // Test helper: drop all cached geo (module state is globalThis-backed). export function resetPoolGeo() { geoCache.clear(); @@ -96,10 +83,12 @@ export function pruneStaleGeo(now = Date.now()) { } // Probe the egress geo of one pool via ipinfo through the pool. Fail-open: -// returns null on any error/timeout. `pool` shape: { proxyUrl, type }. +// returns { ok:true, geo } on success, { ok:false, error } otherwise with +// `error` one of "rate-limit" | "server" | "no-ip" | "network" | "timeout". +// `pool` shape: { proxyUrl, type }. export async function probePoolGeo(pool, timeoutMs = 15000) { const proxyUrl = pool?.proxyUrl; - if (!proxyUrl) return null; + if (!proxyUrl) return { ok: false, error: "network" }; const { proxyAwareFetch } = await import("../utils/proxyFetch.js"); const isRelay = ["vercel", "cloudflare", "deno"].includes(pool?.type); const proxyOptions = isRelay @@ -112,25 +101,29 @@ export async function probePoolGeo(pool, timeoutMs = 15000) { const res = await proxyAwareFetch("https://ipinfo.io/json", { signal: ctrl.signal }, proxyOptions); if (!res.ok) { const txt = await res.text().catch(() => ""); - logProbeFailure(`http ${res.status}`, proxyUrl, txt); - return null; + const error = res.status === 429 || res.status === 403 ? "rate-limit" + : res.status >= 500 ? "server" + : "network"; + return { ok: false, error, detail: `${res.status} ${txt.slice(0, 60)}` }; } const data = await res.json().catch(() => null); if (!data?.ip) { - logProbeFailure("no-ip", proxyUrl, ""); - return null; + return { ok: false, error: "no-ip" }; } return { - ip: data.ip || "", - country: data.country || "", - region: data.region || "", - city: data.city || "", - org: data.org || "", - isDatacenter: /(cloudflare|vercel|amazon|aws|google|microsoft|azure|digitalocean|hetzner|ovh|contabo|leaseweb)/i.test(String(data.org || "")), + ok: true, + geo: { + ip: data.ip || "", + country: data.country || "", + region: data.region || "", + city: data.city || "", + org: data.org || "", + isDatacenter: /(cloudflare|vercel|amazon|aws|google|microsoft|azure|digitalocean|hetzner|ovh|contabo|leaseweb)/i.test(String(data.org || "")), + }, }; } catch (error) { - logProbeFailure("fail", proxyUrl, `${error?.name}: ${error?.message}`); - return null; + const timedOut = ctrl.signal?.aborted && error?.name === "AbortError"; + return { ok: false, error: timedOut ? "timeout" : "network", detail: `${error?.name}: ${error?.message}` }; } finally { clearTimeout(timer); } diff --git a/src/lib/network/poolEgressProbe.js b/src/lib/network/poolEgressProbe.js index cd76ec15..43a84eb6 100644 --- a/src/lib/network/poolEgressProbe.js +++ b/src/lib/network/poolEgressProbe.js @@ -2,6 +2,10 @@ // Fitness / Proxy Pools UI can show each pool's egress IP + country, and // future provider region policies can pre-mark pools unfit. // Fail-open everywhere; never blocks startup or requests. +// +// Rate safety: each probe hits ipinfo.io through the pool; quotas are small, +// so failing pools get a negative backoff instead of being re-probed every +// pass, and per-pass output is a single aggregated summary line. import { getProxyPools } from "@/models"; import { probePoolGeo, setPoolGeo, getPoolGeo } from "open-sse/services/poolGeo.js"; @@ -9,41 +13,68 @@ import { isNonServerRuntime } from "@/sse/services/backgroundTokenRefresh.js"; const PROBE_INTERVAL_MS = 30 * 60 * 1000; const INITIAL_DELAY_MS = 15 * 1000; -const CONCURRENCY = 4; +const CONCURRENCY = 3; // Re-probe a pool when its last sample is older than this — multiple samples // per pool are what let us flag flapping (changing egress) relays. -const GEO_REPROBE_MS = 15 * 60 * 1000; +const GEO_REPROBE_MS = 30 * 60 * 1000; + +// Backoff windows per failure family: server/rate problems are likely to +// persist (broken relay, exhausted quota), network blips recover sooner. +const BACKOFF = { + "rate-limit": 2 * 60 * 60 * 1000, + server: 2 * 60 * 60 * 1000, + network: 30 * 60 * 1000, + timeout: 30 * 60 * 1000, + "no-ip": 30 * 60 * 1000, +}; let started = false; let intervalHandle = null; let initialTimeoutHandle = null; let probing = false; +const probeBackoff = new Map(); // poolId -> retryAfter (ms epoch) -// probe with bounded concurrency; re-probes pools whose sample is stale enough. +function isTruthyEnv(v) { + if (v == null || v === "") return false; + return ["1", "true", "yes", "on"].includes(String(v).trim().toLowerCase()); +} + +// probe with bounded concurrency; skips pools in backoff or with fresh samples. async function probeAll() { if (probing) return; probing = true; try { const pools = await getProxyPools({ isActive: true }); const now = Date.now(); - const targets = (pools || []).filter((p) => { - if (!p?.proxyUrl) return false; + const active = (pools || []).filter((p) => !!p?.proxyUrl); + const targets = active.filter((p) => { + const until = probeBackoff.get(p.id); + if (until && until > now) return false; // waiting out a failure const geo = getPoolGeo(p.id); return !geo || now - geo.ts >= GEO_REPROBE_MS; }); if (targets.length === 0) return; - console.log(`[PoolEgressProbe] probing ${targets.length}/${(pools || []).length} active pools`); + + const failTally = {}; let next = 0; const workers = Array.from({ length: Math.min(CONCURRENCY, Math.max(targets.length, 1)) }, async () => { while (next < targets.length) { const pool = targets[next++]; - const geo = await probePoolGeo(pool); - if (geo) setPoolGeo(pool.id, geo); + const res = await probePoolGeo(pool); + if (res.ok) { + setPoolGeo(pool.id, res.geo); + if (probeBackoff.has(pool.id)) probeBackoff.delete(pool.id); + } else { + failTally[res.error] = (failTally[res.error] || 0) + 1; + probeBackoff.set(pool.id, now + (BACKOFF[res.error] || BACKOFF.network)); + } } }); await Promise.allSettled(workers); - const filled = (pools || []).filter((p) => getPoolGeo(p.id)).length; - console.log(`[PoolEgressProbe] pass done — geo cached for ${filled}/${(pools || []).length} pools`); + + const filled = active.filter((p) => getPoolGeo(p.id)).length; + const failSummary = Object.entries(failTally).map(([k, n]) => `${k}×${n}`).join(", ") || "none"; + console.log(`[PoolEgressProbe] geo ${filled}/${active.length} · fail: ${failSummary}`); } catch (e) { console.log(`[PoolEgressProbe] pass failed: ${e?.message || e}`); } finally { @@ -54,6 +85,7 @@ async function probeAll() { export function startPoolEgressProbe({ intervalMs } = {}) { if (started) return false; if (isNonServerRuntime()) return false; + if (isTruthyEnv(process.env.POOL_GEO_PROBE_DISABLED)) return false; started = true; const period = Number.isFinite(intervalMs) && intervalMs > 0 ? intervalMs : PROBE_INTERVAL_MS; console.log("[PoolEgressProbe] Scheduler started", { intervalMs: period, initialDelayMs: INITIAL_DELAY_MS }); @@ -72,4 +104,4 @@ export function stopPoolEgressProbe() { started = false; } -export const __test__ = { probeAll }; +export const __test__ = { probeAll }; \ No newline at end of file