|
41 | 41 | class AnthropicHandlerMixin: |
42 | 42 | """Mixin providing Anthropic API handler methods for HeadroomProxy.""" |
43 | 43 |
|
| 44 | + async def _count_tokens_offloaded(self, model, messages): # noqa: ANN001, ANN201 |
| 45 | + """Resolve a tokenizer and count messages off the event loop. |
| 46 | +
|
| 47 | + Tokenizer resolution can be expensive on first use (HuggingFace |
| 48 | + backends may download vocab files) and counting a full Claude Code |
| 49 | + conversation is CPU-bound, so both run on the compression executor |
| 50 | + bounded by ``COMPRESSION_TIMEOUT_SECONDS`` (GH #1701: an unbounded |
| 51 | + on-loop load froze the whole server). On timeout or error this |
| 52 | + fails open to character-based estimation. |
| 53 | +
|
| 54 | + Returns: |
| 55 | + Tuple of ``(tokenizer, token_count)``. The tokenizer is fully |
| 56 | + initialized, so later ``count_messages`` calls on it are pure |
| 57 | + CPU work. |
| 58 | + """ |
| 59 | + from headroom.proxy.helpers import COMPRESSION_TIMEOUT_SECONDS |
| 60 | + from headroom.tokenizers import EstimatingTokenCounter, get_tokenizer |
| 61 | + |
| 62 | + def _resolve_and_count(): # noqa: ANN202 |
| 63 | + tokenizer = get_tokenizer(model) |
| 64 | + return tokenizer, tokenizer.count_messages(messages) |
| 65 | + |
| 66 | + try: |
| 67 | + return await self._run_compression_in_executor( |
| 68 | + _resolve_and_count, |
| 69 | + timeout=float(COMPRESSION_TIMEOUT_SECONDS), |
| 70 | + ) |
| 71 | + except Exception as e: # fail open — includes asyncio.TimeoutError |
| 72 | + # Log the downgrade once per model, not per request. |
| 73 | + fallback_models = getattr(self, "_token_count_fallback_models", None) |
| 74 | + if fallback_models is None: |
| 75 | + fallback_models = set() |
| 76 | + self._token_count_fallback_models = fallback_models |
| 77 | + if model not in fallback_models: |
| 78 | + fallback_models.add(model) |
| 79 | + logger.warning( |
| 80 | + f"Token counting for model {model} failed or timed out " |
| 81 | + f"({e.__class__.__name__}); falling back to estimation" |
| 82 | + ) |
| 83 | + estimator = EstimatingTokenCounter() |
| 84 | + return estimator, estimator.count_messages(messages) |
| 85 | + |
44 | 86 | @staticmethod |
45 | 87 | def _resolve_ccr_workspace( |
46 | 88 | request: Any, |
@@ -469,7 +511,6 @@ async def handle_anthropic_messages( |
469 | 511 | read_request_json_with_bytes, |
470 | 512 | ) |
471 | 513 | from headroom.proxy.modes import is_cache_mode, is_token_mode |
472 | | - from headroom.tokenizers import get_tokenizer |
473 | 514 | from headroom.utils import extract_user_query |
474 | 515 |
|
475 | 516 | start_time = time.time() |
@@ -899,9 +940,10 @@ async def _finalize_pre_upstream() -> None: |
899 | 940 | media_type="application/json", |
900 | 941 | ) |
901 | 942 |
|
902 | | - # Count original tokens |
903 | | - tokenizer = get_tokenizer(model) |
904 | | - original_tokens = tokenizer.count_messages(messages) |
| 943 | + # Count original tokens off the event loop: first-use tokenizer |
| 944 | + # resolution may hit the network (HF download) and counting a full |
| 945 | + # conversation is CPU-bound — on-loop it froze the server (#1701). |
| 946 | + tokenizer, original_tokens = await self._count_tokens_offloaded(model, messages) |
905 | 947 |
|
906 | 948 | # Enterprise Security: scan request before compression |
907 | 949 | _security_ctx = None |
@@ -1164,7 +1206,9 @@ def should_skip_ccr_request_compression( |
1164 | 1206 | ) |
1165 | 1207 | if skip_ccr_request_compression: |
1166 | 1208 | optimized_messages = messages |
1167 | | - optimized_tokens = tokenizer.count_messages(optimized_messages) |
| 1209 | + _, optimized_tokens = await self._count_tokens_offloaded( |
| 1210 | + model, optimized_messages |
| 1211 | + ) |
1168 | 1212 | else: |
1169 | 1213 | # Zone 1: Swap cached compressed versions into working copy |
1170 | 1214 | working_messages = comp_cache.apply_cached(messages) |
@@ -3034,7 +3078,6 @@ async def handle_anthropic_batch_create( |
3034 | 3078 | from headroom.ccr import CCRToolInjector |
3035 | 3079 | from headroom.proxy.helpers import MAX_REQUEST_BODY_SIZE, _read_request_json |
3036 | 3080 | from headroom.proxy.modes import is_cache_mode |
3037 | | - from headroom.tokenizers import get_tokenizer |
3038 | 3081 | from headroom.utils import extract_user_query |
3039 | 3082 |
|
3040 | 3083 | start_time = time.time() |
@@ -3142,17 +3185,27 @@ async def handle_anthropic_batch_create( |
3142 | 3185 | ) |
3143 | 3186 | if is_cache_mode(self.config.mode): |
3144 | 3187 | optimized_messages = messages |
3145 | | - original_tokens = get_tokenizer(model).count_messages(messages) |
| 3188 | + _, original_tokens = await self._count_tokens_offloaded(model, messages) |
3146 | 3189 | optimized_tokens = original_tokens |
3147 | 3190 | else: |
3148 | | - result = self.anthropic_pipeline.apply( |
3149 | | - messages=messages, |
3150 | | - model=model, |
3151 | | - model_limit=context_limit, |
3152 | | - context=extract_user_query(messages), |
3153 | | - frozen_message_count=frozen_message_count, |
3154 | | - request_id=request_id, |
3155 | | - **proxy_pipeline_kwargs(self.config), |
| 3191 | + from headroom.proxy.helpers import COMPRESSION_TIMEOUT_SECONDS |
| 3192 | + |
| 3193 | + # Offload off the event loop (#1701): an inline apply() |
| 3194 | + # blocks every other request for the duration; a timeout |
| 3195 | + # here is caught below and passes the item through. |
| 3196 | + result = await self._run_compression_in_executor( |
| 3197 | + lambda messages=messages, model=model, context_limit=context_limit, frozen_message_count=frozen_message_count: ( |
| 3198 | + self.anthropic_pipeline.apply( |
| 3199 | + messages=messages, |
| 3200 | + model=model, |
| 3201 | + model_limit=context_limit, |
| 3202 | + context=extract_user_query(messages), |
| 3203 | + frozen_message_count=frozen_message_count, |
| 3204 | + request_id=request_id, |
| 3205 | + **proxy_pipeline_kwargs(self.config), |
| 3206 | + ) |
| 3207 | + ), |
| 3208 | + timeout=COMPRESSION_TIMEOUT_SECONDS, |
3156 | 3209 | ) |
3157 | 3210 |
|
3158 | 3211 | optimized_messages = result.messages |
|
0 commit comments