feat(proxy-fitness): throttle egress geo probe + aggregate failure log
Refuse re-probing failures too fast: 500/429 failures get a 2h backoff, other errors 30m (flapping relays/quota stops hammering ipinfo every pass, passes every 30 min instead of 15). One-line summary per pass (geo N/M · fail: rate×a server×b) replaces the per-pool error wall. New env POOL_GEO_PROBE_DISABLED=1 turns the feature off.
This commit is contained in:
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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 };
|
||||
Reference in New Issue
Block a user