Skip to content
Closed
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
91 changes: 71 additions & 20 deletions open-sse/executors/antigravity.js
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,26 @@ function sanitizeFunctionName(name) {
}

const MAX_RETRY_AFTER_MS = 10000;
const ANTIGRAVITY_TRANSIENT_RETRY_MAX_MS = 15000;
const MAX_ANTIGRAVITY_OUTPUT_TOKENS = 16384;

const ANTIGRAVITY_TRANSIENT_ERROR_PATTERNS = [
/high\s+traffic/i,
/agent\s+(execution\s+)?terminated\s+due\s+to\s+error/i,
/capacity/i,
/temporarily\s+unavailable/i,
/timeout/i,
/stream\s+(ended|closed|terminated|interrupted)/i,
/empty\s+response/i,
];

const ANTIGRAVITY_TRANSIENT_STATUSES = new Set([
HTTP_STATUS.SERVER_ERROR,
HTTP_STATUS.BAD_GATEWAY,
HTTP_STATUS.SERVICE_UNAVAILABLE,
HTTP_STATUS.GATEWAY_TIMEOUT,
]);

// Fields Google generateContent rejects (Claude/OpenAI/Qwen thinking fields set at body root by thinkingUnified.js)
const ANTIGRAVITY_REQUEST_BLACKLIST = [
"output_config",
Expand Down Expand Up @@ -170,15 +188,22 @@ export class AntigravityExecutor extends BaseExecutor {

if (tools && tools.length > 0) {
// Merge all groups into a single functionDeclarations group (Gemini expects 1 group)
const allDeclarations = tools.flatMap(group =>
(group.functionDeclarations || []).map(fn => ({
...fn,
name: sanitizeFunctionName(fn.name),
parameters: fn.parameters
? cleanJSONSchemaForAntigravity(structuredClone(fn.parameters))
: { type: "object", properties: { reason: { type: "string", description: "Brief explanation" } }, required: ["reason"] }
}))
);
const seenToolNames = new Set();
const allDeclarations = [];
for (const group of tools) {
for (const fn of group.functionDeclarations || []) {
const name = sanitizeFunctionName(fn.name);
if (seenToolNames.has(name)) continue;
seenToolNames.add(name);
allDeclarations.push({
...fn,
name,
parameters: fn.parameters
? cleanJSONSchemaForAntigravity(structuredClone(fn.parameters))
: { type: "object", properties: { reason: { type: "string", description: "Brief explanation" } }, required: ["reason"] }
});
}
}
tools = allDeclarations.length > 0 ? [{ functionDeclarations: allDeclarations }] : [];
}

Expand Down Expand Up @@ -305,23 +330,49 @@ export class AntigravityExecutor extends BaseExecutor {
return totalMs > 0 ? totalMs : null;
}

extractErrorMessage(errorJson, bodyText = "") {
return [
errorJson?.error?.message,
errorJson?.message,
errorJson?.error,
bodyText,
].filter(Boolean).map(v => typeof v === "string" ? v : JSON.stringify(v)).join("\n");
}

isTransientAntigravityError(status, message) {
if (status === HTTP_STATUS.RATE_LIMITED) return true;
if (ANTIGRAVITY_TRANSIENT_STATUSES.has(status)) return true;
return ANTIGRAVITY_TRANSIENT_ERROR_PATTERNS.some(pattern => pattern.test(message || ""));
}

// Hook called by BaseExecutor.tryRetry: derive delay from Retry-After (header → body),
// cap at MAX_RETRY_AFTER_MS, else exponential backoff for 429. Return false to veto (fallback URL).
// cap at MAX_RETRY_AFTER_MS, else retry transient Antigravity failures with backoff.
// Return false to veto (fallback URL / final error).
async computeRetryDelay(response, attempt) {
let bodyText = "";
let errorJson = null;
let retryMs = this.parseRetryHeaders(response.headers);

try {
bodyText = await response.clone().text();
errorJson = bodyText ? JSON.parse(bodyText) : null;
} catch {
// ignore parse errors → fall through to status/message based retry
}

const errorMessage = this.extractErrorMessage(errorJson, bodyText);

if (!retryMs) {
try {
const errorJson = JSON.parse(await response.clone().text());
retryMs = this.parseRetryFromErrorMessage(errorJson?.error?.message || errorJson?.message || "");
} catch {
// ignore parse errors → fall through to backoff
}
retryMs = this.parseRetryFromErrorMessage(errorMessage);
}
if (retryMs) return retryMs <= MAX_RETRY_AFTER_MS ? retryMs : false;
if (response.status === HTTP_STATUS.RATE_LIMITED) {
return Math.min(1000 * (2 ** attempt), MAX_RETRY_AFTER_MS); // exponential backoff
}
return false;

if (!this.isTransientAntigravityError(response.status, errorMessage)) return false;

const cap = response.status === HTTP_STATUS.RATE_LIMITED
? MAX_RETRY_AFTER_MS
: ANTIGRAVITY_TRANSIENT_RETRY_MAX_MS;
return Math.min(1000 * (2 ** attempt), cap); // exponential backoff
}

/**
Expand Down
13 changes: 10 additions & 3 deletions open-sse/handlers/chatCore.js
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import { dedupeTools } from "../utils/toolDeduper.js";
import { injectCaveman } from "../rtk/caveman.js";
import { injectPonytail } from "../rtk/ponytail.js";
import { compressMessages, formatRtkLog } from "../rtk/index.js";
import { compressWithHeadroom, formatHeadroomLog } from "../rtk/headroom.js";
import { compressWithHeadroom, formatHeadroomLog, formatHeadroomSizeLog, isHeadroomPhantomSavings } from "../rtk/headroom.js";
import { getCapabilitiesForModel } from "../providers/capabilities.js";
import { stripUnsupportedModalities } from "../translator/concerns/modality.js";
import { prefetchRemoteImages } from "../translator/concerns/prefetch.js";
Expand Down Expand Up @@ -162,9 +162,16 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
if (rtkLine) console.log(rtkLine);

// Headroom: optional external proxy compression; fail open if proxy is absent.
const headroomStats = await compressWithHeadroom(translatedBody, { enabled: headroomEnabled, url: headroomUrl, model: upstreamModel, format: finalFormat, compressUserMessages: headroomCompressUserMessages });
const headroomDiagnostics = {};
const headroomStats = await compressWithHeadroom(translatedBody, { enabled: headroomEnabled, url: headroomUrl, model: upstreamModel, format: finalFormat, compressUserMessages: headroomCompressUserMessages, diagnostics: headroomDiagnostics });
const headroomLine = formatHeadroomLog(headroomStats);
if (headroomLine) log?.info?.("HEADROOM", headroomLine);
const headroomSizeLine = formatHeadroomSizeLog(headroomDiagnostics);
if (headroomLine) {
log?.info?.("HEADROOM", `${headroomLine}${headroomSizeLine ? ` | ${headroomSizeLine}` : ""}`);
if (isHeadroomPhantomSavings(headroomStats, headroomDiagnostics)) {
log?.warn?.("HEADROOM", `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${headroomSizeLine}`);
}
} else if (headroomEnabled) log?.warn?.("HEADROOM", `skipped: ${headroomDiagnostics.reason || "compression unavailable"}${headroomDiagnostics.endpoint ? ` (${headroomDiagnostics.endpoint})` : ""}`);

// Caveman: inject terse-style system prompt
if (cavemanEnabled && cavemanLevel) {
Expand Down
3 changes: 3 additions & 0 deletions open-sse/providers/registry/antigravity.js
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,9 @@ export default {
"429": {
attempts: 3,
},
"500": {
attempts: 3,
},
"503": {
attempts: 3,
},
Expand Down
156 changes: 136 additions & 20 deletions open-sse/rtk/headroom.js
Original file line number Diff line number Diff line change
Expand Up @@ -3,52 +3,153 @@ import { openaiToClaudeRequest } from "../translator/request/openai-to-claude.js

const DEFAULT_TIMEOUT_MS = 3000;

function jsonBytes(value) {
try {
return new TextEncoder().encode(JSON.stringify(value) || "").length;
} catch {
return 0;
}
}

function messagePayload(body) {
if (Array.isArray(body?.messages)) return body.messages;
if (Array.isArray(body?.input)) return body.input;
return null;
}

function captureSizeSnapshot(body) {
const messages = messagePayload(body);
return {
bodyBytes: jsonBytes(body),
messageBytes: messages ? jsonBytes(messages) : 0,
};
}

function setDiagnostic(diagnostics, reason) {
if (diagnostics && !diagnostics.reason) diagnostics.reason = reason;
}

function scrubSensitiveUrlText(text) {
return String(text)
.replace(/\/\/[^/@\s]+@/g, "//")
.replace(/(https?:\/\/[^\s?#]+)[?#][^\s)]*/g, "$1");
}

function describeFetchError(error) {
const cause = error?.cause;
const code = cause?.code || error?.code;
const message = scrubSensitiveUrlText(cause?.message || error?.message || String(error));
return code ? `${code}: ${message}` : message;
}

function buildCompressEndpoint(url) {
try {
const parsed = new URL(url);
parsed.pathname = `${parsed.pathname.replace(/\/$/, "")}/v1/compress`;
parsed.hash = "";
return parsed.toString();
} catch {
const raw = String(url).replace(/#.*$/, "");
const [base, query = ""] = raw.split("?", 2);
const endpoint = `${base.replace(/\/$/, "")}/v1/compress`;
return query ? `${endpoint}?${query}` : endpoint;
}
}

function maskEndpoint(endpoint) {
try {
const parsed = new URL(endpoint);
parsed.username = "";
parsed.password = "";
parsed.search = "";
parsed.hash = "";
return parsed.toString();
} catch {
return String(endpoint).replace(/\/\/[^/@\s]+@/, "//").replace(/[?#].*$/, "");
}
}

// POST messages to Headroom /v1/compress; returns compressed messages + stats or null.
async function callCompress(url, messages, model, timeoutMs, compressUserMessages) {
const endpoint = `${String(url).replace(/\/$/, "")}/v1/compress`;
async function callCompress(url, messages, model, timeoutMs, compressUserMessages, diagnostics) {
const endpoint = buildCompressEndpoint(url);
diagnostics.endpoint = maskEndpoint(endpoint);
const payload = { messages, model };
if (compressUserMessages) payload.config = { compress_user_messages: true };
const res = await fetch(endpoint, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(payload),
signal: AbortSignal.timeout(timeoutMs),
});
if (!res.ok) return null;
let res;
try {
res = await fetch(endpoint, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(payload),
signal: AbortSignal.timeout(timeoutMs),
});
} catch (error) {
setDiagnostic(diagnostics, `request failed: ${describeFetchError(error)}`);
return null;
}
if (!res.ok) {
setDiagnostic(diagnostics, `proxy returned HTTP ${res.status}`);
return null;
}
const data = await res.json();
if (!Array.isArray(data?.messages)) return null;
if (!Array.isArray(data?.messages)) {
setDiagnostic(diagnostics, "proxy response missing messages[]");
return null;
}
return data;
}

// Compress request body via Headroom proxy. Fail-open: returns null on any error.
// /v1/compress only understands OpenAI shape, so Claude bodies are translated
// to OpenAI, compressed, then translated back using 9Router's own translators.
export async function compressWithHeadroom(body, { enabled, url, model, format, compressUserMessages, timeoutMs = DEFAULT_TIMEOUT_MS } = {}) {
if (!enabled || !url || !body) return null;
export async function compressWithHeadroom(body, { enabled, url, model, format, compressUserMessages, timeoutMs = DEFAULT_TIMEOUT_MS, diagnostics = null } = {}) {
if (!enabled) {
setDiagnostic(diagnostics, "disabled");
return null;
}
if (!url) {
setDiagnostic(diagnostics, "missing proxy URL");
return null;
}
if (!body) {
setDiagnostic(diagnostics, "missing request body");
return null;
}

try {
if (diagnostics) diagnostics.before = captureSizeSnapshot(body);

// Claude shape: translate → OpenAI → compress → translate back.
if (format === "claude") {
const oai = claudeToOpenAIRequest(model, body, false);
if (!Array.isArray(oai?.messages)) return null;
const data = await callCompress(url, oai.messages, model, timeoutMs, compressUserMessages);
if (!Array.isArray(oai?.messages)) {
setDiagnostic(diagnostics, "Claude request did not translate to messages[]");
return null;
}
const data = await callCompress(url, oai.messages, model, timeoutMs, compressUserMessages, diagnostics || {});
if (!data) return null;
const claudeBody = openaiToClaudeRequest(model, { ...oai, messages: data.messages }, false);
if (Array.isArray(claudeBody?.messages)) body.messages = claudeBody.messages;
if (claudeBody?.system !== undefined) body.system = claudeBody.system;
if (diagnostics) diagnostics.after = captureSizeSnapshot(body);
return data;
}

// OpenAI shape: messages/input go straight to the proxy.
const key = Array.isArray(body.messages) ? "messages"
: Array.isArray(body.input) ? "input"
: null;
if (!key) return null;
const data = await callCompress(url, body[key], model, timeoutMs, compressUserMessages);
if (!key) {
setDiagnostic(diagnostics, `unsupported ${format || "unknown"} request shape`);
return null;
}
const data = await callCompress(url, body[key], model, timeoutMs, compressUserMessages, diagnostics || {});
if (!data) return null;
body[key] = data.messages;
if (diagnostics) diagnostics.after = captureSizeSnapshot(body);
return data;
} catch {
} catch (error) {
setDiagnostic(diagnostics, `unexpected error: ${error?.message || String(error)}`);
return null;
}
}
Expand All @@ -57,7 +158,22 @@ export function formatHeadroomLog(stats) {
if (!stats) return null;
const before = stats.tokens_before || 0;
const after = stats.tokens_after || 0;
const saved = stats.tokens_saved || 0;
const pct = before > 0 ? ((saved / before) * 100).toFixed(1) : "0";
return `saved ${saved} tokens / ${before} (${pct}%) ${after ? `after=${after}` : ""}`.trim();
const delta = stats.tokens_saved || 0;
const pct = before > 0 ? ((delta / before) * 100).toFixed(1) : "0";
return `reported token delta=${delta} before=${before}${after ? ` after=${after}` : ""} (${pct}%)`.trim();
}

export function formatHeadroomSizeLog(diagnostics) {
const before = diagnostics?.before;
const after = diagnostics?.after;
if (!before || !after) return "";
return `body=${before.bodyBytes}B→${after.bodyBytes}B messages=${before.messageBytes}B→${after.messageBytes}B`;
}

export function isHeadroomPhantomSavings(stats, diagnostics, minShrinkRatio = 0.05) {
if (!stats?.tokens_saved || stats.tokens_saved <= 0) return false;
const before = diagnostics?.before?.bodyBytes || 0;
const after = diagnostics?.after?.bodyBytes || 0;
if (before <= 0 || after <= 0) return false;
return after >= before * (1 - minShrinkRatio);
}
34 changes: 32 additions & 2 deletions tests/unit/antigravity-retry-hook.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,38 @@ describe("antigravity computeRetryDelay hook (D3)", () => {
expect(await ag.computeRetryDelay(res(429), 3)).toBe(Math.min(1000 * 2 ** 3, MAX));
});

it("503 without retry info → veto (no auto backoff)", async () => {
expect(await ag.computeRetryDelay(res(503), 1)).toBe(false);
it("503 without retry info → transient backoff", async () => {
expect(await ag.computeRetryDelay(res(503), 1)).toBe(2000);
});

it("retries Antigravity agent terminated body even when status is not 429", async () => {
const r = res(500, {}, { error: { message: "Agent execution terminated due to error" } });
expect(await ag.computeRetryDelay(r, 1)).toBe(2000);
});

it("retries high traffic body", async () => {
const r = res(500, {}, { error: { message: "Our servers are experiencing high traffic" } });
expect(await ag.computeRetryDelay(r, 2)).toBe(4000);
});

it("does not retry non-transient 400 errors", async () => {
const r = res(400, {}, { error: { message: "Invalid request" } });
expect(await ag.computeRetryDelay(r, 1)).toBe(false);
});

it("deduplicates sanitized tool names", () => {
const out = ag.transformRequest("claude-opus-4-6-thinking", {
request: {
contents: [{ role: "user", parts: [{ text: "hi" }] }],
tools: [{ functionDeclarations: [
{ name: "read/file", parameters: { type: "object", properties: {} } },
{ name: "read file", parameters: { type: "object", properties: {} } },
{ name: "read/file", parameters: { type: "object", properties: {} } },
] }],
},
}, true, { projectId: "project-1", connectionId: "conn-1" });

expect(out.request.tools[0].functionDeclarations.map(fn => fn.name)).toEqual(["read_file"]);
});

it("buildHeaders includes cached session id after transformRequest", () => {
Expand Down
Loading