Files
OmniRoute/open-sse/handlers/chatCore/nonStreamingSse.ts
Bob.Hou 60580ffeb7 fix(antigravity): streaming passthrough for non-streaming clients (#7408)
* chore(ci): add .mergify.yml to main — Mergify only reads config from the default branch (#7168)

* fix(ci): add the auto-enqueue pull_request_rule to the Mergify config (queue_conditions alone are eligibility-only) (#7179)

* fix(ci): migrate Mergify auto-enqueue to merge_protections_settings.auto_merge_conditions (rules-based path is EOL 2026-07-16) (#7216)

* fix(ci): drop Mergify batch settings (batching is a paid-tier feature; free plan queue is serial) (#7220)

* fix(ci): merge queue tolerates the advisory dast-smoke failure (its GH-hosted build hang dequeued every attempt) (#7225)

* test(ci): make the #6634 selfref guard hermetic — main's copy hard-fails every PR (#7341)

main's copy of this test still does git I/O inside a unit test:

    const baseSrc = git(['show', 'origin/main:' + FILE]);

Runners check out a shallow single ref, so origin/main does not resolve and the
test dies with 'fatal: invalid object name origin/main'. Every PR into main
fails Unit Tests (7/8) on it — today that is #7313, #7315, #7316, #7334, #7336
and #7337, six PRs red on a defect none of them introduced. #7313 has no other
red at all.

release/v3.8.49 already carries a fix (2e42b8efc, #7174: try/catch, fetch
origin/main on demand, t.skip() when unreachable), but it only reaches main at
release time — so main stays broken for the whole cycle. Cherry-picking it would
also import a new problem: PR Test Policy classifies t.skip() as a silenced
assertion, which we watched it correctly catch on #7300 today.

This is the hermetic version instead (ported from #7327, which does the same for
the release branch): read the file straight off disk, compare against an empty
base so baseTaut/baseExtTaut are 0 — the strictest possible comparison point —
and call evaluateMasking() directly. No git ref, no fetch, no skip, nothing the
runner's checkout depth can break.

The #6634 regression stays covered: the guard's logic lives in
SELF_TEST_FIXTURE_RE (check-test-masking.mjs:337), not in the test. Proven both
ways on main before committing — neutralise SELF_TEST_FIXTURE_RE to /$^/ and
the test FAILS; restore it and it passes 2/2, with check-test-masking.mjs left
byte-identical.

Co-authored-by: growab <nekron@icloud.com>

* chore(quality): tighten main's coverage baseline to the CI's real numbers (#7347)

main's ratchet had been failing --require-tighten on every PR: 11 metrics
improved but the baseline was never tightened. Same class as the #6634
selfref guard — an infra fix that lands only on the release branch leaves
main red for the whole cycle, and every PR into main pays for it.

Values are the merged-coverage numbers from a run on main itself (a local
run measures ~68% vs CI's ~80%; the baseline's own note warns about that
gap). Only the 11 coverage values change — gitleaks and semgrepFindings
keep main's own state.

No changelog fragment: #7326 carries it on release/v3.8.49, and a second
one here would double the entry at release time.

* fix(antigravity): remove hardcoded 120s SSE collect timeout

The SSE collection in collectStreamToResponse had a hardcoded 120 s
timeout.  Reasoning-heavy models like gemini-3.1-pro-high on large
prompts (>30 KB) regularly exceed 120 s of generation time, causing
the executor to return a synthetic 504 before the model finishes.

Replace the hardcoded value with FETCH_TIMEOUT_MS (default 600 s,
overridable via FETCH_TIMEOUT_MS env var), which is the standard
upstream-request budget across all OmniRoute providers.

Signed-off-by: Minxi Hou <houminxi@gmail.com>

* fix(antigravity): streaming passthrough for non-streaming clients

When a client sends stream: false to the Antigravity executor
(Gemini models), OmniRoute buffered the entire SSE stream before
responding. Long-thinking models exceeded the 120s timeout.

Remove hardcoded SSE_COLLECT_TIMEOUT_MS. Extract shared
createCreditsExtractionTransform with 16KB buffer cap and abort
handling for client disconnect. Add parseSSEToGeminiResponse for
the non-streaming drain path. Fix hasGeminiTerminalFinishReason
to check top-level candidates (no response wrapper). Add signal
null guards for credits retry path. Return 499 on early abort
instead of piping cancelled body.

Also remove duplicate SKILLS_SANDBOX_RUNTIME from .env.example
and clarify .artifacts/ vs _artifacts/ in .gitignore.

Signed-off-by: Minxi Hou <houminxi@gmail.com>

* refactor(antigravity): extract streaming passthrough to module (file-size cap)

#7408 added the non-streaming SSE pass-through (createCreditsExtractionTransform
plus its two call sites: the credits-retry path and the main non-streaming
path) inline in antigravity.ts, growing it to 1806 lines. Combined with two
other authorized PRs touching the same file (#6979 +11, #7290 +30), the
projected total exceeds the frozen file-size gate (1813).

Extract the new streaming-passthrough logic verbatim into
open-sse/executors/antigravity/streamingPassthrough.ts
(createCreditsExtractionTransform + a new buildSsePassthroughResult that
deduplicates the two near-identical call sites), following the existing
sseCollect.ts submodule pattern -- pure, no host state, no fetch/auth.
antigravity.ts keeps a thin wrapper for createCreditsExtractionTransform
(same public signature the existing unit tests import) that injects
updateAntigravityRemainingCredits so the two modules don't import each
other.

No behavior change: same abort handling, same 499-on-early-disconnect,
same 16KB credits sliding-window cap. antigravity.ts: 1806 -> 1693 lines
(under the 1755 pre-PR baseline, with margin). New module: 176 lines
(cap 800).

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

* refactor(antigravity): split incremental parser + move new tests to own file (file-size caps)

Two remaining frozen file-size violations from #7408, resolved by
extraction/move with zero behavior or assert changes:

- open-sse/handlers/sseParser.ts (979 > frozen 830): the PR appended
  parseSSEToGeminiResponse (+153, the Gemini buffered-SSE ->
  chat.completion parser). Moved verbatim to
  open-sse/handlers/sseParser/geminiResponse.ts, following the handlers
  submodule pattern (chatCore/, responseSanitizer/). sseParser.ts is now
  byte-identical to its pre-PR content (825 lines; PR delta 0). Importers
  (chatCore/nonStreamingSse.ts, tests) point at the new module.

- tests/unit/executor-antigravity.test.ts (1058 > testFrozen 942): the
  PR's new streaming-passthrough tests moved verbatim (same tests, same
  asserts) to tests/unit/antigravity-streaming-passthrough.test.ts:
  the 3 createCreditsExtractionTransform tests plus the non-streaming
  passthrough drain test ("auto-retries short 429 ... collects SSE for
  non-stream clients"), which the PR rewired onto the new raw-SSE path.
  The frozen file drops to 888 lines (below its pre-PR 941).

New files: geminiResponse.ts 156 lines, passthrough test 202 lines (caps
800). Also fixes the stale sseParser.ts path in collectStreamToResponse's
deprecation note.

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

* refactor(antigravity): decompose execute + gemini parser below complexity gate

executeOnce() (complexity 127, 436 lines) and parseSSEToGeminiResponse()
(complexity 39, 117 lines) were both over the check-complexity.mjs gate
(complexity>15, max-lines-per-function>80). Decomposed each into small
named helpers, no behavior change:

- geminiResponse.ts: split into pure per-concern functions (markdown
  shortcut, candidate-parts walk, finishReason, usageMetadata, final
  response assembly).
- antigravity.ts: extracted the per-url-index attempt pipeline
  (runAntigravityAttempt, handleAntigravityRateLimit,
  tryResolveRetryFromErrorBody, shouldAutoRetryTransient) and moved the
  request/result-building helpers (send, credits-retry, embed-retry,
  non-streaming/streaming result builders) into a new
  antigravity/executeAttempt.ts submodule, mirroring the existing
  streamingPassthrough.ts/sseCollect.ts pattern. Also fixes the
  antigravity.ts file-size cap (was pushed to 2084 lines > 1813 frozen
  ceiling by the decomposition itself; now 1428).

check-complexity.mjs: 2054 violations (baseline 2058) — net improvement.
execute/executeOnce/parseSSEToGeminiResponse no longer appear with
ruleId complexity or max-lines-per-function.

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

* Merge branch 'release/v3.8.49' into fix/antigravity-streaming-passthrough

Resolves conflict in open-sse/executors/antigravity.ts between this
branch's streaming-passthrough decomposition and #7290's fallback-chain
decomposition (already merged into release/v3.8.49) — both sides added
imports from the same new antigravity/ submodule files, kept both.

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

* fix(antigravity): keep buffered JSON contract for non-streaming callers

#3786's Pro-family fallback-chain retry loop (execute()) calls executeOnce()
per candidate and inspects result.response directly, expecting a
synthesized chat.completion JSON body on success. The streaming-passthrough
migration made ALL non-streaming (stream: false) responses a raw SSE
pass-through instead, so a successful retry candidate's response.json()
threw ("data: {...}" is not valid JSON) — breaking the fallback chain
(tests/unit/agy-pro-fallback-chain-3786.test.ts, 3 of 13 red).

Route non-streaming (stream: false) responses back through
collectStreamToResponse (buffered collect-to-JSON), which already uses
FETCH_TIMEOUT_MS with no hardcoded 120s ceiling, so long-thinking models
are not penalized. Passthrough is reserved for actual streaming clients
(stream: true), which was the PR's real target scenario.

Extracted the branch into buildAntigravityAttemptResult() to keep
runAntigravityAttempt under the 80-line ratchet cap.

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

---------

Signed-off-by: Minxi Hou <houminxi@gmail.com>
Co-authored-by: Diego Rodrigues de Sa e Souza <8016841+diegosouzapw@users.noreply.github.com>
Co-authored-by: growab <nekron@icloud.com>
Co-authored-by: HouMinXi <1000+HouMinXi@users.noreply.github.com>
Co-authored-by: HouMinXi <19586012+HouMinXi@users.noreply.github.com>
2026-07-18 21:19:20 -03:00

203 lines
7.0 KiB
TypeScript

import { FORMATS } from "../../translator/formats.ts";
import {
parseSSEToResponsesOutput,
parseSSEToClaudeResponse,
parseSSEToOpenAIResponse,
} from "../sseParser.ts";
import { parseSSEToGeminiResponse } from "../sseParser/geminiResponse.ts";
import { getHeaderValueCaseInsensitive } from "./headers.ts";
export function parseNonStreamingSSEPayload(
rawBody: string,
preferredFormat: string,
fallbackModel: string
): { body: Record<string, unknown>; format: string } | null {
const formatsToTry: string[] = [];
const seen = new Set<string>();
const queueFormat = (format: string) => {
if (!format || seen.has(format)) return;
seen.add(format);
formatsToTry.push(format);
};
queueFormat(preferredFormat);
queueFormat(FORMATS.GEMINI);
queueFormat(FORMATS.OPENAI_RESPONSES);
queueFormat(FORMATS.CLAUDE);
queueFormat(FORMATS.OPENAI);
for (const format of formatsToTry) {
const parsed =
format === FORMATS.OPENAI_RESPONSES
? parseSSEToResponsesOutput(rawBody, fallbackModel)
: format === FORMATS.CLAUDE
? parseSSEToClaudeResponse(rawBody, fallbackModel)
: format === FORMATS.GEMINI || format === FORMATS.ANTIGRAVITY
? parseSSEToGeminiResponse(rawBody, fallbackModel)
: parseSSEToOpenAIResponse(rawBody, fallbackModel);
if (parsed && typeof parsed === "object") {
return {
body: parsed as Record<string, unknown>,
format,
};
}
}
return null;
}
export function convertNDJSONToSSE(rawBody: string): string {
const chunks = String(rawBody || "")
.split(/\r?\n/)
.map((line) => line.trim())
.filter((line) => line.length > 0);
if (chunks.length === 0) return rawBody;
return `${chunks.map((chunk) => `data: ${chunk}\n`).join("\n")}\n`;
}
export function normalizeNonStreamingEventPayload(rawBody: string, contentType: string): string {
if (contentType.includes("application/x-ndjson")) {
return convertNDJSONToSSE(rawBody);
}
return rawBody;
}
export function isTruthyStreamBody(body: unknown): boolean {
return !!body && typeof body === "object" && (body as { stream?: unknown }).stream === true;
}
export function isEventStreamAccepted(
headers: Record<string, unknown> | Headers | null | undefined
) {
return (getHeaderValueCaseInsensitive(headers, "accept") || "")
.toLowerCase()
.includes("text/event-stream");
}
export function shouldTreatBufferedEventResponseAsExpected(
upstreamStream: boolean,
providerHeaders: Record<string, unknown> | Headers | null | undefined,
finalBody: unknown
): boolean {
return upstreamStream || isEventStreamAccepted(providerHeaders) || isTruthyStreamBody(finalBody);
}
const NON_STREAMING_SSE_TERMINAL_TYPES = new Set([
"message_stop",
"response.completed",
"response.done",
"response.cancelled",
"response.canceled",
"response.failed",
"response.incomplete",
]);
function isNonStreamingSseTerminalType(eventType: string): boolean {
return NON_STREAMING_SSE_TERMINAL_TYPES.has(eventType);
}
export type NonStreamingSseTerminalState = {
currentEvent: string;
pendingLine: string;
};
function hasClaudeTerminalMessageDelta(parsed: unknown, eventType: string): boolean {
if (eventType !== "message_delta" || !parsed || typeof parsed !== "object") return false;
const delta = (parsed as { delta?: unknown }).delta;
if (!delta || typeof delta !== "object") return false;
const stopReason = (delta as { stop_reason?: unknown }).stop_reason;
return typeof stopReason === "string" ? stopReason.length > 0 : stopReason != null;
}
// Non-empty finishReason is terminal. Gemini SSE payloads from
// streamGenerateContent have candidates at the top level (no
// "response" wrapper). Any non-empty string signals stream end.
function hasGeminiTerminalFinishReason(parsed: unknown): boolean {
if (!parsed || typeof parsed !== "object") return false;
// Top-level candidates (streamGenerateContent?alt=sse)
const obj = parsed as Record<string, unknown>;
const candidates = obj.candidates as unknown[] | undefined;
if (!Array.isArray(candidates) || candidates.length === 0) return false;
const candidate = candidates[0] as Record<string, unknown> | undefined;
if (!candidate || typeof candidate !== "object") return false;
const finishReason = candidate.finishReason;
return typeof finishReason === "string" && finishReason.length > 0;
}
function processNonStreamingSseTerminalLine(
state: NonStreamingSseTerminalState,
rawLine: string
): boolean {
const trimmed = rawLine.trim();
if (!trimmed || trimmed.startsWith(":")) {
const terminalEventOnly = !trimmed && isNonStreamingSseTerminalType(state.currentEvent);
if (!trimmed) state.currentEvent = "";
return terminalEventOnly;
}
if (trimmed.startsWith("event:")) {
state.currentEvent = trimmed.slice(6).trim();
return false;
}
if (!trimmed.startsWith("data:")) return false;
const data = trimmed.slice(5).trim();
if (data === "[DONE]") return true;
if (!data) return false;
// Hot-path optimization: the terminal SSE events we look for (message_stop,
// response.completed, …) all carry a top-level "type" field, OR are signalled by a
// preceding `event:` line (Claude). Gemini signals completion via
// "finishReason" inside response.candidates[0]. OpenAI chat.completion chunks
// carry none of these and terminate with `[DONE]` (handled above), so parsing
// every one of them here is pure waste that compounds into the CPU-runaway on
// large buffered responses. Skip the JSON.parse unless the line could actually
// be a typed terminal.
if (
!data.includes('"type"') &&
// NOTE: "finishReason" is a superset match -- it triggers JSON.parse on
// every Gemini chunk that happens to contain the string (e.g. partial
// candidate payloads), not just the terminal one. This is intentional:
// the extra parses are cheap compared to the CPU-runaway we'd get from
// parsing ALL chunks unconditionally on large buffered responses, and
// the superset is safe (false positives just parse a non-terminal chunk
// and fall through to `return false`).
!data.includes('"finishReason"') &&
!(state.currentEvent === "message_delta" && data.includes("stop_reason"))
) {
return isNonStreamingSseTerminalType(state.currentEvent);
}
try {
const parsed = JSON.parse(data);
const eventType =
parsed && typeof parsed === "object" && typeof parsed.type === "string"
? parsed.type
: state.currentEvent;
return (
isNonStreamingSseTerminalType(eventType) ||
hasClaudeTerminalMessageDelta(parsed, eventType) ||
hasGeminiTerminalFinishReason(parsed)
);
} catch {
// Keep reading malformed data so the parser can report a useful upstream error.
return false;
}
}
export function appendNonStreamingSseTerminalSignal(
state: NonStreamingSseTerminalState,
chunk: string
): boolean {
const lines = `${state.pendingLine}${chunk}`.split(/\r?\n/);
state.pendingLine = lines.pop() ?? "";
for (const rawLine of lines) {
if (processNonStreamingSseTerminalLine(state, rawLine)) return true;
}
return false;
}