diff --git a/open-sse/services/poolGeo.js b/open-sse/services/poolGeo.js new file mode 100644 index 00000000..319934d3 --- /dev/null +++ b/open-sse/services/poolGeo.js @@ -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); + } +} diff --git a/src/lib/network/poolEgressProbe.js b/src/lib/network/poolEgressProbe.js new file mode 100644 index 00000000..cd76ec15 --- /dev/null +++ b/src/lib/network/poolEgressProbe.js @@ -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 }; diff --git a/src/lib/network/stateSweeper.js b/src/lib/network/stateSweeper.js new file mode 100644 index 00000000..df33c673 --- /dev/null +++ b/src/lib/network/stateSweeper.js @@ -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 }; \ No newline at end of file diff --git a/src/shared/services/initializeApp.js b/src/shared/services/initializeApp.js index 5b27f774..e627c560 100644 --- a/src/shared/services/initializeApp.js +++ b/src/shared/services/initializeApp.js @@ -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) { diff --git a/tests/unit/pool-geo.test.js b/tests/unit/pool-geo.test.js new file mode 100644 index 00000000..71f4eef2 --- /dev/null +++ b/tests/unit/pool-geo.test.js @@ -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({}); + }); +}); \ No newline at end of file