Files
pxpipe/src/worker.ts
T
Steven Chong 47b5e412b3 feat: cache-stable history-collapse imaging for GPT (OpenAI) + Anthropic
Collapse old conversation history into rendered PNG sections so the model
reads a compact image instead of re-billed text, while preserving prompt
caching and tool-call behavior. Measures real vs compressed token/cost.

Core:
- GPT history collapse (openai-history.ts): append-only, o200k token-length
  sectioning. Sections seal only at a tool-closed boundary (open call-id set
  empty), so a function_call and its function_call_output never split across
  the collapse cut. Fixes the OpenAI 400 "No tool call found for function
  call output with call_id ..." that hit long Responses-API sessions.
- Anthropic cache contract (history.ts): append-only per-chunk rendering;
  cache_control markers are preserved/moved, never added; chunk boundaries
  align with caller marker seams for byte-stable prefix caching.
- GPT image budget (openai.ts): detail:'original' for gpt-5.x, flagship
  vision-multiplier fix, patch cap; schema-strip preserves real descriptions.
- Savings accounting (openai-savings.ts): cached_tokens + vision-token basis.

Model scope (applicability.ts):
- Default imaged scope = claude-fable-5 + gpt-5.6.
- gpt-5.5 and claude-opus-4-8 stay opt-in: same pipeline, but they degrade
  reading dense imaged history (gist drift), so silently imaging them by
  default is wrong. Promotion is gated on an OCR/recall threshold.

Dashboard: GPT + Anthropic rendering, per-family model toggles, persisted
metrics, thumbnail-expired session UI, reflow/newline handling.

Tests: cache-alignment (GPT + Anthropic), history sectioning + tool-boundary
invariants, savings, dashboard, sessions/restart-restore. 452 passing.
2026-06-20 21:11:24 -04:00

104 lines
4.5 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
* Cloudflare Workers entrypoint. Identical proxy logic to the Node build,
* just wired up through the Worker `fetch` export.
*
* Deploy:
* npx wrangler deploy
*
* Dev:
* npx wrangler dev
*
* Config lives in wrangler.toml.
*/
import { createProxy, type ProxyConfig } from './core/proxy.js';
import type { TransformOptions } from './core/transform.js';
import { toTrackEvent, JsonLogTracker, noopTracker, type Tracker } from './core/tracker.js';
export interface Env {
/** Optional single upstream base for every API family. Family-specific env vars override it. */
PXPIPE_UPSTREAM?: string;
ANTHROPIC_UPSTREAM?: string;
/** Optional override — if set, replaces whatever x-api-key the client sent. */
ANTHROPIC_API_KEY?: string;
OPENAI_UPSTREAM?: string;
/** Optional override — if set, replaces whatever Authorization the client sent. */
OPENAI_API_KEY?: string;
COMPRESS?: string;
COMPRESS_TOOLS?: string;
COMPRESS_SCHEMAS?: string;
COMPRESS_REMINDERS?: string;
COMPRESS_TOOL_RESULTS?: string;
MIN_COMPRESS_CHARS?: string;
MIN_REMINDER_CHARS?: string;
MIN_TOOL_RESULT_CHARS?: string;
COLS?: string;
/** R2 multi-column packing — default 1 (off). 2 squeezes ~2× source rows
* per image; OCR-verify before flipping in production. */
MULTI_COL?: string;
/** When "0" / "false", disable per-request event JSON logs. Default-on.
* Cloudflare ingests console.log as Workers Logs; pipe via Logpush to
* R2/S3 for the same JSONL shape Node writes to disk. */
PXPIPE_TRACK?: string;
}
const truthy = (v: string | undefined, fallback: boolean): boolean =>
v == null ? fallback : v === '1' || v.toLowerCase() === 'true';
export default {
async fetch(req: Request, env: Env, _ctx: ExecutionContext): Promise<Response> {
const transform: TransformOptions = {
compress: truthy(env.COMPRESS, true),
compressTools: truthy(env.COMPRESS_TOOLS, true),
compressSchemas: truthy(env.COMPRESS_SCHEMAS, true),
compressReminders: truthy(env.COMPRESS_REMINDERS, true),
compressToolResults: truthy(env.COMPRESS_TOOL_RESULTS, true),
minCompressChars: env.MIN_COMPRESS_CHARS ? Number(env.MIN_COMPRESS_CHARS) : 2000,
// 500 chars — CPU/latency floor only, not a correctness guard. The
// No floors — the content-aware `isCompressionProfitable()` gate
// decides per-block based on actual pixel cost vs text cost. Host
// can still set a floor via env if they want observability buckets
// (e.g. MIN_TOOL_RESULT_CHARS=200 to skip absurdly small dumps).
minReminderChars: env.MIN_REMINDER_CHARS ? Number(env.MIN_REMINDER_CHARS) : 0,
minToolResultChars: env.MIN_TOOL_RESULT_CHARS ? Number(env.MIN_TOOL_RESULT_CHARS) : 0,
cols: env.COLS ? Number(env.COLS) : 100,
// R2 multi-column ON (2 cols) — single-col drops below break-even on
// real tool-doc slabs. Override via MULTI_COL=1 if OCR misreads layout.
multiCol: env.MULTI_COL ? Math.max(1, Number(env.MULTI_COL) | 0) : 2,
};
const trackingOn = truthy(env.PXPIPE_TRACK, true);
// Workers Logs ingests stdout as separate log lines. Emit one JSON line
// per event so downstream (Logpush → R2/S3) reads the same JSONL shape
// the Node host writes to disk.
const tracker: Tracker = trackingOn ? new JsonLogTracker((s) => console.log(s)) : noopTracker;
const sharedUpstream = env.PXPIPE_UPSTREAM;
const config: ProxyConfig = {
upstream: env.ANTHROPIC_UPSTREAM ?? sharedUpstream ?? 'https://api.anthropic.com',
apiKey: env.ANTHROPIC_API_KEY,
openAIUpstream: env.OPENAI_UPSTREAM ?? sharedUpstream ?? 'https://api.openai.com',
openAIApiKey: env.OPENAI_API_KEY,
transform,
onRequest: (e) => {
// Terse human-readable line (separate from the JSON event below;
// shows up in `wrangler tail`).
const tag = e.info?.compressed
? `compressed ${e.info.origChars}ch → ${e.info.imageCount}img/${e.info.imageBytes}B`
: (e.info?.reason ?? '');
const cacheRead = e.usage?.cache_read_input_tokens ?? 0;
console.log(`${e.method} ${e.path}${e.status} (${e.durationMs}ms) ${tag} cache_read=${cacheRead}`);
if (e.info?.unknownStaticTags && e.info.unknownStaticTags.length > 0) {
console.warn(
`[pxpipe warn] unknown tag(s) in static slab: ${e.info.unknownStaticTags.join(', ')}`,
);
}
tracker.emit(toTrackEvent(e));
},
};
const handle = createProxy(config);
return handle(req);
},
};