[eric] attach: 9router bypass + Anthropic SSE translator for OpenAI/OR PDFs + 2-req OOM guard

This commit is contained in:
ciregenz
2026-05-23 00:50:28 -07:00
parent a30229bbda
commit cbe5372471
5 changed files with 542 additions and 24 deletions
+8 -7
View File
@@ -1187,15 +1187,16 @@ class AgentManager:
# - Anthropic: native document blocks pass through cleanly.
# - Gemini: anthropic_proxy rewrites document → image_url with
# data:application/pdf base64; 9router translates to Gemini
# inlineData natively. VERIFIED empirically May 2026 (47K
# prompt tokens, real PaLM content summarized).
# inlineData natively.
# - OpenRouter: file-parser plugin injected in anthropic-proxy.
# - OpenAI: REFUSED. OpenAI's image_url only accepts image/*
# mime; the type:file shape gets stringified by 9router.
# Workaround: route through `openrouter/openai/gpt-5` which
# uses OR's file-parser plugin.
# - OpenAI direct (GPT-5.x non-codex): anthropic_proxy detects
# document block + bypasses 9router entirely, translating
# to OpenAI Chat Completions and streaming response back
# via anthropic_to_openai.py. Requires openai_api_key.
# - Codex (cx/): models don't support PDFs.
supports_pdf = api in ("anthropic", "gemini", "gemini-cli", "openrouter")
supports_pdf = api in ("anthropic", "gemini", "gemini-cli", "openrouter", "openai")
if api == "openai" and isinstance(model, str) and ("codex" in model.lower() or model.lower().startswith("cx/")):
supports_pdf = False
# Per-file inline caps (raw bytes, before base64). Going over
# means the request would 4xx, blow our 64MB SDK buffer, or
+39
View File
@@ -387,6 +387,45 @@ async def proxy(rest: str, request: Request):
except Exception:
pass
# 9router-bypass paths for PDF-bearing requests on providers where
# 9router 0.3.60 strips or mangles the relevant content/plugin
# fields. We translate + POST directly to the provider's API and
# convert the streaming response back to Anthropic SSE so the
# bundled Claude CLI subprocess consumes it unchanged.
try:
parsed_for_bypass = json.loads(body) if body else None
except Exception:
parsed_for_bypass = None
if isinstance(parsed_for_bypass, dict):
from backend.apps.agents.anthropic_to_openai import (
should_bypass_9router as _should_bypass_oai,
should_bypass_9router_for_openrouter as _should_bypass_or,
forward_to_openai as _forward_oai,
forward_to_openrouter as _forward_or,
)
from backend.apps.settings.settings import load_settings as _load
_s = _load()
if _is_openai_max_completion_tokens_model(model):
_oak = (getattr(_s, "openai_api_key", "") or "").strip()
if _should_bypass_oai(parsed_for_bypass, _oak):
status, body_stream, hdrs = await _forward_oai(
parsed_for_bypass, _oak, dict(request.headers),
)
return StreamingResponse(
body_stream, status_code=status, headers=hdrs,
media_type=hdrs.get("content-type", "text/event-stream"),
)
if _is_openrouter_model(model):
_ork = (getattr(_s, "openrouter_api_key", "") or "").strip()
if _should_bypass_or(parsed_for_bypass, _ork):
status, body_stream, hdrs = await _forward_or(
parsed_for_bypass, _ork, dict(request.headers),
)
return StreamingResponse(
body_stream, status_code=status, headers=hdrs,
media_type=hdrs.get("content-type", "text/event-stream"),
)
if _is_gemini_model(model):
body = _scrub_request_for_gemini(body)
if _is_openai_max_completion_tokens_model(model):
+404
View File
@@ -0,0 +1,404 @@
"""Anthropic Messages API ↔ OpenAI Chat Completions translator.
Used to bypass 9router 0.3.60 for OpenAI requests that include document
blocks (PDFs). 9router's content-block filter strips any block that
isn't `text` or `image_url`, so PDFs sent as OpenAI's native `type:file`
shape never reach the model. This translator POSTs directly to
api.openai.com using the user's OpenAI API key and converts the
streaming response back to Anthropic Messages SSE format so the
bundled Claude CLI subprocess can consume it unchanged.
Scope: PDFs + images + text. Tool-use translation is NOT covered (the
PDF-attach flow does not need tools in the same turn). If a request has
both tools and documents, fall back to the 9router path which handles
tool-use but strips PDFs (refused upstream in agent_manager).
"""
from __future__ import annotations
import asyncio
import json
import logging
import time
import uuid
from typing import AsyncIterator
import httpx
logger = logging.getLogger(__name__)
_OPENAI_UPSTREAM = "https://api.openai.com/v1"
_OPENROUTER_UPSTREAM = "https://openrouter.ai/api/v1"
# Concurrency cap for bypass-route requests. Each in-flight request holds
# the base64'd PDF (raw_bytes * 1.33) in memory across httpx's request
# pipeline + our SSE translator's chunk buffer + the response body. A
# 30MB PDF is ~40MB base64; four in-flight = ~160MB of httpx buffers
# plus Python overhead, which OOM-killed the dev backend on macOS during
# concurrent probes. Cap at 2 so a Mehmet-style multi-PDF attach in one
# session can't take the whole backend down. Requests above the cap
# queue rather than fail.
_BYPASS_CONCURRENCY = 2
_bypass_sema = asyncio.Semaphore(_BYPASS_CONCURRENCY)
# Hard per-request body size ceiling. Anthropic API caps at 32MB,
# OpenAI Chat Completions at 50MB, OpenRouter at whatever underlying
# model accepts. We refuse anything over 40MB raw (≈53MB base64) before
# we even build the request body, so a malicious or accidental huge
# attach never reaches the in-memory pipeline.
_BYPASS_MAX_RAW_BYTES = 40 * 1024 * 1024
def _has_document_block(parsed: dict) -> bool:
msgs = parsed.get("messages")
if not isinstance(msgs, list):
return False
for m in msgs:
content = m.get("content") if isinstance(m, dict) else None
if not isinstance(content, list):
continue
for block in content:
if isinstance(block, dict) and block.get("type") == "document":
return True
return False
def should_bypass_9router(parsed: dict, api_key: str | None) -> bool:
"""True iff request is a GPT-5.x Chat Completions with at least one
document block AND user has an OpenAI API key. Anything else falls
through the normal 9router path."""
if not api_key:
return False
model = (parsed.get("model") or "").lower()
if not any(model.startswith(p) for p in ("gpt-5", "openai/gpt-5", "cp-openai/gpt-5")):
return False
if "codex" in model:
return False
if parsed.get("tools"):
return False
return _has_document_block(parsed)
def should_bypass_9router_for_openrouter(parsed: dict, api_key: str | None) -> bool:
"""True iff request is bound for OpenRouter AND has document blocks
AND user has an OpenRouter API key. 9router 0.3.60 doesn't know
about OR's `plugins` field and silently strips it; we bypass to
inject the file-parser plugin and POST directly to openrouter.ai."""
if not api_key:
return False
model = (parsed.get("model") or "").lower()
if not (model.startswith("openrouter/") or model.startswith("or:")):
return False
if parsed.get("tools"):
return False
return _has_document_block(parsed)
def _content_blocks_to_openai(content) -> list[dict]:
"""Convert Anthropic content blocks → OpenAI Chat Completions parts."""
if isinstance(content, str):
return [{"type": "text", "text": content}]
if not isinstance(content, list):
return [{"type": "text", "text": str(content)}]
out: list[dict] = []
file_counter = 0
for block in content:
if not isinstance(block, dict):
continue
btype = block.get("type")
if btype == "text":
txt = block.get("text") or ""
if txt:
out.append({"type": "text", "text": txt})
elif btype == "image":
src = block.get("source") or {}
if src.get("type") == "base64" and src.get("data"):
mt = src.get("media_type") or "image/png"
out.append({
"type": "image_url",
"image_url": {"url": f"data:{mt};base64,{src['data']}"},
})
elif btype == "document":
src = block.get("source") or {}
if src.get("type") == "base64" and src.get("data"):
file_counter += 1
mt = src.get("media_type") or "application/pdf"
out.append({
"type": "file",
"file": {
"filename": f"attachment_{file_counter}.pdf",
"file_data": f"data:{mt};base64,{src['data']}",
},
})
if not out:
out.append({"type": "text", "text": ""})
return out
def translate_request(parsed: dict) -> dict:
"""Anthropic Messages request → OpenAI Chat Completions request."""
model = parsed.get("model") or ""
if "/" in model:
model = model.split("/", 1)[1]
openai_body: dict = {"model": model, "stream": True}
sys = parsed.get("system")
msgs_out: list[dict] = []
if sys:
if isinstance(sys, str):
msgs_out.append({"role": "system", "content": sys})
elif isinstance(sys, list):
sys_text = "\n".join(
b.get("text", "") for b in sys
if isinstance(b, dict) and b.get("type") == "text"
)
if sys_text:
msgs_out.append({"role": "system", "content": sys_text})
for m in (parsed.get("messages") or []):
if not isinstance(m, dict):
continue
role = m.get("role")
if role not in ("user", "assistant"):
continue
msgs_out.append({"role": role, "content": _content_blocks_to_openai(m.get("content"))})
openai_body["messages"] = msgs_out
mt = parsed.get("max_tokens")
if isinstance(mt, int) and mt > 0:
openai_body["max_completion_tokens"] = mt
if isinstance(parsed.get("temperature"), (int, float)):
openai_body["temperature"] = parsed["temperature"]
# OpenAI omits usage from streamed chunks unless explicitly asked.
# Without this, our Anthropic message_delta would always report 0
# tokens, breaking cost tracking + the context meter for bypass-route
# turns. OpenRouter respects the same flag.
openai_body["stream_options"] = {"include_usage": True}
return openai_body
def _sse_event(event: str, data: dict) -> bytes:
"""Encode an Anthropic-format SSE event."""
return f"event: {event}\ndata: {json.dumps(data)}\n\n".encode("utf-8")
async def _translate_response_stream(
upstream: httpx.Response, model: str,
) -> AsyncIterator[bytes]:
"""Convert OpenAI Chat Completions SSE → Anthropic Messages SSE.
Emits message_start, content_block_start (text block at index 0),
content_block_delta per chunk, then content_block_stop +
message_delta + message_stop on completion.
"""
msg_id = f"msg_{uuid.uuid4().hex[:24]}"
started = False
block_opened = False
output_tokens = 0
input_tokens = 0
stop_reason = "end_turn"
buffer = b""
try:
async for chunk in upstream.aiter_bytes():
if not chunk:
continue
buffer += chunk
while b"\n\n" in buffer:
raw_event, buffer = buffer.split(b"\n\n", 1)
line = raw_event.decode("utf-8", errors="replace").strip()
if not line:
continue
for ln in line.split("\n"):
# SSE comments (`:` prefix) are keep-alives, e.g.
# OpenRouter emits `: OPENROUTER PROCESSING` while
# its file-parser plugin works. Drop them.
if ln.startswith(":"):
continue
if not ln.startswith("data:"):
continue
payload = ln[5:].strip()
if payload == "[DONE]":
continue
try:
ev = json.loads(payload)
except Exception:
continue
if not started:
usage = (ev.get("usage") or {})
input_tokens = int(usage.get("prompt_tokens") or 0)
yield _sse_event("message_start", {
"type": "message_start",
"message": {
"id": msg_id,
"type": "message",
"role": "assistant",
"content": [],
"model": model,
"stop_reason": None,
"stop_sequence": None,
"usage": {
"input_tokens": input_tokens,
"output_tokens": 0,
},
},
})
started = True
choices = ev.get("choices") or []
if not choices:
usage = ev.get("usage") or {}
if usage:
output_tokens = int(usage.get("completion_tokens") or output_tokens)
input_tokens = int(usage.get("prompt_tokens") or input_tokens)
continue
choice = choices[0]
delta = choice.get("delta") or {}
delta_text = delta.get("content")
if isinstance(delta_text, str) and delta_text:
if not block_opened:
yield _sse_event("content_block_start", {
"type": "content_block_start",
"index": 0,
"content_block": {"type": "text", "text": ""},
})
block_opened = True
yield _sse_event("content_block_delta", {
"type": "content_block_delta",
"index": 0,
"delta": {"type": "text_delta", "text": delta_text},
})
finish = choice.get("finish_reason")
if finish:
if finish == "length":
stop_reason = "max_tokens"
elif finish == "tool_calls":
stop_reason = "tool_use"
else:
stop_reason = "end_turn"
finally:
if started:
if block_opened:
yield _sse_event("content_block_stop", {
"type": "content_block_stop", "index": 0,
})
yield _sse_event("message_delta", {
"type": "message_delta",
"delta": {"stop_reason": stop_reason, "stop_sequence": None},
"usage": {"input_tokens": input_tokens, "output_tokens": output_tokens},
})
yield _sse_event("message_stop", {"type": "message_stop"})
async def forward_to_openai(
parsed: dict, api_key: str, headers_in: dict[str, str],
) -> tuple[int, AsyncIterator[bytes], dict[str, str]]:
"""Translate + forward an Anthropic request to OpenAI Chat Completions.
Returns (status, body_stream, response_headers)."""
openai_body = translate_request(parsed)
return await _forward(openai_body, api_key, f"{_OPENAI_UPSTREAM}/chat/completions")
async def forward_to_openrouter(
parsed: dict, api_key: str, headers_in: dict[str, str],
) -> tuple[int, AsyncIterator[bytes], dict[str, str]]:
"""Translate + forward to OpenRouter, injecting the file-parser plugin
so any OR model parses the attached PDFs server-side."""
openai_body = translate_request(parsed)
model = (parsed.get("model") or "").lower()
bare = model
for prefix in ("openrouter/", "or:"):
if bare.startswith(prefix):
bare = bare[len(prefix):]
break
openai_body["model"] = bare
openai_body["plugins"] = [{"id": "file-parser", "pdf": {"engine": "pdf-text"}}]
return await _forward(openai_body, api_key, f"{_OPENROUTER_UPSTREAM}/chat/completions")
def _estimate_body_bytes(body_json: dict) -> int:
"""Sum the base64 payload bytes across content blocks. Used as a
cheap pre-flight check before httpx serializes the body."""
total = 0
for m in body_json.get("messages") or []:
content = m.get("content") if isinstance(m, dict) else None
if not isinstance(content, list):
continue
for block in content:
if not isinstance(block, dict):
continue
if block.get("type") == "image_url":
url = (block.get("image_url") or {}).get("url") or ""
if "," in url:
total += len(url.split(",", 1)[1])
elif block.get("type") == "file":
fd = (block.get("file") or {}).get("file_data") or ""
if "," in fd:
total += len(fd.split(",", 1)[1])
return total
async def _forward(
body_json: dict, api_key: str, url: str,
) -> tuple[int, AsyncIterator[bytes], dict[str, str]]:
# Pre-flight size check. base64 expands ~4/3 so 40MB raw → 53MB b64.
raw_estimate = int(_estimate_body_bytes(body_json) * 0.75)
if raw_estimate > _BYPASS_MAX_RAW_BYTES:
async def reject():
payload = json.dumps({
"type": "error",
"error": {
"type": "invalid_request_error",
"message": (
f"Attached files total ~{raw_estimate // (1024*1024)} MB, "
f"over the {_BYPASS_MAX_RAW_BYTES // (1024*1024)} MB per-request "
"cap on this provider lane. Detach a file or split across "
"separate turns."
),
},
}).encode("utf-8")
yield payload
return 413, reject(), {"content-type": "application/json"}
headers = {
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
"Accept": "text/event-stream",
}
# Acquire the bypass-concurrency semaphore before opening a streaming
# connection. Without this, N simultaneous PDF attaches each hold a
# ~40MB request body + a streaming response buffer, and the OS
# OOM-kills the backend (observed on macOS during a 3-PDF probe
# burst). Semaphore serializes excess requests instead of failing.
await _bypass_sema.acquire()
client = httpx.AsyncClient(timeout=httpx.Timeout(600.0, connect=30.0))
try:
req = client.build_request("POST", url, json=body_json, headers=headers)
upstream = await client.send(req, stream=True)
except Exception:
await client.aclose()
_bypass_sema.release()
raise
async def streamer():
try:
if upstream.status_code >= 400:
raw = await upstream.aread()
yield raw
return
async for chunk in _translate_response_stream(upstream, body_json["model"]):
yield chunk
finally:
try:
await upstream.aclose()
finally:
try:
await client.aclose()
finally:
_bypass_sema.release()
return upstream.status_code, streamer(), {
"content-type": "text/event-stream" if upstream.status_code < 400 else "application/json",
}
+78 -11
View File
@@ -1677,13 +1677,11 @@ def test_resolve_attachments_anthropic_emits_native_document():
os.unlink(path)
def test_resolve_attachments_openai_refuses_pdf_with_openrouter_hint():
"""Empirical probe May 2026: OpenAI image_url rejects non-image mime
types with HTTP 400 'Invalid MIME type'. The type:file shape gets
stringified by 9router 0.3.60. Until we write a 9router-bypass
direct-API translator, refuse OpenAI PDFs with switch hint to
openrouter/openai/gpt-5 which has working PDF support via OR's
file-parser plugin."""
def test_resolve_attachments_openai_accepts_pdf_via_bypass_translator():
"""OpenAI GPT-5.x non-codex accepts PDFs because anthropic_proxy
detects document blocks + bypasses 9router via anthropic_to_openai.
No refusal at agent_manager level; the bypass kicks in at proxy
time when openai_api_key is set."""
import tempfile, os
from backend.apps.agents.agent_manager import AgentManager
mgr = AgentManager()
@@ -1694,14 +1692,83 @@ def test_resolve_attachments_openai_refuses_pdf_with_openrouter_hint():
_text, native, refusals = mgr._resolve_attachments(
[{"path": path, "type": "file"}], api_type="openai", model="gpt-5.5",
)
assert not native
assert refusals
joined = " ".join(refusals).lower()
assert "openrouter" in joined or "claude" in joined
assert native and native[0]["type"] == "document"
assert not refusals
finally:
os.unlink(path)
def test_bypass_estimate_body_bytes_sums_image_url_and_file_blocks():
"""The size estimator must sum payload bytes across BOTH content
types the translator emits (image_url with data: URL, file with
file_data) so the pre-flight reject can fire before httpx serializes."""
from backend.apps.agents.anthropic_to_openai import _estimate_body_bytes
body = {"messages": [{"role": "user", "content": [
{"type": "text", "text": "hi"},
{"type": "image_url", "image_url": {"url": "data:image/png;base64,QUJDRA=="}}, # 8 bytes b64
{"type": "file", "file": {"file_data": "data:application/pdf;base64,RUZHSA=="}}, # 8 bytes b64
]}]}
assert _estimate_body_bytes(body) == 16
def test_bypass_concurrency_semaphore_serializes_excess_requests():
"""The semaphore caps in-flight bypass requests to prevent OOM. Cap=2
means a third concurrent request waits rather than allocating another
~40MB buffer."""
from backend.apps.agents.anthropic_to_openai import _bypass_sema, _BYPASS_CONCURRENCY
assert _BYPASS_CONCURRENCY == 2
# Initial value matches the cap (no in-flight at import time).
assert _bypass_sema._value == _BYPASS_CONCURRENCY
def test_anthropic_to_openai_should_bypass_fires_for_gpt5_pdf():
"""The bypass only fires for: GPT-5.x non-codex + has document block
+ openai_api_key set. Misses any of those: 9router path."""
from backend.apps.agents.anthropic_to_openai import should_bypass_9router
body = {"model": "gpt-5.5", "messages": [{"role": "user", "content": [
{"type": "document", "source": {"type": "base64", "media_type": "application/pdf", "data": "x"}},
]}]}
assert should_bypass_9router(body, "sk-abc")
assert not should_bypass_9router(body, None)
assert not should_bypass_9router(body, "")
assert not should_bypass_9router({**body, "model": "gpt-5.3-codex"}, "sk-abc")
body_no_doc = {"model": "gpt-5.5", "messages": [{"role": "user", "content": "hi"}]}
assert not should_bypass_9router(body_no_doc, "sk-abc")
body_with_tools = {**body, "tools": [{"name": "x"}]}
assert not should_bypass_9router(body_with_tools, "sk-abc")
def test_anthropic_to_openai_request_translation_shape():
"""Translator must produce a valid OpenAI Chat Completions body."""
from backend.apps.agents.anthropic_to_openai import translate_request
body = {
"model": "gpt-5.5",
"max_tokens": 200,
"system": "You are helpful.",
"messages": [{
"role": "user",
"content": [
{"type": "text", "text": "summarize"},
{"type": "document", "source": {
"type": "base64", "media_type": "application/pdf", "data": "JVBERi0=",
}},
{"type": "image", "source": {
"type": "base64", "media_type": "image/png", "data": "iVBOR=",
}},
],
}],
}
out = translate_request(body)
assert out["model"] == "gpt-5.5"
assert out["stream"] is True
assert out["max_completion_tokens"] == 200
assert out["messages"][0] == {"role": "system", "content": "You are helpful."}
user_content = out["messages"][1]["content"]
assert any(p["type"] == "text" and p["text"] == "summarize" for p in user_content)
assert any(p["type"] == "file" and p["file"]["file_data"].startswith("data:application/pdf;base64,") for p in user_content)
assert any(p["type"] == "image_url" and p["image_url"]["url"].startswith("data:image/png;base64,") for p in user_content)
def test_resolve_attachments_gemini_emits_native_document_after_translator_fix():
"""After fixing the 9router 0.3.60 block-stripping bug via
anthropic_proxy._rewrite_document_to_image (now rewrites both
+13 -6
View File
@@ -801,12 +801,15 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
// document→image rewrite); refused on OpenAI/OpenRouter/custom until
// we land file-parser plugin / type:file translation.
// Mirrors backend agent_manager._resolve_attachments support matrix.
// PDFs: Anthropic, Gemini, OpenRouter (all empirically verified May
// 2026 via probe-pdf-roundtrip.py). OpenAI direct refused because
// OpenAI's image_url rejects non-image mime; user should switch to
// openrouter/openai/gpt-5 for OpenAI-via-OR with PDF support.
// PDFs: Anthropic, Gemini, OpenRouter (file-parser plugin), and
// OpenAI direct on GPT-5.x non-Codex (anthropic_proxy bypasses
// 9router and POSTs to api.openai.com via anthropic_to_openai.py).
// Images: every provider via 9router image_url translation.
const pdfSupported = ['anthropic', 'gemini', 'gemini-cli', 'openrouter'].includes(currentModelApi);
const isCodexModel = typeof model === 'string' && (model.toLowerCase().includes('codex') || model.toLowerCase().startsWith('cx/'));
const pdfSupported = (
['anthropic', 'gemini', 'gemini-cli', 'openrouter'].includes(currentModelApi) ||
(currentModelApi === 'openai' && !isCodexModel)
);
const imageSupported = ['anthropic', 'gemini', 'gemini-cli', 'openai', 'openrouter'].includes(currentModelApi);
const pendingPayloadEstimate = useMemo(() => {
@@ -2255,7 +2258,11 @@ const ChatInput = forwardRef<ChatInputHandle, Props>(({ onSend, disabled, mode,
{(() => {
const win = (opt.context_window as number) || 0;
const api = (opt.api as string || 'anthropic').toLowerCase();
const optSupportsPdf = ['anthropic', 'gemini', 'gemini-cli', 'openrouter'].includes(api);
const optIsCodex = typeof opt.value === 'string' && (opt.value.toLowerCase().includes('codex') || opt.value.toLowerCase().startsWith('cx/'));
const optSupportsPdf = (
['anthropic', 'gemini', 'gemini-cli', 'openrouter'].includes(api) ||
(api === 'openai' && !optIsCodex)
);
const optSupportsImage = ['anthropic', 'gemini', 'gemini-cli', 'openai', 'openrouter'].includes(api);
const cannotPdf = pendingKinds.has('pdf') && !optSupportsPdf;
const cannotImg = pendingKinds.has('image') && !optSupportsImage;