From 285b96b137abbd7e8c376e25372b97141509a6d1 Mon Sep 17 00:00:00 2001 From: Dizzle <112548150+maxmad64bis@users.noreply.github.com> Date: Fri, 2 Oct 2026 03:11:38 +0200 Subject: [PATCH] fix(proxies): share the refusal store across duplicated server module copies --- .../fixes/15316-refusal-store-shared-state.md | 1 + open-sse/utils/proxyRefusalMemory.ts | 33 ++++----- open-sse/utils/proxyTransitionListeners.ts | 71 +++++++++++++++++-- src/lib/proxyEvents/proxyTransitionBridge.ts | 2 +- ...oxy-refusal-memory-single-instance.test.ts | 63 ++++++++++++++++ ...roxy-transition-bridge-single-emit.test.ts | 51 +++++++++++++ 6 files changed, 197 insertions(+), 24 deletions(-) create mode 100644 changelog.d/fixes/15316-refusal-store-shared-state.md create mode 100644 tests/unit/proxy-refusal-memory-single-instance.test.ts create mode 100644 tests/unit/proxy-transition-bridge-single-emit.test.ts diff --git a/changelog.d/fixes/15316-refusal-store-shared-state.md b/changelog.d/fixes/15316-refusal-store-shared-state.md new file mode 100644 index 00000000000..fe53f00c177 --- /dev/null +++ b/changelog.d/fixes/15316-refusal-store-shared-state.md @@ -0,0 +1 @@ +- **fix(proxies):** Share the set-aside store across duplicated server module copies, so a `connection reset` or `429` recorded through one copy keeps that egress set aside for every reader ([#15316](https://github.com/diegosouzapw/OmniRoute/pull/15316)) — thanks @maxmad64bis diff --git a/open-sse/utils/proxyRefusalMemory.ts b/open-sse/utils/proxyRefusalMemory.ts index 0b26de765c4..df5f0a2f610 100644 --- a/open-sse/utils/proxyRefusalMemory.ts +++ b/open-sse/utils/proxyRefusalMemory.ts @@ -10,7 +10,14 @@ * (callers gate writes and decisions on it) and it stays free of the proxy dispatcher, so * the DB layer can consult it without loading undici or the SOCKS connector. */ -import { notifyProxyTransition } from "./proxyTransitionListeners.ts"; +import { notifyProxyTransition, getSharedRefusalStore } from "./proxyTransitionListeners.ts"; +import type { + RefusalState, + SharedRefusalStore, + SlowOverrun, + TransportFailure, + TransportSuccess, +} from "./proxyTransitionListeners.ts"; import { stripIpv6Brackets } from "./proxyFamily.ts"; // Field trends show a refused egress rarely recovers within minutes, so the @@ -81,8 +88,10 @@ export const REFUSAL_POLICIES: { export type ProxyRefusalKind = keyof typeof REFUSAL_POLICIES; -// `seq` orders set-aside events so a cache can tell whether it already saw this one. -type RefusalState = { streak: number; until: number; seq: number }; +// Shared across duplicated server module copies (see getSharedRefusalStore): +// rebind on each module evaluation so HMR keeps the same object. +const store: SharedRefusalStore = getSharedRefusalStore(); +const memory: Map = store.memory; const MAX_ENTRIES = 1000; const REFUSAL_KINDS = Object.keys(REFUSAL_POLICIES) as ProxyRefusalKind[]; @@ -91,9 +100,6 @@ const DEFAULT_PORTS: Record = { http: "8080", https: "443", sock const RELAY_TYPES = new Set(["vercel", "deno", "cloudflare"]); const FAMILY_MARKER = /\?family=(ipv4|ipv6)$/; -const memory = new Map(); -let refusalSeq = 0; - const textField = (value: unknown): string => (typeof value === "string" ? value : ""); // The port as proxyConfigToUrl() normalizes it: the scheme default when unset, null if invalid. @@ -201,7 +207,7 @@ function insertState(key: string, kind: ProxyRefusalKind, nowMs: number): number const periodMs = Math.min(policy.baseMs * 2 ** (streak - 1), policy.maxMs); const id = entryId(key, kind); memory.delete(id); - memory.set(id, { streak, until: nowMs + periodMs, seq: ++refusalSeq }); + memory.set(id, { streak, until: nowMs + periodMs, seq: ++store.seq.value }); if (memory.size > MAX_ENTRIES) { const oldest = memory.keys().next().value; if (oldest !== undefined) memory.delete(oldest); @@ -391,7 +397,7 @@ export function listEntryMembers(entryKey: string | null): string[] { /** Sequence number of the last set-aside event recorded in this process (0 = none yet). */ export function getProxyRefusalSeq(): number { - return refusalSeq; + return store.seq.value; } /** @@ -488,11 +494,8 @@ export const TRANSPORT_EVIDENCE_WINDOW_MS = 300_000; export const TRANSPORT_EVIDENCE_THRESHOLD = 3; const MAX_TRANSPORT_EVIDENCE = 1000; -type TransportFailure = { key: string; destination: string; at: number }; -type TransportSuccess = { destination: string; key: string; at: number }; - -const transportFailures: TransportFailure[] = []; -const transportSuccesses: TransportSuccess[] = []; +const transportFailures: TransportFailure[] = store.transportFailures; +const transportSuccesses: TransportSuccess[] = store.transportSuccesses; // Lazy purge mirrors readState: entries older than the evidence window plus // twice the transport cap can no longer contribute, so drop them on read. @@ -577,9 +580,7 @@ export const SLOW_OVERRUN_WINDOW_MS = 300_000; export const SLOW_OVERRUN_THRESHOLD = 3; const MAX_SLOW_OVERRUNS = 1000; -type SlowOverrun = { key: string; at: number }; - -const slowOverruns: SlowOverrun[] = []; +const slowOverruns: SlowOverrun[] = store.slowOverruns; // Lazy purge mirrors readState: entries older than the evidence window plus // twice the slow cap can no longer contribute, so drop them on record. diff --git a/open-sse/utils/proxyTransitionListeners.ts b/open-sse/utils/proxyTransitionListeners.ts index 6c79788e6ab..e0ca474edf2 100644 --- a/open-sse/utils/proxyTransitionListeners.ts +++ b/open-sse/utils/proxyTransitionListeners.ts @@ -18,17 +18,72 @@ export interface ProxyTransition { export type ProxyTransitionListener = (transition: ProxyTransition) => void; -const listeners = new Set(); +export interface SharedRefusalStore { + memory: Map; + seq: { value: number }; + transportFailures: TransportFailure[]; + transportSuccesses: TransportSuccess[]; + slowOverruns: SlowOverrun[]; + listeners: Set; + listenerKeys: Map; +} + +export type RefusalState = { streak: number; until: number; seq: number }; +export type TransportFailure = { key: string; destination: string; at: number }; +export type TransportSuccess = { destination: string; key: string; at: number }; +export type SlowOverrun = { key: string; at: number }; + +const SHARED_STORE_KEY = Symbol.for("omniroute.proxyRefusalMemory"); + +/** + * Process-wide refusal store, shared across duplicated server module copies. + * Lazy `??=` init mirrors getPatchState in proxyFetch.ts: data only, no + * closures, so re-evaluation (HMR) rebinds the same object. + */ +export function getSharedRefusalStore(): SharedRefusalStore { + const holder = globalThis as unknown as Record; + let store = holder[SHARED_STORE_KEY]; + if (!store) { + store = { + memory: new Map(), + seq: { value: 0 }, + transportFailures: [], + transportSuccesses: [], + slowOverruns: [], + listeners: new Set(), + listenerKeys: new Map(), + }; + holder[SHARED_STORE_KEY] = store; + } + return store; +} -export function onProxyTransition(listener: ProxyTransitionListener): () => void { - listeners.add(listener); +export function onProxyTransition(listener: ProxyTransitionListener, key?: string): () => void { + const store = getSharedRefusalStore(); + if (key !== undefined) { + const existing = store.listenerKeys.get(key); + if (existing !== undefined) { + const current = existing; + return () => { + store.listeners.delete(current); + if (store.listenerKeys.get(key) === current) store.listenerKeys.delete(key); + }; + } + store.listeners.add(listener); + store.listenerKeys.set(key, listener); + return () => { + store.listeners.delete(listener); + if (store.listenerKeys.get(key) === listener) store.listenerKeys.delete(key); + }; + } + store.listeners.add(listener); return () => { - listeners.delete(listener); + store.listeners.delete(listener); }; } export function notifyProxyTransition(transition: ProxyTransition): void { - for (const listener of listeners) { + for (const listener of getSharedRefusalStore().listeners) { try { listener(transition); } catch (err) { @@ -39,10 +94,12 @@ export function notifyProxyTransition(transition: ProxyTransition): void { /** Test-only: number of registered listeners. */ export function __listenerCountForTesting(): number { - return listeners.size; + return getSharedRefusalStore().listeners.size; } /** Test-only: forget all listeners. */ export function __resetProxyTransitionListenersForTesting(): void { - listeners.clear(); + const store = getSharedRefusalStore(); + store.listeners.clear(); + store.listenerKeys.clear(); } diff --git a/src/lib/proxyEvents/proxyTransitionBridge.ts b/src/lib/proxyEvents/proxyTransitionBridge.ts index 417677445a3..ed90c95ce67 100644 --- a/src/lib/proxyEvents/proxyTransitionBridge.ts +++ b/src/lib/proxyEvents/proxyTransitionBridge.ts @@ -153,7 +153,7 @@ export function registerProxyTransitionBridge(): void { if (unsubscribe !== null) return; unsubscribe = onProxyTransition((transition) => { emitSetAside(transition); - }); + }, "proxyTransitionBridge"); } /** Test-only: forget emission state and the subscription. */ diff --git a/tests/unit/proxy-refusal-memory-single-instance.test.ts b/tests/unit/proxy-refusal-memory-single-instance.test.ts new file mode 100644 index 00000000000..e3ef5e32166 --- /dev/null +++ b/tests/unit/proxy-refusal-memory-single-instance.test.ts @@ -0,0 +1,63 @@ +/** + * Two copies of the refusal store share one process-wide state. + * + * The server bundle holds several copies of this store (one per module + * graph): transport failures recorded through one copy were invisible to + * the readers of another. All mutable state lives on `globalThis` so + * every copy reads and writes the same store. + */ +import test from "node:test"; +import assert from "node:assert/strict"; + +const copyA = await import("../../open-sse/utils/proxyRefusalMemory.ts?copy=a"); +const copyB = await import("../../open-sse/utils/proxyRefusalMemory.ts?copy=b"); + +const KEY = "http://user@host.example:8080"; +const DEST = "target.example"; +const OTHER = "http://user@other.example:8080"; + +function resetBoth(): void { + copyA.__resetProxyRefusalMemoryForTesting(); + copyB.__resetProxyRefusalMemoryForTesting(); + copyA.__resetTransportEvidenceForTesting(); + copyB.__resetTransportEvidenceForTesting(); + copyA.__resetSlowOverrunsForTesting(); + copyB.__resetSlowOverrunsForTesting(); +} + +test.beforeEach(() => { + resetBoth(); +}); + +test.afterEach(() => { + resetBoth(); +}); + +test("witness: distinct module instances under distinct import URLs", () => { + assert.notEqual(copyA.noteProxyRefusal, copyB.noteProxyRefusal); +}); + +test("set-aside written through one copy is visible through the other", () => { + const now = 1_000_000; + const periodMs = copyA.noteProxyRefusal(KEY, "transport", now); + assert.ok(typeof periodMs === "number"); + assert.equal(copyB.isProxyAvoided(KEY, now), true); +}); + +test("transport cross-evidence recorded through either copy reads from both", () => { + const now = 2_000_000; + copyA.recordTransportFailure(KEY, DEST, now); + copyB.recordTransportFailure(KEY, DEST, now + 1); + copyA.recordTransportFailure(KEY, DEST, now + 2); + copyB.recordTransportSuccess(DEST, OTHER, now + 3); + assert.equal(copyA.hasTransportCrossEvidence(KEY, DEST, now + 4), true); + assert.equal(copyB.hasTransportCrossEvidence(KEY, DEST, now + 4), true); +}); + +test("slow-overrun evidence recorded through one copy reads through the other", () => { + const now = 3_000_000; + copyA.recordSlowOverrun(KEY, now); + copyA.recordSlowOverrun(KEY, now + 1); + copyA.recordSlowOverrun(KEY, now + 2); + assert.equal(copyB.hasSlowOverrunEvidence(KEY, now + 3), true); +}); diff --git a/tests/unit/proxy-transition-bridge-single-emit.test.ts b/tests/unit/proxy-transition-bridge-single-emit.test.ts new file mode 100644 index 00000000000..ae170f68577 --- /dev/null +++ b/tests/unit/proxy-transition-bridge-single-emit.test.ts @@ -0,0 +1,51 @@ +/** + * One set-aside transition emits one `proxy.set_aside` bus event, even when + * the bridge module is loaded twice (duplicated server module copies would + * otherwise subscribe twice and emit twice). + */ +import test from "node:test"; +import assert from "node:assert/strict"; + +const bridgeA = await import("../../src/lib/proxyEvents/proxyTransitionBridge.ts?copy=a"); +const bridgeB = await import("../../src/lib/proxyEvents/proxyTransitionBridge.ts?copy=b"); +const storeA = await import("../../open-sse/utils/proxyRefusalMemory.ts?copy=a"); +const listenersA = await import("../../open-sse/utils/proxyTransitionListeners.ts?copy=a"); +const bus = await import("../../src/lib/events/eventBus.ts"); + +const KEY = "http://user@bridge.example:8080"; +const NOW = 10_000_000; + +function resetAll(): void { + bridgeA.__resetBridgeForTesting(); + bridgeB.__resetBridgeForTesting(); + storeA.__resetProxyRefusalMemoryForTesting(); + listenersA.__resetProxyTransitionListenersForTesting(); +} + +test.beforeEach(() => { + resetAll(); +}); + +test.afterEach(() => { + resetAll(); + delete process.env.PROXY_WEBHOOK_REBOUND_MS; +}); + +test("transition seen by two bridge copies emits a single bus event", () => { + const clock = () => NOW; + bridgeA.__setBridgeNowForTesting(clock); + bridgeB.__setBridgeNowForTesting(clock); + bridgeA.registerProxyTransitionBridge(); + bridgeB.registerProxyTransitionBridge(); + const received: unknown[] = []; + const off = bus.on("proxy.set_aside", (payload) => { + received.push(payload); + }); + try { + const periodMs = storeA.noteProxyRefusal(KEY, "transport", NOW); + assert.ok(typeof periodMs === "number"); + assert.equal(received.length, 1); + } finally { + off(); + } +});