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/14588-opencode-429-proxy-dedup.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(sse):** the rotation wave no longer replays a proxy that already returned 429 within the same request ([#14588](https://github.com/diegosouzapw/OmniRoute/pull/14588)) — thanks @maxmad64bis
16 changes: 8 additions & 8 deletions open-sse/executors/opencode.ts
Original file line number Diff line number Diff line change
Expand Up @@ -614,10 +614,9 @@ export class OpencodeExecutor extends BaseExecutor {
// through the accounts is the retry). Avoids an unbounded loop on a
// persistently malformed upstream.
const emptyRejectionBudget = accounts.length === 1 ? 1 : 0;
// Tried set: proxy keys already proven unusable for this request's
// model (geo-blocked, or transient 5xx). Request-local only — nothing
// persists past execute().
const geoTriedProxyKeys = new Set<string>();
// Tried sets, request-local only: geo/transient + 429 no-replay keys.
const geoTriedProxyKeys = new Set<string>(),
rateLimitedProxyKeys = new Set<string>();
// Opt-in (PROXY_SKIP_RECENTLY_FAILED, default off): members the provider just refused
// (received refusal or refused TCP probe) are skipped. Off = plain rotation.
const skipRecentlyFailed = isProxySkipRecentlyFailedEnabled();
Expand Down Expand Up @@ -645,7 +644,7 @@ export class OpencodeExecutor extends BaseExecutor {
if (a.proxy === null) return !directTried || geoTriedProxyKeys.size === 0;
if (skipRecentlyFailed && isProxyAvoided(proxyEgressKey(a.proxy))) return false;
const k = proxyKeyOf(a.proxy);
return k !== null && !geoTriedProxyKeys.has(k);
return k !== null && !geoTriedProxyKeys.has(k) && !rateLimitedProxyKeys.has(k);
};
let account = this.pickAccountWith(accounts, isProxiedCandidate);
if (attributionOn) {
Expand Down Expand Up @@ -678,12 +677,11 @@ export class OpencodeExecutor extends BaseExecutor {
if (
!isMonoRetryOwed &&
lastResult !== null &&
geoTriedProxyKeys.size > 0 &&
geoTriedProxyKeys.size + rateLimitedProxyKeys.size > 0 &&
!isProxiedCandidate(account) &&
!(account.proxy === null && !directTried)
) {
// Geo exhaustion (last was 403/451) → surface as-is, no success mark.
// Transient exhaustion (last was 5xx) → same: surface last as-is.
// Geo/transient exhaustion → surface as-is, no success mark.
// Any other last status (e.g. 429 after 403s) → skip without a call.
if (lastWasGeo || lastWasTransient) break;
continue;
Expand Down Expand Up @@ -816,6 +814,8 @@ export class OpencodeExecutor extends BaseExecutor {
const status = result.response.status;
if (status === 429) {
markCooldown(account);
const rateKey = proxyKeyOf(account.proxy);
if (rateKey !== null) rateLimitedProxyKeys.add(rateKey);
const setAsideMs = egressPacing.noteRefusedMember(account.proxy, skipRecentlyFailed);
// Opt-in (#13657): a 429 that names a real rate limit stops the wave and
// the real upstream 429 is returned untouched (body, Retry-After, quota
Expand Down
128 changes: 128 additions & 0 deletions tests/unit/opencode-429-proxy-dedup.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
import { describe, it, beforeEach, afterEach, before, after } from "node:test";
import assert from "node:assert";
import net from "node:net";
import { OpencodeExecutor } from "../../open-sse/executors/opencode.ts";
import type { ExecutorLog, ProviderCredentials } from "../../open-sse/executors/base.ts";
import { resolveProxyForRequest } from "../../open-sse/utils/proxyFetch.ts";

// The park-and-replay path re-selects outside the rotation loop, so the
// per-request dedup guarantee below only holds with it disabled.
const PARK_FLAG = "OPENCODE_PARK_AND_RESUME";
const EARLY_STOP_FLAG = "OPENCODE_RATE_LIMITED_429_EARLY_STOP";

const log: ExecutorLog = { debug() {}, info() {}, warn() {}, error() {} };
const FPS = ["d", "e", "f"].map((c) => c.repeat(32));

const servers: net.Server[] = [];
const ports: number[] = [];

function listen(server: net.Server): Promise<number> {
return new Promise((resolve) => {
server.listen(0, "127.0.0.1", () => resolve((server.address() as net.AddressInfo).port));
});
}

before(async () => {
for (let i = 0; i < FPS.length; i++) {
const server = net.createServer((s) => s.destroy());
servers.push(server);
ports.push(await listen(server));
}
});

after(() => {
servers.forEach((s) => s.close());
});

function credentialsFor(proxies: number[]): ProviderCredentials {
const fingerprints = FPS.slice(0, proxies.length);
return {
apiKey: null,
accessToken: null,
connectionId: "noauth",
providerSpecificData: {
fingerprints,
accountProxies: fingerprints.map((fp, i) => ({
fingerprint: fp,
proxy: { type: "http", host: "127.0.0.1", port: proxies[i] },
})),
},
};
}

const BUSY_BODY = JSON.stringify({ error: { message: "upstream busy, try again" } });

describe("opencode 429 proxy dedup per request", () => {
let originalFetch: typeof globalThis.fetch;
let priorPark: string | undefined;
let priorEarlyStop: string | undefined;
let observed: string[];

beforeEach(() => {
originalFetch = globalThis.fetch;
priorPark = process.env[PARK_FLAG];
priorEarlyStop = process.env[EARLY_STOP_FLAG];
delete process.env[PARK_FLAG];
delete process.env[EARLY_STOP_FLAG];
observed = [];
});

afterEach(() => {
globalThis.fetch = originalFetch;
if (priorPark === undefined) delete process.env[PARK_FLAG];
else process.env[PARK_FLAG] = priorPark;
if (priorEarlyStop === undefined) delete process.env[EARLY_STOP_FLAG];
else process.env[EARLY_STOP_FLAG] = priorEarlyStop;
});

function installFetch(plan: Array<{ status: number; body?: string }>) {
let call = 0;
globalThis.fetch = (async (input: RequestInfo | URL) => {
const url =
typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url;
const resolved = resolveProxyForRequest(url);
observed.push(resolved.proxyUrl ? new URL(resolved.proxyUrl).port : "direct");
const step = plan[Math.min(call, plan.length - 1)];
call++;
return new Response(step.body ?? JSON.stringify({ ok: step.status === 200 }), {
status: step.status,
headers: { "Content-Type": "application/json" },
});
}) as typeof globalThis.fetch;
}

async function run(credentials: ProviderCredentials) {
const exec = new OpencodeExecutor("opencode-zen");
return exec.execute({
model: "muse-spark-1.3-contributor-free",
body: { messages: [{ role: "user", content: "hi" }], stream: false },
stream: false,
signal: null,
credentials,
log,
});
}

it("does not replay a rate-limited proxy within one rotation wave", async () => {
installFetch([{ status: 429, body: BUSY_BODY }, { status: 200 }]);
const result = await run(credentialsFor([ports[0], ports[0]]));
assert.strictEqual(observed.length, 1, `same proxy replayed (saw ${observed.length} calls)`);
assert.strictEqual(observed[0], String(ports[0]));
assert.strictEqual((result as { response: Response }).response.status, 429);
});

it("serves the last 429 when every proxy is rate-limited", async () => {
installFetch([{ status: 429, body: BUSY_BODY }]);
const result = await run(credentialsFor([ports[0], ports[0], ports[0]]));
assert.strictEqual(observed.length, 1, `exhaustion replayed a proxy (saw ${observed.length})`);
assert.strictEqual((result as { response: Response }).response.status, 429);
});

it("does not deduplicate two proxies sharing one host", async () => {
installFetch([{ status: 429, body: BUSY_BODY }, { status: 200 }]);
const result = await run(credentialsFor([ports[0], ports[1]]));
assert.strictEqual(observed.length, 2);
assert.notStrictEqual(observed[0], observed[1]);
assert.strictEqual((result as { response: Response }).response.status, 200);
});
});
Loading