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/14221-rank-pool-candidates.md
Original file line number Diff line number Diff line change
@@ -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
124 changes: 111 additions & 13 deletions src/lib/db/proxies/rotation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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++) {
Expand All @@ -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<string, string> = { 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<string, unknown>;
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<T>(candidates: T[], signals?: Partial<PoolRankSignals>): 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.
Expand All @@ -175,25 +251,47 @@ function pickFromCandidates<T>(
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;
Expand All @@ -207,20 +305,20 @@ function pickFromCandidates<T>(
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.
Expand Down
151 changes: 151 additions & 0 deletions tests/unit/proxy-pool-rank-candidates.test.ts
Original file line number Diff line number Diff line change
@@ -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<Member[]> {
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);
});
10 changes: 6 additions & 4 deletions tests/unit/proxy-pool-skips-refused-member.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand Down
Loading