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. diff --git a/open-sse/executors/antigravity/executeAttempt.ts b/open-sse/executors/antigravity/executeAttempt.ts index 276412efaae..793bdb436d1 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; @@ -330,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, @@ -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) { 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; + } +});