Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
e726708
fix(combo): skip provider cooldown for per-model-quota providers on 500
hartmark Jul 3, 2026
0af0ded
fix(combo): don't model-lockout Gemini on500 — sibling retry
hartmark Jul 3, 2026
f8cc028
fix(combo): pinned model must fall through to retry/sibling on500
hartmark Jul 3, 2026
efdb37e
fix(combo): pinned model must validate quality + fall through on failure
hartmark Jul 3, 2026
a3790bf
fix(combo): locked stream quality check must fail, not pass
hartmark Jul 3, 2026
8644635
fix: preserve X-Correlation-Id through earlyStreamKeepalive slow path
hartmark Jul 4, 2026
bc1e9d7
feat(combo): invalidate context-cache pins on combo edit + visual pin…
hartmark Jul 4, 2026
e63ec40
fix(combo): clone response before quality check to prevent locked-str…
hartmark Jul 4, 2026
431b149
fix(stream): only emit error frame if no content forwarded yet
hartmark Jul 4, 2026
7cde704
fix(gemini): map malformed_response finish reason to content_filter
hartmark Jul 4, 2026
760072d
fix(api): correlationId filter uses substring matching
hartmark Jul 4, 2026
b503472
feat(dashboard): request log UI improvements
hartmark Jul 4, 2026
d14197f
test(gemini): improve live test reliability and diagnostics
hartmark Jul 4, 2026
0836da4
test(gemini): add malformed_response finish reason detection test
hartmark Jul 4, 2026
8773d61
test: remove stale gemini-malformed-response test (function moved to …
hartmark Jul 4, 2026
ecbf74d
fix(combo): resolve rawModel reference error in round-robin cooldown …
hartmark Jul 4, 2026
9ec3019
test(combo): fix context-cache pin test — seed anthropic connection f…
hartmark Jul 4, 2026
6042536
fix(combo): streaming 500s, SSE corruption, Gemini malformed-response…
hartmark Jul 5, 2026
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.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

### 🐛 Bug Fixes

- **combo/streaming: fix intermittent 500s, corrupted SSE, and Gemini malformed-response handling ([#5976](https://github.com/diegosouzapw/OmniRoute/issues/5976)).** Five related fixes on the combo streaming path: (1) the quality check now reads a **clone** of the response so the original stream stays unlocked (kills the random `ERR_INVALID_STATE: ReadableStream is locked` 500s), and the abandoned clone branch is cancelled so it no longer buffers the whole body per request; (2) `withEarlyStreamKeepalive` only emits an in-band error frame when **no** bytes were forwarded yet, so a mid-flight upstream drop no longer corrupts a partially-delivered SSE stream; (3) Gemini `finishReason: "MALFORMED_RESPONSE"` maps to OpenAI `content_filter` and triggers combo failover instead of returning broken-but-"successful" output; (4) `/api/usage/call-logs?correlationId=` uses parameterized substring (`LIKE`) matching so partial IDs resolve; (5) per-model 500s skip model-lockout/cooldown for per-model-quota providers. Plus request-logger UI detail improvements. Regression guards: `tests/unit/{earlyStreamKeepalive,finishReason,combo-provider-cooldown-sibling,combo-context-relay,call-logs-correlation-substring,save-call-log-persistence,validate-response-quality}.test.ts`. ([#6216](https://github.com/diegosouzapw/OmniRoute/pull/6216) — thanks @hartmark)
- **chatcore (tools): stop the default 128-tool cap from silently dropping opencode's `task`/MCP tools.** opencode (used as an MCP/agent host) sends a large tool list; when it exceeds the speculative `MAX_TOOLS_LIMIT` (128) default, `truncateToolList` did a blind `tools.slice(0, 128)`, dropping every tool past index 128 — including opencode's built-in `task` tool (subagent launch) and many MCP tools, so models routed through OmniRoute could no longer spawn subagents or reach part of their tools. The cap exists to avoid upstream `400`s for providers with real hard limits (e.g. grok-cli 200), so it is kept for those: detection of the opencode client (`isOpencodeClient` — any `x-opencode-*` header, or `opencode` in the user-agent) now only bypasses the **speculative 128 default**, never a known provider ceiling. Precedence is explicit — a proactive/detected provider limit always truncates (even for opencode); otherwise opencode forwards its full tool list; otherwise the unchanged 128 default applies to every other client. Refactors `getEffectiveToolLimit` into `getKnownToolLimit(provider) ?? DEFAULT_LIMIT` (byte-identical for existing callers) and fixes a cosmetic debug-log that reported the truncated count instead of the original. Regression guard: `tests/unit/tool-limit-detector.test.ts`.

- **fix(mitm):** the macOS MITM-cert install check now matches the system keychain again. `security find-certificate -a -Z` prints the SHA-1 as a colon-less hex string, but the installed-check compared it against `getCertFingerprint()`'s colon-separated form, so the substring match never hit — the cert was reported as not-installed and re-prompted for the sudo install on every run. Fingerprints are now normalized (colons stripped, upper-cased) on both sides via the extracted `macCertOutputHasFingerprint` helper. Regression guard: `tests/unit/mitm-cert-mac-fingerprint.test.ts`. ([#6204](https://github.com/diegosouzapw/OmniRoute/pull/6204), closes [#6134](https://github.com/diegosouzapw/OmniRoute/issues/6134) — thanks @rianonehub)
Expand Down
2 changes: 2 additions & 0 deletions open-sse/handlers/chatCore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -374,6 +374,7 @@ export async function handleChatCore({
skipUpstreamRetry = false,
createPiiTransform = null,
correlationId = null,
modelPinned = false,
}) {
let { provider, model, extendedContext } = modelInfo;
// ── Memory pressure guard ────────────────────────────────────────────
Expand Down Expand Up @@ -806,6 +807,7 @@ export async function handleChatCore({
apiKeyInfo,
noLogEnabled,
correlationId,
modelPinned,
});

// Primary path: merge client model id + alias target so config on either key applies; resolved
Expand Down
3 changes: 3 additions & 0 deletions open-sse/handlers/chatCore/attemptLogging.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ export type PersistAttemptLogsContext = {
apiKeyInfo: { id?: string | null; name?: string | null } | null | undefined;
noLogEnabled: unknown;
correlationId?: string | null;
modelPinned?: boolean;
};

function toConnectionId(value: unknown): string | null {
Expand Down Expand Up @@ -108,6 +109,7 @@ export function persistAttemptLogs(args: PersistAttemptLogsArgs, ctx: PersistAtt
apiKeyInfo,
noLogEnabled,
correlationId,
modelPinned,
} = ctx;
const initialConnectionId = toConnectionId(connectionId);
const finalConnectionId = toConnectionId(credentials?.connectionId) || initialConnectionId;
Expand Down Expand Up @@ -203,5 +205,6 @@ export function persistAttemptLogs(args: PersistAttemptLogsArgs, ctx: PersistAtt
noLog: noLogEnabled,
pipelinePayloads,
correlationId,
modelPinned: modelPinned || false,
}).catch(() => {});
}
104 changes: 90 additions & 14 deletions open-sse/services/combo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
formatRetryAfter,
getModelLockoutInfo,
getRuntimeProviderProfile,
hasPerModelQuota,
isModelLocked,
recordModelLockoutFailure,
recordProviderFailure,
Expand Down Expand Up @@ -107,7 +108,11 @@ import {
getStickyWeightedExecutionKey,
recordStickyWeightedSuccess,
} from "./combo/rrState.ts";
import { validateResponseQuality, toRetryAfterDisplayValue } from "./combo/validateQuality.ts";
import {
validateResponseQuality,
releaseQualityClone,
toRetryAfterDisplayValue,
} from "./combo/validateQuality.ts";
import { resolveComboCooldownWaitDecision } from "./combo/comboCooldownRetry.ts";
import {
computeClosestRetryAfter,
Expand Down Expand Up @@ -705,15 +710,59 @@ export async function handleComboChat({
"COMBO",
`Bypassing strategy — routing directly to pinned context model: ${pinnedModel}`
);
return handleSingleModelWithTimeout(body, pinnedModel);
let pinnedResult: Response | null = null;
try {
pinnedResult = await handleSingleModelWithTimeout(body, pinnedModel, {
modelPinned: true,
} as SingleModelTarget);
} catch (pinErr) {
log.warn(
"COMBO",
`Pinned model ${pinnedModel} threw error: ${pinErr instanceof Error ? pinErr.message : String(pinErr)}, falling through to combo retry/fallback`
);
}
if (pinnedResult) {
if (pinnedResult.ok) {
let pinnedClone: Response;
try {
pinnedClone = pinnedResult.clone();
} catch {
pinnedClone = pinnedResult;
}
const pinnedQuality = await validateResponseQuality(
pinnedClone,
clientRequestedStream,
log,
config.responseValidation
);
releaseQualityClone(pinnedClone, pinnedResult, pinnedQuality);
if (pinnedQuality.valid) return pinnedResult;
log.warn(
"COMBO",
`Pinned model ${pinnedModel} returned 200 but failed quality check: ${pinnedQuality.reason}, falling through to combo retry/fallback`
);
} else {
const pinnedStatus = pinnedResult.status || 500;
if (![408, 429, 500, 502, 503, 504].includes(pinnedStatus)) {
return pinnedResult;
}
log.warn(
"COMBO",
`Pinned model ${pinnedModel} failed (${pinnedStatus}), falling through to combo retry/fallback`
);
}
}
// Fall through to the target iteration loop below — retries and sibling
// models will be tried via the normal combo machinery.
}
log.warn(
"COMBO",
pinInCombo
? `Context-cache pin "${pinnedModel}" provider durably unhealthy — dropping pin, using strategy`
: `Stale context-cache pin "${pinnedModel}" not in combo "${combo.name}" targets — dropping pin, using strategy`
);
return handleSingleModelWithTimeout(body, pinnedModel);
// Fall through to the normal target iteration loop below — the pin is
// dropped, so the combo strategy picks the best available target.
}

// Fusion strategy: parallel panel + judge synthesis. Handled in a separate module
Expand Down Expand Up @@ -1504,12 +1553,22 @@ export async function handleComboChat({
undefined;
const effectiveConnectionId = selectedConnectionId || target.connectionId || "";

// Clone BEFORE quality check — validateResponseQuality reads the body
// via getReader() which locks the stream. The clone's body is consumed
// by the quality check; the original stays unlocked for piping.
let qualityClone: Response;
try {
qualityClone = result.clone();
} catch {
qualityClone = result;
}
const quality = await validateResponseQuality(
result,
qualityClone,
clientRequestedStream,
log,
config.responseValidation
);
releaseQualityClone(qualityClone, result, quality);
if (!quality.valid) {
log.warn(
"COMBO",
Expand Down Expand Up @@ -1744,7 +1803,7 @@ export async function handleComboChat({
})();
}

return { ok: true, response: quality.clonedResponse ?? result };
return { ok: true, response: result };
}

// Extract error info from response
Expand Down Expand Up @@ -2023,7 +2082,16 @@ export async function handleComboChat({
}
log.warn("COMBO", `Model ${modelStr} failed, trying next`, { status: result.status });

if (resilienceSettings.providerCooldown.enabled && provider && provider !== "unknown") {
// #5976: per-model-quota providers (Gemini, GitHub, etc.) multiplex models
// behind one connection. A model-level 500 must NOT cool down the entire
// provider — sibling models may still succeed. Skip cooldown recording for
// these providers on 500 errors so the next target can try.
if (
resilienceSettings.providerCooldown.enabled &&
provider &&
provider !== "unknown" &&
!(result.status === 500 && hasPerModelQuota(provider, rawModel))
) {
recordProviderCooldown(
provider,
targetWithConnection.connectionId ?? undefined,
Expand Down Expand Up @@ -2572,12 +2640,19 @@ async function handleRoundRobinCombo({

// Success — validate response quality before returning
if (result.ok) {
let rrClone: Response;
try {
rrClone = result.clone();
} catch {
rrClone = result;
}
const quality = await validateResponseQuality(
result,
rrClone,
clientRequestedStream,
log,
config.responseValidation
);
releaseQualityClone(rrClone, result, quality);
if (!quality.valid) {
log.warn(
"COMBO-RR",
Expand Down Expand Up @@ -2665,12 +2740,8 @@ async function handleRoundRobinCombo({
}
})();
}
// validateResponseQuality peeks streaming bodies via getReader(),
// which locks `result.body`. It returns a clonedResponse that replays
// the buffered prefix and forwards the rest. Returning the original
// (now-locked) `result` makes Next.js throw "ReadableStream is locked"
// → 500. Mirror the priority strategy and return the replay response.
return quality.clonedResponse ?? result;
// Clone is consumed by quality check; original stays unlocked.
return result;
}

// Extract error info
Expand Down Expand Up @@ -2847,7 +2918,12 @@ async function handleRoundRobinCombo({
if (offset > 0) fallbackCount++;
log.warn("COMBO-RR", `${modelStr} failed, trying next model`, { status: result.status });

if (resilienceSettings.providerCooldown.enabled && provider && provider !== "unknown") {
if (
resilienceSettings.providerCooldown.enabled &&
provider &&
provider !== "unknown" &&
!(result.status === 500 && hasPerModelQuota(provider, parseModel(modelStr).model || modelStr))
) {
recordProviderCooldown(
provider,
targetWithConnection.connectionId ?? undefined,
Expand Down
13 changes: 10 additions & 3 deletions open-sse/services/combo/runtimeUnits.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import { errorResponse } from "../../utils/error.ts";
import { recordComboRequest } from "../comboMetrics.ts";
import { resolveDelayMs } from "./comboPredicates.ts";
import { validateResponseQuality } from "./validateQuality.ts";
import { validateResponseQuality, releaseQualityClone } from "./validateQuality.ts";
import type { ResponseValidationConfig } from "./responseValidation.ts";
import type {
ComboCollectionLike,
Expand Down Expand Up @@ -228,12 +228,19 @@ export async function executeRuntimeUnitCombo(args: {
});
return { response, unit };
}
let unitClone: Response;
try {
unitClone = response.clone();
} catch {
unitClone = response;
}
const quality = await validateResponseQuality(
response,
unitClone,
clientRequestedStream,
args.log,
args.config.responseValidation as ResponseValidationConfig | undefined
);
releaseQualityClone(unitClone, response, quality);
if (quality.valid) {
recordComboRequest(args.combo.name, unit.modelStr, {
success: true,
Expand All @@ -242,7 +249,7 @@ export async function executeRuntimeUnitCombo(args: {
strategy: effectiveStrategy,
target: { executionKey: unit.executionKey, stepId: unit.stepId, label: unit.label },
});
return { response: quality.clonedResponse ?? response, unit };
return { response, unit };
}
}
if (![408, 429, 500, 502, 503, 504].includes(response.status)) break;
Expand Down
2 changes: 2 additions & 0 deletions open-sse/services/combo/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,8 @@ export type SingleModelTarget =
allowRateLimitedConnection?: boolean;
effectiveComboStrategy?: string | null;
modelAbortSignal?: AbortSignal | null;
/** True when this target was selected via context-cache session pinning. */
modelPinned?: boolean;
})
| { modelAbortSignal: AbortSignal };

Expand Down
66 changes: 58 additions & 8 deletions open-sse/services/combo/validateQuality.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,7 @@ import {
createSSEDataLineNormalizer,
isKnownNonClaudeStreamPayload,
} from "../../utils/streamHelpers.ts";
import {
evaluateResponseValidation,
type ResponseValidationConfig,
} from "./responseValidation.ts";
import { evaluateResponseValidation, type ResponseValidationConfig } from "./responseValidation.ts";
import { getReasoningTokens } from "../../../src/lib/usage/tokenAccounting.ts";
import type { ComboRetryAfter } from "./types.ts";

Expand Down Expand Up @@ -99,6 +96,7 @@ export async function validateResponseQuality(
let hasMessageStart = false;
let hasContentBlock = false;
let hasLifecycleEnd = false;
let anyContentFound = false;
const sseLineNormalizer = createSSEDataLineNormalizer();
let pendingEventType = "";

Expand Down Expand Up @@ -236,6 +234,18 @@ export async function validateResponseQuality(
return { valid: false, reason: "streaming empty content block" };
}

// Non-Claude stream with no recognizable content at all — the stream
// ended without any content deltas (e.g. Gemini returning HTTP 200
// with an empty body or only metadata chunks). Mark as invalid for
// combo failover so the sibling model gets tried.
if (!anyContentFound && !hasContentBlock) {
log.warn?.(
"COMBO",
"Streaming response ended with no recognized content — marking as invalid for combo failover"
);
return { valid: false, reason: "streaming no recognized content" };
}

// Incomplete lifecycle or non-Claude stream — replay all buffered
// bytes. The reader is exhausted so the forwarding reader will
// immediately signal done.
Expand All @@ -251,6 +261,7 @@ export async function validateResponseQuality(
const foundContent = parseAccumulatedSse();

if (foundContent) {
anyContentFound = true;
// A content_block_* event was found — stop peeking. Return a
// clonedResponse that replays all buffered bytes (the current chunk
// is already in bufferedChunks) and then forwards the remainder of
Expand All @@ -259,9 +270,23 @@ export async function validateResponseQuality(
return { valid: true, clonedResponse };
}
}
} catch {
// If reading the stream fails, pass through — other mechanisms
// (stream readiness timeout) will catch truly broken streams.
} catch (streamErr) {
// If reading the stream fails due to a locked stream or pipe error,
// the content cannot be verified — mark as invalid for combo failover.
// A locked ReadableStream means the response body is already consumed
// or corrupted (e.g. "Invalid state: The ReadableStream is locked").
// Broad match: Chrome/V8 throws "body used already", Firefox throws
// "ReadableStream is locked", etc.
const errMsg = streamErr instanceof Error ? streamErr.message : String(streamErr);
if (
streamErr instanceof TypeError &&
(errMsg.includes("locked") ||
errMsg.includes("disturbed") ||
errMsg.includes("used already"))
) {
return { valid: false, reason: "stream locked or disturbed" };
}
// Other read errors — pass through (stream readiness timeout will catch truly broken streams)
return { valid: true };
}
}
Expand Down Expand Up @@ -308,7 +333,8 @@ export async function validateResponseQuality(

const choices = json?.choices;
if (json?.object === "response") {
if (!responsesApiOutputHasContent(json.output)) return { valid: false, reason: "empty_choices" };
if (!responsesApiOutputHasContent(json.output))
return { valid: false, reason: "empty_choices" };
const status = typeof json.status === "string" ? json.status : "";
if (status && !["completed", "done"].includes(status)) {
return { valid: false, reason: "no_terminal" };
Expand Down Expand Up @@ -389,3 +415,27 @@ export async function validateResponseQuality(
}),
};
}

/**
* Release the peek-and-abandon clone used by {@link validateResponseQuality}.
*
* The quality check clones the upstream response, reads the clone only until the
* first content block, then hands back a `clonedResponse` that callers on the
* streaming path DISCARD (they forward the original, untouched response). Because
* a `Response.clone()` tees the body, that abandoned branch would otherwise buffer
* the entire remaining body in memory until the original finishes streaming.
*
* Cancelling the abandoned branch releases that buffer. Per the ReadableStream tee
* contract, cancelling one branch does NOT cancel the shared source while the other
* branch (the original response being streamed to the client) is still active, so
* this is safe. No-op when the clone fell back to the original (clone unsupported)
* or when quality reading already exhausted the body (no `clonedResponse`).
*/
export function releaseQualityClone(
clone: Response,
original: Response,
quality: { clonedResponse?: Response }
): void {
if (clone === original) return;
void quality.clonedResponse?.body?.cancel().catch(() => {});
}
Loading
Loading