Files
9router/open-sse/rtk/headroom.js
Ankit d4d11357ab fix(headroom): translate openai-responses input through OpenAI for compression
Codex (openai-responses) body.input holds Responses items, not OpenAI
messages. Translate input -> OpenAI -> compress -> back to input so the
Responses contract is preserved. Fixes #1998

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-06-26 11:12:10 +07:00

204 lines
7.3 KiB
JavaScript

import { claudeToOpenAIRequest } from "../translator/request/claude-to-openai.js";
import { openaiToClaudeRequest } from "../translator/request/openai-to-claude.js";
import {
openaiResponsesToOpenAIRequest,
openaiToOpenAIResponsesRequest,
} from "../translator/request/openai-responses.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, diagnostics) {
const endpoint = buildCompressEndpoint(url);
diagnostics.endpoint = maskEndpoint(endpoint);
const payload = { messages, model };
if (compressUserMessages) payload.config = { compress_user_messages: true };
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)) {
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, 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)) {
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 Responses shape (Codex): body.input holds Responses items, NOT OpenAI
// messages. Translate input -> OpenAI -> compress -> translate back to input so
// body.input keeps the Responses contract (the proxy only understands OpenAI). (#1998)
if (format === "openai-responses") {
const oai = openaiResponsesToOpenAIRequest(model, body, false);
if (!Array.isArray(oai?.messages)) return null;
const data = await callCompress(url, oai.messages, model, timeoutMs, compressUserMessages, diagnostics || {});
if (!data) return null;
// input: undefined so the translator rebuilds input from the compressed
// messages instead of returning the original input unchanged.
const responsesBody = openaiToOpenAIResponsesRequest(
model,
{ ...oai, input: undefined, messages: data.messages },
false
);
if (Array.isArray(responsesBody?.input)) body.input = responsesBody.input;
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) {
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 (error) {
setDiagnostic(diagnostics, `unexpected error: ${error?.message || String(error)}`);
return null;
}
}
export function formatHeadroomLog(stats) {
if (!stats) return null;
const before = stats.tokens_before || 0;
const after = stats.tokens_after || 0;
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);
}