Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
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
6 changes: 6 additions & 0 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,12 @@ LABEL org.opencontainers.image.title="omniroute" \
ENV NODE_ENV=production
ENV PORT=20128
ENV HOSTNAME=0.0.0.0
# Runtime heap ceiling. 1024MB is enough for normal traffic but can be tight
# for large fusion-combo panels (many models fanned out in parallel, each
# response buffered in full — see open-sse/services/fusion.ts::FUSION_DEFAULTS
# .maxPanel, issue #1905). Override at `docker run` time with
# `-e OMNIROUTE_MEMORY_MB=2048` (or higher) if you raise fusionTuning.maxPanel
# above the default cap.
ENV OMNIROUTE_MEMORY_MB=1024
ENV NODE_OPTIONS="--max-old-space-size=${OMNIROUTE_MEMORY_MB}"

Expand Down
1 change: 1 addition & 0 deletions changelog.d/fixes/1382-streaming-empty-content-block.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(combo):** streaming Claude responses whose content block opens (`content_block_start`) and closes with no usable text/tool_use — a shape some upstreams return for tool-heavy requests on HTTP 200 — are now detected by `validateResponseQuality`'s SSE peek and trigger combo failover instead of being forwarded to the client as a silent empty completion (thanks @heishen6).
1 change: 1 addition & 0 deletions changelog.d/fixes/1809-mitm-stop-dns-before-kill.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(cli):** `stopMitm()` now removes /etc/hosts DNS-spoof entries before killing the MITM server process, closing the window where a client's DNS still resolved a target host to `127.0.0.1` while nothing was listening there — the cause of `connect ECONNREFUSED 127.0.0.1:443` right after stopping the MITM proxy (thanks @dionisius95).
1 change: 1 addition & 0 deletions changelog.d/fixes/1905-fusion-panel-oom.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(combos):** fusion combos now reject an oversized panel (>40 models by default, tunable via `fusionTuning.maxPanel`) with a clean 400 before fanning out, instead of buffering dozens of concurrent full responses in memory and OOM-crashing the whole container. (thanks @fontvu)
29 changes: 24 additions & 5 deletions open-sse/services/autoCombo/__tests__/speedRanking.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,15 +165,34 @@ describe("rankBySpeed — factor breakdown", () => {
});

it("falls back to 0.5 per missing metric so new providers are not crushed", () => {
const ranked = rankBySpeed([candidate({ provider: "fresh", model: "m" })]);
const ranked = rankBySpeed([
candidate({
provider: "fresh",
model: "m",
p95LatencyMs: undefined,
latencyStdDev: undefined,
}),
]);
expect(ranked).toHaveLength(1);
// No telemetry at all → weighted sum lands near 0.5 with reliability multiplier 1
expect(ranked[0].factors.reliability).toBe(1);
expect(ranked[0].factors.health).toBe(1);
expect(ranked[0].factors.ttft).toBe(0.5);
expect(ranked[0].factors.tps).toBe(0.5);
});
});
expect(ranked[0].factors.tps).toBe(0.5);
});

it("uses p95 latency when TTFT and E2E telemetry are unavailable", () => {
const ranked = rankBySpeed([
candidate({ provider: "slow-tail", model: "m", p95LatencyMs: 4000 }),
candidate({ provider: "fast-tail", model: "m", p95LatencyMs: 1000 }),
]);
const fast = ranked.find((entry) => entry.provider === "fast-tail");
const slow = ranked.find((entry) => entry.provider === "slow-tail");

expect(fast?.factors.ttft).toBeGreaterThan(slow?.factors.ttft ?? 1);
expect(fast?.factors.e2e).toBeGreaterThan(slow?.factors.e2e ?? 1);
});
});

describe("rankBySpeed — weight overrides", () => {
it("respects caller weight overrides (e.g. heavy TTFT bias)", () => {
Expand Down Expand Up @@ -223,4 +242,4 @@ describe("pickFastest", () => {
const winner = pickFastest([slow, fast]);
expect(winner?.provider).toBe("fast");
});
});
});
10 changes: 8 additions & 2 deletions open-sse/services/autoCombo/speedRanking.ts
Original file line number Diff line number Diff line change
Expand Up @@ -211,9 +211,15 @@ function speedFactorsFor(
failureRate: number
): SpeedFactors {
return {
ttft: lowerIsBetter(positiveFinite(candidate.avgTtftMs), maxima.ttft),
ttft: lowerIsBetter(
positiveFinite(candidate.avgTtftMs) ?? positiveFinite(candidate.p95LatencyMs),
maxima.ttft
),
tps: higherIsBetter(positiveFinite(candidate.avgTokensPerSecond), maxima.tps),
e2e: lowerIsBetter(positiveFinite(candidate.avgE2ELatencyMs), maxima.e2e),
e2e: lowerIsBetter(
positiveFinite(candidate.avgE2ELatencyMs) ?? positiveFinite(candidate.p95LatencyMs),
maxima.e2e
),
p95: lowerIsBetter(positiveFinite(candidate.p95LatencyMs), maxima.p95),
health: healthScoreFor(candidate.circuitBreakerState),
reliability: clamp01(1 - failureRate),
Expand Down
151 changes: 116 additions & 35 deletions open-sse/services/combo/validateQuality.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,91 @@ function extractEnvelopeErrorText(json: Record<string, unknown>): string | null
return parts.length > 0 ? parts.join(" ") : null;
}

/** Mutable lifecycle flags threaded through {@link applySseLifecycleEvent}. */
interface SseLifecycleFlags {
hasMessageStart: boolean;
hasContentBlock: boolean;
hasRealContent: boolean;
hasLifecycleEnd: boolean;
}

/** Read `parsed.<key>` as a nested object bag, or null when absent/not an object. */
function asObject(parsed: Record<string, unknown>, key: string): Record<string, unknown> | null {
const value = parsed[key];
return value && typeof value === "object" ? (value as Record<string, unknown>) : null;
}

/**
* A content_block_start is real signal only for tool_use / redacted_thinking —
* a tool call is meaningful even before its input_json_delta arrives. text and
* thinking blocks routinely open empty; keep peeking for a delta instead.
*/
function contentBlockStartIsRealSignal(parsed: Record<string, unknown>): boolean {
const blockType = asObject(parsed, "content_block")?.type;
return blockType === "tool_use" || blockType === "redacted_thinking";
}

/**
* A content_block_delta is real signal when it carries non-empty text/thinking,
* or any input_json_delta fragment — even an empty-string first chunk proves a
* tool_use block is actively streaming its arguments.
*/
function contentBlockDeltaIsRealSignal(parsed: Record<string, unknown>): boolean {
const delta = asObject(parsed, "delta");
if (!delta) return false;
const deltaType = typeof delta.type === "string" ? delta.type : "";
if (deltaType === "input_json_delta") return true;
if (deltaType !== "text_delta" && deltaType !== "thinking_delta") return false;
const text = delta.text ?? delta.thinking;
return typeof text === "string" && text.length > 0;
}

/** A message_delta closes the lifecycle once it carries a stop_reason. */
function messageDeltaEndsLifecycle(parsed: Record<string, unknown>): boolean {
return asObject(parsed, "delta")?.stop_reason != null;
}

/**
* Apply a single parsed Claude SSE event to the peeked lifecycle `flags`
* (mutated in place). Extracted from `parseAccumulatedSse`'s inline switch to
* keep that function under the complexity/line ratchets — logic unchanged.
*
* Returns true once REAL content (not just an empty content_block_start) is
* detected — the caller should stop peeking and treat the stream as non-empty.
*/
function applySseLifecycleEvent(
eventType: string,
parsed: Record<string, unknown>,
flags: SseLifecycleFlags
): boolean {
switch (eventType) {
case "message_start":
flags.hasMessageStart = true;
return false;
case "content_block_start":
flags.hasContentBlock = true;
if (!contentBlockStartIsRealSignal(parsed)) return false;
flags.hasRealContent = true;
return true;
case "content_block_delta":
flags.hasContentBlock = true;
if (!contentBlockDeltaIsRealSignal(parsed)) return false;
flags.hasRealContent = true;
return true;
case "content_block_stop":
flags.hasContentBlock = true;
return false;
case "message_stop":
flags.hasLifecycleEnd = true;
return false;
case "message_delta":
if (messageDeltaEndsLifecycle(parsed)) flags.hasLifecycleEnd = true;
return false;
default:
return false;
}
}

function responsesApiOutputHasContent(output: unknown): boolean {
return (
Array.isArray(output) &&
Expand Down Expand Up @@ -125,9 +210,22 @@ export async function validateResponseQuality(
let decodedSoFar = "";

// SSE lifecycle state.
let hasMessageStart = false;
let hasContentBlock = false;
let hasLifecycleEnd = false;
//
// #1382: hasContentBlock only means "a content_block_* event was observed"
// — it does NOT mean the block carried usable content. A content_block_start
// for a text/thinking block routinely opens with empty text (real content
// arrives via subsequent content_block_delta events); some upstreams
// (reported: DeepSeek/GLM via claude→openai translation on tool-heavy
// requests) open and close such a block without ever emitting a delta.
// hasRealContent tracks whether we've actually seen usable output: a
// tool_use/redacted_thinking block start (self-evidently real, even before
// any delta), or a delta carrying non-empty text/thinking/tool-input.
const sse: SseLifecycleFlags = {
hasMessageStart: false,
hasContentBlock: false,
hasRealContent: false,
hasLifecycleEnd: false,
};
let anyContentFound = false;
let sawAnyBytes = false;
const sseLineNormalizer = createSSEDataLineNormalizer();
Expand All @@ -138,8 +236,9 @@ export async function validateResponseQuality(
* flags in the closure. The last (potentially incomplete) line is kept in
* `decodedSoFar` for the next iteration.
*
* Returns true when a content_block_* event is detected — the caller
* should stop peeking and treat the stream as non-empty.
* Returns true once REAL content (not just an empty content_block_start)
* is detected — the caller should stop peeking and treat the stream as
* non-empty.
*/
function parseAccumulatedSse(): boolean {
const lines = decodedSoFar.split(/\r?\n/);
Expand Down Expand Up @@ -177,32 +276,8 @@ export async function validateResponseQuality(
return true;
}

switch (eventType) {
case "message_start":
hasMessageStart = true;
break;
case "content_block_start":
case "content_block_delta":
case "content_block_stop":
hasContentBlock = true;
// Signal caller to stop buffering immediately.
return true;
case "message_stop":
hasLifecycleEnd = true;
break;
case "message_delta": {
const delta = parsed.delta;
if (
delta &&
typeof delta === "object" &&
(delta as Record<string, unknown>).stop_reason != null
) {
hasLifecycleEnd = true;
}
break;
}
default:
break;
if (applySseLifecycleEvent(eventType, parsed, sse)) {
return true;
}
}
return false;
Expand Down Expand Up @@ -258,11 +333,17 @@ export async function validateResponseQuality(
if (decodedSoFar.trim()) decodedSoFar += "\n\n";
parseAccumulatedSse();

if (hasMessageStart && hasLifecycleEnd && !hasContentBlock) {
// Complete Claude lifecycle with zero content blocks → failover.
if (sse.hasMessageStart && sse.hasLifecycleEnd && !sse.hasRealContent) {
// Complete Claude lifecycle with zero content blocks, or with
// content_block_start/stop pairs that never carried real text/
// thinking/tool_use content (#1382 — tool-heavy claude→openai
// requests against upstreams like DeepSeek/GLM can "complete" a
// lifecycle around an empty block) → failover.
log.warn?.(
"COMBO",
"Streaming Claude response has complete lifecycle but zero content blocks (content_filter?) — marking as invalid for combo failover"
sse.hasContentBlock
? "Streaming Claude response has complete lifecycle but its content block(s) carried no usable text/tool_use — marking as invalid for combo failover"
: "Streaming Claude response has complete lifecycle but zero content blocks (content_filter?) — marking as invalid for combo failover"
);
return { valid: false, reason: "streaming empty content block" };
}
Expand All @@ -273,7 +354,7 @@ export async function validateResponseQuality(
// (an explicit `data: [DONE]`, ping/metadata events, an incomplete
// Claude lifecycle) keep the pass-through contract (#3399/#3685):
// those are handled by the stream-readiness timeout, not failover.
if (!anyContentFound && !hasContentBlock && !sawAnyBytes) {
if (!anyContentFound && !sse.hasContentBlock && !sawAnyBytes) {
log.warn?.(
"COMBO",
"Streaming response ended with no recognized content — marking as invalid for combo failover"
Expand Down
23 changes: 23 additions & 0 deletions open-sse/services/fusion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,12 +27,20 @@ export const FUSION_DEFAULTS = {
minPanel: 2, // answers needed before stragglers get a grace window
stragglerGraceMs: 8000, // wait this long for laggards once quorum is reached
panelHardTimeoutMs: 90000, // absolute cap so one hung model can't stall forever
// Hard cap on panel size (issue #1905). Every panel member is fanned out in
// parallel and its full response text buffered in memory simultaneously —
// with the runtime heap capped (Dockerfile OMNIROUTE_MEMORY_MB, default
// 1024MB), a large panel (reported: ~73 models) with sizable concurrent
// responses can exceed the heap ceiling and OOM-crash the whole process.
// Reject oversized panels up front with a clean 400 instead.
maxPanel: 40,
} as const;

export type FusionTuning = {
minPanel?: number;
stragglerGraceMs?: number;
panelHardTimeoutMs?: number;
maxPanel?: number;
};

type Body = Record<string, unknown>;
Expand Down Expand Up @@ -246,6 +254,21 @@ export async function handleFusionChat({
return handleSingleModel(body, panel[0]);
}

// Reject an oversized panel BEFORE fan-out (issue #1905): fanning out N
// parallel calls and buffering N full response bodies at once is what
// drives the process into an OOM crash, not any one call in isolation.
const maxPanel = tuning?.maxPanel ?? FUSION_DEFAULTS.maxPanel;
if (panel.length > maxPanel) {
log.warn(
"FUSION",
`Combo "${comboName ?? ""}" panel=${panel.length} exceeds maxPanel=${maxPanel} — rejecting before fan-out (#1905)`
);
return errorResponse(
400,
`Fusion panel too large (${panel.length} models, max ${maxPanel}) — reduce the combo's target count or raise fusionTuning.maxPanel`
);
}

const cfg = {
minPanel: tuning?.minPanel ?? FUSION_DEFAULTS.minPanel,
stragglerGraceMs: tuning?.stragglerGraceMs ?? FUSION_DEFAULTS.stragglerGraceMs,
Expand Down
Loading