Spaces:
Build error
Build error
Commit ·
72a87b4
1
Parent(s): 0302796
feat: make token_headroom the default mode + fix Gemini handler bug
Browse files- Change default HEADROOM_MODE from cost_savings to token_headroom
across server.py, cli/proxy.py, and mcp_server.py. Prefix caching
is native to providers; Headroom's value-add is compression.
- Fix undefined _compression_failed variable in Gemini handler
(ruff + mypy error).
- Apply ruff format fixes.
- headroom/ccr/mcp_server.py +1 -1
- headroom/cli/proxy.py +2 -2
- headroom/proxy/server.py +160 -32
headroom/ccr/mcp_server.py
CHANGED
|
@@ -76,7 +76,7 @@ def _format_session_summary(summary: dict[str, Any], local_stats: dict[str, Any]
|
|
| 76 |
lines.append("Headroom Session Summary")
|
| 77 |
lines.append("=" * 40)
|
| 78 |
|
| 79 |
-
mode = summary.get("mode", "
|
| 80 |
api_reqs = summary.get("api_requests", 0)
|
| 81 |
model = summary.get("primary_model", "unknown")
|
| 82 |
lines.append(f"Mode: {mode} | {api_reqs} API requests | {model}")
|
|
|
|
| 76 |
lines.append("Headroom Session Summary")
|
| 77 |
lines.append("=" * 40)
|
| 78 |
|
| 79 |
+
mode = summary.get("mode", "token_headroom")
|
| 80 |
api_reqs = summary.get("api_requests", 0)
|
| 81 |
model = summary.get("primary_model", "unknown")
|
| 82 |
lines.append(f"Mode: {mode} | {api_reqs} API requests | {model}")
|
headroom/cli/proxy.py
CHANGED
|
@@ -14,7 +14,7 @@ from .main import main
|
|
| 14 |
"--mode",
|
| 15 |
default=None,
|
| 16 |
type=click.Choice(["cost_savings", "token_headroom"]),
|
| 17 |
-
help="Optimization mode:
|
| 18 |
)
|
| 19 |
@click.option("--no-optimize", is_flag=True, help="Disable optimization (passthrough mode)")
|
| 20 |
@click.option("--no-cache", is_flag=True, help="Disable semantic caching")
|
|
@@ -184,7 +184,7 @@ def proxy(
|
|
| 184 |
effective_anyllm_provider = os.environ.get("HEADROOM_ANYLLM_PROVIDER") or anyllm_provider
|
| 185 |
|
| 186 |
# Resolve mode: CLI flag > env var > default
|
| 187 |
-
effective_mode = mode or os.environ.get("HEADROOM_MODE", "
|
| 188 |
|
| 189 |
config = ProxyConfig(
|
| 190 |
host=host,
|
|
|
|
| 14 |
"--mode",
|
| 15 |
default=None,
|
| 16 |
type=click.Choice(["cost_savings", "token_headroom"]),
|
| 17 |
+
help="Optimization mode: token_headroom (compress for session extension) or cost_savings (preserve prefix cache). Default: token_headroom. Env: HEADROOM_MODE",
|
| 18 |
)
|
| 19 |
@click.option("--no-optimize", is_flag=True, help="Disable optimization (passthrough mode)")
|
| 20 |
@click.option("--no-cache", is_flag=True, help="Disable semantic caching")
|
|
|
|
| 184 |
effective_anyllm_provider = os.environ.get("HEADROOM_ANYLLM_PROVIDER") or anyllm_provider
|
| 185 |
|
| 186 |
# Resolve mode: CLI flag > env var > default
|
| 187 |
+
effective_mode = mode or os.environ.get("HEADROOM_MODE", "token_headroom")
|
| 188 |
|
| 189 |
config = ProxyConfig(
|
| 190 |
host=host,
|
headroom/proxy/server.py
CHANGED
|
@@ -527,6 +527,21 @@ def _build_session_summary(
|
|
| 527 |
# Maximum request body size (100MB - increased to support image-heavy requests)
|
| 528 |
MAX_REQUEST_BODY_SIZE = 100 * 1024 * 1024
|
| 529 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 530 |
|
| 531 |
# =============================================================================
|
| 532 |
# Data Models
|
|
@@ -605,10 +620,10 @@ class ProxyConfig:
|
|
| 605 |
bedrock_profile: str | None = None # AWS profile (optional)
|
| 606 |
anyllm_provider: str = "openai" # any-llm provider (openai, mistral, groq, etc.)
|
| 607 |
|
| 608 |
-
# Optimization mode: "
|
| 609 |
-
# cost_savings: preserve prefix cache for cost reduction
|
| 610 |
# token_headroom: compress older messages for session extension
|
| 611 |
-
|
|
|
|
| 612 |
|
| 613 |
# Optimization
|
| 614 |
optimize: bool = True
|
|
@@ -872,6 +887,19 @@ class TokenBucketRateLimiter:
|
|
| 872 |
)
|
| 873 |
self._lock = asyncio.Lock()
|
| 874 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 875 |
def _refill(self, state: RateLimitState, rate_per_minute: float) -> float:
|
| 876 |
"""Refill bucket based on elapsed time."""
|
| 877 |
now = time.time()
|
|
@@ -884,6 +912,9 @@ class TokenBucketRateLimiter:
|
|
| 884 |
async def check_request(self, key: str = "default") -> tuple[bool, float]:
|
| 885 |
"""Check if request is allowed. Returns (allowed, wait_seconds)."""
|
| 886 |
async with self._lock:
|
|
|
|
|
|
|
|
|
|
| 887 |
state = self._request_buckets[key]
|
| 888 |
available = self._refill(state, self.requests_per_minute)
|
| 889 |
|
|
@@ -1747,6 +1778,20 @@ class HeadroomProxy:
|
|
| 1747 |
if session_id not in self._compression_caches:
|
| 1748 |
from headroom.cache.compression_cache import CompressionCache
|
| 1749 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1750 |
self._compression_caches[session_id] = CompressionCache()
|
| 1751 |
return self._compression_caches[session_id]
|
| 1752 |
|
|
@@ -1834,7 +1879,7 @@ class HeadroomProxy:
|
|
| 1834 |
logger.warning(
|
| 1835 |
f"Unknown HEADROOM_MODE '{self.config.mode}', falling back to 'cost_savings'"
|
| 1836 |
)
|
| 1837 |
-
self.config.mode = "
|
| 1838 |
logger.info(f"Mode: {self.config.mode}")
|
| 1839 |
if self.config.mode == "token_headroom":
|
| 1840 |
logger.info(" Prefix freeze: re-freeze after compression")
|
|
@@ -2104,6 +2149,21 @@ class HeadroomProxy:
|
|
| 2104 |
)
|
| 2105 |
model = body.get("model", "unknown")
|
| 2106 |
messages = body.get("messages", [])
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2107 |
stream = body.get("stream", False)
|
| 2108 |
|
| 2109 |
# Image compression (before text optimization)
|
|
@@ -2127,7 +2187,9 @@ class HeadroomProxy:
|
|
| 2127 |
|
| 2128 |
# Rate limiting
|
| 2129 |
if self.rate_limiter:
|
| 2130 |
-
|
|
|
|
|
|
|
| 2131 |
allowed, wait_seconds = await self.rate_limiter.check_request(rate_key)
|
| 2132 |
if not allowed:
|
| 2133 |
await self.metrics.record_rate_limited()
|
|
@@ -2211,6 +2273,7 @@ class HeadroomProxy:
|
|
| 2211 |
prefix_tracker = self.session_tracker_store.get_or_create(session_id, "anthropic")
|
| 2212 |
frozen_message_count = prefix_tracker.get_frozen_message_count()
|
| 2213 |
|
|
|
|
| 2214 |
if self.config.optimize and messages:
|
| 2215 |
try:
|
| 2216 |
context_limit = self.anthropic_provider.get_context_limit(model)
|
|
@@ -2229,12 +2292,17 @@ class HeadroomProxy:
|
|
| 2229 |
# Re-freeze boundary: consecutive stable messages from start
|
| 2230 |
frozen_message_count = comp_cache.compute_frozen_count(messages)
|
| 2231 |
|
| 2232 |
-
result =
|
| 2233 |
-
|
| 2234 |
-
|
| 2235 |
-
|
| 2236 |
-
|
| 2237 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2238 |
)
|
| 2239 |
|
| 2240 |
# Cache newly compressed messages (index-aligned diff)
|
|
@@ -2250,12 +2318,17 @@ class HeadroomProxy:
|
|
| 2250 |
# original_tokens was set at line ~2183 from uncompressed messages.
|
| 2251 |
optimized_tokens = result.tokens_after
|
| 2252 |
else:
|
| 2253 |
-
result =
|
| 2254 |
-
|
| 2255 |
-
|
| 2256 |
-
|
| 2257 |
-
|
| 2258 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2259 |
)
|
| 2260 |
|
| 2261 |
if result.messages != messages:
|
|
@@ -2269,6 +2342,8 @@ class HeadroomProxy:
|
|
| 2269 |
waste_signals_dict = result.waste_signals.to_dict()
|
| 2270 |
except Exception as e:
|
| 2271 |
logger.warning(f"Optimization failed: {e}")
|
|
|
|
|
|
|
| 2272 |
|
| 2273 |
tokens_saved = max(0, original_tokens - optimized_tokens)
|
| 2274 |
optimization_latency = (time.time() - start_time) * 1000
|
|
@@ -2811,6 +2886,8 @@ class HeadroomProxy:
|
|
| 2811 |
response_headers["x-headroom-transforms"] = ",".join(transforms_applied)
|
| 2812 |
if cache_hit:
|
| 2813 |
response_headers["x-headroom-cached"] = "true"
|
|
|
|
|
|
|
| 2814 |
|
| 2815 |
return Response(
|
| 2816 |
content=response.content,
|
|
@@ -4234,7 +4311,7 @@ class HeadroomProxy:
|
|
| 4234 |
)
|
| 4235 |
|
| 4236 |
async def generate():
|
| 4237 |
-
nonlocal body # May need to modify for continuation requests
|
| 4238 |
|
| 4239 |
# For memory mode, we buffer the response to check for tool calls
|
| 4240 |
buffered_chunks: list[bytes] = []
|
|
@@ -4255,10 +4332,27 @@ class HeadroomProxy:
|
|
| 4255 |
chunk_str = chunk.decode("utf-8", errors="ignore")
|
| 4256 |
stream_state["sse_buffer"] += chunk_str
|
| 4257 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 4258 |
if memory_enabled:
|
| 4259 |
# Buffer for memory tool detection
|
| 4260 |
buffered_chunks.append(chunk)
|
| 4261 |
full_sse_data += chunk_str
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 4262 |
else:
|
| 4263 |
# Immediate streaming when memory not enabled
|
| 4264 |
yield chunk
|
|
@@ -4696,6 +4790,21 @@ class HeadroomProxy:
|
|
| 4696 |
)
|
| 4697 |
model = body.get("model", "unknown")
|
| 4698 |
messages = body.get("messages", [])
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 4699 |
stream = body.get("stream", False)
|
| 4700 |
|
| 4701 |
# Image compression (before text optimization)
|
|
@@ -4778,6 +4887,7 @@ class HeadroomProxy:
|
|
| 4778 |
)
|
| 4779 |
openai_frozen_count = openai_prefix_tracker.get_frozen_message_count()
|
| 4780 |
|
|
|
|
| 4781 |
if self.config.optimize and messages:
|
| 4782 |
try:
|
| 4783 |
context_limit = self.openai_provider.get_context_limit(model)
|
|
@@ -4791,12 +4901,17 @@ class HeadroomProxy:
|
|
| 4791 |
# Re-freeze boundary
|
| 4792 |
openai_frozen_count = comp_cache.compute_frozen_count(messages)
|
| 4793 |
|
| 4794 |
-
result =
|
| 4795 |
-
|
| 4796 |
-
|
| 4797 |
-
|
| 4798 |
-
|
| 4799 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 4800 |
)
|
| 4801 |
|
| 4802 |
if result.messages != working_messages:
|
|
@@ -4810,12 +4925,17 @@ class HeadroomProxy:
|
|
| 4810 |
# so tokens_saved captures both Zone 1 + Zone 2 savings.
|
| 4811 |
optimized_tokens = result.tokens_after
|
| 4812 |
else:
|
| 4813 |
-
result =
|
| 4814 |
-
|
| 4815 |
-
|
| 4816 |
-
|
| 4817 |
-
|
| 4818 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 4819 |
)
|
| 4820 |
|
| 4821 |
if result.messages != messages:
|
|
@@ -4829,6 +4949,8 @@ class HeadroomProxy:
|
|
| 4829 |
waste_signals_dict = result.waste_signals.to_dict()
|
| 4830 |
except Exception as e:
|
| 4831 |
logger.warning(f"Optimization failed: {e}")
|
|
|
|
|
|
|
| 4832 |
|
| 4833 |
tokens_saved = max(0, original_tokens - optimized_tokens)
|
| 4834 |
optimization_latency = (time.time() - start_time) * 1000
|
|
@@ -5052,6 +5174,8 @@ class HeadroomProxy:
|
|
| 5052 |
response_headers["x-headroom-transforms"] = ",".join(transforms_applied)
|
| 5053 |
if cache_read_tokens > 0:
|
| 5054 |
response_headers["x-headroom-cached"] = "true"
|
|
|
|
|
|
|
| 5055 |
|
| 5056 |
return Response(
|
| 5057 |
content=response.content,
|
|
@@ -5941,6 +6065,7 @@ class HeadroomProxy:
|
|
| 5941 |
optimized_messages = messages
|
| 5942 |
optimized_tokens = original_tokens
|
| 5943 |
|
|
|
|
| 5944 |
if self.config.optimize and messages:
|
| 5945 |
try:
|
| 5946 |
# Use OpenAI pipeline (similar message format)
|
|
@@ -5959,6 +6084,7 @@ class HeadroomProxy:
|
|
| 5959 |
if result.waste_signals:
|
| 5960 |
waste_signals_dict = result.waste_signals.to_dict()
|
| 5961 |
except Exception as e:
|
|
|
|
| 5962 |
logger.warning(f"[{request_id}] Gemini optimization failed: {e}")
|
| 5963 |
|
| 5964 |
tokens_saved = max(0, original_tokens - optimized_tokens)
|
|
@@ -6074,6 +6200,8 @@ class HeadroomProxy:
|
|
| 6074 |
response_headers["x-headroom-transforms"] = ",".join(transforms_applied)
|
| 6075 |
if cache_read_tokens > 0:
|
| 6076 |
response_headers["x-headroom-cached"] = "true"
|
|
|
|
|
|
|
| 6077 |
|
| 6078 |
return Response(
|
| 6079 |
content=response.content,
|
|
@@ -6541,7 +6669,7 @@ def create_app(config: ProxyConfig | None = None) -> FastAPI:
|
|
| 6541 |
"total_tokens_saved": total_tokens_saved,
|
| 6542 |
}
|
| 6543 |
else:
|
| 6544 |
-
compression_cache_stats = {"mode": "
|
| 6545 |
|
| 6546 |
# Build unified savings summary (all layers)
|
| 6547 |
compression_tokens = m.tokens_saved_total
|
|
@@ -7846,7 +7974,7 @@ if __name__ == "__main__":
|
|
| 7846 |
max_keepalive_connections=_get_env_int("HEADROOM_MAX_KEEPALIVE", args.max_keepalive),
|
| 7847 |
http2=not args.no_http2 and _get_env_bool("HEADROOM_HTTP2", True),
|
| 7848 |
tool_profiles=tool_profiles if tool_profiles else None,
|
| 7849 |
-
mode=_get_env_str("HEADROOM_MODE", "
|
| 7850 |
)
|
| 7851 |
|
| 7852 |
# Get worker and concurrency settings
|
|
|
|
| 527 |
# Maximum request body size (100MB - increased to support image-heavy requests)
|
| 528 |
MAX_REQUEST_BODY_SIZE = 100 * 1024 * 1024
|
| 529 |
|
| 530 |
+
# Maximum SSE buffer size (10MB - prevents memory exhaustion from malformed streams)
|
| 531 |
+
MAX_SSE_BUFFER_SIZE = 10 * 1024 * 1024
|
| 532 |
+
|
| 533 |
+
# Maximum message array length (prevents DoS from deeply nested payloads)
|
| 534 |
+
MAX_MESSAGE_ARRAY_LENGTH = 10000
|
| 535 |
+
|
| 536 |
+
# Maximum compression cache sessions (prevents unbounded memory growth)
|
| 537 |
+
MAX_COMPRESSION_CACHE_SESSIONS = 500
|
| 538 |
+
|
| 539 |
+
# Maximum rate limiter buckets (prevents DoS via spoofed API keys)
|
| 540 |
+
MAX_RATE_LIMITER_BUCKETS = 1000
|
| 541 |
+
|
| 542 |
+
# Compression pipeline timeout in seconds
|
| 543 |
+
COMPRESSION_TIMEOUT_SECONDS = 30
|
| 544 |
+
|
| 545 |
|
| 546 |
# =============================================================================
|
| 547 |
# Data Models
|
|
|
|
| 620 |
bedrock_profile: str | None = None # AWS profile (optional)
|
| 621 |
anyllm_provider: str = "openai" # any-llm provider (openai, mistral, groq, etc.)
|
| 622 |
|
| 623 |
+
# Optimization mode: "token_headroom" (default) or "cost_savings"
|
|
|
|
| 624 |
# token_headroom: compress older messages for session extension
|
| 625 |
+
# cost_savings: preserve prefix cache for cost reduction
|
| 626 |
+
mode: str = "token_headroom"
|
| 627 |
|
| 628 |
# Optimization
|
| 629 |
optimize: bool = True
|
|
|
|
| 887 |
)
|
| 888 |
self._lock = asyncio.Lock()
|
| 889 |
|
| 890 |
+
async def _cleanup_stale_buckets(self) -> None:
|
| 891 |
+
"""Remove buckets that haven't been used in the last 10 minutes."""
|
| 892 |
+
now = time.time()
|
| 893 |
+
stale_threshold = now - 600 # 10 minutes
|
| 894 |
+
stale_keys = [
|
| 895 |
+
k for k, v in self._request_buckets.items() if v.last_update < stale_threshold
|
| 896 |
+
]
|
| 897 |
+
for k in stale_keys:
|
| 898 |
+
del self._request_buckets[k]
|
| 899 |
+
self._token_buckets.pop(k, None)
|
| 900 |
+
if stale_keys:
|
| 901 |
+
logger.debug(f"Cleaned up {len(stale_keys)} stale rate limiter buckets")
|
| 902 |
+
|
| 903 |
def _refill(self, state: RateLimitState, rate_per_minute: float) -> float:
|
| 904 |
"""Refill bucket based on elapsed time."""
|
| 905 |
now = time.time()
|
|
|
|
| 912 |
async def check_request(self, key: str = "default") -> tuple[bool, float]:
|
| 913 |
"""Check if request is allowed. Returns (allowed, wait_seconds)."""
|
| 914 |
async with self._lock:
|
| 915 |
+
# Prevent unbounded bucket growth from spoofed keys
|
| 916 |
+
if len(self._request_buckets) > MAX_RATE_LIMITER_BUCKETS:
|
| 917 |
+
await self._cleanup_stale_buckets()
|
| 918 |
state = self._request_buckets[key]
|
| 919 |
available = self._refill(state, self.requests_per_minute)
|
| 920 |
|
|
|
|
| 1778 |
if session_id not in self._compression_caches:
|
| 1779 |
from headroom.cache.compression_cache import CompressionCache
|
| 1780 |
|
| 1781 |
+
# Evict oldest caches if at capacity
|
| 1782 |
+
if len(self._compression_caches) >= MAX_COMPRESSION_CACHE_SESSIONS:
|
| 1783 |
+
# Remove oldest quarter to amortize cleanup cost
|
| 1784 |
+
oldest_keys = list(self._compression_caches.keys())[
|
| 1785 |
+
: MAX_COMPRESSION_CACHE_SESSIONS // 4
|
| 1786 |
+
]
|
| 1787 |
+
for key in oldest_keys:
|
| 1788 |
+
del self._compression_caches[key]
|
| 1789 |
+
logger.info(
|
| 1790 |
+
"Evicted %d compression caches (exceeded %d max sessions)",
|
| 1791 |
+
len(oldest_keys),
|
| 1792 |
+
MAX_COMPRESSION_CACHE_SESSIONS,
|
| 1793 |
+
)
|
| 1794 |
+
|
| 1795 |
self._compression_caches[session_id] = CompressionCache()
|
| 1796 |
return self._compression_caches[session_id]
|
| 1797 |
|
|
|
|
| 1879 |
logger.warning(
|
| 1880 |
f"Unknown HEADROOM_MODE '{self.config.mode}', falling back to 'cost_savings'"
|
| 1881 |
)
|
| 1882 |
+
self.config.mode = "token_headroom"
|
| 1883 |
logger.info(f"Mode: {self.config.mode}")
|
| 1884 |
if self.config.mode == "token_headroom":
|
| 1885 |
logger.info(" Prefix freeze: re-freeze after compression")
|
|
|
|
| 2149 |
)
|
| 2150 |
model = body.get("model", "unknown")
|
| 2151 |
messages = body.get("messages", [])
|
| 2152 |
+
|
| 2153 |
+
# Validate message array size
|
| 2154 |
+
if len(messages) > MAX_MESSAGE_ARRAY_LENGTH:
|
| 2155 |
+
return JSONResponse(
|
| 2156 |
+
status_code=400,
|
| 2157 |
+
content={
|
| 2158 |
+
"type": "error",
|
| 2159 |
+
"error": {
|
| 2160 |
+
"type": "invalid_request_error",
|
| 2161 |
+
"message": f"Message array too large ({len(messages)} messages). "
|
| 2162 |
+
f"Maximum is {MAX_MESSAGE_ARRAY_LENGTH}.",
|
| 2163 |
+
},
|
| 2164 |
+
},
|
| 2165 |
+
)
|
| 2166 |
+
|
| 2167 |
stream = body.get("stream", False)
|
| 2168 |
|
| 2169 |
# Image compression (before text optimization)
|
|
|
|
| 2187 |
|
| 2188 |
# Rate limiting
|
| 2189 |
if self.rate_limiter:
|
| 2190 |
+
api_key = headers.get("x-api-key", "")
|
| 2191 |
+
client_ip = request.client.host if request.client else "unknown"
|
| 2192 |
+
rate_key = f"{api_key[:16]}:{client_ip}" if api_key else client_ip
|
| 2193 |
allowed, wait_seconds = await self.rate_limiter.check_request(rate_key)
|
| 2194 |
if not allowed:
|
| 2195 |
await self.metrics.record_rate_limited()
|
|
|
|
| 2273 |
prefix_tracker = self.session_tracker_store.get_or_create(session_id, "anthropic")
|
| 2274 |
frozen_message_count = prefix_tracker.get_frozen_message_count()
|
| 2275 |
|
| 2276 |
+
_compression_failed = False
|
| 2277 |
if self.config.optimize and messages:
|
| 2278 |
try:
|
| 2279 |
context_limit = self.anthropic_provider.get_context_limit(model)
|
|
|
|
| 2292 |
# Re-freeze boundary: consecutive stable messages from start
|
| 2293 |
frozen_message_count = comp_cache.compute_frozen_count(messages)
|
| 2294 |
|
| 2295 |
+
result = await asyncio.wait_for(
|
| 2296 |
+
asyncio.to_thread(
|
| 2297 |
+
lambda: self.anthropic_pipeline.apply(
|
| 2298 |
+
messages=working_messages,
|
| 2299 |
+
model=model,
|
| 2300 |
+
model_limit=context_limit,
|
| 2301 |
+
frozen_message_count=frozen_message_count,
|
| 2302 |
+
biases=biases,
|
| 2303 |
+
)
|
| 2304 |
+
),
|
| 2305 |
+
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
| 2306 |
)
|
| 2307 |
|
| 2308 |
# Cache newly compressed messages (index-aligned diff)
|
|
|
|
| 2318 |
# original_tokens was set at line ~2183 from uncompressed messages.
|
| 2319 |
optimized_tokens = result.tokens_after
|
| 2320 |
else:
|
| 2321 |
+
result = await asyncio.wait_for(
|
| 2322 |
+
asyncio.to_thread(
|
| 2323 |
+
lambda: self.anthropic_pipeline.apply(
|
| 2324 |
+
messages=messages,
|
| 2325 |
+
model=model,
|
| 2326 |
+
model_limit=context_limit,
|
| 2327 |
+
frozen_message_count=frozen_message_count,
|
| 2328 |
+
biases=biases,
|
| 2329 |
+
)
|
| 2330 |
+
),
|
| 2331 |
+
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
| 2332 |
)
|
| 2333 |
|
| 2334 |
if result.messages != messages:
|
|
|
|
| 2342 |
waste_signals_dict = result.waste_signals.to_dict()
|
| 2343 |
except Exception as e:
|
| 2344 |
logger.warning(f"Optimization failed: {e}")
|
| 2345 |
+
# Flag compression failure for observability
|
| 2346 |
+
_compression_failed = True
|
| 2347 |
|
| 2348 |
tokens_saved = max(0, original_tokens - optimized_tokens)
|
| 2349 |
optimization_latency = (time.time() - start_time) * 1000
|
|
|
|
| 2886 |
response_headers["x-headroom-transforms"] = ",".join(transforms_applied)
|
| 2887 |
if cache_hit:
|
| 2888 |
response_headers["x-headroom-cached"] = "true"
|
| 2889 |
+
if _compression_failed:
|
| 2890 |
+
response_headers["x-headroom-compression-failed"] = "true"
|
| 2891 |
|
| 2892 |
return Response(
|
| 2893 |
content=response.content,
|
|
|
|
| 4311 |
)
|
| 4312 |
|
| 4313 |
async def generate():
|
| 4314 |
+
nonlocal body, memory_enabled # May need to modify for continuation requests
|
| 4315 |
|
| 4316 |
# For memory mode, we buffer the response to check for tool calls
|
| 4317 |
buffered_chunks: list[bytes] = []
|
|
|
|
| 4332 |
chunk_str = chunk.decode("utf-8", errors="ignore")
|
| 4333 |
stream_state["sse_buffer"] += chunk_str
|
| 4334 |
|
| 4335 |
+
# Safety: prevent unbounded buffer growth
|
| 4336 |
+
if len(stream_state["sse_buffer"]) > MAX_SSE_BUFFER_SIZE:
|
| 4337 |
+
logger.error(
|
| 4338 |
+
"SSE buffer exceeded maximum size (%d bytes), "
|
| 4339 |
+
"truncating to prevent memory exhaustion",
|
| 4340 |
+
MAX_SSE_BUFFER_SIZE,
|
| 4341 |
+
)
|
| 4342 |
+
stream_state["sse_buffer"] = stream_state["sse_buffer"][
|
| 4343 |
+
-MAX_SSE_BUFFER_SIZE // 2 :
|
| 4344 |
+
]
|
| 4345 |
+
|
| 4346 |
if memory_enabled:
|
| 4347 |
# Buffer for memory tool detection
|
| 4348 |
buffered_chunks.append(chunk)
|
| 4349 |
full_sse_data += chunk_str
|
| 4350 |
+
if len(full_sse_data) > MAX_SSE_BUFFER_SIZE:
|
| 4351 |
+
logger.warning(
|
| 4352 |
+
"Memory-mode SSE buffer exceeded maximum size, "
|
| 4353 |
+
"disabling memory detection for this request"
|
| 4354 |
+
)
|
| 4355 |
+
memory_enabled = False
|
| 4356 |
else:
|
| 4357 |
# Immediate streaming when memory not enabled
|
| 4358 |
yield chunk
|
|
|
|
| 4790 |
)
|
| 4791 |
model = body.get("model", "unknown")
|
| 4792 |
messages = body.get("messages", [])
|
| 4793 |
+
|
| 4794 |
+
# Validate message array size
|
| 4795 |
+
if len(messages) > MAX_MESSAGE_ARRAY_LENGTH:
|
| 4796 |
+
return JSONResponse(
|
| 4797 |
+
status_code=400,
|
| 4798 |
+
content={
|
| 4799 |
+
"error": {
|
| 4800 |
+
"message": f"Message array too large ({len(messages)} messages). "
|
| 4801 |
+
f"Maximum is {MAX_MESSAGE_ARRAY_LENGTH}.",
|
| 4802 |
+
"type": "invalid_request_error",
|
| 4803 |
+
"code": "invalid_request",
|
| 4804 |
+
}
|
| 4805 |
+
},
|
| 4806 |
+
)
|
| 4807 |
+
|
| 4808 |
stream = body.get("stream", False)
|
| 4809 |
|
| 4810 |
# Image compression (before text optimization)
|
|
|
|
| 4887 |
)
|
| 4888 |
openai_frozen_count = openai_prefix_tracker.get_frozen_message_count()
|
| 4889 |
|
| 4890 |
+
_compression_failed = False
|
| 4891 |
if self.config.optimize and messages:
|
| 4892 |
try:
|
| 4893 |
context_limit = self.openai_provider.get_context_limit(model)
|
|
|
|
| 4901 |
# Re-freeze boundary
|
| 4902 |
openai_frozen_count = comp_cache.compute_frozen_count(messages)
|
| 4903 |
|
| 4904 |
+
result = await asyncio.wait_for(
|
| 4905 |
+
asyncio.to_thread(
|
| 4906 |
+
lambda: self.openai_pipeline.apply(
|
| 4907 |
+
messages=working_messages,
|
| 4908 |
+
model=model,
|
| 4909 |
+
model_limit=context_limit,
|
| 4910 |
+
frozen_message_count=openai_frozen_count,
|
| 4911 |
+
biases=_hook_biases,
|
| 4912 |
+
)
|
| 4913 |
+
),
|
| 4914 |
+
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
| 4915 |
)
|
| 4916 |
|
| 4917 |
if result.messages != working_messages:
|
|
|
|
| 4925 |
# so tokens_saved captures both Zone 1 + Zone 2 savings.
|
| 4926 |
optimized_tokens = result.tokens_after
|
| 4927 |
else:
|
| 4928 |
+
result = await asyncio.wait_for(
|
| 4929 |
+
asyncio.to_thread(
|
| 4930 |
+
lambda: self.openai_pipeline.apply(
|
| 4931 |
+
messages=messages,
|
| 4932 |
+
model=model,
|
| 4933 |
+
model_limit=context_limit,
|
| 4934 |
+
frozen_message_count=openai_frozen_count,
|
| 4935 |
+
biases=_hook_biases,
|
| 4936 |
+
)
|
| 4937 |
+
),
|
| 4938 |
+
timeout=COMPRESSION_TIMEOUT_SECONDS,
|
| 4939 |
)
|
| 4940 |
|
| 4941 |
if result.messages != messages:
|
|
|
|
| 4949 |
waste_signals_dict = result.waste_signals.to_dict()
|
| 4950 |
except Exception as e:
|
| 4951 |
logger.warning(f"Optimization failed: {e}")
|
| 4952 |
+
# Flag compression failure for observability
|
| 4953 |
+
_compression_failed = True
|
| 4954 |
|
| 4955 |
tokens_saved = max(0, original_tokens - optimized_tokens)
|
| 4956 |
optimization_latency = (time.time() - start_time) * 1000
|
|
|
|
| 5174 |
response_headers["x-headroom-transforms"] = ",".join(transforms_applied)
|
| 5175 |
if cache_read_tokens > 0:
|
| 5176 |
response_headers["x-headroom-cached"] = "true"
|
| 5177 |
+
if _compression_failed:
|
| 5178 |
+
response_headers["x-headroom-compression-failed"] = "true"
|
| 5179 |
|
| 5180 |
return Response(
|
| 5181 |
content=response.content,
|
|
|
|
| 6065 |
optimized_messages = messages
|
| 6066 |
optimized_tokens = original_tokens
|
| 6067 |
|
| 6068 |
+
_compression_failed = False
|
| 6069 |
if self.config.optimize and messages:
|
| 6070 |
try:
|
| 6071 |
# Use OpenAI pipeline (similar message format)
|
|
|
|
| 6084 |
if result.waste_signals:
|
| 6085 |
waste_signals_dict = result.waste_signals.to_dict()
|
| 6086 |
except Exception as e:
|
| 6087 |
+
_compression_failed = True
|
| 6088 |
logger.warning(f"[{request_id}] Gemini optimization failed: {e}")
|
| 6089 |
|
| 6090 |
tokens_saved = max(0, original_tokens - optimized_tokens)
|
|
|
|
| 6200 |
response_headers["x-headroom-transforms"] = ",".join(transforms_applied)
|
| 6201 |
if cache_read_tokens > 0:
|
| 6202 |
response_headers["x-headroom-cached"] = "true"
|
| 6203 |
+
if _compression_failed:
|
| 6204 |
+
response_headers["x-headroom-compression-failed"] = "true"
|
| 6205 |
|
| 6206 |
return Response(
|
| 6207 |
content=response.content,
|
|
|
|
| 6669 |
"total_tokens_saved": total_tokens_saved,
|
| 6670 |
}
|
| 6671 |
else:
|
| 6672 |
+
compression_cache_stats = {"mode": "token_headroom"}
|
| 6673 |
|
| 6674 |
# Build unified savings summary (all layers)
|
| 6675 |
compression_tokens = m.tokens_saved_total
|
|
|
|
| 7974 |
max_keepalive_connections=_get_env_int("HEADROOM_MAX_KEEPALIVE", args.max_keepalive),
|
| 7975 |
http2=not args.no_http2 and _get_env_bool("HEADROOM_HTTP2", True),
|
| 7976 |
tool_profiles=tool_profiles if tool_profiles else None,
|
| 7977 |
+
mode=_get_env_str("HEADROOM_MODE", "token_headroom"),
|
| 7978 |
)
|
| 7979 |
|
| 7980 |
# Get worker and concurrency settings
|