feat(pxpipe): PXPIPE token saver — multimodal prompt compression (#2465)
Add pxpipe as an experimental fifth Token Saver: Claude-format request bodies above a configurable size threshold are rendered as dense PNGs via the pxpipe-proxy library API (transformAnthropicMessages) before dispatch, cutting estimated input tokens by ~35-60% on token-dense contexts. Integration follows the Headroom pattern: applied to the final body in chatCore just before dispatch, fail-open on any error/timeout. Managed npm install into DATA_DIR/pxpipe, dynamic loader with per-version cache-bust, JSONL event log with rotation, /api/pxpipe/* endpoints, Token Saver card (marked experimental) + /dashboard/pxpipe page, and per-request Activated/Skipped annotation in Request Details. Disabled by default.
This commit is contained in:
committed by
decolua
parent
e1f3399b73
commit
dcf1927f22
@@ -24,6 +24,7 @@ import { injectCaveman } from "../rtk/caveman.js";
|
||||
import { injectPonytail } from "../rtk/ponytail.js";
|
||||
import { compressMessages, formatRtkLog } from "../rtk/index.js";
|
||||
import { compressWithHeadroom, formatHeadroomLog, formatHeadroomSizeLog, isHeadroomPhantomSavings } from "../rtk/headroom.js";
|
||||
import { compressWithPxpipe, formatPxpipeLog } from "../rtk/pxpipe.js";
|
||||
import { getCapabilitiesForModel } from "../providers/capabilities.js";
|
||||
import { stripUnsupportedModalities } from "../translator/concerns/modality.js";
|
||||
import { prefetchRemoteImages } from "../translator/concerns/prefetch.js";
|
||||
@@ -35,7 +36,7 @@ import { prefetchRemoteImages } from "../translator/concerns/prefetch.js";
|
||||
* @param {object} options.credentials - Provider credentials
|
||||
* @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses")
|
||||
*/
|
||||
export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, headroomEnabled, headroomUrl, headroomCompressUserMessages, cavemanEnabled, cavemanLevel, ponytailEnabled, ponytailLevel, sourceFormatOverride, providerThinking }) {
|
||||
export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, headroomEnabled, headroomUrl, headroomCompressUserMessages, cavemanEnabled, cavemanLevel, ponytailEnabled, ponytailLevel, pxpipeEnabled, pxpipeMinChars, pxpipeTimeoutMs, pxpipeTransform, onPxpipeEvent, sourceFormatOverride, providerThinking }) {
|
||||
const { provider, model } = modelInfo;
|
||||
const requestStartTime = Date.now();
|
||||
|
||||
@@ -186,6 +187,21 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
log?.debug?.("PONYTAIL", `${ponytailLevel} | ${finalFormat}`);
|
||||
}
|
||||
|
||||
// PXPIPE: image bulky context (Claude-format bodies only), last saver before dispatch
|
||||
let pxpipeSummary = null;
|
||||
if (pxpipeEnabled) {
|
||||
const pxpipeResult = await compressWithPxpipe(translatedBody, {
|
||||
enabled: true, format: finalFormat, model: upstreamModel,
|
||||
minChars: pxpipeMinChars, timeoutMs: pxpipeTimeoutMs, transform: pxpipeTransform,
|
||||
});
|
||||
pxpipeSummary = pxpipeResult.summary;
|
||||
if (pxpipeResult.body) translatedBody = pxpipeResult.body;
|
||||
const pxpipeLine = formatPxpipeLog(pxpipeSummary);
|
||||
if (pxpipeLine) log?.info?.("PXPIPE", pxpipeLine);
|
||||
else log?.debug?.("PXPIPE", `skipped: ${pxpipeSummary.reason}${pxpipeSummary.detail ? ` (${pxpipeSummary.detail})` : ""}`);
|
||||
try { onPxpipeEvent?.({ provider, model, ...pxpipeSummary }); } catch { /* stats must not break requests */ }
|
||||
}
|
||||
|
||||
const executor = getExecutor(provider);
|
||||
trackPendingRequest(model, provider, connectionId, true);
|
||||
appendRequestLog({ model, provider, connectionId, status: "PENDING" }).catch(() => { });
|
||||
@@ -254,6 +270,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
request: extractRequestConfig(body, stream),
|
||||
providerRequest: translatedBody || null,
|
||||
response: { error: error.message || String(error), status: error.name === "AbortError" ? 499 : 502, thinking: null },
|
||||
pxpipe: pxpipeSummary,
|
||||
status: "error"
|
||||
})).catch(() => { });
|
||||
|
||||
@@ -300,6 +317,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
request: extractRequestConfig(body, stream),
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
response: { error: message, status: statusCode, thinking: null },
|
||||
pxpipe: pxpipeSummary,
|
||||
status: "error"
|
||||
})).catch(() => { });
|
||||
|
||||
@@ -309,7 +327,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
return createErrorResult(statusCode, errMsg, resetsAtMs);
|
||||
}
|
||||
|
||||
const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess };
|
||||
const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, pxpipe: pxpipeSummary };
|
||||
const appendLog = (extra) => appendRequestLog({ model, provider, connectionId, ...extra }).catch(() => { });
|
||||
const trackDone = () => trackPendingRequest(model, provider, connectionId, false);
|
||||
|
||||
|
||||
@@ -198,7 +198,7 @@ export function translateNonStreamingResponse(responseBody, targetFormat, source
|
||||
/**
|
||||
* Handle non-streaming response from provider.
|
||||
*/
|
||||
export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, trackDone, appendLog }) {
|
||||
export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, trackDone, appendLog, pxpipe }) {
|
||||
trackDone();
|
||||
const contentType = providerResponse.headers.get("content-type") || "";
|
||||
let responseBody;
|
||||
@@ -296,6 +296,7 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
thinking: translatedResponse?.choices?.[0]?.message?.reasoning_content || translatedResponse?.reasoning_content || null,
|
||||
finish_reason: translatedResponse?.choices?.[0]?.finish_reason || "unknown"
|
||||
},
|
||||
pxpipe,
|
||||
status: "success"
|
||||
}, { endpoint: clientRawRequest?.endpoint || null })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to save:", err.message);
|
||||
|
||||
@@ -69,6 +69,7 @@ export function buildRequestDetail(base, overrides = {}) {
|
||||
providerRequest: base.providerRequest || null,
|
||||
providerResponse: base.providerResponse || null,
|
||||
response: base.response || {},
|
||||
pxpipe: base.pxpipe || undefined,
|
||||
status: base.status || "success",
|
||||
...overrides
|
||||
};
|
||||
|
||||
@@ -43,7 +43,7 @@ function buildTransformStream({ provider, sourceFormat, targetFormat, userAgent,
|
||||
/**
|
||||
* Handle streaming response — pipe provider SSE through transform stream to client.
|
||||
*/
|
||||
export async function handleStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, userAgent, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, streamController, onStreamComplete, streamDetailId }) {
|
||||
export async function handleStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, userAgent, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, streamController, onStreamComplete, streamDetailId, pxpipe }) {
|
||||
if (onRequestSuccess) {
|
||||
Promise.resolve()
|
||||
.then(onRequestSuccess)
|
||||
@@ -94,6 +94,7 @@ export async function handleStreamingResponse({ providerResponse, provider, mode
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
providerResponse: "[Streaming - raw response not captured]",
|
||||
response: { content: "[Streaming in progress...]", thinking: null, type: "streaming" },
|
||||
pxpipe,
|
||||
status: "success"
|
||||
}, { id: streamDetailId })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to save streaming request:", err.message);
|
||||
@@ -108,7 +109,7 @@ export async function handleStreamingResponse({ providerResponse, provider, mode
|
||||
/**
|
||||
* Build onStreamComplete callback for streaming usage tracking.
|
||||
*/
|
||||
export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest }) {
|
||||
export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest, pxpipe }) {
|
||||
const streamDetailId = `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`;
|
||||
|
||||
const onStreamComplete = (contentObj, usage, ttftAt) => {
|
||||
@@ -127,6 +128,7 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
providerResponse: safeContent,
|
||||
response: { content: safeContent, thinking: safeThinking, type: "streaming" },
|
||||
pxpipe,
|
||||
status: "success"
|
||||
}, { id: streamDetailId })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to update streaming content:", err.message);
|
||||
|
||||
104
open-sse/rtk/pxpipe.js
Normal file
104
open-sse/rtk/pxpipe.js
Normal file
@@ -0,0 +1,104 @@
|
||||
// PXPIPE: render bulky Claude-format context as dense PNGs via pxpipe-proxy's
|
||||
// library API (transformAnthropicMessages). Fail-open like every token saver:
|
||||
// any error/timeout returns { body: null, summary } and leaves the request untouched.
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
|
||||
const DEFAULT_TIMEOUT_MS = 15000;
|
||||
const DEFAULT_MIN_CHARS = 25000;
|
||||
// pxpipe's own profitability gate assumes ~4 chars/token; reuse it for the
|
||||
// estimated before/after numbers surfaced in stats (marked "estimated" in UI).
|
||||
const EST_CHARS_PER_TOKEN = 4;
|
||||
|
||||
function bodyChars(body) {
|
||||
try {
|
||||
return JSON.stringify(body)?.length || 0;
|
||||
} catch {
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
function estTokens(chars) {
|
||||
return Math.round(chars / EST_CHARS_PER_TOKEN);
|
||||
}
|
||||
|
||||
function skipped(reason, extra = {}) {
|
||||
return { body: null, summary: { applied: false, reason, ...extra } };
|
||||
}
|
||||
|
||||
// Transform a Claude-format request body through pxpipe. Returns
|
||||
// { body: <new body object> | null, summary } — body is null when nothing changed.
|
||||
// opts.transform is injected by the host (src side) so open-sse stays free of
|
||||
// filesystem/install concerns and remains usable standalone.
|
||||
export async function compressWithPxpipe(body, { enabled, format, model, minChars, timeoutMs, transform } = {}) {
|
||||
if (!enabled) return skipped("disabled");
|
||||
if (typeof transform !== "function") return skipped("not_installed");
|
||||
if (!body) return skipped("missing_body");
|
||||
if (format !== FORMATS.CLAUDE) return skipped("unsupported_format", { detail: format });
|
||||
|
||||
const startedAt = Date.now();
|
||||
const originalChars = bodyChars(body);
|
||||
const threshold = Number(minChars) > 0 ? Number(minChars) : DEFAULT_MIN_CHARS;
|
||||
if (originalChars < threshold) {
|
||||
return skipped("below_threshold", { originalChars, threshold });
|
||||
}
|
||||
|
||||
try {
|
||||
const encoded = new TextEncoder().encode(JSON.stringify(body));
|
||||
const budget = Number(timeoutMs) > 0 ? Number(timeoutMs) : DEFAULT_TIMEOUT_MS;
|
||||
// transformAnthropicMessages is local CPU work and can't be aborted; race a
|
||||
// timer and discard the result if it loses (input body is never mutated).
|
||||
const result = await Promise.race([
|
||||
transform({
|
||||
body: encoded,
|
||||
model,
|
||||
options: { minCompressChars: threshold },
|
||||
}),
|
||||
new Promise((resolve) => setTimeout(() => resolve(null), budget)),
|
||||
]);
|
||||
if (!result) return skipped("timeout", { originalChars, durationMs: Date.now() - startedAt });
|
||||
if (!result.applied) {
|
||||
return skipped(result.reason || "passthrough", {
|
||||
detail: result.detail,
|
||||
originalChars,
|
||||
durationMs: Date.now() - startedAt,
|
||||
});
|
||||
}
|
||||
|
||||
const newBody = JSON.parse(new TextDecoder().decode(result.body));
|
||||
const compressedBodyChars = bodyChars(newBody);
|
||||
const info = result.info || {};
|
||||
const imagedChars = info.compressedChars || 0;
|
||||
// The transformed body is BIGGER in bytes (base64 PNGs) but cheaper in tokens:
|
||||
// images bill by pixels (Anthropic: pixels/750), not by encoded length. So the
|
||||
// after-estimate is remaining-text tokens + image tokens — never chars/4 of the
|
||||
// new body. Provider-billed usage recorded per request stays the ground truth.
|
||||
const imageTokensEst = info.imageTokens
|
||||
|| (info.imagePixels ? Math.round(info.imagePixels / 750) : (info.imageCount || 0) * 4761);
|
||||
const summary = {
|
||||
applied: true,
|
||||
reason: "applied",
|
||||
originalChars,
|
||||
compressedBodyChars,
|
||||
imagedChars,
|
||||
imageCount: info.imageCount || 0,
|
||||
imageBytes: info.imageBytes || 0,
|
||||
tokensBeforeEst: info.baselineTokens || estTokens(originalChars),
|
||||
tokensAfterEst: estTokens(Math.max(0, originalChars - imagedChars)) + imageTokensEst,
|
||||
durationMs: Date.now() - startedAt,
|
||||
cacheOwnsControl: result.cache?.ownsCacheControl === true,
|
||||
};
|
||||
summary.tokensSavedEst = Math.max(0, summary.tokensBeforeEst - summary.tokensAfterEst);
|
||||
summary.savedPct = summary.tokensBeforeEst > 0
|
||||
? +((summary.tokensSavedEst / summary.tokensBeforeEst) * 100).toFixed(2)
|
||||
: 0;
|
||||
return { body: newBody, summary };
|
||||
} catch (e) {
|
||||
return skipped("transform_error", { detail: e?.message || String(e), originalChars, durationMs: Date.now() - startedAt });
|
||||
}
|
||||
}
|
||||
|
||||
export function formatPxpipeLog(summary) {
|
||||
if (!summary) return null;
|
||||
if (!summary.applied) return null;
|
||||
return `imaged ${summary.imagedChars}ch → ${summary.imageCount} image(s) | est ${summary.tokensBeforeEst}→${summary.tokensAfterEst} tokens (-${summary.savedPct}%) | ${summary.durationMs}ms`;
|
||||
}
|
||||
Reference in New Issue
Block a user