Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog.d/fixes/antigravity-replayable-request-body.md
Original file line number Diff line number Diff line change
@@ -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.
29 changes: 7 additions & 22 deletions open-sse/executors/antigravity/executeAttempt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -330,7 +315,7 @@ export async function sendAntigravityRequest(
headers: Record<string, string>,
transformedBody: Record<string, unknown>,
credentials: AntigravityCredentials,
stream: boolean,
_stream: boolean,
signal: AbortSignal | null | undefined,
log: SafeAntigravityLog,
retryAttempt: number,
Expand Down Expand Up @@ -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,
});

Expand All @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down
108 changes: 108 additions & 0 deletions tests/unit/antigravity-replayable-request-body.test.ts
Original file line number Diff line number Diff line change
@@ -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;
}
});
Loading