Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog.d/fixes/15316-refusal-store-shared-state.md
Original file line number Diff line number Diff line change
@@ -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
33 changes: 17 additions & 16 deletions open-sse/utils/proxyRefusalMemory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<string, RefusalState> = store.memory;

const MAX_ENTRIES = 1000;
const REFUSAL_KINDS = Object.keys(REFUSAL_POLICIES) as ProxyRefusalKind[];
Expand All @@ -91,9 +100,6 @@ const DEFAULT_PORTS: Record<string, string> = { http: "8080", https: "443", sock
const RELAY_TYPES = new Set(["vercel", "deno", "cloudflare"]);
const FAMILY_MARKER = /\?family=(ipv4|ipv6)$/;

const memory = new Map<string, RefusalState>();
let refusalSeq = 0;

/**
* Credential-free label for one egress key (`scheme://user@host:port` as
* proxyEgressKey writes it): `scheme://host:port`. String surgery only, no
Expand Down Expand Up @@ -250,7 +256,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);
Expand Down Expand Up @@ -457,7 +463,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;
}

/**
Expand Down Expand Up @@ -554,11 +560,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.
Expand Down Expand Up @@ -675,9 +678,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.
Expand Down
71 changes: 64 additions & 7 deletions open-sse/utils/proxyTransitionListeners.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,17 +18,72 @@ export interface ProxyTransition {

export type ProxyTransitionListener = (transition: ProxyTransition) => void;

const listeners = new Set<ProxyTransitionListener>();
export interface SharedRefusalStore {
memory: Map<string, RefusalState>;
seq: { value: number };
transportFailures: TransportFailure[];
transportSuccesses: TransportSuccess[];
slowOverruns: SlowOverrun[];
listeners: Set<ProxyTransitionListener>;
listenerKeys: Map<string, ProxyTransitionListener>;
}

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<symbol, SharedRefusalStore | undefined>;
let store = holder[SHARED_STORE_KEY];
if (!store) {
store = {
memory: new Map<string, RefusalState>(),
seq: { value: 0 },
transportFailures: [],
transportSuccesses: [],
slowOverruns: [],
listeners: new Set<ProxyTransitionListener>(),
listenerKeys: new Map<string, ProxyTransitionListener>(),
};
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) {
Expand All @@ -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();
}
2 changes: 1 addition & 1 deletion src/lib/proxyEvents/proxyTransitionBridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand Down
63 changes: 63 additions & 0 deletions tests/unit/proxy-refusal-memory-single-instance.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
51 changes: 51 additions & 0 deletions tests/unit/proxy-transition-bridge-single-emit.test.ts
Original file line number Diff line number Diff line change
@@ -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();
}
});
Loading