From 4b8ec9c3d22a4c94be123b6efb7e1535f1764cd9 Mon Sep 17 00:00:00 2001 From: mdigitalbh81 Date: Tue, 29 Sep 2026 00:41:46 -0300 Subject: [PATCH 1/4] fix(antigravity): keep streaming request bodies replayable --- .../executors/antigravity/executeAttempt.ts | 27 +++++-------------- 1 file changed, 6 insertions(+), 21 deletions(-) diff --git a/open-sse/executors/antigravity/executeAttempt.ts b/open-sse/executors/antigravity/executeAttempt.ts index 276412efaae..b7709995be2 100644 --- a/open-sse/executors/antigravity/executeAttempt.ts +++ b/open-sse/executors/antigravity/executeAttempt.ts @@ -196,21 +196,6 @@ export function buildAntigravity429ErrorMessage(errorJson: unknown): string { return errorMessage; } -function getChunkedOrFixedBody(bodyStr: string, stream: boolean): BodyInit { - if (stream) { - return new ReadableStream( - { - async start(controller) { - controller.enqueue(new TextEncoder().encode(bodyStr)); - controller.close(); - }, - }, - { highWaterMark: 16384 } - ); - } - return bodyStr; -} - function cloneAntigravityRequestBody(body: unknown): unknown { if (!body || typeof body !== "object") { return body; @@ -363,11 +348,13 @@ export async function sendAntigravityRequest( "TELEMETRY", `[Antigravity] PhysicalSend - RequestId: ${correlationId ?? "none"}, URL: ${url}, Model: ${model}, PhysicalSend: ${physicalSendOrdinal}, RetryAttempt: ${retryAttempt}` ); + // The Antigravity request payload is finite JSON even when the response is streamed. + // Keep the upload replayable instead of wrapping it in a one-shot ReadableStream; proxyFetch + // can then use its normal replay/fallback path without retaining a duplex upload stream. let response = await fetchAntigravityWithReadinessTimeout(url, { method: "POST", headers: finalHeaders, - body: getChunkedOrFixedBody(serializedRequest.bodyString, stream), - ...(stream ? { duplex: "half" } : {}), + body: serializedRequest.bodyString, signal, }); @@ -384,8 +371,7 @@ export async function sendAntigravityRequest( response = await fetchAntigravityWithReadinessTimeout(url, { method: "POST", headers: retryHeaders, - body: getChunkedOrFixedBody(serializedRequest.bodyString, stream), - ...(stream ? { duplex: "half" } : {}), + body: serializedRequest.bodyString, signal, }); finalHeaders = retryHeaders; @@ -455,8 +441,7 @@ export async function tryCreditsRetry( const creditsResp = await fetchAntigravityWithReadinessTimeout(url, { method: "POST", headers: finalCreditsHeaders, - body: getChunkedOrFixedBody(serializedCreditsRequest.bodyString, stream), - ...(stream ? { duplex: "half" } : {}), + body: serializedCreditsRequest.bodyString, signal, }); if (creditsResp.ok || creditsResp.status !== HTTP_STATUS.RATE_LIMITED) { From f52efd8438903c369530c7f3de74f5b8b41cd20c Mon Sep 17 00:00:00 2001 From: mdigitalbh81 Date: Tue, 29 Sep 2026 00:42:10 -0300 Subject: [PATCH 2/4] test(antigravity): cover replayable streaming request uploads --- tests/unit/executor-antigravity.test.ts | 102 ++++++++++++++++++++++++ 1 file changed, 102 insertions(+) diff --git a/tests/unit/executor-antigravity.test.ts b/tests/unit/executor-antigravity.test.ts index 9040aa70fa9..0d9ace60d9b 100644 --- a/tests/unit/executor-antigravity.test.ts +++ b/tests/unit/executor-antigravity.test.ts @@ -12,6 +12,10 @@ import { } from "../../open-sse/services/antigravityVersion.ts"; import { clearAntigravityProjectCache } from "../../open-sse/services/antigravityProjectBootstrap.ts"; import { runWithCapture } from "../../open-sse/utils/providerRequestLogging.ts"; +import { + sendAntigravityRequest, + tryCreditsRetry, +} from "../../open-sse/executors/antigravity/executeAttempt.ts"; type AntigravityTransformResult = Exclude< Awaited>, @@ -1140,3 +1144,101 @@ test("AntigravityExecutor.transformRequest maps Claude models through Gemini con assert.equal(result.request.temperature, undefined); assert.equal(result.request.toolConfig, undefined); }); + + +test("Antigravity streaming sends finite JSON as a replayable fixed request body", async () => { + const originalFetch = globalThis.fetch; + const sends: Array<{ body: BodyInit | null | undefined; duplex: unknown }> = []; + + globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { + sends.push({ + body: init?.body, + duplex: (init as RequestInit & { duplex?: unknown } | undefined)?.duplex, + }); + return new Response( + 'data: {"response":{"candidates":[{"content":{"parts":[{"text":"ok"}]},"finishReason":"STOP"}]}}\n\n', + { status: sends.length === 1 ? 403 : 200, headers: { "Content-Type": "text/event-stream" } } + ); + }) as typeof fetch; + + try { + const body = { + project: "project-1", + requestId: "agent-test", + request: { contents: [{ role: "user", parts: [{ text: "hello" }] }] }, + model: "gemini-2.5-flash", + userAgent: "antigravity", + requestType: "agent", + }; + const result = await sendAntigravityRequest( + "antigravity", + "https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse", + "gemini-2.5-flash", + { "Content-Type": "application/json" }, + body, + { accessToken: "token", projectId: "project-1" }, + true, + null, + { debug() {}, info() {}, warn() {}, error() {} }, + 0, + { value: 0 }, + "agent-test" + ); + + assert.equal(result.response.status, 200); + assert.equal(sends.length, 2, "403 retry must perform a second physical send"); + for (const send of sends) { + assert.equal(typeof send.body, "string"); + assert.equal(send.duplex, undefined); + } + assert.equal(String(sends[0]?.body), String(sends[1]?.body)); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("Antigravity credits retry also keeps a streaming request body replayable", async () => { + const originalFetch = globalThis.fetch; + let capturedBody: BodyInit | null | undefined; + let capturedDuplex: unknown; + + globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { + capturedBody = init?.body; + capturedDuplex = (init as RequestInit & { duplex?: unknown } | undefined)?.duplex; + return new Response( + 'data: {"response":{"candidates":[{"content":{"parts":[{"text":"ok"}]},"finishReason":"STOP"}]}}\n\n', + { status: 200, headers: { "Content-Type": "text/event-stream" } } + ); + }) as typeof fetch; + + try { + const result = await tryCreditsRetry( + "antigravity", + "https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse", + { "Content-Type": "application/json" }, + { + project: "project-1", + requestId: "agent-test-credits", + request: { contents: [{ role: "user", parts: [{ text: "hello" }] }] }, + model: "gemini-2.5-flash", + userAgent: "antigravity", + requestType: "agent", + }, + { accessToken: "token", projectId: "project-1" }, + true, + null, + { debug() {}, info() {}, warn() {}, error() {} }, + "account-1", + () => {}, + { value: 0 }, + "agent-test-credits" + ); + + assert.equal(result?.response.status, 200); + assert.equal(typeof capturedBody, "string"); + assert.equal(capturedDuplex, undefined); + assert.match(String(capturedBody), /"enabledCreditTypes":\["GOOGLE_ONE_AI"\]/); + } finally { + globalThis.fetch = originalFetch; + } +}); From 4c8bde6e7fc4a5a4e92437c7c76fa537263addcf Mon Sep 17 00:00:00 2001 From: mdigitalbh81 Date: Tue, 29 Sep 2026 00:42:19 -0300 Subject: [PATCH 3/4] docs(changelog): note Antigravity replayable request body fix --- changelog.d/fixes/antigravity-replayable-request-body.md | 1 + 1 file changed, 1 insertion(+) create mode 100644 changelog.d/fixes/antigravity-replayable-request-body.md diff --git a/changelog.d/fixes/antigravity-replayable-request-body.md b/changelog.d/fixes/antigravity-replayable-request-body.md new file mode 100644 index 00000000000..c642e07a4c0 --- /dev/null +++ b/changelog.d/fixes/antigravity-replayable-request-body.md @@ -0,0 +1 @@ +- **fix(antigravity):** send the finite JSON upload as a replayable fixed request body even when the upstream response is streamed, avoiding the one-shot `ReadableStream`/`duplex` upload path and keeping retries replay-safe under large Codex/Responses payloads. From 1789b1dc97d97093f5c9315fd6ff5fa7242e7762 Mon Sep 17 00:00:00 2001 From: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com> Date: Tue, 29 Sep 2026 03:58:36 -0300 Subject: [PATCH 4/4] test(antigravity): move replayable-body regressions to their own file The two new regressions pushed tests/unit/executor-antigravity.test.ts past the 1200-line test-file cap; they now live in tests/unit/antigravity-replayable-request-body.test.ts. The now-unused stream parameter of sendAntigravityRequest is renamed to _stream so eslint passes. --- .../executors/antigravity/executeAttempt.ts | 2 +- ...ntigravity-replayable-request-body.test.ts | 108 ++++++++++++++++++ tests/unit/executor-antigravity.test.ts | 102 ----------------- 3 files changed, 109 insertions(+), 103 deletions(-) create mode 100644 tests/unit/antigravity-replayable-request-body.test.ts diff --git a/open-sse/executors/antigravity/executeAttempt.ts b/open-sse/executors/antigravity/executeAttempt.ts index b7709995be2..793bdb436d1 100644 --- a/open-sse/executors/antigravity/executeAttempt.ts +++ b/open-sse/executors/antigravity/executeAttempt.ts @@ -315,7 +315,7 @@ export async function sendAntigravityRequest( headers: Record, transformedBody: Record, credentials: AntigravityCredentials, - stream: boolean, + _stream: boolean, signal: AbortSignal | null | undefined, log: SafeAntigravityLog, retryAttempt: number, diff --git a/tests/unit/antigravity-replayable-request-body.test.ts b/tests/unit/antigravity-replayable-request-body.test.ts new file mode 100644 index 00000000000..c44751e8993 --- /dev/null +++ b/tests/unit/antigravity-replayable-request-body.test.ts @@ -0,0 +1,108 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; + +import { + sendAntigravityRequest, + tryCreditsRetry, +} from "../../open-sse/executors/antigravity/executeAttempt.ts"; + +// The Antigravity request payload is finite JSON even when the response is streamed, so every +// physical send (first attempt, 403 retry, Google One AI credits retry) must carry a replayable +// string body with no `duplex` upload stream (#5770 follow-up). + +test("Antigravity streaming sends finite JSON as a replayable fixed request body", async () => { + const originalFetch = globalThis.fetch; + const sends: Array<{ body: BodyInit | null | undefined; duplex: unknown }> = []; + + globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { + sends.push({ + body: init?.body, + duplex: (init as (RequestInit & { duplex?: unknown }) | undefined)?.duplex, + }); + return new Response( + 'data: {"response":{"candidates":[{"content":{"parts":[{"text":"ok"}]},"finishReason":"STOP"}]}}\n\n', + { status: sends.length === 1 ? 403 : 200, headers: { "Content-Type": "text/event-stream" } } + ); + }) as typeof fetch; + + try { + const body = { + project: "project-1", + requestId: "agent-test", + request: { contents: [{ role: "user", parts: [{ text: "hello" }] }] }, + model: "gemini-2.5-flash", + userAgent: "antigravity", + requestType: "agent", + }; + const result = await sendAntigravityRequest( + "antigravity", + "https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse", + "gemini-2.5-flash", + { "Content-Type": "application/json" }, + body, + { accessToken: "token", projectId: "project-1" }, + true, + null, + { debug() {}, info() {}, warn() {}, error() {} }, + 0, + { value: 0 }, + "agent-test" + ); + + assert.equal(result.response.status, 200); + assert.equal(sends.length, 2, "403 retry must perform a second physical send"); + for (const send of sends) { + assert.equal(typeof send.body, "string"); + assert.equal(send.duplex, undefined); + } + assert.equal(String(sends[0]?.body), String(sends[1]?.body)); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("Antigravity credits retry also keeps a streaming request body replayable", async () => { + const originalFetch = globalThis.fetch; + let capturedBody: BodyInit | null | undefined; + let capturedDuplex: unknown; + + globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { + capturedBody = init?.body; + capturedDuplex = (init as (RequestInit & { duplex?: unknown }) | undefined)?.duplex; + return new Response( + 'data: {"response":{"candidates":[{"content":{"parts":[{"text":"ok"}]},"finishReason":"STOP"}]}}\n\n', + { status: 200, headers: { "Content-Type": "text/event-stream" } } + ); + }) as typeof fetch; + + try { + const result = await tryCreditsRetry( + "antigravity", + "https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse", + { "Content-Type": "application/json" }, + { + project: "project-1", + requestId: "agent-test-credits", + request: { contents: [{ role: "user", parts: [{ text: "hello" }] }] }, + model: "gemini-2.5-flash", + userAgent: "antigravity", + requestType: "agent", + }, + { accessToken: "token", projectId: "project-1" }, + true, + null, + { debug() {}, info() {}, warn() {}, error() {} }, + "account-1", + () => {}, + { value: 0 }, + "agent-test-credits" + ); + + assert.equal(result?.response.status, 200); + assert.equal(typeof capturedBody, "string"); + assert.equal(capturedDuplex, undefined); + assert.match(String(capturedBody), /"enabledCreditTypes":\["GOOGLE_ONE_AI"\]/); + } finally { + globalThis.fetch = originalFetch; + } +}); diff --git a/tests/unit/executor-antigravity.test.ts b/tests/unit/executor-antigravity.test.ts index 0d9ace60d9b..9040aa70fa9 100644 --- a/tests/unit/executor-antigravity.test.ts +++ b/tests/unit/executor-antigravity.test.ts @@ -12,10 +12,6 @@ import { } from "../../open-sse/services/antigravityVersion.ts"; import { clearAntigravityProjectCache } from "../../open-sse/services/antigravityProjectBootstrap.ts"; import { runWithCapture } from "../../open-sse/utils/providerRequestLogging.ts"; -import { - sendAntigravityRequest, - tryCreditsRetry, -} from "../../open-sse/executors/antigravity/executeAttempt.ts"; type AntigravityTransformResult = Exclude< Awaited>, @@ -1144,101 +1140,3 @@ test("AntigravityExecutor.transformRequest maps Claude models through Gemini con assert.equal(result.request.temperature, undefined); assert.equal(result.request.toolConfig, undefined); }); - - -test("Antigravity streaming sends finite JSON as a replayable fixed request body", async () => { - const originalFetch = globalThis.fetch; - const sends: Array<{ body: BodyInit | null | undefined; duplex: unknown }> = []; - - globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { - sends.push({ - body: init?.body, - duplex: (init as RequestInit & { duplex?: unknown } | undefined)?.duplex, - }); - return new Response( - 'data: {"response":{"candidates":[{"content":{"parts":[{"text":"ok"}]},"finishReason":"STOP"}]}}\n\n', - { status: sends.length === 1 ? 403 : 200, headers: { "Content-Type": "text/event-stream" } } - ); - }) as typeof fetch; - - try { - const body = { - project: "project-1", - requestId: "agent-test", - request: { contents: [{ role: "user", parts: [{ text: "hello" }] }] }, - model: "gemini-2.5-flash", - userAgent: "antigravity", - requestType: "agent", - }; - const result = await sendAntigravityRequest( - "antigravity", - "https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse", - "gemini-2.5-flash", - { "Content-Type": "application/json" }, - body, - { accessToken: "token", projectId: "project-1" }, - true, - null, - { debug() {}, info() {}, warn() {}, error() {} }, - 0, - { value: 0 }, - "agent-test" - ); - - assert.equal(result.response.status, 200); - assert.equal(sends.length, 2, "403 retry must perform a second physical send"); - for (const send of sends) { - assert.equal(typeof send.body, "string"); - assert.equal(send.duplex, undefined); - } - assert.equal(String(sends[0]?.body), String(sends[1]?.body)); - } finally { - globalThis.fetch = originalFetch; - } -}); - -test("Antigravity credits retry also keeps a streaming request body replayable", async () => { - const originalFetch = globalThis.fetch; - let capturedBody: BodyInit | null | undefined; - let capturedDuplex: unknown; - - globalThis.fetch = (async (_url: string | URL | Request, init?: RequestInit) => { - capturedBody = init?.body; - capturedDuplex = (init as RequestInit & { duplex?: unknown } | undefined)?.duplex; - return new Response( - 'data: {"response":{"candidates":[{"content":{"parts":[{"text":"ok"}]},"finishReason":"STOP"}]}}\n\n', - { status: 200, headers: { "Content-Type": "text/event-stream" } } - ); - }) as typeof fetch; - - try { - const result = await tryCreditsRetry( - "antigravity", - "https://cloudcode-pa.googleapis.com/v1internal:streamGenerateContent?alt=sse", - { "Content-Type": "application/json" }, - { - project: "project-1", - requestId: "agent-test-credits", - request: { contents: [{ role: "user", parts: [{ text: "hello" }] }] }, - model: "gemini-2.5-flash", - userAgent: "antigravity", - requestType: "agent", - }, - { accessToken: "token", projectId: "project-1" }, - true, - null, - { debug() {}, info() {}, warn() {}, error() {} }, - "account-1", - () => {}, - { value: 0 }, - "agent-test-credits" - ); - - assert.equal(result?.response.status, 200); - assert.equal(typeof capturedBody, "string"); - assert.equal(capturedDuplex, undefined); - assert.match(String(capturedBody), /"enabledCreditTypes":\["GOOGLE_ONE_AI"\]/); - } finally { - globalThis.fetch = originalFetch; - } -});