Spaces:
Build error
Build error
tudragon154203 commited on
Commit ·
a5d7c82
1
Parent(s): d546552
Avoid blocking proxy health routes during compression
Browse files
headroom/proxy/handlers/anthropic.py
CHANGED
|
@@ -714,7 +714,11 @@ class AnthropicHandlerMixin:
|
|
| 714 |
):
|
| 715 |
compressor = _get_image_compressor()
|
| 716 |
if compressor and compressor.has_images(messages):
|
| 717 |
-
messages =
|
|
|
|
|
|
|
|
|
|
|
|
|
| 718 |
if compressor.last_result:
|
| 719 |
logger.info(
|
| 720 |
f"Image compression: {compressor.last_result.technique.value} "
|
|
@@ -2070,7 +2074,8 @@ class AnthropicHandlerMixin:
|
|
| 2070 |
original_tokens = get_tokenizer(model).count_messages(messages)
|
| 2071 |
optimized_tokens = original_tokens
|
| 2072 |
else:
|
| 2073 |
-
result =
|
|
|
|
| 2074 |
messages=messages,
|
| 2075 |
model=model,
|
| 2076 |
model_limit=context_limit,
|
|
|
|
| 714 |
):
|
| 715 |
compressor = _get_image_compressor()
|
| 716 |
if compressor and compressor.has_images(messages):
|
| 717 |
+
messages = await asyncio.to_thread(
|
| 718 |
+
compressor.compress,
|
| 719 |
+
messages,
|
| 720 |
+
provider="anthropic",
|
| 721 |
+
)
|
| 722 |
if compressor.last_result:
|
| 723 |
logger.info(
|
| 724 |
f"Image compression: {compressor.last_result.technique.value} "
|
|
|
|
| 2074 |
original_tokens = get_tokenizer(model).count_messages(messages)
|
| 2075 |
optimized_tokens = original_tokens
|
| 2076 |
else:
|
| 2077 |
+
result = await asyncio.to_thread(
|
| 2078 |
+
self.anthropic_pipeline.apply,
|
| 2079 |
messages=messages,
|
| 2080 |
model=model,
|
| 2081 |
model_limit=context_limit,
|
headroom/proxy/handlers/batch.py
CHANGED
|
@@ -145,7 +145,8 @@ class BatchHandlerMixin:
|
|
| 145 |
)
|
| 146 |
|
| 147 |
# Use OpenAI pipeline (similar message format after conversion)
|
| 148 |
-
result =
|
|
|
|
| 149 |
messages=messages,
|
| 150 |
model=model,
|
| 151 |
model_limit=context_limit,
|
|
@@ -904,7 +905,8 @@ class BatchHandlerMixin:
|
|
| 904 |
if self.config.optimize:
|
| 905 |
try:
|
| 906 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 907 |
-
result =
|
|
|
|
| 908 |
messages=messages,
|
| 909 |
model=model,
|
| 910 |
model_limit=context_limit,
|
|
|
|
| 145 |
)
|
| 146 |
|
| 147 |
# Use OpenAI pipeline (similar message format after conversion)
|
| 148 |
+
result = await asyncio.to_thread(
|
| 149 |
+
self.openai_pipeline.apply,
|
| 150 |
messages=messages,
|
| 151 |
model=model,
|
| 152 |
model_limit=context_limit,
|
|
|
|
| 905 |
if self.config.optimize:
|
| 906 |
try:
|
| 907 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 908 |
+
result = await asyncio.to_thread(
|
| 909 |
+
self.openai_pipeline.apply,
|
| 910 |
messages=messages,
|
| 911 |
model=model,
|
| 912 |
model_limit=context_limit,
|
headroom/proxy/handlers/gemini.py
CHANGED
|
@@ -277,7 +277,8 @@ class GeminiHandlerMixin:
|
|
| 277 |
try:
|
| 278 |
# Use OpenAI pipeline (similar message format)
|
| 279 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 280 |
-
result =
|
|
|
|
| 281 |
messages=messages,
|
| 282 |
model=model,
|
| 283 |
model_limit=context_limit,
|
|
@@ -537,7 +538,8 @@ class GeminiHandlerMixin:
|
|
| 537 |
if self.config.optimize and messages and _license_ok:
|
| 538 |
try:
|
| 539 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 540 |
-
result =
|
|
|
|
| 541 |
messages=messages,
|
| 542 |
model=model,
|
| 543 |
model_limit=context_limit,
|
|
@@ -744,7 +746,8 @@ class GeminiHandlerMixin:
|
|
| 744 |
if self.config.optimize and messages:
|
| 745 |
try:
|
| 746 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 747 |
-
result =
|
|
|
|
| 748 |
messages=messages,
|
| 749 |
model=model,
|
| 750 |
model_limit=context_limit,
|
|
|
|
| 277 |
try:
|
| 278 |
# Use OpenAI pipeline (similar message format)
|
| 279 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 280 |
+
result = await asyncio.to_thread(
|
| 281 |
+
self.openai_pipeline.apply,
|
| 282 |
messages=messages,
|
| 283 |
model=model,
|
| 284 |
model_limit=context_limit,
|
|
|
|
| 538 |
if self.config.optimize and messages and _license_ok:
|
| 539 |
try:
|
| 540 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 541 |
+
result = await asyncio.to_thread(
|
| 542 |
+
self.openai_pipeline.apply,
|
| 543 |
messages=messages,
|
| 544 |
model=model,
|
| 545 |
model_limit=context_limit,
|
|
|
|
| 746 |
if self.config.optimize and messages:
|
| 747 |
try:
|
| 748 |
context_limit = self.openai_provider.get_context_limit(model)
|
| 749 |
+
result = await asyncio.to_thread(
|
| 750 |
+
self.openai_pipeline.apply,
|
| 751 |
messages=messages,
|
| 752 |
model=model,
|
| 753 |
model_limit=context_limit,
|
headroom/proxy/handlers/openai.py
CHANGED
|
@@ -241,7 +241,11 @@ class OpenAIHandlerMixin:
|
|
| 241 |
|
| 242 |
compressor = _get_image_compressor()
|
| 243 |
if compressor and compressor.has_images(messages):
|
| 244 |
-
messages =
|
|
|
|
|
|
|
|
|
|
|
|
|
| 245 |
if compressor.last_result:
|
| 246 |
logger.info(
|
| 247 |
f"[{request_id}] Image: {compressor.last_result.technique.value} "
|
|
@@ -2602,7 +2606,8 @@ class OpenAIHandlerMixin:
|
|
| 2602 |
if compress_tagged_content is not None:
|
| 2603 |
pipeline_kwargs["compress_tagged_content"] = bool(compress_tagged_content)
|
| 2604 |
|
| 2605 |
-
result =
|
|
|
|
| 2606 |
messages=messages,
|
| 2607 |
model=model,
|
| 2608 |
**pipeline_kwargs,
|
|
|
|
| 241 |
|
| 242 |
compressor = _get_image_compressor()
|
| 243 |
if compressor and compressor.has_images(messages):
|
| 244 |
+
messages = await asyncio.to_thread(
|
| 245 |
+
compressor.compress,
|
| 246 |
+
messages,
|
| 247 |
+
provider="openai",
|
| 248 |
+
)
|
| 249 |
if compressor.last_result:
|
| 250 |
logger.info(
|
| 251 |
f"[{request_id}] Image: {compressor.last_result.technique.value} "
|
|
|
|
| 2606 |
if compress_tagged_content is not None:
|
| 2607 |
pipeline_kwargs["compress_tagged_content"] = bool(compress_tagged_content)
|
| 2608 |
|
| 2609 |
+
result = await asyncio.to_thread(
|
| 2610 |
+
self.openai_pipeline.apply,
|
| 2611 |
messages=messages,
|
| 2612 |
model=model,
|
| 2613 |
**pipeline_kwargs,
|
tests/test_proxy/test_proxy_healthchecks.py
CHANGED
|
@@ -206,3 +206,62 @@ def test_shutdown_tolerates_stubbed_memory_handler():
|
|
| 206 |
response = client.get("/health")
|
| 207 |
|
| 208 |
assert response.status_code == 200
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 206 |
response = client.get("/health")
|
| 207 |
|
| 208 |
assert response.status_code == 200
|
| 209 |
+
@pytest.mark.asyncio
|
| 210 |
+
async def test_readyz_remains_responsive_during_slow_compress(monkeypatch):
|
| 211 |
+
import asyncio
|
| 212 |
+
import time
|
| 213 |
+
|
| 214 |
+
from httpx import ASGITransport, AsyncClient
|
| 215 |
+
|
| 216 |
+
config = ProxyConfig(
|
| 217 |
+
optimize=True,
|
| 218 |
+
cache_enabled=False,
|
| 219 |
+
rate_limit_enabled=False,
|
| 220 |
+
cost_tracking_enabled=False,
|
| 221 |
+
)
|
| 222 |
+
app = create_app(config)
|
| 223 |
+
app.state.ready = True
|
| 224 |
+
app.state.startup_error = None
|
| 225 |
+
app.state.proxy.http_client = object()
|
| 226 |
+
|
| 227 |
+
def slow_apply(*, messages, model, **kwargs): # noqa: ANN003
|
| 228 |
+
time.sleep(0.25)
|
| 229 |
+
return SimpleNamespace(
|
| 230 |
+
messages=messages,
|
| 231 |
+
transforms_applied=[],
|
| 232 |
+
transforms_summary=[],
|
| 233 |
+
markers_inserted=[],
|
| 234 |
+
tokens_before=16,
|
| 235 |
+
tokens_after=16,
|
| 236 |
+
skip_reason="no_change",
|
| 237 |
+
)
|
| 238 |
+
|
| 239 |
+
app.state.proxy.openai_pipeline.apply = slow_apply
|
| 240 |
+
|
| 241 |
+
transport = ASGITransport(app=app)
|
| 242 |
+
async with AsyncClient(transport=transport, base_url="http://testserver") as client:
|
| 243 |
+
compress_task = asyncio.create_task(
|
| 244 |
+
client.post(
|
| 245 |
+
"/v1/compress",
|
| 246 |
+
json={
|
| 247 |
+
"model": "gpt-4o-mini",
|
| 248 |
+
"messages": [
|
| 249 |
+
{"role": "system", "content": "You are helpful."},
|
| 250 |
+
{"role": "user", "content": "Compress this."},
|
| 251 |
+
],
|
| 252 |
+
},
|
| 253 |
+
)
|
| 254 |
+
)
|
| 255 |
+
|
| 256 |
+
await asyncio.sleep(0.05)
|
| 257 |
+
|
| 258 |
+
started = time.perf_counter()
|
| 259 |
+
readyz_response = await client.get("/readyz")
|
| 260 |
+
readyz_elapsed = time.perf_counter() - started
|
| 261 |
+
|
| 262 |
+
compress_response = await compress_task
|
| 263 |
+
|
| 264 |
+
assert readyz_response.status_code == 200
|
| 265 |
+
assert readyz_response.json()["ready"] is True
|
| 266 |
+
assert readyz_elapsed < 0.15
|
| 267 |
+
assert compress_response.status_code == 200
|