feat(proxy-fitness): egress geo probe + unstable detection + state sweeper
Background probe fetches ipinfo through each pool itself (provider-agnostic transport), fills an egress IP/country cache (TTL 1h, 8-IP history) and flags flapping relays as unstable. Periodic state sweeper (10 min) prunes expired fitness marks, stale geo/ip-history and Freebuff session/cooldown state; schedulers skip non-server runtimes. Unit tests for the geo cache.
This commit is contained in:
@@ -0,0 +1,137 @@
|
||||
// Pool egress geo — engine layer, shared by the dashboard UI (Proxy Fitness /
|
||||
// Proxy Pools egress column) and any provider region policy that wants to
|
||||
// pre-mark pools unfit by egress region.
|
||||
//
|
||||
// The probe transport is provider-agnostic: it fetches ipinfo.io THROUGH the
|
||||
// pool itself (relay headers for vercel/cloudflare/deno, proxy URL for
|
||||
// socks/http), so it works for every pool type and every provider.
|
||||
|
||||
// The probe transport is provider-agnostic: it fetches ipinfo.io THROUGH the
|
||||
// pool itself (relay headers for vercel/cloudflare/deno, proxy URL for
|
||||
// socks/http), so it works for every pool type and every provider.
|
||||
//
|
||||
// State lives on globalThis (same reason as proxyPoolFitness): the background
|
||||
// probe and the /api/proxy-pools reader must share ONE cache across Next dev
|
||||
// bundles.
|
||||
|
||||
const GEO_STATE_KEY = "__9routerPoolGeo__";
|
||||
const geoCache = (globalThis[GEO_STATE_KEY] ??= new Map()); // poolId -> { ip, country, ..., ts, ipHistory }
|
||||
|
||||
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();
|
||||
}
|
||||
|
||||
// Attach stability classification: >=2 distinct egress IPs observed = flapping
|
||||
// (typical for serverless relays — Vercel/Cloudflare egress varies per colo).
|
||||
function withStability(entry) {
|
||||
const ips = new Set([entry?.ip, ...(entry?.ipHistory || []).map((h) => h.ip)]);
|
||||
ips.delete("");
|
||||
const ipCount = ips.size;
|
||||
return { ...entry, ipCount, isUnstable: ipCount >= 2 };
|
||||
}
|
||||
|
||||
export function getPoolGeo(poolId) {
|
||||
const entry = geoCache.get(poolId);
|
||||
if (!entry) return null;
|
||||
if (entry.ts + POOL_GEO_TTL_MS < Date.now()) {
|
||||
geoCache.delete(poolId);
|
||||
return null;
|
||||
}
|
||||
return withStability(entry);
|
||||
}
|
||||
|
||||
export function setPoolGeo(poolId, geo) {
|
||||
if (!poolId || !geo?.ip) return;
|
||||
const prev = geoCache.get(poolId);
|
||||
const ipHistory = prev?.ipHistory ? [...prev.ipHistory] : [];
|
||||
if (prev?.ip && prev.ip !== geo.ip) {
|
||||
// Record the IP we are leaving — the history tracks past egress IPs.
|
||||
ipHistory.push({ ip: prev.ip, ts: Date.now() });
|
||||
if (ipHistory.length > POOL_GEO_IP_HISTORY_MAX) ipHistory.shift();
|
||||
}
|
||||
geoCache.set(poolId, { ...geo, ts: Date.now(), ipHistory });
|
||||
}
|
||||
|
||||
export function poolGeoSnapshot(now = Date.now()) {
|
||||
const out = {};
|
||||
for (const [poolId, entry] of geoCache) {
|
||||
if (entry.ts + POOL_GEO_TTL_MS <= now) {
|
||||
geoCache.delete(poolId);
|
||||
continue;
|
||||
}
|
||||
out[poolId] = withStability(entry);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// Sweep geo entries past their TTL (ipHistory rides along with the entry).
|
||||
// Returns how many entries were removed.
|
||||
export function pruneStaleGeo(now = Date.now()) {
|
||||
let removed = 0;
|
||||
for (const [poolId, entry] of geoCache) {
|
||||
if (entry.ts + POOL_GEO_TTL_MS <= now) {
|
||||
geoCache.delete(poolId);
|
||||
removed += 1;
|
||||
}
|
||||
}
|
||||
return removed;
|
||||
}
|
||||
|
||||
// Probe the egress geo of one pool via ipinfo through the pool. Fail-open:
|
||||
// returns null on any error/timeout. `pool` shape: { proxyUrl, type }.
|
||||
export async function probePoolGeo(pool, timeoutMs = 15000) {
|
||||
const proxyUrl = pool?.proxyUrl;
|
||||
if (!proxyUrl) return null;
|
||||
const { proxyAwareFetch } = await import("../utils/proxyFetch.js");
|
||||
const isRelay = ["vercel", "cloudflare", "deno"].includes(pool?.type);
|
||||
const proxyOptions = isRelay
|
||||
? { vercelRelayUrl: proxyUrl }
|
||||
: { connectionProxyEnabled: true, connectionProxyUrl: proxyUrl };
|
||||
|
||||
const ctrl = new AbortController();
|
||||
const timer = setTimeout(() => ctrl.abort(new Error("geo probe timeout")), timeoutMs);
|
||||
try {
|
||||
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 data = await res.json().catch(() => null);
|
||||
if (!data?.ip) {
|
||||
logProbeFailure("no-ip", proxyUrl, "");
|
||||
return null;
|
||||
}
|
||||
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 || "")),
|
||||
};
|
||||
} catch (error) {
|
||||
logProbeFailure("fail", proxyUrl, `${error?.name}: ${error?.message}`);
|
||||
return null;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
// Background pool egress geo probe — fills the poolGeo cache so the Proxy
|
||||
// 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.
|
||||
|
||||
import { getProxyPools } from "@/models";
|
||||
import { probePoolGeo, setPoolGeo, getPoolGeo } from "open-sse/services/poolGeo.js";
|
||||
import { isNonServerRuntime } from "@/sse/services/backgroundTokenRefresh.js";
|
||||
|
||||
const PROBE_INTERVAL_MS = 30 * 60 * 1000;
|
||||
const INITIAL_DELAY_MS = 15 * 1000;
|
||||
const CONCURRENCY = 4;
|
||||
// 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;
|
||||
|
||||
let started = false;
|
||||
let intervalHandle = null;
|
||||
let initialTimeoutHandle = null;
|
||||
let probing = false;
|
||||
|
||||
// probe with bounded concurrency; re-probes pools whose sample is stale enough.
|
||||
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 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`);
|
||||
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);
|
||||
}
|
||||
});
|
||||
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`);
|
||||
} catch (e) {
|
||||
console.log(`[PoolEgressProbe] pass failed: ${e?.message || e}`);
|
||||
} finally {
|
||||
probing = false;
|
||||
}
|
||||
}
|
||||
|
||||
export function startPoolEgressProbe({ intervalMs } = {}) {
|
||||
if (started) return false;
|
||||
if (isNonServerRuntime()) 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 });
|
||||
initialTimeoutHandle = setTimeout(() => { probeAll().catch(() => {}); }, INITIAL_DELAY_MS);
|
||||
if (initialTimeoutHandle.unref) initialTimeoutHandle.unref();
|
||||
intervalHandle = setInterval(() => { probeAll().catch(() => {}); }, period);
|
||||
if (intervalHandle.unref) intervalHandle.unref();
|
||||
return true;
|
||||
}
|
||||
|
||||
export function stopPoolEgressProbe() {
|
||||
if (initialTimeoutHandle) clearTimeout(initialTimeoutHandle);
|
||||
if (intervalHandle) clearInterval(intervalHandle);
|
||||
initialTimeoutHandle = null;
|
||||
intervalHandle = null;
|
||||
started = false;
|
||||
}
|
||||
|
||||
export const __test__ = { probeAll };
|
||||
@@ -0,0 +1,47 @@
|
||||
// Periodic in-memory state sweeper — a cron-like job that prunes expired data
|
||||
// so here long-running servers never accumulate stale entries:
|
||||
// • pool fitness marks (proxyPoolFitness.pruneExpired)
|
||||
// • pool egress geo cache incl. ipHistory (poolGeo.pruneStaleGeo)
|
||||
// • freebuff session cache + cooldowns (freebuff.pruneSessionState)
|
||||
// Fail-open everywhere; never blocks startup or requests.
|
||||
|
||||
import { pruneExpired } from "open-sse/services/proxyPoolFitness.js";
|
||||
import { pruneStaleGeo } from "open-sse/services/poolGeo.js";
|
||||
import { isNonServerRuntime } from "@/sse/services/backgroundTokenRefresh.js";
|
||||
|
||||
const SWEEP_INTERVAL_MS = 10 * 60 * 1000;
|
||||
|
||||
let started = false;
|
||||
let handle = null;
|
||||
|
||||
async function sweep() {
|
||||
try {
|
||||
const fitness = pruneExpired();
|
||||
const geo = pruneStaleGeo();
|
||||
const { pruneSessionState } = await import("open-sse/executors/freebuff.js");
|
||||
const sessions = pruneSessionState();
|
||||
if (fitness || geo || sessions) {
|
||||
console.log(`[StateSweeper] pruned ${fitness} fitness, ${geo} geo, ${sessions} session/cooldown entries`);
|
||||
}
|
||||
} catch {
|
||||
// fail-open: next tick retries
|
||||
}
|
||||
}
|
||||
|
||||
export function startStateSweeper({ intervalMs } = {}) {
|
||||
if (started) return false;
|
||||
if (isNonServerRuntime()) return false;
|
||||
started = true;
|
||||
const period = Number.isFinite(intervalMs) && intervalMs > 0 ? intervalMs : SWEEP_INTERVAL_MS;
|
||||
handle = setInterval(() => { sweep().catch(() => {}); }, period);
|
||||
if (handle.unref) handle.unref();
|
||||
return true;
|
||||
}
|
||||
|
||||
export function stopStateSweeper() {
|
||||
if (handle) clearInterval(handle);
|
||||
handle = null;
|
||||
started = false;
|
||||
}
|
||||
|
||||
export const __test__ = { sweep };
|
||||
@@ -118,6 +118,16 @@ async function runHeavyStartup() {
|
||||
import("@/sse/services/backgroundTokenRefresh.js")
|
||||
.then(({ startBackgroundTokenRefresh }) => startBackgroundTokenRefresh())
|
||||
.catch((e) => console.log("[BackgroundTokenRefresh] scheduler start failed:", e.message));
|
||||
|
||||
// Pool egress geo probe — fills the Proxy Fitness / Proxy Pools egress column.
|
||||
import("@/lib/network/poolEgressProbe.js")
|
||||
.then(({ startPoolEgressProbe }) => startPoolEgressProbe())
|
||||
.catch((e) => console.log("[PoolEgressProbe] scheduler start failed:", e.message));
|
||||
|
||||
// Periodic in-memory state sweeper — prunes expired fitness/geo/session state.
|
||||
import("@/lib/network/stateSweeper.js")
|
||||
.then(({ startStateSweeper }) => startStateSweeper())
|
||||
.catch((e) => console.log("[StateSweeper] scheduler start failed:", e.message));
|
||||
}
|
||||
|
||||
function hasQuotaAutoPingEnabled(settings) {
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
import { describe, it, expect, beforeEach } from "vitest";
|
||||
import { setPoolGeo, getPoolGeo, poolGeoSnapshot, pruneStaleGeo, resetPoolGeo } from "open-sse/services/poolGeo.js";
|
||||
|
||||
describe("pool egress geo cache", () => {
|
||||
beforeEach(() => resetPoolGeo());
|
||||
|
||||
it("stores geo and classifies stability from egress changes", () => {
|
||||
setPoolGeo("p1", { ip: "1.1.1.1", country: "US" });
|
||||
setPoolGeo("p1", { ip: "1.1.1.1", country: "US" }); // same IP twice
|
||||
expect(getPoolGeo("p1").isUnstable).toBe(false);
|
||||
expect(getPoolGeo("p1").ipCount).toBe(1);
|
||||
|
||||
setPoolGeo("p1", { ip: "2.2.2.2", country: "US" }); // changed
|
||||
expect(getPoolGeo("p1").isUnstable).toBe(true);
|
||||
expect(getPoolGeo("p1").ipCount).toBe(2);
|
||||
});
|
||||
|
||||
it("prunes TTL-stale entries (ipHistory rides along)", () => {
|
||||
setPoolGeo("p1", { ip: "1.1.1.1", country: "US" });
|
||||
setPoolGeo("p1", { ip: "2.2.2.2", country: "US" });
|
||||
const cache = globalThis["__9routerPoolGeo__"];
|
||||
cache.get("p1").ts = Date.now() - 2 * 60 * 60 * 1000;
|
||||
expect(pruneStaleGeo()).toBe(1);
|
||||
expect(poolGeoSnapshot()).toEqual({});
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user