diff --git a/.env.example b/.env.example index 34744d642d58..5f52deb964ec 100644 --- a/.env.example +++ b/.env.example @@ -1813,6 +1813,7 @@ CURSOR_USER_AGENT="Cursor/3.4" # OPENCODE_PARK_AND_RESUME=false # #13924 feature flag (Settings → Feature Flags wins): park the request with a heartbeat after repeated transient 429s, then replay one capped leg of up to 3 accounts #OPENCODE_POOL_STRAIN_MARKER_PATH=/tmp/opencode-pool-strain.json # #13924: pool-strain marker path (JSON {since, reason, ttl_s}); fresh marker parks without recounting # RESPONSES_FIRST_BYTE_TIMEOUT_MS=15000 # #13484: OpenCode Responses first-byte window, only used when the OPENCODE_RESPONSES_STALL_ROTATION flag is on (0 disables) +# FLUSH_EMPTY_RETRY_ENABLED=false # #14213 feature flag (Settings → Feature Flags wins): retry empty translated streaming turns through the normal credential path (up to STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX retries) # ── API Bridge (/v1 proxy server) ── # API_BRIDGE_PROXY_TIMEOUT_MS=600000 # Proxy hop timeout (default: 10min) diff --git a/changelog.d/fixes/14213-flush-empty-retry.md b/changelog.d/fixes/14213-flush-empty-retry.md new file mode 100644 index 000000000000..049b203247dc --- /dev/null +++ b/changelog.d/fixes/14213-flush-empty-retry.md @@ -0,0 +1 @@ +- **fix(sse):** retry empty translated streaming turns through the normal credential path (up to `STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX` retries) instead of exposing an empty 200 or an empty-content 502, and a stream that drops before anything reaches the client takes the same retry path; a stream that answers and then stops producing without closing is replayed the same way once its stall budget runs out, instead of buffering forever; off by default behind `FLUSH_EMPTY_RETRY_ENABLED` ([#14213](https://github.com/diegosouzapw/OmniRoute/pull/14213)) diff --git a/config/quality/file-size-baseline.json b/config/quality/file-size-baseline.json index d70232bc224b..c788de22119d 100644 --- a/config/quality/file-size-baseline.json +++ b/config/quality/file-size-baseline.json @@ -3,6 +3,7 @@ "_rebaseline_2026_09_21_14250_member_egress_lines": "PR #14250 own growth: src/app/(dashboard)/dashboard/settings/components/ProxyRegistryManager.tsx 1477->1479 (+2 = the PoolMemberEgressLines import and its one-line mount under the pool members label, next to PoolEgressObservation). The member-egress observation itself lives outside the frozen file, all under cap: PoolMemberEgressLines.tsx, the dedicated GET /api/settings/proxies/pool/member-egress route, readPoolMemberEgressObservation in src/lib/proxyPoolEgressObservation.ts and getRecentEgressIpForProxy in src/lib/db/proxyLogs.ts. Only the mount point is irreducible. Covered by tests/unit/proxy-pool-member-egress-route.test.ts and tests/unit/ui/PoolMemberEgressLines.test.tsx.", "_rebaseline_2026_09_22_opencode_train10c_drift": "Release-tip drift: open-sse/executors/opencode.ts 1247->1251, left by the train-10c merge wave (2026-09-22) — the PR->release fast-gates do not run check:file-size, so the tip went red for every train boarding afterwards. Absorbed once at the tip under the owner-approved train-rebaseline policy (see _rebaseline_2026_09_18_merge_train_8_frozen_growth). Structural shrink tracked in #3501.", "_rebaseline_2026_09_21_14290_stack_handover_growth": "Stacked on #13924 (rewritten/squash-merged as 893fef9c on release/v3.8.51). PR #14290's own growth on top of that parent: open-sse/executors/opencode.ts 1233->1247 (+14 irreducible seam: settle429Arm park arm forcing burstStreak to threshold + handover comment; park/replay block owned by #13924, throttle leaf opencodeEgressThrottle.ts 543 lines under cap). Covered by tests/unit/opencode-429-park-resume.test.ts handover case (8/8) + tests/unit/opencode-egress-throttle.test.ts (22/22).", + "_rebaseline_2026_09_21_14213_empty_turn_retry_growth": "PR #14213 own growth: empty translated streaming turn retry through the normal credential path (classifier + replay parity + bounded reader in new open-sse/utils/emptyTurnRetry.ts, hook wiring in open-sse/handlers/chatCore.ts, streamEmptyChoices extraction). open-sse/handlers/chatCore.ts 6287->6400 (+113 gate count on the reconciled release tip 30f7088c, whose own chatCore.ts already sits at the frozen 6287; irreducible call-site wiring at the insertion point; logic lives in the new 422-line module under the cap). Irreducible spec wiring, additive and flag-off inert. Covered by flush-empty-retry.test.ts + flush-empty-retry-hook.test.ts (46 tests).", "_rebaseline_2026_09_20_14226_core_overrides_column": "own growth 1788->1800 (+12) src/lib/db/core.ts: rate_limit_overrides_json column in reimport INSERT + preservation SELECT, irreducible spec growth, covered by rate-limit-overrides-startup/reimport tests", "_rebaseline_2026_09_18_merge_train_8_frozen_growth": "Owner-approved train rebaseline (2026-09-18, /merge-prs; precedent _rebaseline_2026_07_23_v3849_merge_train_15). Own growth of 32 merge-ready contributor PRs that each add irreducible call-site/plumbing lines to an already-frozen file, measured on the combined merge-train tip 04cf8095 (release tip green before boarding). Per-file (old->new, contributing PRs): src/app/(dashboard)/dashboard/combos/page.tsx 5080->5091 (#13951); src/app/(dashboard)/dashboard/endpoint/EndpointPageClient.tsx 2491->2493 (#13533); src/app/api/v1/models/catalog.ts 2117->2127 (#13994); src/lib/db/apiKeys.ts 1659->1671 (#12952, #13861); src/lib/tokenHealthCheck.ts 1221->1254 (#13444, #13874); src/shared/components/RequestLoggerDetail.tsx 1200->1210 (#13373); src/shared/components/RequestLoggerV2.tsx 1718->1748 (#13373); src/shared/middleware/chatBodyAdmission.ts 1200->1206 (#13823); src/sse/handlers/chat.ts 2541->2547 (combined growth); src/sse/services/auth.ts 3582->3592 (combined growth); open-sse/executors/antigravity.ts 1665->1717 (#13125, #13318, #13659); open-sse/executors/codex.ts 1553->1570 (#13708); open-sse/executors/cursor.ts 1847->1868 (#13125); open-sse/executors/deepseek-web.ts 1200->1224 (#13226); open-sse/executors/default.ts 1200->1205 (#11828); open-sse/handlers/imageGeneration.ts 3304->3334 (#12982); open-sse/services/accountFallback.ts 2507->2515 (#13008); open-sse/services/combo/executeTargetAttempt.ts 1228->1258 (#12235); open-sse/services/combo/roundRobinCombo.ts 1221->1261 (#12235); open-sse/services/rateLimitManager.ts 1200->1329 (#13895); open-sse/translator/response/openai-responses.ts 1466->1518 (#12841, #13956); open-sse/utils/cursorAgentProtobuf.ts 1547->1588 (#13125); open-sse/utils/stream.ts 3140->3239 (#12688, #12855); tests/integration/chat-pipeline.test.ts 1740->1756 (#12966); tests/unit/account-fallback-service.test.ts 2056->2072 (#13040); tests/unit/chatcore-translation-paths.test.ts 3449->3546 (#13856, #13972); tests/unit/token-refresh-service.test.ts 1407->1408 (combined growth); tests/unit/translator-openai-to-gemini.test.ts 1625->1809 (#13318, #13848). Files previously under the 1200 cap that crossed it are frozen at the measured size. Structural shrink of these god-files stays tracked in #3501; the ceilings never move up again outside a documented entry. ADJUST (train 8d re-measure after #13548 ejection and #14101 landing): open-sse/services/accountFallback.ts 2515->2517 (#13008 +22, #13350 +7, #13984 +2, #13040 +2 on a tip at 2499).", "_rebaseline_2026_09_18_14065_codex_reasoning_whitelist": "Release-tip drift: open-sse/executors/codex.ts 1552->1553 (+1) from #14065 (fix(codex): whitelist reasoning object keys before the wire, #13643), merged 2026-09-18 without its own rebaseline — the PR->release fast-gates do not run check:file-size, so the tip went red for every train boarding afterwards. Absorbed at the release tip by the /merge-prs captain session (owner-approved train-rebaseline policy, 2026-09-18). Structural shrink tracked in #3501.", @@ -480,7 +481,7 @@ "open-sse/executors/codex.ts": 1570, "open-sse/executors/cursor.ts": 1868, "open-sse/executors/muse-spark-web.ts": 1405, - "open-sse/handlers/chatCore.ts": 6287, + "open-sse/handlers/chatCore.ts": 6400, "open-sse/handlers/imageGeneration.ts": 3334, "open-sse/handlers/search.ts": 1789, "open-sse/mcp-server/schemas/tools.ts": 1621, diff --git a/docs/reference/ENVIRONMENT.md b/docs/reference/ENVIRONMENT.md index c9ee58231cda..931fda33c22b 100644 --- a/docs/reference/ENVIRONMENT.md +++ b/docs/reference/ENVIRONMENT.md @@ -1146,6 +1146,7 @@ Anthropic-compatible provider instead. | `PROXY_HEALTH_TEST_STAGGER_MS` | `100` | `src/lib/proxyHealth/probeTarget.ts` | Delay in ms between two probe departures inside a batch. Without it the whole batch leaves at the same moment and a shared egress IP can trip a rate-limited target. Set to `0` to disable the spacing; capped at 5000. | | `PROXY_HEALTH_USE_PROVIDER_TARGET` | `true` | `src/lib/proxyHealth/providerProbeTarget.ts` | Set "false" to stop probing the real host of a proxy's assigned provider (`GET /models`, no API key) and always use `PROXY_HEALTH_TEST_URL` instead. | | `PROXY_HEALTH_AUTO_DEACTIVATE` | `false` | `src/lib/proxyHealth/statusPolicy.ts` | When `false` (default), automated reachability probes (the scheduler + the `/api/settings/proxies/auto-test` "Test All" button) are **read-only** and never write a proxy's status — only the operator sets active/inactive, so a flaky probe can't strand an assigned proxy (#6246). Set `true` to restore the legacy test-and-set behaviour. | +| `FLUSH_EMPTY_RETRY_ENABLED` | `false` | `src/shared/utils/featureFlags.ts` | Opt-in feature flag (see [FEATURE_FLAGS.md](./FEATURE_FLAGS.md); a dashboard DB override wins). `true` (or `1`, `yes`) retries empty translated streaming turns through the normal credential path (up to `STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX` retries) instead of exposing an empty 200 or an empty-content 502. | | `PROXY_POOL_EGRESS_OBSERVATION` | `false` | `src/shared/utils/featureFlags.ts` | Opt-in feature flag (see [FEATURE_FLAGS.md](./FEATURE_FLAGS.md); a dashboard DB override wins). `true` (or `1`, `yes`) shows the read-only pool egress observation under a proxy pool in the dashboard (distinct egress IPs, connections and the most seen behind one IP over the last 24 h, from the proxy log). Never used for routing. | | `PROXY_AUTO_REMOVE` | `false` | `src/lib/proxyHealth/scheduler.ts` | Set `true` to let the scheduler auto-remove proxies after repeated consecutive failures. | | `PROXY_AUTO_REMOVE_AFTER` | `3` | `src/lib/proxyHealth/scheduler.ts` | Consecutive failures before the scheduler auto-removes a proxy (when `PROXY_AUTO_REMOVE=true`). | diff --git a/docs/reference/FEATURE_FLAGS.md b/docs/reference/FEATURE_FLAGS.md index b353dddf5d1b..b6e63a716a17 100644 --- a/docs/reference/FEATURE_FLAGS.md +++ b/docs/reference/FEATURE_FLAGS.md @@ -46,7 +46,7 @@ A boolean flag is considered **enabled** when its effective value is `"true"`, ## Flag Catalog -75 flags across 6 categories. **Default** is the definition default — the value +76 flags across 6 categories. **Default** is the definition default — the value used when neither a DB override nor an environment variable is present. ### Security (10) @@ -64,7 +64,7 @@ used when neither a DB override nor an environment variable is present. | `AUTH_LOG_INCLUDE_ACCOUNT_ID` | boolean | `false` | Include account prefix in AUTH log lines (e.g. "Using account: abc12345..."). Disabled by default so account identifiers are redacted from shared/multi-tenant process logs. Independent from Debug Mode; flipping Debug Mode does not reveal this. | | `OMNIROUTE_OIDC_DISABLE_PASSWORD_LOGIN` | boolean | `false` | When OIDC is enabled, disable password login so users can only authenticate via OIDC Single Sign-On. When disabled (default), both password login and OIDC are available. | -### Network (17) +### Network (18) | Key | Type | Default | Restart | Description | | ----------------------------------------------- | ------- | ------- | ------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | @@ -80,6 +80,7 @@ used when neither a DB override nor an environment variable is present. | `OPENCODE_USER_BLOCKED_ROTATION` | boolean | `false` | | OpenCode executor: on a 403/451 carrying a `user_blocked` refusal (not geo, not a Cloudflare fingerprint rejection), cool the refused account down and rotate to the next account at most once per request; a second refusal is returned as-is, without a success mark. Off by default: routing around an upstream user block can look like evasion and spread the flag across the fleet. | | `OPENCODE_TRANSIENT_FAILOVER_BACKOFF` | boolean | `false` | | OpenCode rotation: after two consecutive transient upstream failures (5xx or an empty 400), pause before the next account — 1.5s doubling per further failure, capped at 6s per pause and 10s per request, skipped on client disconnect; the failed body is released before waiting. Off by default: failover stays immediate. | | `OPENCODE_PARK_AND_RESUME` | boolean | `false` | | OpenCode rotation: park the request after repeated transient 429s (or a fresh pool-strain marker) with a heartbeat, then replay one capped leg of up to 3 sequential accounts instead of fanning out the whole fleet. Off by default: every 429 rotates to the next account exactly as before. | +| `FLUSH_EMPTY_RETRY_ENABLED` | boolean | `false` | | On translated streaming turns, when the upstream turn carries no usable content (reasoning-only completion or zero valuable chunks), issue bounded retries through the normal credential path (up to `STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX`) before anything is exposed to the client. Off by default: empty turns keep the current behavior (empty 200 or empty-content 502). | | `OPENCODE_RATE_LIMITED_429_EARLY_STOP` | boolean | `false` | | OpenCode rotation: stop the account wave at the first 429 classified as a real rate limit (parseable `Retry-After`, or a body naming a rate/usage limit) and return that upstream 429 unchanged. Unclassified 429s keep rotating. Off by default: the free tier is limited per egress IP (#9611), so every 429 rotates and an exhausted wave returns the last upstream 429. | | `MITM_DISABLE_TLS_VERIFY` | boolean | `false` | ✓ | Disable TLS certificate verification for the MITM proxy. **Danger.** | | `OMNIROUTE_ALLOW_PRIVATE_PROVIDER_URLS` | boolean | `false` | | Allow provider URLs pointing to private/internal networks. | diff --git a/open-sse/config/constants.ts b/open-sse/config/constants.ts index d3d49d8eaad3..a0f9be98806e 100644 --- a/open-sse/config/constants.ts +++ b/open-sse/config/constants.ts @@ -365,11 +365,15 @@ export const CREDENTIAL_HEALTH_CACHE_TTL = (() => { * as soon as this many bytes accumulate, regardless of the timer. * - EARLY_RETRY_MAX: max transparent re-opens of the upstream stream while the * holdback is still uncommitted (free-claude-code uses 5 total attempts = 4 retries). + * - EMPTY_TURN_RETRY_MAX: max bounded retries of a translated stream turn that ends + * with no usable content (same family: bounded retries of a failing stream + * before anything is exposed to the client). */ export const STREAM_RECOVERY = { HOLDBACK_MS: 750, BUFFER_MAX_BYTES: 65536, EARLY_RETRY_MAX: 4, + EMPTY_TURN_RETRY_MAX: 4, /** * Minimum character overlap `trimContinuationOverlap` must find between the * already-emitted text and a mid-stream continuation for the continuation to be diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index b4e0e329e54f..38feddfc7de5 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -44,6 +44,11 @@ import { buildNonStreamingResponseHeaders } from "./chatCore/nonStreamingRespons import { maybeWrapForcedNonStreamingResponsesJson } from "./chatCore/responsesJsonToSse.ts"; import { enforceOutputTokenBudget } from "./chatCore/outputTokenBudget.ts"; import { maybeConvertJsonBodyToSse } from "./chatCore/jsonBodyToSse.ts"; +import { + judgeBufferedTurn, + readBoundedResponseOutcome, + FLUSH_EMPTY_RETRY_MAX_BYTES, +} from "../utils/emptyTurnRetry.ts"; import { assembleStreamingResponseHeaders } from "./chatCore/streamingResponseHeaders.ts"; import { storeStreamingSemanticCacheResponse } from "./chatCore/streamingSemanticCacheStore.ts"; import { assembleStreamingPipeline } from "./chatCore/streamingPipeline.ts"; @@ -5823,6 +5828,114 @@ export async function handleChatCore({ } providerResponse = streamReadiness.response; + // Flush-empty retry (opt-in `FLUSH_EMPTY_RETRY_ENABLED`, default off): when the + // upstream turn carries no usable content (reasoning-only 200, or a + // zero-valuable-chunk turn that the empty-stream guard would turn into a 502), + // issue bounded retries through the normal credential path BEFORE anything is + // exposed to the client — in particular before `onRequestSuccess` below. + // Empty turns are stochastic upstream misses, not account faults, so no + // cooldown and no forced exclusion: the round-robin picker may rotate + // fingerprint slots opportunistically, a single slot simply replays the same + // account. Budget: `STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX` retries, then fall + // back to the current behavior. Translate-path streams only (mirror of the + // empty-stream guard); flag off = byte-for-byte unchanged. Bounded reader + // (abandon past the cap, never a full `text()` read); the original + // reconstructed response is piped, only the bounded copy is classified. + // Known TTFT cost when armed: a small valid turn under the cap is fully + // buffered before the first client byte (flag off by default, so the + // streaming path is untouched unless opted in). + if (stream && providerResponse.ok && providerResponse.body) { + let flushEmptyRetryArmed = false; + try { + flushEmptyRetryArmed = isFeatureFlagEnabled("FLUSH_EMPTY_RETRY_ENABLED"); + } catch { + flushEmptyRetryArmed = false; + } + const isTranslatePath = + targetFormat === FORMATS.OPENAI_RESPONSES || + needsTranslation(targetFormat, clientResponseFormat); + if (flushEmptyRetryArmed && isTranslatePath) { + for ( + let emptyTurnRetries = 0; + emptyTurnRetries <= STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX; + emptyTurnRetries++ + ) { + const verdict = judgeBufferedTurn( + await readBoundedResponseOutcome( + providerResponse, + FLUSH_EMPTY_RETRY_MAX_BYTES, + streamReadinessPolicy.timeoutMs + ), + targetFormat, + clientResponseFormat, + clientRawRequest?.signal?.aborted === true + ); + if (verdict.kind === "pass") { + log?.debug?.("FLUSH_EMPTY_RETRY", `passing the turn through: ${verdict.why}`); + break; + } + if (emptyTurnRetries >= STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX) { + log?.warn?.( + "FLUSH_EMPTY_RETRY", + "retry budget exhausted, falling back to current behavior" + ); + break; + } + log?.warn?.( + "FLUSH_EMPTY_RETRY", + `${verdict.reason}, bounded retry through the normal credential path` + ); + const nextCreds = await getProviderCredentials( + provider, + null, + null, + currentModel + ).catch(() => null); + if (!nextCreds?.connectionId) break; + const retryConnectionId = String(nextCreds.connectionId); + Object.assign(credentials, nextCreds); + log?.info?.("FLUSH_EMPTY_RETRY", `retrying on ${retryConnectionId}`); + await providerResponse.body?.cancel().catch(() => {}); + let retryResult: unknown = null; + try { + retryResult = await executeProviderRequest(currentModel, false); + } catch { + break; + } + const retryResponse = (retryResult as { response?: Response })?.response; + if (!retryResponse?.ok || !retryResponse.body) { + if (retryResponse) await retryResponse.body?.cancel().catch(() => {}); + break; + } + const prepared = await maybeConvertJsonBodyToSse(retryResponse, { + log, + provider, + model, + }); + const ready = prepared.ok + ? await ensureStreamReadiness(prepared, { + timeoutMs: streamReadinessPolicy.timeoutMs, + maxTimeoutMs: streamReadinessPolicy.maxTimeoutMs, + provider, + model, + log, + }) + : null; + const preparedStream = ready && ready.ok ? ready.response : null; + if (!preparedStream) { + await retryResponse.body?.cancel().catch(() => {}); + break; + } + // Swap BEFORE re-classifying so the next loop iteration reads the retry. + providerResponse = preparedStream; + finalBody = providerRequestCapture.body( + (retryResult as { transformedBody?: unknown })?.transformedBody ?? translatedBody + ); + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); + } + } + } + // Notify success - caller can clear error status if needed if (onRequestSuccess) { await onRequestSuccess(); diff --git a/open-sse/utils/emptyTurnRetry.ts b/open-sse/utils/emptyTurnRetry.ts new file mode 100644 index 000000000000..28a5f8c60b1f --- /dev/null +++ b/open-sse/utils/emptyTurnRetry.ts @@ -0,0 +1,422 @@ +import { translateResponse, initState } from "../translator/index.ts"; +import { FORMATS } from "../translator/formats.ts"; +import { hasValidUsage } from "./usageTracking.ts"; +import { parseSSELine, hasValuableContent } from "./streamHelpers.ts"; +import { isEmptyTurnCore } from "./streamEmptyChoices.ts"; +import { sanitizeStreamingChunk } from "../handlers/responseSanitizer.ts"; +import { getAnyReasoningValue, getReadableReasoningValue } from "./reasoningFields.ts"; + +/** Max upstream bytes buffered for flush-empty-retry classification (flag-gated). */ +export const FLUSH_EMPTY_RETRY_MAX_BYTES = 256_000; + +/** Legit empty stops (same set as `LEGIT_EMPTY_OPENAI_FINISH` in errorClassifier): a + * turn truncated at the token limit (`length`), a tool-call turn (`tool_calls`), + * or a filtered turn (`content_filter`) is a valid completion, not an empty-turn + * failure — never a retry trigger. */ +const LEGIT_EMPTY_TURN_FINISH = new Set(["length", "tool_calls", "content_filter"]); + +export type EmptyTurnSummary = { + finishReason: string; + contentText: string; + reasoningText: string; + forwardedValuableChunk: boolean; + hasValidUsage: boolean; + toolCallsPresent: boolean; +}; + +/** + * Unified "turn with no usable content" classifier (one mechanism, two arms: + * reasoning-only turn with stop + empty text + non-empty reasoning, and + * zero-valuable-chunk turn via shared `isEmptyTurnCore`). Legit empty stops + * (length/tool_calls/content_filter) and tool-call turns are never empty. + */ +export function isUselessEmptyTurn(summary: EmptyTurnSummary): boolean { + if (LEGIT_EMPTY_TURN_FINISH.has(summary.finishReason)) return false; + if (summary.toolCallsPresent) return false; + if (summary.contentText.length > 0) return false; + // Reasoning-only arm: non-empty reasoning with no content. The stop gate + // applies to chat turns (finish=stop disambiguates from mid-stream deltas); + // Responses translators never set a probe finish reason, so Responses + // reasoning-only turns (reasoning deltas, no completed/output) take the + // same arm without the stop requirement. + if ( + summary.reasoningText.length > 0 && + (summary.finishReason === "stop" || summary.finishReason === "") + ) + return true; + return isEmptyTurnCore(summary.forwardedValuableChunk, summary.hasValidUsage); +} + +/** A read that outlived the idle budget, kept distinct from a real chunk. */ +const IDLE_READ = Symbol("idle-read"); + +/** + * One read under an idle budget. The budget covers the gap between chunks, not + * the whole turn, so a long generation that keeps producing is never cut short. + * `idleMs <= 0` keeps the plain unbounded read. + */ +async function readWithinIdleBudget( + reader: ReadableStreamDefaultReader, + idleMs: number +): Promise | typeof IDLE_READ> { + if (idleMs <= 0) return reader.read(); + let timer: ReturnType | undefined; + const expiry = new Promise((resolve) => { + timer = setTimeout(() => resolve(IDLE_READ), idleMs); + }); + try { + return await Promise.race([reader.read(), expiry]); + } finally { + clearTimeout(timer); + } +} + +async function drainBoundedChunks( + reader: ReadableStreamDefaultReader, + maxBytes: number, + idleMs: number +): Promise<{ chunks: Uint8Array[]; total: number; over: boolean; idle: boolean }> { + const chunks: Uint8Array[] = []; + let total = 0; + for (;;) { + const read = await readWithinIdleBudget(reader, idleMs); + if (read === IDLE_READ) return { chunks, total, over: false, idle: true }; + const { done, value } = read; + if (done) break; + if (!value) continue; + total += value.byteLength; + if (total > maxBytes) return { chunks, total, over: true, idle: false }; + chunks.push(value); + } + return { chunks, total, over: false, idle: false }; +} + +function concatChunks(chunks: Uint8Array[], total: number): string | null { + try { + const out = new Uint8Array(total); + let offset = 0; + for (const chunk of chunks) { + out.set(chunk, offset); + offset += chunk.byteLength; + } + return new TextDecoder().decode(out); + } catch { + return null; + } +} + +export type BoundedReadOutcome = + | { kind: "text"; text: string } + // Over the byte cap, no body, or undecodable: not classifiable, pass it through. + | { kind: "skipped" } + // The body threw while being read (e.g. the upstream dropped the connection + // after its headers): nothing usable was delivered. + | { kind: "error" } + // The body stopped producing without closing or erroring: buffering can never + // end on its own. `text` is whatever had been buffered when the budget ran + // out, so a stalled turn is judged on its content like any other. + | { kind: "idle"; text: string }; + +/** + * Bounded read of a `Response` body: streams chunks through a reader with a + * byte counter and abandons past `maxBytes` (`skipped` = fall back to the + * normal path). Never a full `text()` read: a large valid turn is abandoned + * without ever being fully buffered. A body that throws while being read is + * reported as `error`, distinct from `skipped`, so the caller can treat a + * dropped stream like an empty turn. `idleMs` bounds the gap between chunks — + * without it a stream that stops producing without closing buffers forever, + * since this read owns no other deadline and runs before the client pipe (and + * its idle watchdog) exists. The consumed clone is discarded by the caller; the + * piped original is untouched. + */ +export async function readBoundedResponseOutcome( + response: Response, + maxBytes: number, + idleMs = 0 +): Promise { + const clone = response.clone(); + if (!clone.body) return { kind: "skipped" }; + const reader = clone.body.getReader(); + try { + const { chunks, total, over, idle } = await drainBoundedChunks(reader, maxBytes, idleMs); + if (idle) { + // Never awaited: this branch exists because the stream stopped answering. + void reader.cancel().catch(() => undefined); + return { kind: "idle", text: concatChunks(chunks, total) ?? "" }; + } + if (over) { + try { + await reader.cancel(); + } catch { + // best-effort + } + return { kind: "skipped" }; + } + const text = concatChunks(chunks, total); + return text === null ? { kind: "skipped" } : { kind: "text", text }; + } catch { + return { kind: "error" }; + } finally { + try { + reader.releaseLock(); + } catch { + // best-effort + } + } +} + +// Discriminated by a string, not a boolean literal: a boolean discriminant does not +// narrow under every tsconfig in this repo (the API-route check is one of them). +export type BufferedTurnVerdict = { kind: "retry"; reason: string } | { kind: "pass"; why: string }; + +/** + * Decide from a bounded read whether the buffered turn deserves a retry: an + * empty turn, or a stream that dropped before anything reached the client + * (unless the client itself went away). Every other outcome passes through. + */ +export function judgeBufferedTurn( + read: BoundedReadOutcome, + targetFormat: string, + clientFormat: string, + clientAborted: boolean +): BufferedTurnVerdict { + if (read.kind === "skipped") { + return { kind: "pass", why: "not classified (over the buffer cap or unreadable)" }; + } + if (read.kind === "idle") { + if (clientAborted) { + return { kind: "pass", why: "stream stalled after the client went away" }; + } + const stalled = summarizeReplayedUpstreamTurn(read.text, targetFormat, clientFormat); + // Content already produced is worth keeping: the client pipe forwards it and + // owns the rest. Nothing usable means the turn is as empty as a silent one. + return stalled && !isUselessEmptyTurn(stalled) + ? { kind: "pass", why: "stalled turn already carries usable content" } + : { kind: "retry", reason: "stream stalled before any usable output" }; + } + if (read.kind === "error") { + return clientAborted + ? { kind: "pass", why: "stream dropped after the client went away" } + : { kind: "retry", reason: "stream dropped before any output" }; + } + const summary = summarizeReplayedUpstreamTurn(read.text, targetFormat, clientFormat); + if (!summary) return { kind: "pass", why: "turn could not be summarized" }; + return isUselessEmptyTurn(summary) + ? { kind: "retry", reason: "empty turn" } + : { kind: "pass", why: "turn has usable content" }; +} + +/** Text of a bounded read, or null when the body was skipped or failed. */ +export async function readBoundedResponseText( + response: Response, + maxBytes: number, + idleMs = 0 +): Promise { + const outcome = await readBoundedResponseOutcome(response, maxBytes, idleMs); + return outcome.kind === "text" ? outcome.text : null; +} + +type ProbeAccum = { + state: Record; + forwardedValuableChunk: boolean; + finishReason: string; + toolCallsPresent: boolean; +}; + +function asRecord(value: unknown): Record | null { + return !!value && typeof value === "object" && !Array.isArray(value) + ? (value as Record) + : null; +} + +function firstChoiceOf(rec: Record): Record | null { + if (!Array.isArray(rec.choices)) return null; + return asRecord(rec.choices[0]); +} + +function choiceDelta(rec: Record): Record | null { + const choice = firstChoiceOf(rec); + if (!choice) return null; + return asRecord(choice.delta); +} + +function createProbeState(sourceFormat: string): Record | null { + try { + return { + ...(initState(sourceFormat) as Record), + accumulatedContent: "", + accumulatedReasoning: "", + }; + } catch { + return null; + } +} + +function appendText(state: Record, key: string, text: string): void { + if (state[key] === undefined || !text) return; + state[key] = String(state[key]) + text; +} + +function accumulateRawChunk(parsed: Record, probe: ProbeAccum): void { + const rawDelta = choiceDelta(parsed); + const content = rawDelta?.content; + if (typeof content === "string" && content) { + appendText(probe.state, "accumulatedContent", content); + } + const reasoning = getReadableReasoningValue(rawDelta ?? {}); + if (reasoning) appendText(probe.state, "accumulatedReasoning", reasoning); +} + +function sanitizeForVerdict(item: unknown, sourceFormat: string): Record | null { + const rec = asRecord(item); + if (!rec) return null; + const isResponsesEvent = + typeof rec?.event === "string" && (rec.event as string).startsWith("response."); + if (sourceFormat === FORMATS.OPENAI && !isResponsesEvent) { + return sanitizeStreamingChunk(rec) as Record; + } + return rec; +} + +function accumulateHubContent(rec: Record, probe: ProbeAccum): string { + const delta = choiceDelta(rec); + const content = delta && typeof delta.content === "string" ? delta.content : ""; + if (content) appendText(probe.state, "accumulatedContent", content); + return content; +} + +function accumulateHubReasoning(rec: Record, probe: ProbeAccum): string { + const readable = getReadableReasoningValue(rec); + const delta = choiceDelta(rec); + const anyReasoning = readable || getAnyReasoningValue(delta ?? {}); + if (anyReasoning) appendText(probe.state, "accumulatedReasoning", anyReasoning); + return readable; +} + +function scanSiblingReasoning(translated: unknown[], item: unknown, probe: ProbeAccum): void { + for (const sib of translated) { + if (!sib || typeof sib !== "object" || Array.isArray(sib) || sib === item) continue; + const sibDelta = choiceDelta(sib as Record); + const sibReasoning = getReadableReasoningValue(sibDelta ?? {}); + if (sibReasoning) appendText(probe.state, "accumulatedReasoning", sibReasoning); + if (sibReasoning) probe.forwardedValuableChunk = true; + } +} + +function recordValuableItem(rec: Record, probe: ProbeAccum): void { + probe.forwardedValuableChunk = true; + const choice = firstChoiceOf(rec); + const finish = choice?.finish_reason; + if (typeof finish === "string" && finish) probe.finishReason = finish; + const delta = choice ? asRecord(choice.delta) : null; + const toolCalls = delta?.tool_calls; + if (Array.isArray(toolCalls) && toolCalls.length > 0) probe.toolCallsPresent = true; +} + +function classifyTranslatedItem( + item: unknown, + translated: unknown[], + sourceFormat: string, + probe: ProbeAccum +): void { + const rec = sanitizeForVerdict(item, sourceFormat); + if (!rec) return; + accumulateHubContent(rec, probe); + const hubReasoning = accumulateHubReasoning(rec, probe); + if (!hubReasoning && Array.isArray(translated)) { + scanSiblingReasoning(translated, item, probe); + } + if (!hasValuableContent(rec, sourceFormat)) return; + recordValuableItem(rec, probe); +} + +function replayParseLine( + line: string, + targetFormat: string, + sourceFormat: string, + probe: ProbeAccum +): void { + const trimmed = line.trim(); + if (!trimmed) return; + const parsedLine = parseSSELine(trimmed); + if (!parsedLine || (parsedLine as Record).done) return; + const parsed = parsedLine as Record; + accumulateRawChunk(parsed, probe); + let translated: unknown; + try { + translated = translateResponse( + targetFormat, + sourceFormat, + parsed as Record, + probe.state + ); + } catch { + return; + } + if (!Array.isArray(translated)) return; + for (const item of translated) { + classifyTranslatedItem(item, translated, sourceFormat, probe); + } +} + +function replayFlush(targetFormat: string, sourceFormat: string, probe: ProbeAccum): void { + try { + const flushed = translateResponse(targetFormat, sourceFormat, null, probe.state); + if (!Array.isArray(flushed)) return; + for (const item of flushed) { + const rec = asRecord(item); + if (rec && hasValuableContent(rec, sourceFormat)) { + probe.forwardedValuableChunk = true; + } + } + } catch { + // Flush failure is conservative: keep what the chunks already told us + } +} + +function buildProbeSummary(probe: ProbeAccum): EmptyTurnSummary { + return { + finishReason: + probe.finishReason || + (typeof probe.state.finishReason === "string" ? probe.state.finishReason : ""), + contentText: + typeof probe.state.accumulatedContent === "string" ? probe.state.accumulatedContent : "", + reasoningText: + typeof probe.state.accumulatedReasoning === "string" ? probe.state.accumulatedReasoning : "", + forwardedValuableChunk: probe.forwardedValuableChunk, + hasValidUsage: hasValidUsage(probe.state.usage as never), + toolCallsPresent: + probe.toolCallsPresent || + (probe.state.toolCalls instanceof Map && probe.state.toolCalls.size > 0), + }; +} + +/** + * Replay buffered upstream SSE bytes through a disposable translator and + * summarize the turn for `isUselessEmptyTurn`. Same translator entries as the + * live transform (`translateResponse` + `initState`), no client output. + * Returns null when the body is not classifiable (conservative: no retry). + */ +export function summarizeReplayedUpstreamTurn( + text: string, + targetFormat: string, + sourceFormat: string +): EmptyTurnSummary | null { + const state = createProbeState(sourceFormat); + if (!state) return null; + const probe: ProbeAccum = { + state, + forwardedValuableChunk: false, + finishReason: "", + toolCallsPresent: false, + }; + try { + for (const line of text.split("\n")) { + replayParseLine(line, targetFormat, sourceFormat, probe); + } + replayFlush(targetFormat, sourceFormat, probe); + } catch { + return null; + } + return buildProbeSummary(probe); +} diff --git a/open-sse/utils/streamEmptyChoices.ts b/open-sse/utils/streamEmptyChoices.ts index 05f8d8e4b9b7..aa072758aade 100644 --- a/open-sse/utils/streamEmptyChoices.ts +++ b/open-sse/utils/streamEmptyChoices.ts @@ -44,32 +44,45 @@ type EmptyChoicesRejectContext = { targetFormat?: string; model?: string | null; usage?: unknown; - onFailure?: ((payload: { - status: number; - message: string; - code?: string; - type?: string; - }) => boolean | void | Promise) | null; - onComplete?: ((payload: { - status: number; - usage: unknown; - responseBody?: unknown; - providerPayload?: unknown; - clientPayload?: unknown; - error?: string | null; - errorCode?: string | null; - }) => void) | null; + onFailure?: + | ((payload: { + status: number; + message: string; + code?: string; + type?: string; + }) => boolean | void | Promise) + | null; + onComplete?: + | ((payload: { + status: number; + usage: unknown; + responseBody?: unknown; + providerPayload?: unknown; + clientPayload?: unknown; + error?: string | null; + errorCode?: string | null; + }) => void) + | null; clearPendingRequestFromStream?: () => void; }; +/** + * Shared empty-turn core (extracted literally from the #9268 guard below): + * a turn is empty when nothing valuable was forwarded AND no valid usage + * was accumulated. Imported by the flush-empty-retry classifier so both + * sites share one implementation. + */ +export function isEmptyTurnCore(forwardedValuableChunk: boolean, hasValidUsage: boolean): boolean { + return !forwardedValuableChunk && !hasValidUsage; +} + /** * Returns `true` when the empty-stream condition was detected and the caller * must abort the stream (controller.error + early return); `false` when the * stream legitimately forwarded content/usage and should complete normally. */ export function rejectEmptyChoicesStream(ctx: EmptyChoicesRejectContext): boolean { - if (ctx.forwardedValuableChunk || ctx.hasValidUsage) return false; - + if (!isEmptyTurnCore(ctx.forwardedValuableChunk, ctx.hasValidUsage)) return false; const error = new Error( "Provider returned empty content — stream forwarded no valuable chunks" ) as Error & { statusCode: number; code: string }; diff --git a/src/shared/constants/featureFlagDefinitions.ts b/src/shared/constants/featureFlagDefinitions.ts index 87c776a2ab73..219f61415406 100644 --- a/src/shared/constants/featureFlagDefinitions.ts +++ b/src/shared/constants/featureFlagDefinitions.ts @@ -263,6 +263,18 @@ export const FEATURE_FLAG_DEFINITIONS: FeatureFlagDefinition[] = [ requiresRestart: false, warningLevel: "caution", }, + { + key: "FLUSH_EMPTY_RETRY_ENABLED", + label: "Flush Empty Turn Retry", + description: + "On translated streaming turns, when the upstream turn carries no usable content (reasoning-only completion or zero valuable chunks), issue bounded retries through the normal credential path (up to `STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX`) before anything is exposed to the client. Off by default: empty turns keep the current behavior (empty 200 or empty-content 502).", + descriptionI18nKey: "featureFlagFlushEmptyRetryEnabledDescription", + category: "network", + defaultValue: "false", + type: "boolean", + requiresRestart: false, + warningLevel: "caution", + }, { key: "OPENCODE_RATE_LIMITED_429_EARLY_STOP", label: "OpenCode Rate-Limited 429 Early Stop", diff --git a/stryker.conf.json b/stryker.conf.json index 2cc8d8aa7f8b..37626e961068 100644 --- a/stryker.conf.json +++ b/stryker.conf.json @@ -167,6 +167,7 @@ "tests/unit/circuit-breaker-resolved-5xx-12254.test.ts", "tests/unit/circuit-breaker-stream-controller-4602.test.ts", "tests/unit/overloaded-not-provider-breaker.test.ts", + "tests/unit/flush-empty-retry-hook.test.ts", "tests/unit/claude-code-parity.test.ts", "tests/unit/claude-effort-suffix-strip.test.ts", "tests/unit/claude-oauth-provider.test.ts", diff --git a/tests/unit/feature-flags-settings.test.ts b/tests/unit/feature-flags-settings.test.ts index 0e9ba3045334..4c553cb60922 100644 --- a/tests/unit/feature-flags-settings.test.ts +++ b/tests/unit/feature-flags-settings.test.ts @@ -40,7 +40,9 @@ const { // the dead ONEPROXY_ENABLED (readerless since the 1proxy purge, #12091) // brought it back to 53. UNIVERSAL_CONTEXT_HANDOFF_ENABLED bumped it to 54. // #13641 added SEARCH_STATS_HIDE_DELETED_CONNECTIONS, bumping the count to 56. -const EXPECTED_FEATURE_FLAG_COUNT = 75; +// 893fef9c added OPENCODE_PARK_AND_RESUME (74 -> 75); FLUSH_EMPTY_RETRY_ENABLED +// (flush empty-turn retry, default off) bumps it to 76. +const EXPECTED_FEATURE_FLAG_COUNT = 76; // ────────────────────────────────────────────────────── // Test group 1 — Flag definitions registry diff --git a/tests/unit/flush-empty-retry-hook.test.ts b/tests/unit/flush-empty-retry-hook.test.ts new file mode 100644 index 000000000000..95a2fa105da1 --- /dev/null +++ b/tests/unit/flush-empty-retry-hook.test.ts @@ -0,0 +1,266 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +// Hook-level tests for the flush empty-turn retry (translate-path streams). +// +// A Gemini upstream speaks a non-OpenAI format, so the OpenAI-speaking client +// path always goes through the translate transform — the hook's +// `isTranslatePath` gate is armed. Two Gemini accounts isolate rotation: +// the first serves an empty turn, the second serves content. +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-flush-empty-retry-")); +process.env.DATA_DIR = TEST_DATA_DIR; +process.env.REQUIRE_API_KEY = "false"; +process.env.DASHBOARD_PASSWORD = ""; +process.env.INITIAL_PASSWORD = ""; +delete process.env.JWT_SECRET; +if (!process.env.API_KEY_SECRET) { + process.env.API_KEY_SECRET = `test-flush-empty-retry-${Date.now()}`; +} +process.env.STREAM_READINESS_TIMEOUT_MS = "1000"; + +const core = await import("../../src/lib/db/core.ts"); +const providersDb = await import("../../src/lib/db/providers.ts"); +const { handleChat } = await import("../../src/sse/handlers/chat.ts"); +const { initTranslators } = await import("../../open-sse/translator/index.ts"); +const { clearInflight } = await import("../../open-sse/services/requestDedup.ts"); +const { resetAllCircuitBreakers } = await import("../../src/shared/utils/circuitBreaker.ts"); + +const originalFetch = globalThis.fetch; + +const FLAG = "FLUSH_EMPTY_RETRY_ENABLED"; +const ORIGINAL_FLAG = process.env[FLAG]; + +function setFlag(enabled: boolean) { + if (enabled) process.env[FLAG] = "true"; + else delete process.env[FLAG]; +} + +async function resetStorage() { + clearInflight(); + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); + fs.mkdirSync(TEST_DATA_DIR, { recursive: true }); + resetAllCircuitBreakers(); + initTranslators(); +} + +async function seedGemini(name: string, apiKey: string) { + const row = (await providersDb.createProviderConnection({ + provider: "gemini", + authType: "apikey", + name, + apiKey, + isActive: true, + testStatus: "active", + })) as { id: string }; + return { id: row.id, apiKey }; +} + +// Anonymized replay of the 19/09 reasoning-only payload shape: 74 tokens of +// thinking text (native thought part), empty message text, terminal stop. +function reasoningOnlyStreamResponse(): Response { + const thinking = "r".repeat(74); + return new Response( + `data: ${JSON.stringify({ + candidates: [{ content: { parts: [{ text: thinking, thought: true }] } }], + })}\n\ndata: ${JSON.stringify({ + candidates: [{ finishReason: "STOP" }], + })}\n\n`, + { status: 200, headers: { "Content-Type": "text/event-stream" } } + ); +} + +// Anonymized replay of the 19/09 zero-chunk payload shape: no candidates at +// all (the upstream turn emitted nothing translatable). +function zeroChunkStreamResponse(): Response { + return new Response( + `data: ${JSON.stringify({ candidates: [] })}\n\n` + + `data: ${JSON.stringify({ candidates: [] })}\n\n`, + { status: 200, headers: { "Content-Type": "text/event-stream" } } + ); +} + +// The upstream sent a first frame with nothing usable in it, then the +// connection dropped (undici surfaces this as a TypeError "terminated" on the +// body). The first frame lets the stream readiness gate pass, so the failure +// happens while the flush hook is buffering the turn. +function droppedStreamResponse(): Response { + let sent = false; + const stream = new ReadableStream({ + pull(controller) { + if (!sent) { + sent = true; + controller.enqueue( + new TextEncoder().encode(`data: ${JSON.stringify({ candidates: [] })}\n\n`) + ); + return; + } + controller.error(new TypeError("terminated")); + }, + }); + return new Response(stream, { status: 200, headers: { "Content-Type": "text/event-stream" } }); +} + +function contentStreamResponse(text: string): Response { + return new Response( + `data: ${JSON.stringify({ + candidates: [{ content: { parts: [{ text }] } }], + })}\n\ndata: ${JSON.stringify({ + candidates: [{ finishReason: "STOP" }], + })}\n\ndata: [DONE]\n\n`, + { status: 200, headers: { "Content-Type": "text/event-stream" } } + ); +} + +function streamRequest() { + const nonce = `flush-empty-retry-${Date.now()}-${Math.random().toString(36).slice(2)}`; + return new Request("http://localhost/v1/chat/completions", { + method: "POST", + headers: { "Content-Type": "application/json", Accept: "text/event-stream" }, + body: JSON.stringify({ + model: "gemini/gemini-2.5-flash", + messages: [{ role: "user", content: `Reply with OK only. ${nonce}` }], + max_tokens: 64, + stream: true, + temperature: 0, + }), + }); +} + +function stubFetch(dispatches: string[], handler: (auth: string, callIndex: number) => Response) { + globalThis.fetch = (async (_url: unknown, init: { headers?: unknown }) => { + const headers = new Headers((init?.headers ?? {}) as HeadersInit); + const auth = headers.get("authorization") ?? headers.get("x-goog-api-key") ?? ""; + const callIndex = dispatches.length; + dispatches.push(auth); + return handler(auth, callIndex); + }) as typeof fetch; +} + +async function drainText(response: Response): Promise { + const text = await response.text().catch(() => ""); + return text; +} + +test.beforeEach(async () => { + globalThis.fetch = originalFetch; + setFlag(true); + await resetStorage(); +}); + +test.afterEach(async () => { + await new Promise((r) => setTimeout(r, 50)); + globalThis.fetch = originalFetch; + if (ORIGINAL_FLAG === undefined) delete process.env[FLAG]; + else process.env[FLAG] = ORIGINAL_FLAG; + await resetStorage(); +}); + +test.after(async () => { + globalThis.fetch = originalFetch; + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); +}); + +test("flag off: empty reasoning-only turn exposes current behavior with zero retry", async () => { + setFlag(false); + await seedGemini("gemini-off-a", "sk-flush-off-a"); + await seedGemini("gemini-off-b", "sk-flush-off-b"); + const dispatches: string[] = []; + stubFetch(dispatches, () => reasoningOnlyStreamResponse()); + const response = await handleChat(streamRequest()); + await drainText(response); + assert.equal( + dispatches.length, + 1, + `flag off must issue exactly 1 dispatch, got ${dispatches.length}` + ); +}); + +test("reasoning-only turn retries once through the normal credential path and serves content", async () => { + await seedGemini("gemini-empty-a", "sk-flush-empty-a"); + await seedGemini("gemini-content-b", "sk-flush-content-b"); + const dispatches: string[] = []; + stubFetch(dispatches, (_auth, callIndex) => + callIndex === 0 ? reasoningOnlyStreamResponse() : contentStreamResponse("served-after-retry") + ); + const response = await handleChat(streamRequest()); + const bodyText = await drainText(response); + assert.equal(dispatches.length, 2, `expected initial + 1 retry, got ${dispatches.length}`); + assert.match(bodyText, /served-after-retry/, "client must receive the retry content"); +}); + +test("persistent empty turns exhaust the budget with exactly 5 dispatches", async () => { + await seedGemini("gemini-empty2-a", "sk-flush-empty2-a"); + await seedGemini("gemini-empty2-b", "sk-flush-empty2-b"); + const dispatches: string[] = []; + stubFetch(dispatches, () => zeroChunkStreamResponse()); + const response = await handleChat(streamRequest()); + await drainText(response); + assert.equal( + dispatches.length, + 5, + `budget 4 means 1 initial + 4 retries, got ${dispatches.length}` + ); +}); + +test("single slot retries the same account and serves content", async () => { + await seedGemini("gemini-lone", "sk-flush-lone"); + const dispatches: string[] = []; + stubFetch(dispatches, (_auth, callIndex) => + callIndex === 0 ? reasoningOnlyStreamResponse() : contentStreamResponse("served-after-retry") + ); + const response = await handleChat(streamRequest()); + const bodyText = await drainText(response); + assert.equal( + dispatches.length, + 2, + `single slot replays the same account, got ${dispatches.length}` + ); + assert.match(bodyText, /served-after-retry/, "client must receive the retry content"); +}); + +test("a stream that drops before anything reaches the client retries and serves content", async () => { + await seedGemini("gemini-drop-a", "sk-flush-drop-a"); + await seedGemini("gemini-drop-b", "sk-flush-drop-b"); + const dispatches: string[] = []; + stubFetch(dispatches, (_auth, callIndex) => + callIndex === 0 ? droppedStreamResponse() : contentStreamResponse("served-after-drop") + ); + const response = await handleChat(streamRequest()); + const bodyText = await drainText(response); + assert.equal(dispatches.length, 2, `expected initial + 1 retry, got ${dispatches.length}`); + assert.match(bodyText, /served-after-drop/, "client must receive the retry content"); +}); + +test("persistent stream drops exhaust the same budget as empty turns", async () => { + await seedGemini("gemini-drop2-a", "sk-flush-drop2-a"); + await seedGemini("gemini-drop2-b", "sk-flush-drop2-b"); + const dispatches: string[] = []; + stubFetch(dispatches, () => droppedStreamResponse()); + const response = await handleChat(streamRequest()); + await drainText(response); + assert.equal( + dispatches.length, + 5, + `budget 4 means 1 initial + 4 retries, got ${dispatches.length}` + ); +}); + +test("flag off: a dropped stream is not retried", async () => { + setFlag(false); + await seedGemini("gemini-drop3-a", "sk-flush-drop3-a"); + await seedGemini("gemini-drop3-b", "sk-flush-drop3-b"); + const dispatches: string[] = []; + stubFetch(dispatches, () => droppedStreamResponse()); + const response = await handleChat(streamRequest()); + await drainText(response); + assert.equal( + dispatches.length, + 1, + `flag off must issue exactly 1 dispatch, got ${dispatches.length}` + ); +}); diff --git a/tests/unit/flush-empty-retry.test.ts b/tests/unit/flush-empty-retry.test.ts new file mode 100644 index 000000000000..a159b99ba50f --- /dev/null +++ b/tests/unit/flush-empty-retry.test.ts @@ -0,0 +1,263 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +const { + isUselessEmptyTurn, + summarizeReplayedUpstreamTurn, + readBoundedResponseText, + readBoundedResponseOutcome, + judgeBufferedTurn, + FLUSH_EMPTY_RETRY_MAX_BYTES, +} = await import("../../open-sse/utils/emptyTurnRetry.ts"); +const { isEmptyTurnCore } = await import("../../open-sse/utils/streamEmptyChoices.ts"); +const { FORMATS } = await import("../../open-sse/translator/formats.ts"); + +const base = { + finishReason: "stop", + contentText: "", + reasoningText: "", + forwardedValuableChunk: false, + hasValidUsage: false, + toolCallsPresent: false, +}; + +test("reasoning-only turn: stop + empty text + non-empty reasoning is an empty turn", () => { + assert.equal(isUselessEmptyTurn({ ...base, reasoningText: "some thinking trace" }), true); +}); + +test("zero-valuable-chunk turn: no valuable chunks and no valid usage is an empty turn", () => { + assert.equal(isUselessEmptyTurn({ ...base, finishReason: "" }), true); +}); + +test("non-empty content text is not an empty turn", () => { + assert.equal(isUselessEmptyTurn({ ...base, contentText: "hello" }), false); +}); + +test("valid usage alone is not an empty turn", () => { + assert.equal(isUselessEmptyTurn({ ...base, hasValidUsage: true }), false); +}); + +test("forwarded valuable chunk alone is not an empty turn", () => { + assert.equal(isUselessEmptyTurn({ ...base, forwardedValuableChunk: true }), false); +}); + +test("legit finish reasons are not empty turns", () => { + for (const finishReason of ["length", "tool_calls", "content_filter"]) { + assert.equal(isUselessEmptyTurn({ ...base, finishReason }), false); + } +}); + +test("tool calls present are not an empty turn", () => { + assert.equal(isUselessEmptyTurn({ ...base, toolCallsPresent: true }), false); +}); + +test("shared core matches the frozen guard condition", () => { + assert.equal(isEmptyTurnCore(false, false), true); + assert.equal(isEmptyTurnCore(true, false), false); + assert.equal(isEmptyTurnCore(false, true), false); +}); + +test("memory cap is a positive bound", () => { + assert.ok(FLUSH_EMPTY_RETRY_MAX_BYTES > 0); +}); + +function sse(...frames: string[]): string { + return frames.map((f) => `data: ${f}\n\n`).join("") + "data: [DONE]\n\n"; +} + +const chatChunk = (delta: Record, finish: string | null = null) => + JSON.stringify({ + id: "chatcmpl-probe", + object: "chat.completion.chunk", + model: "probe", + choices: [{ delta, finish_reason: finish }], + }); + +test("replay: reasoning-only chat turn classifies empty (parity with live transform)", () => { + const body = sse(chatChunk({ reasoning_content: "thinking trace here" }), chatChunk({}, "stop")); + const summary = summarizeReplayedUpstreamTurn(body, FORMATS.OPENAI, FORMATS.OPENAI); + assert.ok(summary, "replay must produce a summary"); + assert.equal(summary.reasoningText.length > 0, true); + assert.equal(summary.contentText, ""); + assert.equal(isUselessEmptyTurn(summary), true); +}); + +test("replay: all-empty choices turn classifies empty", () => { + const body = sse( + JSON.stringify({ + id: "chatcmpl-probe", + object: "chat.completion.chunk", + model: "probe", + choices: [], + }) + ); + const summary = summarizeReplayedUpstreamTurn(body, FORMATS.OPENAI, FORMATS.OPENAI); + assert.ok(summary, "replay must produce a summary"); + assert.equal(isUselessEmptyTurn(summary), true); +}); + +test("replay: turn with real content classifies useful", () => { + const body = sse(chatChunk({ content: "hello" }), chatChunk({}, "stop")); + const summary = summarizeReplayedUpstreamTurn(body, FORMATS.OPENAI, FORMATS.OPENAI); + assert.ok(summary, "replay must produce a summary"); + assert.equal(isUselessEmptyTurn(summary), false); +}); + +test("replay: retry with content after empty first turn serves content", () => { + const attempts = [ + sse(chatChunk({ reasoning_content: "only thinking" }), chatChunk({}, "stop")), + sse(chatChunk({ content: "real answer" }), chatChunk({}, "stop")), + ]; + let calls = 0; + for (const body of attempts) { + calls++; + const summary = summarizeReplayedUpstreamTurn(body, FORMATS.OPENAI, FORMATS.OPENAI); + assert.ok(summary); + if (calls === 1) { + assert.equal(isUselessEmptyTurn(summary), true); + continue; + } + assert.equal(isUselessEmptyTurn(summary), false); + } + assert.equal(calls, 2, "exactly 2 attempts: initial + 1 retry"); +}); + +test("replay: retry also empty stays empty (fall back to current behavior)", () => { + const body = sse( + JSON.stringify({ + id: "chatcmpl-probe", + object: "chat.completion.chunk", + model: "probe", + choices: [], + }) + ); + for (let attempt = 0; attempt < 2; attempt++) { + const summary = summarizeReplayedUpstreamTurn(body, FORMATS.OPENAI, FORMATS.OPENAI); + assert.ok(summary); + assert.equal(isUselessEmptyTurn(summary), true); + } +}); + +test("bounded read abandons past the cap without buffering everything", async () => { + const big = new ReadableStream({ + start(controller) { + controller.enqueue( + new TextEncoder().encode("data: " + "x".repeat(FLUSH_EMPTY_RETRY_MAX_BYTES + 1) + "\n\n") + ); + controller.close(); + }, + }); + const res = new Response(big, { status: 200 }); + const out = await readBoundedResponseText(res, FLUSH_EMPTY_RETRY_MAX_BYTES); + assert.equal(out, null, "past-cap body must be abandoned, not buffered"); +}); + +test("bounded read returns small bodies intact", async () => { + const body = sse(chatChunk({ content: "hi" })); + const res = new Response(body, { status: 200 }); + const out = await readBoundedResponseText(res, FLUSH_EMPTY_RETRY_MAX_BYTES); + assert.equal(out, body); +}); + +test("bounded read outcome tells a read failure apart from an over-cap body", async () => { + const failing = new ReadableStream({ + pull(controller) { + controller.error(new TypeError("terminated")); + }, + }); + const failed = await readBoundedResponseOutcome( + new Response(failing, { status: 200 }), + FLUSH_EMPTY_RETRY_MAX_BYTES + ); + assert.equal(failed.kind, "error", "a stream that throws while being read is a failure"); + + const big = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("x".repeat(FLUSH_EMPTY_RETRY_MAX_BYTES + 1))); + controller.close(); + }, + }); + const skipped = await readBoundedResponseOutcome( + new Response(big, { status: 200 }), + FLUSH_EMPTY_RETRY_MAX_BYTES + ); + assert.equal(skipped.kind, "skipped", "an over-cap body is passed through, not retried"); + + const body = sse(chatChunk({ content: "hi" })); + const ok = await readBoundedResponseOutcome( + new Response(body, { status: 200 }), + FLUSH_EMPTY_RETRY_MAX_BYTES + ); + assert.deepEqual(ok, { kind: "text", text: body }); +}); + +test("buffered turn verdict: a dropped stream retries unless the client went away", () => { + const dropped = { kind: "error" } as const; + const retry = judgeBufferedTurn(dropped, FORMATS.OPENAI_RESPONSES, FORMATS.OPENAI, false); + assert.equal(retry.kind, "retry"); + const gone = judgeBufferedTurn(dropped, FORMATS.OPENAI_RESPONSES, FORMATS.OPENAI, true); + assert.equal(gone.kind, "pass", "a client that disconnected must not cost another upstream call"); + const skipped = judgeBufferedTurn( + { kind: "skipped" }, + FORMATS.OPENAI_RESPONSES, + FORMATS.OPENAI, + false + ); + assert.equal(skipped.kind, "pass", "an unclassifiable turn passes through"); +}); + +test("replay: Responses reasoning-only deltas without completed classify empty", () => { + const body = [ + `data: ${JSON.stringify({ type: "response.reasoning_text.delta", delta: "thinking here" })}\n\n`, + `data: ${JSON.stringify({ type: "response.reasoning_summary_text.done", text: "thinking here" })}\n\n`, + ].join(""); + const summary = summarizeReplayedUpstreamTurn(body, FORMATS.OPENAI_RESPONSES, FORMATS.OPENAI); + assert.ok(summary, "replay must produce a summary"); + assert.equal(isUselessEmptyTurn(summary), true); +}); + +test("a stalled stream with nothing usable is retried, not surfaced as an error", async () => { + // An upstream that answers, sends a keepalive, then goes silent without ever + // closing or erroring: buffering can never end on its own. + const silent = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(": keepalive\n\n")); + }, + }); + const started = Date.now(); + const out = await readBoundedResponseOutcome( + new Response(silent, { status: 200 }), + FLUSH_EMPTY_RETRY_MAX_BYTES, + 50 + ); + assert.equal(out.kind, "idle", "a stream that stops producing must end the buffered read"); + assert.ok(Date.now() - started < 2000, "the read must return on its own budget"); + + const verdict = judgeBufferedTurn(out, FORMATS.OPENAI_RESPONSES, FORMATS.OPENAI, false); + assert.equal(verdict.kind, "retry", "nothing usable was produced: replay, never a bare failure"); + + const gone = judgeBufferedTurn(out, FORMATS.OPENAI_RESPONSES, FORMATS.OPENAI, true); + assert.equal(gone.kind, "pass", "a client that disconnected must not cost another dispatch"); +}); + +test("a stalled stream that already carries content is passed through, not replayed", async () => { + const partial = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(sse(chatChunk({ content: "hello" })))); + }, + }); + const out = await readBoundedResponseOutcome( + new Response(partial, { status: 200 }), + FLUSH_EMPTY_RETRY_MAX_BYTES, + 50 + ); + assert.equal(out.kind, "idle"); + const verdict = judgeBufferedTurn(out, FORMATS.OPENAI, FORMATS.OPENAI, false); + assert.equal(verdict.kind, "pass", "content already produced is worth keeping"); +}); + +test("bounded read keeps no budget when the idle budget is zero", async () => { + const body = sse(chatChunk({ content: "hi" })); + const out = await readBoundedResponseOutcome(new Response(body, { status: 200 }), 256_000, 0); + assert.deepEqual(out, { kind: "text", text: body }, "a disabled budget must not change reads"); +}); diff --git a/tests/unit/server-owned-tool-loop-flag.test.ts b/tests/unit/server-owned-tool-loop-flag.test.ts index 16a28aa15754..87ee22aef331 100644 --- a/tests/unit/server-owned-tool-loop-flag.test.ts +++ b/tests/unit/server-owned-tool-loop-flag.test.ts @@ -68,8 +68,9 @@ describe("isServerOwnedToolLoopEnabled wrapper", () => { describe("feature-flags-settings count update", () => { it("flag count matches updated expected value", () => { - // 893fef9c added OPENCODE_PARK_AND_RESUME (74 -> 75). - assert.equal(FEATURE_FLAG_DEFINITIONS.length, 75); + // 893fef9c added OPENCODE_PARK_AND_RESUME (74 -> 75); FLUSH_EMPTY_RETRY_ENABLED + // (flush empty-turn retry, default off) bumps it to 76. + assert.equal(FEATURE_FLAG_DEFINITIONS.length, 76); }); }); diff --git a/tests/unit/stream-recovery.test.ts b/tests/unit/stream-recovery.test.ts index fa04c55d42e7..f9b6c48c8c51 100644 --- a/tests/unit/stream-recovery.test.ts +++ b/tests/unit/stream-recovery.test.ts @@ -51,6 +51,7 @@ test("STREAM_RECOVERY constants mirror the free-claude-code values", () => { assert.equal(STREAM_RECOVERY.HOLDBACK_MS, 750); assert.equal(STREAM_RECOVERY.BUFFER_MAX_BYTES, 65536); assert.equal(STREAM_RECOVERY.EARLY_RETRY_MAX, 4); + assert.equal(STREAM_RECOVERY.EMPTY_TURN_RETRY_MAX, 4); }); test("HoldbackBuffer holds chunks until flushed, then commits", () => {