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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(proxies):** set aside egress after repeated transport failures with cross-egress success ([#14802](https://github.com/diegosouzapw/OmniRoute/pull/14802)) — thanks @maxmad64bis
3 changes: 2 additions & 1 deletion config/quality/file-size-baseline.json
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"_rebaseline_2026_09_24_14802_transport_setaside": "PR #14802 own growth: open-sse/utils/proxyFetch.ts 1334->1341 (+7 = two-symbol import + proxied-success hook at the single dispatch return + tagged-failure guard at the single final throw; full evidence helpers live in open-sse/utils/proxyRefusalMemory.ts (not frozen), irreducible call-site wiring at the two existing chokepoints). Covered by tests/unit/proxy-transport-setaside-cross-evidence.test.ts (10/10) + tests/unit/proxy-transport-traffic-simulation.test.ts (4/4).",
"_rebaseline_2026_09_24_14708_persisted_cooldown_precedence": "PR #14708 own growth: open-sse/services/combo/roundRobinCombo.ts 1263->1285 (+22 = the persisted-cooldown pre-check at the only round-robin dispatch point with no persisted check (import + async pre-check + skip/continue), plus the sticky-target reflow). Irreducible: the check must live at the dispatch site that bypasses executeTargetAttempt. Covered by tests/unit/combo-transient-persisted-precedence.test.ts.",
"_rebaseline_2026_09_24_14367_synthetic_key_budget_orphan": "PR #14367 (original fix; #14661 re-landed the same diff): src/lib/db/apiKeys.ts 1718->1719 (+1 = the single import of SYNTHETIC_ENV_API_KEY_ID threading the canonical synthetic-identity constant into getApiKeyMetadata; the exclusion list lives in src/shared/constants/apiKeyIdentities.ts). Irreducible. Covered by tests/unit/db/health-check-orphan-budgets.test.ts.",
"_rebaseline_2026_09_23_13593_self_hop": "PR for #13593 own growth: open-sse/utils/proxyFetch.ts +4 (one import, a two-line comment, and one call stamp the process self-hop token onto fetches aimed at this listener). Re-measured with split(\"\\n\").length on the tree reconciled with release/v3.8.51 @ce11dda6: merge-base e693a7e4 1287 -> tip 1330 (tip-only growth: +7 #14620 applied-proxy shared context, frozen at 1294; +36 #14311 loopback no-replay 34cff242, whose ceiling the 2026-09-24 merge wave deferred to its follow-up rebaseline PR) -> 1334 with this PR (+4). The URL match and header live in open-sse/utils/selfHop.ts (under cap); admission already skips the public pressure shed for isInternalAdmissionBypass, which now accepts the token. Irreducible call-site wiring at the single fetch chokepoint. Covered by tests/unit/own-listener-self-hop.test.ts.",
Expand Down Expand Up @@ -510,7 +511,7 @@
"open-sse/services/combo/executeTargetAttempt.ts": 1273,
"open-sse/translator/response/openai-responses.ts": 1536,
"open-sse/utils/cursorAgentProtobuf.ts": 1588,
"open-sse/utils/proxyFetch.ts": 1334,
"open-sse/utils/proxyFetch.ts": 1341,
"open-sse/utils/stream.ts": 3296,
"open-sse/vendor/codex-chatgpt-web/adapters/chatgpt-web/browser-worker.ts": 4398,
"open-sse/vendor/codex-chatgpt-web/bridge.ts": 1335,
Expand Down
9 changes: 8 additions & 1 deletion open-sse/utils/proxyFetch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import tlsClient, { type TlsFetchOptions, guardTlsFirstByte } from "./tlsClient.
import { withUpstreamStatusCapture } from "./upstreamStatusCapture.ts";
import { stampOwnListenerSelfHop } from "./selfHop.ts";
import { describeFallbackFailure, redactProxyDetailsInMessage } from "./proxyFetchRedaction.ts";
import { recordFinalTransportOutcome, recordProxiedSuccess } from "./proxyTransportOutcome.ts";
import { sanitizeTransportError } from "./proxyTransportError.ts";
import { isProxyReachable } from "@/lib/proxyHealth";
import {
Expand Down Expand Up @@ -1171,11 +1172,13 @@ async function patchedFetchUnrecorded(
let lastProxyError: unknown = null;
for (let attempt = 0; attempt < maxProxyAttempts; attempt++) {
try {
return await _undiciProxy(input, {
const response = await _undiciProxy(input, {
...options,
dispatcher:
attempt === 0 ? createProxyDispatcher(proxyUrl) : getProxyRetryDispatcher(proxyUrl),
});
recordProxiedSuccess(proxyUrl, targetUrl); // completed response, any status
return response;
} catch (error) {
if (isCallerAbort(error, getEffectiveSignal(input, options))) throw error;
const msg = error instanceof Error ? error.message : String(error);
Expand Down Expand Up @@ -1207,6 +1210,10 @@ async function patchedFetchUnrecorded(
originalMsg ? `Proxy request failed: ${originalMsg}` : "Proxy request failed",
"PROXY_REQUEST_FAILED"
);
// Read the code off the thrown sanitized error (tag survives the
// sanitize as errorCode passthrough; untagged reads undefined).
if (sanitized.errorCode === "proxy_unreachable")
recordFinalTransportOutcome(proxyUrl, targetUrl);
if (sanitized.causeCode) {
sanitized.message += ` (cause ${sanitized.causeCode})`;
}
Expand Down
114 changes: 114 additions & 0 deletions open-sse/utils/proxyRefusalMemory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,11 +56,18 @@ if (quota429MaxMs < quota429BaseMs) {
export const REFUSAL_POLICIES: {
proxy_unreachable: { baseMs: 60_000; maxMs: 600_000 };
ip_quota_429: RefusalPolicy;
transport: { baseMs: 60_000; maxMs: 600_000 };
} = {
/** The TCP probe could not open a connection to the proxy. */
proxy_unreachable: { baseMs: 60_000, maxMs: 600_000 },
/** The provider refused through this proxy; the member is set aside for a cooldown. */
ip_quota_429: { baseMs: quota429BaseMs, maxMs: quota429MaxMs },
/**
* Repeated tagged transport failures through this egress with cross-evidence:
* the same destination answers through a different egress, so the member —
* not the destination — is at fault. Same curve as a refused probe.
*/
transport: { baseMs: 60_000, maxMs: 600_000 },
};

export type ProxyRefusalKind = keyof typeof REFUSAL_POLICIES;
Expand Down Expand Up @@ -236,6 +243,113 @@ export function __resetProxyRefusalMemoryForTesting(): void {
memory.clear();
}

/**
* The destination a transport outcome is attributed to: the lower-cased host of
* the requested target URL, without IPv6 brackets. Shared by the failure and
* success hooks so cross-evidence always matches on the same normalization.
* Null when the target is not a parseable URL.
*/
export function transportDestinationKey(targetUrl: string): string | null {
try {
const host = new URL(targetUrl).hostname;
if (!host) return null;
return stripIpv6Brackets(host).toLowerCase();
} catch {
return null;
}
}

/**
* Cross-evidence for the `transport` refusal kind. A tagged transport failure
* alone never sets a member aside: it is only evidence, kept per (egress key,
* destination host). A later success to the same destination through a
* *different* egress proves the destination answers and the failing member —
* not the destination — is at fault.
*/
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[] = [];

// Lazy purge mirrors readState: entries older than the evidence window plus
// twice the transport cap can no longer contribute, so drop them on read.
function purgeTransportEvidence(nowMs: number): void {
const cutoff = nowMs - TRANSPORT_EVIDENCE_WINDOW_MS - 2 * REFUSAL_POLICIES.transport.maxMs;
while (transportFailures.length > 0 && transportFailures[0].at < cutoff) {
transportFailures.shift();
}
while (transportSuccesses.length > 0 && transportSuccesses[0].at < cutoff) {
transportSuccesses.shift();
}
}

/** Record one final tagged transport failure. Never throws, never writes refusal memory. */
export function recordTransportFailure(
key: string | null,
destination: string | null,
nowMs: number = Date.now()
): void {
if (key === null || destination === null || destination === "") return;
purgeTransportEvidence(nowMs);
transportFailures.push({ key, destination, at: nowMs });
while (transportFailures.length > MAX_TRANSPORT_EVIDENCE) transportFailures.shift();
}

/** Record one success to a destination through an egress. Never throws. */
export function recordTransportSuccess(
destination: string | null,
key: string | null,
nowMs: number = Date.now()
): void {
if (key === null || destination === null || destination === "") return;
purgeTransportEvidence(nowMs);
transportSuccesses.push({ destination, key, at: nowMs });
while (transportSuccesses.length > MAX_TRANSPORT_EVIDENCE) transportSuccesses.shift();
}

/**
* True when the evidence condemns this egress for this destination: at least
* TRANSPORT_EVIDENCE_THRESHOLD tagged failures through it inside the window
* AND at least one success to the same destination through a different egress
* inside the window. A success through the egress itself is not evidence.
*/
export function hasTransportCrossEvidence(
key: string | null,
destination: string | null,
nowMs: number = Date.now()
): boolean {
if (key === null || destination === null || destination === "") return false;
purgeTransportEvidence(nowMs);
const from = nowMs - TRANSPORT_EVIDENCE_WINDOW_MS;
let failures = 0;
for (const f of transportFailures) {
if (f.key === key && f.destination === destination && f.at >= from) {
failures++;
if (failures >= TRANSPORT_EVIDENCE_THRESHOLD) break;
}
}
if (failures < TRANSPORT_EVIDENCE_THRESHOLD) return false;
return transportSuccesses.some(
(s) => s.destination === destination && s.key !== key && s.at >= from
);
}

/** Test-only: forget transport evidence (refusal memory is separate). */
export function __resetTransportEvidenceForTesting(): void {
transportFailures.length = 0;
transportSuccesses.length = 0;
}

/** Test-only: current evidence store sizes. */
export function __transportEvidenceSizeForTesting(): { failures: number; successes: number } {
return { failures: transportFailures.length, successes: transportSuccesses.length };
}

/** Test-only: number of (key, kind) entries held. */
export function __proxyRefusalMemorySizeForTesting(): number {
return memory.size;
Expand Down
2 changes: 1 addition & 1 deletion open-sse/utils/proxyTransitionListeners.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
* in this file so the store stays pure and cycle-free.
*/

export type ProxyTransitionKind = "ip_quota_429" | "proxy_unreachable";
export type ProxyTransitionKind = "ip_quota_429" | "proxy_unreachable" | "transport";

export interface ProxyTransition {
key: string;
Expand Down
68 changes: 68 additions & 0 deletions open-sse/utils/proxyTransportOutcome.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
/**
* Transport-failure outcome feedback for proxy selection (#14802). The proxy
* dispatcher (proxyFetch) calls the two record hooks at its single success
* return and its single final throw; the evidence itself is kept by the pure
* store in proxyRefusalMemory.ts. Everything here is opt-in: with
* PROXY_SKIP_RECENTLY_FAILED off nothing is recorded (no hot-path cost) and
* nothing is ever set aside.
*/
import { isProxySkipRecentlyFailedEnabled } from "@/shared/utils/featureFlags";
import {
hasTransportCrossEvidence,
noteProxyRefusal,
proxyEgressKey,
recordTransportFailure,
recordTransportSuccess,
transportDestinationKey,
} from "./proxyRefusalMemory.ts";

/**
* Decide a tagged final transport failure (errorCode "proxy_unreachable" on the
* thrown sanitized error). Opt-in: only with PROXY_SKIP_RECENTLY_FAILED on,
* and only with cross-evidence (repeated tagged failures through this egress
* plus a real success to the same destination through a different egress),
* does the member get set aside under the "transport" kind. A null key (direct
* egress, edge relay), an explicit single-member pool, or missing evidence
* writes nothing. The failure itself must already be recorded via
* recordTransportFailure before calling; a retried-then-recovered attempt is
* never recorded, so only final failures count.
*
* The dispatcher does not know the pool size, so on that path a lone member is
* protected by the cross-evidence rule itself: without a success through a
* different egress to the same destination nothing is ever written.
*/
export function noteTransportOutcome(args: {
key: string | null;
destination: string | null;
poolSize?: number;
nowMs?: number;
}): void {
const { key, destination, poolSize, nowMs = Date.now() } = args;
if (key === null) return;
if (typeof poolSize === "number" && poolSize <= 1) return;
if (!isProxySkipRecentlyFailedEnabled()) return;
if (hasTransportCrossEvidence(key, destination, nowMs)) noteProxyRefusal(key, "transport", nowMs);
}

/** Record a final tagged transport failure as evidence, then decide. Best effort, never throws. */
export function recordFinalTransportOutcome(proxyUrl: string | null, targetUrl: string): void {
try {
if (!isProxySkipRecentlyFailedEnabled()) return;
const key = proxyEgressKey(proxyUrl);
const destination = transportDestinationKey(targetUrl);
recordTransportFailure(key, destination);
noteTransportOutcome({ key, destination });
} catch {
/* evidence is best-effort; never break the request path */
}
}

/** Record a proxied success as cross-evidence. Best effort, never throws. */
export function recordProxiedSuccess(proxyUrl: string | null, targetUrl: string): void {
try {
if (!isProxySkipRecentlyFailedEnabled()) return;
recordTransportSuccess(transportDestinationKey(targetUrl), proxyEgressKey(proxyUrl));
} catch {
/* evidence is best-effort; never break the request path */
}
}
2 changes: 1 addition & 1 deletion src/lib/events/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ export interface CredentialHealthChangedPayload {
}

export interface ProxySetAsidePayload {
reason: "ip_quota_429" | "proxy_unreachable";
reason: "ip_quota_429" | "proxy_unreachable" | "transport";
setAsideUntil: string;
durationMs: number;
egressKeyMasked: string;
Expand Down
5 changes: 5 additions & 0 deletions src/sse/handlers/proxyOutcomeMemory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,3 +40,8 @@ export function noteProxyOutcome(
if (inRefusalScope) noteProxyServed(key);
else noteProxyRecovered(key, "proxy_unreachable");
}

// The transport cross-evidence decision lives next to its store in open-sse (the
// proxy dispatcher calls it without reaching into src/sse); re-exported here so
// the outcome feedback for proxies stays discoverable from one place.
export { noteTransportOutcome } from "@omniroute/open-sse/utils/proxyTransportOutcome.ts";
Loading
Loading