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/15455-responses-stall-measure.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(opencode):** count silent streamed replies that would have tripped the 15 s first-byte window while the guard stays off ([#15455](https://github.com/diegosouzapw/OmniRoute/pull/15455)) — thanks @maxmad64bis
4 changes: 3 additions & 1 deletion config/quality/file-size-baseline.json
Original file line number Diff line number Diff line change
Expand Up @@ -397,6 +397,8 @@
"_rebaseline_2026_07_27_3850_relax_filesize_cap": "OWNER-APPROVED TEMPORARY relax for v3.8.50-3.8.54 PREPARE phase (docs/ROADMAP.md). cap 800->900 (+100), testCap 800->900 (+100). Targets: decompose-existing-frozen unchanged (frozen still only-shrink); this only relaxes the cap for NEW files in the decompose/extract-while-PREPARE phase (.51='executor registry in-place' and .52='combo.ts decomposition' create new leaf modules above 800). RE-TIGHTENING MANDATORY in v3.8.51: cap target 850 = 850 once decomposition wave stabilizes. SUPERSEDED by _rebaseline_2026_07_27_3850_relax_filesize_cap_v2_20pct (v1 +20% buffer) — retained for audit. Tracked via same roadmap issue.",
"_rebaseline_2026_07_27_v3849_train1h": "Merge-train 1H (31 PRs) — owner-approved 2026-07-27. Two distinct causes, kept separate on purpose: (1) GENUINE irreducible growth at existing chokepoints — providerLimits/auth (#8632 Kimi quota-reset recovery), rateLimitManager (#8616 idle wedged limiters), models-catalog-route.test (#8610 OpenCode Go effort aliases); (2) COLLISION with #8585, which banked shrinks measured on the pre-train release tip while 30 sibling PRs in the SAME train grew those files again — chat/accountFallback (#8628), chatCore (#8613), videoGeneration (#8581), imageGeneration. The zero-headroom frozen entries cannot absorb either. Ceilings re-pinned to the post-merge tip; #8612 (also in this train) automates shrink-banking so this self-inflicted drift stops recurring. Detail: src/lib/usage/providerLimits.ts 1006->1013 (#8632); src/sse/services/auth.ts 2492->2508 (#8632); open-sse/services/rateLimitManager.ts 1014->1060 (#8616); src/sse/handlers/chat.ts 1842->1845 (#8628); open-sse/handlers/chatCore.ts 4939->4955 (#8613); open-sse/handlers/imageGeneration.ts 3100->3101 ((sem PR — teto do #8585)); open-sse/handlers/videoGeneration.ts 1038->1063 (#8581); open-sse/services/accountFallback.ts 1965->1966 (#8628); tests/unit/models-catalog-route.test.ts 1608->1636 (#8610)",
"_rebaseline_2026_10_02_15159_jonlwheat2_wave": "Wave of jonlwheat2-gif audit PRs (#15159) landed together. src/lib/db/providers.ts 1199->1237 (+38, #15290: findExistingCookieConnection now matches the credential before the editable name, plus the collision guard; own growth at the existing chokepoint) is frozen at its measured size; open-sse/executors/deepseek-web.ts 1224->1228 (+4: prettier line-wrapping of the #15217 sanitization). Re-measured on the release tip after the wave.",
"_rebaseline_2026_10_03_15455_responses_stall_measure": "PR 15455 own growth: open-sse/executors/opencode.ts 1457->1463 (+6 = one import swap plus one guarded-call factory line at the existing stall wiring chokepoint, logic in open-sse/executors/opencodeResponsesStall.ts under no frozen ceiling). Additive and flag-off inert. Covered by tests/unit/opencode-responses-stall-measure.test.ts.",
"_rebaseline_2026_10_06_15455_opencode_silent_replies": "Own growth of #15455 (silent streamed replies counted while the first-byte guard runs): open-sse/executors/opencode.ts 1463->1467 on top of the release tip, wiring at the existing chokepoints; logic stays in sibling modules. Measured on the PR head merged with the tip.",
"frozen": {
"src/lib/db/providers.ts": 1237,
"_rebaseline_2026_09_29_release_drain_wave": "Release-drain merge wave 2026-09-28/29 (~150 PRs into release/v3.8.51) follow-up; each file grew through already-merged PRs, re-measured with split(\"\\n\").length on release/v3.8.51 @b636b6d9: src/app/(dashboard)/dashboard/combos/page.tsx 5099->5110 (#14150/#14280/#14278); src/sse/handlers/chat.ts 2583->2586 (#14209/#13888/#15010); src/sse/services/auth.ts 3617->3622 (#13888 native Claude quota scope); open-sse/executors/default.ts 1207->1209 (#14870 GPT-6 max_completion_tokens); open-sse/handlers/responseSanitizer.ts 1200->1204 (#14813/#14894 (first frozen entry, file crossed the 1200 cap)); open-sse/utils/proxyFetch.ts 1383->1404 (#14315 stale keep-alive dispatcher + #13988 local direct timeout). Covered by the merged PRs' own tests; shrink work tracked with the other frozen files.",
Expand Down Expand Up @@ -576,7 +578,7 @@
"open-sse/executors/deepseek-web.ts": 1228,
"open-sse/executors/default.ts": 1209,
"open-sse/services/rateLimitManager.ts": 1329,
"open-sse/executors/opencode.ts": 1460,
"open-sse/executors/opencode.ts": 1467,
"open-sse/handlers/responseSanitizer.ts": 1204
},
"_rebaseline_2026_09_15_roundrobin_dashboard_events": "Fix #13089 (Combo Studio Live dashboard shows an empty backlog for round-robin combos): open-sse/services/combo/roundRobinCombo.ts 1205->1213. Round-robin is the only combo strategy that bypasses handleComboChat/executeTargetAttempt.ts, the path that publishes the combo.target.attempt/succeeded/failed EventBus events the Live dashboard listens for — so round-robin completions never showed up. The new call-site wiring (createRRDashboardEvents(...) instantiated once per target, one-line .attempt()/.succeeded()/.failed() calls at the 6 existing dispatch/outcome points) is the emitter logic actually extracted into a new module, open-sse/services/combo/rrDashboardEvents.ts — this is the minimum irreducible footprint for wiring 6 required call sites into 6 fixed control-flow points of the frozen file. Covered by tests/unit/issue-13089-roundrobin-live-ws-events.test.ts (2 tests: success + failure paths).",
Expand Down
11 changes: 9 additions & 2 deletions open-sse/executors/opencode.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,8 @@ import { withRequestShapeRetry } from "./opencodeRequestShape.ts";
// contract applies), and existing importers keep resolving it from the executor.
export { isPremiumOpencodeModel };
import {
guardResponsesStall,
isResponsesFirstByteTimeout,
makeStallGuardedCall,
setupStallGuard,
} from "./opencodeResponsesStall.ts";
import { discardResponseBody } from "./opencodeResponseBody.ts";
Expand Down Expand Up @@ -562,7 +562,14 @@ export class OpencodeExecutor extends BaseExecutor {
const hasProxies = accounts.some((a) => a.proxy !== null);
// Opt-in Responses first-byte stall guard; 0 = no-op.
const stallWindowMs = setupStallGuard(input.stream, this._requestFormat, log, cid).windowMs;
const guardStall = <T>(r: T) => guardResponsesStall(r, stallWindowMs, input.signal);
const guardStall = makeStallGuardedCall(
input.stream,
this._requestFormat,
stallWindowMs,
input.signal,
log,
cid
);
const headersWait = headersWaitState(
input,
this._requestFormat,
Expand Down
140 changes: 138 additions & 2 deletions open-sse/executors/opencodeResponsesStall.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,12 @@ import {
getResponsesFirstByteTimeoutMs,
getUpstreamTimeoutConfig,
} from "@/shared/utils/runtimeTimeouts";
import { guardResponsesStreamFirstByte } from "../utils/firstByteWatchdog.ts";
import {
guardResponsesStreamFirstByte,
isResponsesFirstByteTimeout,
} from "../utils/firstByteWatchdog.ts";

export { isResponsesFirstByteTimeout } from "../utils/firstByteWatchdog.ts";
export { isResponsesFirstByteTimeout };

export type StallGuardSetup = {
windowMs: number;
Expand Down Expand Up @@ -81,6 +84,139 @@ export function resolveResponsesStallWindowMs(
return configured > cap ? cap : configured;
}

/**
* Shadow first-byte counters: occurrences and cumulative silent time measured
* while the guard stays off, plus shadow-only failures (clone or watch
* errors, read through the shared debug sink). Three module scalars,
* constant size by construction — nothing grows per request.
*/
let shadowOccurrences = 0;
let shadowTotalMs = 0;
let shadowFailures = 0;

export type StallShadowOptions = {
stream: boolean | undefined;
requestFormat: string | null;
windowMs: number;
signal?: AbortSignal | null;
log?: {
warn?: (tag: string, message: string) => void;
debug?: (tag: string, message: string) => void;
} | null;
cid?: string;
readBound?: () => number;
readConfigured?: () => number;
};

function shadowApplies<T>(
result: T,
stream: boolean | undefined,
requestFormat: string | null,
windowMs: number
): result is T & { response: Response } {
if (!stream || requestFormat !== "openai-responses" || windowMs !== 0) return false;
if (!result || typeof result !== "object" || !("response" in result)) return false;
const response = (result as { response: Response }).response;
return !!response?.ok && !!response.body;
}

function shadowWindowMs(readBound: () => number, readConfigured: () => number): number | null {
const configured = readConfigured();
if (configured === 0) return null;
const cap = readBound();
return cap > 0 ? Math.min(configured, cap) : configured;
}

function reportShadowHit(
log: StallShadowOptions["log"],
cid: string,
windowMs: number,
elapsedMs: number
): void {
shadowOccurrences += 1;
shadowTotalMs += elapsedMs;
log?.warn?.(
"OPENCODE",
`${cid}silent streamed reply past ${windowMs}ms ` +
`(shadow count ${shadowOccurrences}, total ${shadowTotalMs}ms, ` +
`watch failures ${shadowFailures})`
);
}

function reportShadowMiss(log: StallShadowOptions["log"], cid: string): void {
shadowFailures += 1;
log?.debug?.("OPENCODE", `${cid}silent stream shadow watch ended`);
}

/**
* Passive shadow of the stall guard: when the guard does not apply
* (`windowMs` 0) but the request is a streamed Responses reply with a 2xx
* body, watches a clone of the body for the configured window and counts how
* often it would have fired. Never rotates, never cools down, never rejects
* toward the caller: the input result is handed back untouched and every
* shadow failure lands in `shadowFailures` instead.
*/
export function noteStallShadow<T>(result: T, options: StallShadowOptions): T {
const {
stream,
requestFormat,
windowMs,
signal,
log,
cid = "",
readBound = () => getUpstreamTimeoutConfig().streamReadinessTimeoutMs,
readConfigured = getResponsesFirstByteTimeoutMs,
} = options;
if (!shadowApplies(result, stream, requestFormat, windowMs)) return result;
const window = shadowWindowMs(readBound, readConfigured);
if (window === null) return result;
let clone: Response;
try {
clone = result.response.clone();
} catch {
shadowFailures += 1;
return result;
}
const startedAt = Date.now();
void guardResponsesStreamFirstByte(clone, window, signal).then(
() => {
void clone.body?.cancel().catch(() => {
shadowFailures += 1;
});
},
(error: unknown) => {
if (isResponsesFirstByteTimeout(error)) {
reportShadowHit(log, cid, window, Date.now() - startedAt);
return;
}
reportShadowMiss(log, cid);
}
);
return result;
}

/**
* Builds the request-scoped stall wrapper: counts silent streamed replies
* that would have fired past the configured window while the guard stays
* off, then applies the guard itself. Keeps the executor call site a one-line
* wiring hunk; the shadow runs inside `noteStallShadow` next to the guard it
* mirrors and never rotates, cools down, or rejects toward the caller.
*/
export function makeStallGuardedCall(
stream: boolean | undefined,
requestFormat: string | null,
windowMs: number,
signal?: AbortSignal | null,
log?: { warn?: (tag: string, message: string) => void } | null,
cid = ""
): <T>(result: T | Promise<T>) => Promise<T> {
return <T>(result: T | Promise<T>): Promise<T> =>
Promise.resolve(result).then((resolved) => {
noteStallShadow(resolved, { stream, requestFormat, windowMs, signal, log, cid });
return guardResponsesStall(resolved, windowMs, signal);
});
}

/**
* Returns `result` itself when `windowMs` is 0 or it carries no 2xx body;
* otherwise resolves once the first body byte arrives, or throws
Expand Down
123 changes: 123 additions & 0 deletions tests/unit/opencode-responses-stall-measure.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
import { describe, it } from "node:test";
import assert from "node:assert/strict";
import { noteStallShadow } from "../../open-sse/executors/opencodeResponsesStall.ts";

function silentBody(): ReadableStream<Uint8Array> {
return new ReadableStream<Uint8Array>({ pull() {} });
}

function talkingBody(): ReadableStream<Uint8Array> {
const text =
'event: response.created\ndata: {"type":"response.created","response":{"id":"r1"}}\n\n';
return new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode(text));
controller.close();
},
});
}

function watchedResult(body: ReadableStream<Uint8Array>): {
response: Response;
} {
return { response: new Response(body, { status: 200 }) };
}

async function settle(): Promise<void> {
await new Promise((resolve) => setTimeout(resolve, 60));
}

describe("silent stream shadow measurement", () => {
it("shadow counts a silent streamed reply while the guard stays off", async () => {
const seen: string[] = [];
const result = watchedResult(silentBody());
const returned = noteStallShadow(result, {
stream: true,
requestFormat: "openai-responses",
windowMs: 0,
log: { warn: (_tag, message) => seen.push(message) },
cid: "",
readBound: () => 80_000,
readConfigured: () => 30,
});
assert.equal(returned, result);
await settle();
assert.equal(seen.length, 1);
assert.match(seen[0] as string, /silent streamed reply past 30ms/);
assert.match(seen[0] as string, /shadow count 1/);
await result.response.body?.cancel().catch(() => {});
});

it("shadow stays silent on a talking body", async () => {
const seen: string[] = [];
const result = watchedResult(talkingBody());
noteStallShadow(result, {
stream: true,
requestFormat: "openai-responses",
windowMs: 0,
log: { warn: (_tag, message) => seen.push(message) },
readBound: () => 80_000,
readConfigured: () => 30,
});
await settle();
assert.deepEqual(seen, []);
const text = await result.response.text();
assert.match(text, /response\.created/);
});

it("shadow ignores non-streamed and non-responses requests", async () => {
const seen: string[] = [];
noteStallShadow(watchedResult(silentBody()), {
stream: false,
requestFormat: "openai-responses",
windowMs: 0,
log: { warn: (_tag, message) => seen.push(message) },
readBound: () => 80_000,
readConfigured: () => 30,
});
noteStallShadow(watchedResult(silentBody()), {
stream: true,
requestFormat: "openai-chat",
windowMs: 0,
log: { warn: (_tag, message) => seen.push(message) },
readBound: () => 80_000,
readConfigured: () => 30,
});
await settle();
assert.deepEqual(seen, []);
});

it("shadow stays off while the guard is armed", async () => {
const seen: string[] = [];
const result = watchedResult(silentBody());
noteStallShadow(result, {
stream: true,
requestFormat: "openai-responses",
windowMs: 50,
log: { warn: (_tag, message) => seen.push(message) },
readBound: () => 80_000,
readConfigured: () => 50,
});
await settle();
assert.deepEqual(seen, []);
await result.response.body?.cancel().catch(() => {});
});

it("shadow storage stays constant after several requests", async () => {
const seen: string[] = [];
for (let i = 0; i < 5; i++) {
const result = watchedResult(silentBody());
noteStallShadow(result, {
stream: true,
requestFormat: "openai-responses",
windowMs: 0,
log: { warn: (_tag, message) => seen.push(message) },
readBound: () => 80_000,
readConfigured: () => 20,
});
await result.response.body?.cancel().catch(() => {});
}
await settle();
assert.equal(seen.length, 5);
});
});
Loading