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.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ _In development — bullets added per PR; finalized at release._
- **fix(combo): round-robin members fail over faster under concurrency saturation via a configurable queue depth** — when a round-robin combo member was saturated, requests sat in the per-model semaphore's **unbounded** queue and only failed over to the next member after the full `queueTimeoutMs` (default 30s) elapsed — so a burst of agentic requests deep-queued one hot member instead of spilling to healthy ones. The per-model semaphore now accepts a bounded queue depth and emits `SEMAPHORE_QUEUE_FULL` once it is full (the round-robin loop already cascades on that code), so a configured low depth fails over immediately. A new `queueDepth` combo-config knob (global default / provider override / per-combo, default **20** for backward compatibility; **0** = never queue → fail over now) is exposed in Settings → Combo Defaults. ([#3872](https://github.com/diegosouzapw/OmniRoute/issues/3872) — thanks @KooshaPari)
- **fix(pricing): align Claude Code (`cc`) pricing with current Anthropic per-MTok rates** — the `cc` provider block in the default pricing table had stale numbers across every Claude 4.x family entry — most visibly, `claude-opus-4-5-20251101` was billed at the deprecated Opus 4.1 rate (`input $15` / `output $75`), and `claude-haiku-4-5-20251001` was at half the current Haiku 4.5 rate. The `cached` (cache hit) and `cache_creation` (5-minute cache write) multipliers were also off across Opus 4.6/4.7/4.8, Sonnet 4.5/4.6, Haiku 4.5, and Fable 5. All eight entries now match the rates Anthropic publishes (input, 5m cache write at 1.25x input, cache hit at 0.1x input, output; reasoning billed at the output rate), so cost accounting on the dashboard and per-request usage events stop under- or over-reporting Claude Code spend. (thanks @chulanpro5)
- **fix(executors): sanitize Anthropic-shape content parts before GitHub Copilot `/chat/completions`** — Claude models on GitHub Copilot driven from clients like Cursor IDE (e.g. `gh/claude-sonnet-4.6`) failed with `Provider returned error: type has to be either 'image_url' or 'text' (reset after 30s)` because the client passed through Anthropic-shape content parts (`tool_use`, `tool_result`, `thinking`) untouched, and the Copilot chat-completions endpoint only accepts `text`/`image_url`. `GithubExecutor.transformRequest` now serializes any unsupported part type as `text` (preserving the model's context), drops empty parts, and collapses to `null` when an assistant message's only content was tool_calls — `tool_calls` ride alongside untouched. Codex-family models still route through `/responses` unchanged. (thanks @cngznNN)
- **fix(sse):** refactor stall detection to reduce false positives on slow but progressing streams. (thanks @zakirkun)

---

Expand Down
137 changes: 130 additions & 7 deletions open-sse/utils/streamHandler.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,22 @@
import { trackPendingRequest } from "@/lib/usageDb";
import { STREAM_IDLE_TIMEOUT_MS } from "../config/constants.ts";
import { FORMATS } from "../translator/formats.ts";
import { PENDING_REQUEST_CLEARED_MARKER } from "./stream.ts";

// Stream handler with disconnect detection - shared for all providers

const DISCONNECT_ABORT_DELAY_MS = 2_000;

// Default budget for the pipeWithDisconnect raw-upstream stall watchdog.
// Inherits STREAM_IDLE_TIMEOUT_MS so a single env knob still governs the
// max time we tolerate silence from upstream. Reasoning models (Claude
// thinking, Kiro EventStream binary frames) emit zero post-transform
// output for long stretches while raw bytes keep arriving — measuring
// stall on the transform output false-positives on those streams, so
// the watchdog must track upstream byte activity instead. Ported from
// decolua/9router#1243.
const DEFAULT_STREAM_STALL_TIMEOUT_MS = STREAM_IDLE_TIMEOUT_MS;

type StreamDisconnectEvent = {
reason: string;
duration: number;
Expand Down Expand Up @@ -391,19 +402,131 @@ export function createDisconnectAwareStream(transformStream, streamController) {
}

/**
* Pipe provider response through transform with disconnect detection
* @param {Response} providerResponse - Response from provider
* @param {TransformStream} transformStream - Transform stream for SSE
* @param {object} streamController - Stream controller from createStreamController
* Pipe provider response through transform with disconnect detection.
*
* Stall watchdog tracks raw upstream byte activity, not transform output.
* Reasoning models (Claude thinking via Kiro, etc.) can produce zero SSE
* output for long stretches while partial EventStream frames keep arriving;
* measuring stall on the transform output caused false stalls. Any upstream
* chunk resets the timer. If no bytes arrive for `stallTimeoutMs`, the
* stream surfaces a "stream stall timeout" error and aborts.
*
* Ported from decolua/9router#1243 by @zakirkun.
*
* @param providerResponse - Response from provider
* @param transformStream - Transform stream for SSE
* @param streamController - Stream controller from createStreamController
* @param opts.stallTimeoutMs - Override the stall budget (defaults to
* STREAM_IDLE_TIMEOUT_MS / DEFAULT_STREAM_STALL_TIMEOUT_MS). `0` disables
* the watchdog.
*/
export function pipeWithDisconnect(
providerResponse: Response,
transformStream: TransformStream<Uint8Array, Uint8Array>,
streamController: StreamController
streamController: StreamController,
opts: { stallTimeoutMs?: number } = {}
) {
const transformedBody = providerResponse.body.pipeThrough(transformStream);
const stallTimeoutMs = opts.stallTimeoutMs ?? DEFAULT_STREAM_STALL_TIMEOUT_MS;

// Watchdog disabled — preserve legacy behavior verbatim.
if (!stallTimeoutMs || stallTimeoutMs <= 0) {
const transformedBody = providerResponse.body.pipeThrough(transformStream);
return createDisconnectAwareStream(
{ readable: transformedBody, writable: { getWriter: () => ({ abort: () => {} }) } },
streamController
);
}

let stallTimer: ReturnType<typeof setTimeout> | null = null;
// Captured on the upstream tap's `start`, used by the watchdog to error the
// pipeline so the downstream reader unblocks and emits a clean SSE error
// event. Without this, aborting the AbortController alone does not unblock
// a `reader.read()` already suspended on the transform pipe — the request
// would hang until the upstream finally closed the socket.
let upstreamTapController: TransformStreamDefaultController<Uint8Array> | null = null;
// Set when the watchdog fires so the downstream pull() catch (which sees
// the same error propagated through the pipeline) does not call
// handleError a second time — pending-cleanup is idempotent but onError
// callbacks should fire once per error.
let stallFired = false;

const clearStall = () => {
if (stallTimer) {
clearTimeout(stallTimer);
stallTimer = null;
}
};
const armStall = () => {
clearStall();
stallTimer = setTimeout(() => {
stallTimer = null;
stallFired = true;
const stallError = new Error("stream stall timeout");
// Notify the controller (onError callback + pending-request cleanup).
try {
streamController.handleError?.(stallError);
} catch {}
// Error the pipeline so the downstream reader unblocks. createDisconnect-
// AwareStream's catch block translates this into buildStreamErrorChunks
// (sanitized SSE error event with finish_reason:"error", per the format).
try {
upstreamTapController?.error(stallError);
} catch {}
// Abort the underlying fetch so upstream releases the connection.
try {
streamController.abort?.();
} catch {}
}, stallTimeoutMs);
};

// Wrap controller so every termination path clears the stall timer.
// Without this, abort/complete/error/disconnect paths leave the timer armed
// and a stale abort could fire after the request has already ended.
const wrappedController: StreamController = {
...streamController,
handleComplete: () => {
clearStall();
streamController.handleComplete();
},
handleError: (e: unknown) => {
clearStall();
// Watchdog already fired its own handleError — the inner pull() catch
// sees the same error propagated through the pipeline; suppress the
// duplicate to keep onError callbacks single-fire.
if (stallFired) return;
streamController.handleError(e);
},
handleDisconnect: (reason?: string) => {
clearStall();
streamController.handleDisconnect(reason);
},
abort: () => {
clearStall();
streamController.abort();
},
};

// Inert tap that resets the stall timer on every raw upstream byte chunk.
// Sits between the provider body and the SSE transform so reasoning models
// that buffer many raw bytes into a single emitted event do not look
// stalled to the watchdog.
const upstreamTap = new TransformStream<Uint8Array, Uint8Array>({
start(controller) {
upstreamTapController = controller;
armStall();
},
transform(chunk, controller) {
armStall();
controller.enqueue(chunk);
},
flush() {
clearStall();
},
});

const transformedBody = providerResponse.body.pipeThrough(upstreamTap).pipeThrough(transformStream);
return createDisconnectAwareStream(
{ readable: transformedBody, writable: { getWriter: () => ({ abort: () => {} }) } },
streamController
wrappedController
);
}
128 changes: 128 additions & 0 deletions tests/unit/stream-handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -468,3 +468,131 @@ test("pipeWithDisconnect does not double-clear transform errors already accounte
assert.equal(pending.byModel[modelKey], 1);
assert.equal(pending.byAccount[connectionId][modelKey], 1);
});

// Stall detection: tied to RAW upstream byte activity, not transform output.
// Ports decolua/9router#1243 — reasoning models (Claude thinking, Kiro
// EventStream binary frames) can stream raw bytes for long stretches while
// the SSE transform produces zero output as it accumulates a frame. The
// stall watchdog must NOT fire on those slow-but-progressing streams.
test("pipeWithDisconnect does NOT flag a slow but progressing upstream as stalled (no false positive)", async () => {
// Upstream emits 3 small chunks 30ms apart (90ms total). The transform
// never forwards any output (simulates a translator buffering a frame
// boundary that has not yet completed). The stall budget is 200ms — well
// above the 30ms gap between upstream bytes, so a byte-activity watchdog
// should never fire. A transform-output-activity watchdog would
// false-stall here.
const source = new ReadableStream({
async start(controller) {
controller.enqueue(encoder.encode("a"));
await new Promise((r) => setTimeout(r, 30));
controller.enqueue(encoder.encode("b"));
await new Promise((r) => setTimeout(r, 30));
controller.enqueue(encoder.encode("c"));
await new Promise((r) => setTimeout(r, 30));
controller.close();
},
});

// Black-hole transform — consumes every byte, emits nothing until flush.
const swallowingTransform = new TransformStream({
transform() {
/* drop chunk — output stream is silent */
},
flush(controller) {
controller.enqueue(encoder.encode("done"));
},
});

let onErrorCalled = false;
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});

const stream = pipeWithDisconnect(
new Response(source),
swallowingTransform,
streamController,
{ stallTimeoutMs: 200 }
);

const text = await readStreamText(stream);

// No stall error — final flush output reaches the client cleanly.
assert.equal(text, "done");
assert.equal(onErrorCalled, false, "stall watchdog must NOT fire on a slow but progressing upstream");
assert.doesNotMatch(text, /stall/i);
assert.doesNotMatch(text, /"finish_reason":"error"/);
});

test("pipeWithDisconnect flags a truly stalled upstream (no bytes for the full stall budget)", async () => {
// Upstream emits one byte and then goes silent forever. Stall budget is
// 80ms — the watchdog must fire and surface a stream-stall error.
const source = new ReadableStream({
start(controller) {
controller.enqueue(encoder.encode("x"));
// never enqueue again, never close — simulate a truly hung upstream
},
cancel() {
// upstream cancel hook so the stall abort path can release the source
},
});

let onErrorEvent = null;
const streamController = createStreamController({
onError(event) {
onErrorEvent = event;
return true;
},
});

const stream = pipeWithDisconnect(
new Response(source),
new TransformStream(),
streamController,
{ stallTimeoutMs: 80 }
);

const text = await readStreamText(stream);

assert.ok(onErrorEvent !== null, "stall watchdog must fire when upstream stops sending bytes");
assert.match(onErrorEvent.message, /stall/i);
assert.match(text, /stall/i);
assert.match(text, /"finish_reason":"error"/);
});

test("pipeWithDisconnect stall watchdog does not fire after normal stream completion", async () => {
// Upstream completes quickly. The stall timer must be cleared on
// completion so a stale abort cannot fire after the request has ended.
const source = new ReadableStream({
start(controller) {
controller.enqueue(encoder.encode("ok"));
controller.close();
},
});

let onErrorCalled = false;
const streamController = createStreamController({
onError() {
onErrorCalled = true;
return true;
},
});

const stream = pipeWithDisconnect(
new Response(source),
new TransformStream(),
streamController,
{ stallTimeoutMs: 50 }
);

const text = await readStreamText(stream);

// Wait past the stall budget — no late stall error must surface.
await new Promise((r) => setTimeout(r, 120));

assert.equal(text, "ok");
assert.equal(onErrorCalled, false, "stall watchdog must be cleared on stream completion");
});
Loading