diff --git a/changelog.d/fixes/14221-rank-pool-candidates.md b/changelog.d/fixes/14221-rank-pool-candidates.md new file mode 100644 index 00000000000..dcd26aa05a0 --- /dev/null +++ b/changelog.d/fixes/14221-rank-pool-candidates.md @@ -0,0 +1 @@ +- **fix(proxies):** stop re-serving a pool member that just failed at the head of the pool: candidates are ordered by crossed short-memory health signals (opt-in, off by default), so a request lands on the first useful try instead of paying the failed head first ([#14221](https://github.com/diegosouzapw/OmniRoute/pull/14221)) — thanks @maxmad64bis diff --git a/src/lib/db/proxies/rotation.ts b/src/lib/db/proxies/rotation.ts index 853bfd4f984..12645f55a35 100644 --- a/src/lib/db/proxies/rotation.ts +++ b/src/lib/db/proxies/rotation.ts @@ -15,6 +15,7 @@ import { proxyEgressKey, } from "@omniroute/open-sse/utils/proxyRefusalMemory.ts"; import { isProxySkipRecentlyFailedEnabled } from "@/shared/utils/featureFlags"; +import { getCachedProxyHealth } from "@/lib/proxyHealth"; import type { JsonRecord, ProxyScope, ProxyRotationStrategy } from "./types"; import { PROXY_ROTATION_STRATEGIES, DEFAULT_PROXY_ROTATION_STRATEGY } from "./types"; import { @@ -149,6 +150,14 @@ function eligibleMemberIndexes(candidates: unknown[]): number[] | null { return isProxySkipRecentlyFailedEnabled() ? eligible : null; } +// True once the sticky window elapsed (or never started): the held member is due +// for rotation. Shared by the pre-rank bypass (held member served untouched) and +// the sticky branch below (advance on expiry) — same `state`, no extra DB read. +function isStickyExpired(state: { stickyWindowMinutes: number; rotatedAt: string | null }): boolean { + const lastRotated = state.rotatedAt ? Date.parse(state.rotatedAt) : NaN; + return !Number.isFinite(lastRotated) || Date.now() - lastRotated >= state.stickyWindowMinutes * 60_000; +} + // First eligible index at or after `start`, going round the pool. function firstEligibleFrom(start: number, eligible: number[], size: number): number { for (let step = 0; step < size; step++) { @@ -158,6 +167,73 @@ function firstEligibleFrom(start: number, eligible: number[], size: number): num return start; } +/** Health signals read from short-lived process memory, injectable for tests. */ +export interface PoolRankSignals { + isAvoided: (key: string | null) => boolean; + probeHealth: (url: string) => boolean | null; +} + +const DEFAULT_POOL_RANK_SIGNALS: PoolRankSignals = { + isAvoided: isProxyAvoided, + probeHealth: getCachedProxyHealth, +}; + +// Relay entries carry the relay URL in `host` and no dispatcher: not rankable. +const RELAY_TYPES = new Set(["vercel", "deno", "cloudflare"]); +// Same scheme defaults as proxyConfigToUrl() so the rebuilt URL hits the probe cache key. +const DEFAULT_PORTS: Record = { http: "8080", https: "443", socks5: "1080" }; + +// Rebuild the probe URL for a pool row the way the dispatcher builds it +// (`proxyConfigToUrl`, read-only replica): scheme + encoded auth + host + port, +// with the `?family=` marker when set. Null when the row cannot egress. +function candidateProbeUrl(row: unknown): string | null { + if (!row || typeof row !== "object" || Array.isArray(row)) return null; + const record = row as Record; + const host = typeof record.host === "string" ? record.host : ""; + if (!host) return null; + const type = String(record.type || "http").toLowerCase(); + if (RELAY_TYPES.has(type) || !(type in DEFAULT_PORTS)) return null; + const parsed = Number(record.port); + const port = + record.port && Number.isInteger(parsed) && parsed >= 1 && parsed <= 65535 + ? String(parsed) + : DEFAULT_PORTS[type]; + const bracketed = host.includes(":") && !host.startsWith("[") ? `[${host}]` : host; + const username = typeof record.username === "string" ? record.username : ""; + const password = typeof record.password === "string" ? record.password : ""; + const auth = + username || password ? `${encodeURIComponent(username)}:${encodeURIComponent(password)}@` : ""; + const family = typeof record.family === "string" ? record.family : ""; + const marker = family === "ipv4" || family === "ipv6" ? `?family=${family}` : ""; + return `${type}://${auth}${bracketed}:${port}${marker}`; +} + +/** + * Order pool candidates by crossed short-memory health signals without removing + * anyone: a member just set aside ranks last, then a member whose last cached + * probe verdict was negative. Unknown (no signal, unreconstructible URL) keeps + * the current position order. Stable: health ties keep their relative order, so + * an all-clear or all-set-aside pool returns its input order unchanged. + */ +export function rankPoolCandidates(candidates: T[], signals?: Partial): T[] { + if (candidates.length < 2) return [...candidates]; + const { isAvoided, probeHealth } = { + ...DEFAULT_POOL_RANK_SIGNALS, + isAvoided: signals?.isAvoided ?? DEFAULT_POOL_RANK_SIGNALS.isAvoided, + probeHealth: signals?.probeHealth ?? DEFAULT_POOL_RANK_SIGNALS.probeHealth, + }; + const scored = candidates.map((candidate, index) => { + if (isAvoided(proxyEgressKey(candidate))) return { candidate, index, score: 2 }; + const url = candidateProbeUrl(candidate); + if (url !== null && probeHealth(url) === false) return { candidate, index, score: 1 }; + return { candidate, index, score: 0 }; + }); + if (scored.every((entry) => entry.score === scored[0].score)) return [...candidates]; + return scored + .sort((a, b) => a.score - b.score || a.index - b.index) + .map((entry) => entry.candidate); +} + /** * Pick one member from an already-alive candidate list according to the scope's * rotation strategy. Assumes `candidates` is non-empty and ordered by position. @@ -175,25 +251,47 @@ function pickFromCandidates( if (candidates.length === 1) return candidates[0]; const state = getOrCreateRotationRow(db, normalizedScope, rotationScopeId); - const eligible = eligibleMemberIndexes(candidates); + + if (state.strategy === "sticky") { + const expired = isStickyExpired(state); + if (!expired) { + const idx = ((state.cursor % candidates.length) + candidates.length) % candidates.length; + const eligible = eligibleMemberIndexes(candidates); + return candidates[eligible ? firstEligibleFrom(idx, eligible, candidates.length) : idx]; + } + } + + // Order by crossed short-memory health signals (opt-in, PROXY_SKIP_RECENTLY_FAILED): + // stops re-serving at the head a proxy that just failed, without removing anyone. + // Sticky past its window and every other strategy rank normally; a held sticky + // member returns above, untouched. The eligible-skip below still applies on the + // ranked list, so a set-aside member stays skipped while another is eligible and + // the cursor advances past the member actually served. + // NOTE: ranking changes which member the persisted cursor lands on. After a + // set-aside, the next pick serves the healthiest member at-or-after the cursor + // (not the cursor member itself when it was set aside) — the cursor then + // advances past the member served, preserving rotation without re-serving the + // failed head first. + const ranked = isProxySkipRecentlyFailedEnabled() + ? rankPoolCandidates(candidates) + : [...candidates]; + const eligible = eligibleMemberIndexes(ranked); if (state.strategy === "random") { // crypto.randomInt (unbiased, uniform in [0, length)) instead of Math.random — // CodeQL js/insecure-randomness flags Math.random flowing into the selected proxy's // credentials (a "security context"). Load-balancing selection is not a secret, but // crypto.randomInt silences the alert at the source and is unbiased (#6365 follow-up). - if (eligible) return candidates[eligible[randomInt(eligible.length)]]; - return candidates[randomInt(candidates.length)]; + if (eligible) return ranked[eligible[randomInt(eligible.length)]]; + return ranked[randomInt(ranked.length)]; } if (state.strategy === "latency") { - return pickByLatency(db, eligible ? eligible.map((index) => candidates[index]) : candidates); + return pickByLatency(db, eligible ? eligible.map((index) => ranked[index]) : ranked); } if (state.strategy === "sticky") { - const windowMs = state.stickyWindowMinutes * 60_000; - const lastRotated = state.rotatedAt ? Date.parse(state.rotatedAt) : NaN; - const expired = !Number.isFinite(lastRotated) || Date.now() - lastRotated >= windowMs; + const expired = isStickyExpired(state); let cursor = state.cursor; if (expired) { cursor = state.cursor + 1; @@ -207,20 +305,20 @@ function pickFromCandidates( rotationScopeId ); } - const idx = ((cursor % candidates.length) + candidates.length) % candidates.length; + const idx = ((cursor % ranked.length) + ranked.length) % ranked.length; // A held member set aside is replaced for this pick only: no extra write. - return candidates[eligible ? firstEligibleFrom(idx, eligible, candidates.length) : idx]; + return ranked[eligible ? firstEligibleFrom(idx, eligible, ranked.length) : idx]; } // round-robin (default): pick at the current cursor, then advance it monotonically, // past any member skipped so the next pick starts after the one actually served. - const idx = ((state.cursor % candidates.length) + candidates.length) % candidates.length; - const served = eligible ? firstEligibleFrom(idx, eligible, candidates.length) : idx; - const skipped = (served - idx + candidates.length) % candidates.length; + const idx = ((state.cursor % ranked.length) + ranked.length) % ranked.length; + const served = eligible ? firstEligibleFrom(idx, eligible, ranked.length) : idx; + const skipped = (served - idx + ranked.length) % ranked.length; db.prepare( "UPDATE proxy_scope_rotation SET cursor = ?, updated_at = ? WHERE scope = ? AND scope_id IS ?" ).run(state.cursor + skipped + 1, new Date().toISOString(), normalizedScope, rotationScopeId); - return candidates[served]; + return ranked[served]; } // Fetch the alive, position-ordered candidate rows for a (scope, scope_id) pool. diff --git a/tests/unit/proxy-pool-rank-candidates.test.ts b/tests/unit/proxy-pool-rank-candidates.test.ts new file mode 100644 index 00000000000..5d87152a284 --- /dev/null +++ b/tests/unit/proxy-pool-rank-candidates.test.ts @@ -0,0 +1,151 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +// Pool candidate ranking: with PROXY_SKIP_RECENTLY_FAILED on, the pool orders +// candidates by crossed short-memory health signals instead of serving position +// order, without removing anyone. The pool stops re-serving at the head a proxy +// that just failed, so the request lands on the first useful try. With the flag +// off (the default) the order stays exactly the plain rotation. + +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-pool-rank-")); +process.env.DATA_DIR = TEST_DATA_DIR; +process.env.API_KEY_SECRET = "test-secret"; + +const core = await import("../../src/lib/db/core.ts"); +const proxiesDb = await import("../../src/lib/db/proxies.ts"); +const rotation = await import("../../src/lib/db/proxies/rotation.ts"); +const memory = await import("../../open-sse/utils/proxyRefusalMemory.ts"); + +function resetStorage() { + memory.__resetProxyRefusalMemoryForTesting(); + process.env.PROXY_SKIP_RECENTLY_FAILED = "true"; + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); + fs.mkdirSync(TEST_DATA_DIR, { recursive: true }); +} + +test.beforeEach(() => { + resetStorage(); +}); + +test.after(() => { + delete process.env.PROXY_SKIP_RECENTLY_FAILED; + memory.__resetProxyRefusalMemoryForTesting(); + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); +}); + +type Member = { id: string; host: string; port: number }; +let seq = 0; + +async function pool(size: number, scope = "provider", scopeId = "openai"): Promise { + const members: Member[] = []; + for (let i = 0; i < size; i++) { + seq++; + const host = `10.9.0.${seq}`; + const port = 9200 + seq; + const proxy = await proxiesDb.createProxy({ name: `member ${seq}`, type: "http", host, port }); + await proxiesDb.addProxyToScopePool(scope, scopeId, proxy.id); + members.push({ id: proxy.id, host, port }); + } + return members; +} + +function keyOf(member: Member) { + return memory.proxyEgressKey({ type: "http", host: member.host, port: member.port }); +} + +function setAside(member: Member) { + memory.noteProxyRefusal(keyOf(member), "proxy_unreachable"); +} + +async function pick(scope = "provider", scopeId = "openai") { + const resolved = await proxiesDb.resolveProxyForScopeFromRegistry(scope, scopeId); + return (resolved as { proxy: { host: string } }).proxy.host; +} + +test("a member set aside ranks last: B, C before A", async () => { + const [a, b, c] = await pool(3); + setAside(a); + const ranked = rotation.rankPoolCandidates([a, b, c]); + assert.deepEqual( + ranked.map((m) => m.host), + [b.host, c.host, a.host] + ); +}); + +test("no signal keeps position order", async () => { + const [a, b, c] = await pool(3); + assert.deepEqual( + rotation.rankPoolCandidates([a, b, c]).map((m) => m.host), + [a.host, b.host, c.host] + ); +}); + +test("with every member set aside the pool keeps its order and size", async () => { + const members = await pool(3); + for (const member of members) setAside(member); + const ranked = rotation.rankPoolCandidates(members); + assert.equal(ranked.length, 3); + assert.deepEqual( + ranked.map((m) => m.host), + members.map((m) => m.host) + ); +}); + +test("a negative probe verdict ranks the member last", async () => { + const [a, b] = await pool(2); + const ranked = rotation.rankPoolCandidates([a, b], { + probeHealth: (url) => (url.includes(a.host) ? false : null), + }); + assert.deepEqual( + ranked.map((m) => m.host), + [b.host, a.host] + ); +}); + +test("the probe URL carries auth and the port like the dispatcher", async () => { + const row = { type: "http", host: "10.9.9.9", port: 9201, username: "u", password: "p" }; + const ranked = rotation.rankPoolCandidates([row, { ...row, host: "10.9.9.8" }], { + probeHealth: (url) => (url.startsWith("http://u:p@10.9.9.9:9201") ? false : null), + }); + assert.deepEqual( + ranked.map((m) => (m as { host: string }).host), + ["10.9.9.8", "10.9.9.9"] + ); +}); + +test("with the flag off the order stays the plain rotation", async () => { + const [a] = await pool(2); + setAside(a); + delete process.env.PROXY_SKIP_RECENTLY_FAILED; + try { + assert.equal(await pick(), a.host); + } finally { + process.env.PROXY_SKIP_RECENTLY_FAILED = "true"; + } +}); + +test("functional: the pool serves the healthy member first after the head fails", async () => { + const [a, b] = await pool(2); + setAside(a); + assert.equal(await pick(), b.host); + assert.equal(await pick(), b.host); +}); + +test("sticky holds its member while its window runs, even with a rank signal", async () => { + const members = await pool(3); + await proxiesDb.setScopeRotationStrategy("provider", "openai", "sticky", { + stickyWindowMinutes: 30, + }); + const held = await pick(); + const heldMember = members.find((m) => m.host === held) ?? null; + assert.ok(heldMember); + const other = members.find((m) => m.host !== held) ?? null; + assert.ok(other); + setAside(other); + assert.equal(await pick(), held); +}); diff --git a/tests/unit/proxy-pool-skips-refused-member.test.ts b/tests/unit/proxy-pool-skips-refused-member.test.ts index 56e2b6ec4e7..4bf49f07ded 100644 --- a/tests/unit/proxy-pool-skips-refused-member.test.ts +++ b/tests/unit/proxy-pool-skips-refused-member.test.ts @@ -179,15 +179,17 @@ test("a connection's chat-path resolution skips a member set aside", async () => assert.equal((first as { proxy: { host: string } }).proxy.host, a.host); setAside(b); + // Ranked (b last) + skip: cursor 1 lands on c, served past the set-aside member. const next = await settingsDb.resolveProxyForConnection("conn-pool"); assert.equal((next as { proxy: { host: string } }).proxy.host, c.host); memory.noteProxyRecovered(keyOf(b), "proxy_unreachable"); assert.equal(memory.isProxyAvoided(keyOf(b)), false); - assert.deepEqual( - [await pick("account", "conn-pool"), await pick("account", "conn-pool")], - [a.host, b.host] - ); + // Rank + skip moves the cursor past c (served at cursor 1, cursor now 2), so the + // next pick serves c again (cursor 2 in the restored position order [a,b,c]), + // then rotation resumes at a. + const after = [await pick("account", "conn-pool"), await pick("account", "conn-pool")]; + assert.deepEqual(after, [c.host, a.host]); }); test("with the flag off a connection's chat-path resolution still rotates normally", async () => {