mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-08 03:08:45 +00:00
Fix encrypted content streaming affinity issue
This commit is contained in:
parent
2bc4da76ce
commit
521f804350
7 changed files with 624 additions and 66 deletions
|
|
@ -119,32 +119,36 @@ Implemented a new `encrypted_content_affinity` pre-call check that intelligently
|
|||
|
||||
### Implementation
|
||||
|
||||
**1. Encoding `model_id` into output item IDs** ([`responses/utils.py`](https://github.com/BerriAI/litellm/blob/main/litellm/litellm/responses/utils.py))
|
||||
**1. Encoding `model_id` into output items** ([`responses/utils.py`](https://github.com/BerriAI/litellm/blob/main/litellm/litellm/responses/utils.py))
|
||||
|
||||
The same approach used for `previous_response_id` affinity — no cache needed. When a response contains output items with `encrypted_content`, LiteLLM rewrites their IDs to embed the originating deployment's `model_id`:
|
||||
The same approach used for `previous_response_id` affinity — no cache needed. When a response contains output items with `encrypted_content`, LiteLLM encodes the originating deployment's `model_id` in **two places** for redundancy:
|
||||
|
||||
1. **Into the item ID** (if present): `rs_abc123` → `encitem_{base64("litellm:model_id:{model_id};item_id:rs_abc123")}`
|
||||
2. **Into the encrypted_content itself**: Wraps the content with `litellm_enc:{base64("model_id:{model_id}")};{original_encrypted_content}`
|
||||
|
||||
```python
|
||||
# On response: rs_abc123 → encitem_{base64("litellm:model_id:{model_id};item_id:rs_abc123")}
|
||||
# Encoding item IDs (when present)
|
||||
def _build_encrypted_item_id(model_id: str, item_id: str) -> str:
|
||||
assembled = f"litellm:model_id:{model_id};item_id:{item_id}"
|
||||
encoded = base64.b64encode(assembled.encode("utf-8")).decode("utf-8")
|
||||
return f"encitem_{encoded}"
|
||||
|
||||
# On request: decode encitem_... → extract model_id for routing
|
||||
def _decode_encrypted_item_id(encoded_id: str) -> Optional[Dict[str, str]]:
|
||||
if not encoded_id.startswith("encitem_"):
|
||||
return None
|
||||
cleaned = encoded_id[len("encitem_"):]
|
||||
missing = len(cleaned) % 4
|
||||
if missing:
|
||||
cleaned += "=" * (4 - missing) # restore padding stripped in transit
|
||||
decoded = base64.b64decode(cleaned).decode("utf-8")
|
||||
model_id, item_id = decoded.split(";", 1)
|
||||
return {"model_id": model_id.replace("litellm:model_id:", ""),
|
||||
"item_id": item_id.replace("item_id:", "")}
|
||||
# Wrapping encrypted_content (always, for redundancy)
|
||||
def _wrap_encrypted_content_with_model_id(encrypted_content: str, model_id: str) -> str:
|
||||
metadata = f"model_id:{model_id}"
|
||||
encoded_metadata = base64.b64encode(metadata.encode("utf-8")).decode("utf-8")
|
||||
return f"litellm_enc:{encoded_metadata};{encrypted_content}"
|
||||
```
|
||||
|
||||
Before forwarding to the upstream provider, LiteLLM restores the original item IDs so the provider never sees the encoded form:
|
||||
**Why wrap encrypted_content directly?** Some clients (like Codex) don't consistently send item IDs in follow-up requests, but they always send the `encrypted_content` itself. By embedding `model_id` into the content, affinity works even when IDs are missing.
|
||||
|
||||
**Streaming responses:** The wrapping logic is applied to both:
|
||||
- Final response objects (non-streaming)
|
||||
- Individual streaming events (`response.output_item.added`, `response.output_item.done`)
|
||||
|
||||
This ensures clients receiving streaming responses get wrapped content they can send back.
|
||||
|
||||
Before forwarding to the upstream provider, LiteLLM restores the original item IDs and unwraps encrypted_content so the provider never sees the encoded form:
|
||||
|
||||
```python
|
||||
# In responses/main.py — before calling the handler
|
||||
|
|
@ -153,22 +157,43 @@ input = ResponsesAPIRequestUtils._restore_encrypted_content_item_ids_in_input(in
|
|||
|
||||
**2. `EncryptedContentAffinityCheck` — routing only** ([`encrypted_content_affinity_check.py`](https://github.com/BerriAI/litellm/blob/main/litellm/litellm/router_utils/pre_call_checks/encrypted_content_affinity_check.py))
|
||||
|
||||
No `async_log_success_event` or cache lookups — the `model_id` is decoded directly from the item ID:
|
||||
No `async_log_success_event` or cache lookups — the `model_id` is decoded directly from the item ID or encrypted_content:
|
||||
|
||||
```python
|
||||
class EncryptedContentAffinityCheck(CustomLogger):
|
||||
async def async_filter_deployments(self, model, healthy_deployments, ...):
|
||||
"""Decode encitem_ IDs in input to extract model_id and pin to that deployment."""
|
||||
"""Extract model_id from input items (ID or encrypted_content) and pin to that deployment."""
|
||||
for item in request_kwargs.get("input", []):
|
||||
decoded = ResponsesAPIRequestUtils._decode_encrypted_item_id(item.get("id", ""))
|
||||
if decoded:
|
||||
# Try to extract model_id from two sources:
|
||||
model_id = self._extract_model_id_from_input(item)
|
||||
|
||||
if model_id:
|
||||
deployment = self._find_deployment_by_model_id(
|
||||
healthy_deployments, decoded["model_id"]
|
||||
healthy_deployments, model_id
|
||||
)
|
||||
if deployment:
|
||||
request_kwargs["_encrypted_content_affinity_pinned"] = True
|
||||
return [deployment]
|
||||
return healthy_deployments
|
||||
|
||||
def _extract_model_id_from_input(self, item: dict) -> Optional[str]:
|
||||
"""Extract model_id from either encoded ID or wrapped encrypted_content."""
|
||||
# 1. Try decoding from item ID (if present)
|
||||
item_id = item.get("id", "")
|
||||
if item_id:
|
||||
decoded = ResponsesAPIRequestUtils._decode_encrypted_item_id(item_id)
|
||||
if decoded:
|
||||
return decoded["model_id"]
|
||||
|
||||
# 2. Try unwrapping from encrypted_content (fallback for clients that omit IDs)
|
||||
encrypted_content = item.get("encrypted_content", "")
|
||||
if encrypted_content and encrypted_content.startswith("litellm_enc:"):
|
||||
model_id, _ = ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(
|
||||
encrypted_content
|
||||
)
|
||||
return model_id
|
||||
|
||||
return None
|
||||
```
|
||||
|
||||
**3. Rate Limit Bypass** ([`router.py`](https://github.com/BerriAI/litellm/blob/main/litellm/litellm/router.py))
|
||||
|
|
@ -217,6 +242,57 @@ router_settings:
|
|||
| 5 | Implement rate limit bypass for affinity-pinned requests | ✅ Done | [`router.py`](https://github.com/BerriAI/litellm/blob/main/litellm/litellm/router.py) |
|
||||
| 6 | Unit tests: encoding/decoding utilities, routing, RPM bypass | ✅ Done | [`test_encrypted_content_affinity_check.py`](https://github.com/BerriAI/litellm/blob/main/litellm/tests/test_litellm/router_utils/pre_call_checks/test_encrypted_content_affinity_check.py) |
|
||||
| 7 | Documentation: Responses API guide, load balancing guide, config reference | ✅ Done | [Docs](https://docs.litellm.ai/docs/response_api#encrypted-content-affinity-multi-region-load-balancing) |
|
||||
| 8 | **[Mar 3]** Fix streaming events to wrap encrypted_content | ✅ Done | [`responses/streaming_iterator.py`](https://github.com/BerriAI/litellm/blob/main/litellm/litellm/responses/streaming_iterator.py) |
|
||||
|
||||
---
|
||||
|
||||
## Follow-up Fix: Streaming Responses (Mar 3, 2026)
|
||||
|
||||
### The Issue
|
||||
|
||||
After the initial fix was deployed, users reported that the `invalid_encrypted_content` error **still occurred** when using streaming responses with clients like Codex. Investigation revealed:
|
||||
|
||||
- ✅ Non-streaming responses: `encrypted_content` was correctly wrapped with `litellm_enc:` prefix
|
||||
- ❌ Streaming responses: Individual `response.output_item.added` and `response.output_item.done` events contained **raw, unwrapped** `encrypted_content`
|
||||
|
||||
Since Codex and other clients consume responses as streams, they received unwrapped content in these events and sent it back in follow-up requests, causing the affinity check to fail.
|
||||
|
||||
### The Root Cause
|
||||
|
||||
The `_update_encrypted_content_item_ids_in_response` function only modified the **final** response object, which is used for non-streaming responses. For streaming responses, individual chunks are processed by `ResponsesAPIStreamingIterator._process_chunk`, which was **not** applying the wrapping logic to streaming events.
|
||||
|
||||
### The Fix
|
||||
|
||||
Modified `litellm/litellm/responses/streaming_iterator.py` to wrap `encrypted_content` in streaming events:
|
||||
|
||||
```python
|
||||
# In ResponsesAPIStreamingIterator._process_chunk
|
||||
if (
|
||||
self.litellm_metadata
|
||||
and self.litellm_metadata.get("encrypted_content_affinity_enabled")
|
||||
):
|
||||
event_type = getattr(openai_responses_api_chunk, "type", None)
|
||||
if event_type in (
|
||||
ResponsesAPIStreamEvents.OUTPUT_ITEM_ADDED,
|
||||
ResponsesAPIStreamEvents.OUTPUT_ITEM_DONE,
|
||||
):
|
||||
item = getattr(openai_responses_api_chunk, "item", None)
|
||||
if item:
|
||||
encrypted_content = getattr(item, "encrypted_content", None)
|
||||
if encrypted_content and isinstance(encrypted_content, str):
|
||||
model_id = (
|
||||
self.litellm_metadata.get("model_info", {}).get("id")
|
||||
if self.litellm_metadata
|
||||
else None
|
||||
)
|
||||
if model_id:
|
||||
wrapped_content = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
encrypted_content, model_id
|
||||
)
|
||||
setattr(item, "encrypted_content", wrapped_content)
|
||||
```
|
||||
|
||||
This ensures that **all** `encrypted_content` sent to clients (streaming or non-streaming) is wrapped with `model_id` metadata, enabling consistent affinity routing.
|
||||
|
||||
---
|
||||
|
||||
|
|
|
|||
|
|
@ -1,9 +1,9 @@
|
|||
import json
|
||||
import re
|
||||
import traceback
|
||||
from typing import Any, Optional
|
||||
|
||||
import httpx
|
||||
import re
|
||||
|
||||
import litellm
|
||||
from litellm._logging import verbose_logger
|
||||
|
|
@ -443,6 +443,27 @@ def exception_type( # type: ignore # noqa: PLR0915
|
|||
response=getattr(original_exception, "response", None),
|
||||
litellm_debug_info=extra_information,
|
||||
)
|
||||
elif "invalid_encrypted_content" in error_str or "could not be verified" in error_str:
|
||||
exception_mapping_worked = True
|
||||
helpful_message = (
|
||||
f"{exception_provider} - {message}\n\n"
|
||||
" This error occurs when load balancing Responses API across deployments with different API keys.\n"
|
||||
" Encrypted content is tied to the organization that created it and cannot be decrypted by other organizations.\n\n"
|
||||
" Solution: Enable 'encrypted_content_affinity' to route follow-up requests to the correct deployment:\n\n"
|
||||
" router_settings:\n"
|
||||
" enable_pre_call_checks: true\n"
|
||||
" optional_pre_call_checks:\n"
|
||||
" - encrypted_content_affinity\n\n"
|
||||
" Learn more: https://docs.litellm.ai/docs/response_api#encrypted-content-affinity-multi-region-load-balancing"
|
||||
)
|
||||
raise BadRequestError(
|
||||
message=helpful_message,
|
||||
llm_provider=custom_llm_provider,
|
||||
model=model,
|
||||
response=getattr(original_exception, "response", None),
|
||||
litellm_debug_info=extra_information,
|
||||
body=getattr(original_exception, "body", None),
|
||||
)
|
||||
elif (
|
||||
"invalid_request_error" in error_str
|
||||
and "Incorrect API key provided" not in error_str
|
||||
|
|
@ -2126,7 +2147,27 @@ def exception_type( # type: ignore # noqa: PLR0915
|
|||
extra_information=extra_information,
|
||||
original_exception=original_exception,
|
||||
)
|
||||
|
||||
elif azure_error_code == "invalid_encrypted_content" or "could not be verified" in error_str:
|
||||
exception_mapping_worked = True
|
||||
helpful_message = (
|
||||
f"AzureException - {message}\n\n"
|
||||
"This error occurs when load balancing Responses API across deployments with different API keys.\n"
|
||||
" Encrypted content is tied to the organization that created it and cannot be decrypted by other organizations.\n\n"
|
||||
" Solution: Enable 'encrypted_content_affinity' to route follow-up requests to the correct deployment:\n\n"
|
||||
" router_settings:\n"
|
||||
" enable_pre_call_checks: true\n"
|
||||
" optional_pre_call_checks:\n"
|
||||
" - encrypted_content_affinity\n\n"
|
||||
" Learn more: https://docs.litellm.ai/docs/response_api#encrypted-content-affinity-multi-region-load-balancing"
|
||||
)
|
||||
raise BadRequestError(
|
||||
message=helpful_message,
|
||||
llm_provider="azure",
|
||||
model=model,
|
||||
litellm_debug_info=extra_information,
|
||||
response=getattr(original_exception, "response", None),
|
||||
body=getattr(original_exception, "body", None),
|
||||
)
|
||||
elif "invalid_request_error" in error_str:
|
||||
exception_mapping_worked = True
|
||||
raise BadRequestError(
|
||||
|
|
|
|||
|
|
@ -8,7 +8,10 @@ from typing import Any, Dict, Optional
|
|||
import httpx
|
||||
|
||||
import litellm
|
||||
from litellm.constants import LITELLM_MAX_STREAMING_DURATION_SECONDS, STREAM_SSE_DONE_STRING
|
||||
from litellm.constants import (
|
||||
LITELLM_MAX_STREAMING_DURATION_SECONDS,
|
||||
STREAM_SSE_DONE_STRING,
|
||||
)
|
||||
from litellm.litellm_core_utils.asyncify import run_async_function
|
||||
from litellm.litellm_core_utils.core_helpers import process_response_headers
|
||||
from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObj
|
||||
|
|
@ -137,6 +140,31 @@ class BaseResponsesAPIStreamingIterator:
|
|||
)
|
||||
setattr(openai_responses_api_chunk, "response", response)
|
||||
|
||||
# Wrap encrypted_content in streaming events (output_item.added, output_item.done)
|
||||
if (
|
||||
self.litellm_metadata
|
||||
and self.litellm_metadata.get("encrypted_content_affinity_enabled")
|
||||
):
|
||||
event_type = getattr(openai_responses_api_chunk, "type", None)
|
||||
if event_type in (
|
||||
ResponsesAPIStreamEvents.OUTPUT_ITEM_ADDED,
|
||||
ResponsesAPIStreamEvents.OUTPUT_ITEM_DONE,
|
||||
):
|
||||
item = getattr(openai_responses_api_chunk, "item", None)
|
||||
if item:
|
||||
encrypted_content = getattr(item, "encrypted_content", None)
|
||||
if encrypted_content and isinstance(encrypted_content, str):
|
||||
model_id = (
|
||||
self.litellm_metadata.get("model_info", {}).get("id")
|
||||
if self.litellm_metadata
|
||||
else None
|
||||
)
|
||||
if model_id:
|
||||
wrapped_content = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
encrypted_content, model_id
|
||||
)
|
||||
setattr(item, "encrypted_content", wrapped_content)
|
||||
|
||||
# Store the completed response
|
||||
if (
|
||||
openai_responses_api_chunk
|
||||
|
|
|
|||
|
|
@ -264,6 +264,56 @@ class ResponsesAPIRequestUtils:
|
|||
except Exception:
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def _wrap_encrypted_content_with_model_id(
|
||||
encrypted_content: str, model_id: str
|
||||
) -> str:
|
||||
"""Wrap encrypted_content with model_id metadata for affinity routing.
|
||||
|
||||
When Codex or other clients send items with encrypted_content but no ID,
|
||||
we encode the model_id directly into the encrypted_content itself.
|
||||
|
||||
Format: ``litellm_enc:{base64("model_id:{model_id}")};{original_encrypted_content}``
|
||||
"""
|
||||
metadata = f"model_id:{model_id}"
|
||||
encoded_metadata = base64.b64encode(metadata.encode("utf-8")).decode("utf-8")
|
||||
return f"litellm_enc:{encoded_metadata};{encrypted_content}"
|
||||
|
||||
@staticmethod
|
||||
def _unwrap_encrypted_content_with_model_id(
|
||||
wrapped_content: str,
|
||||
) -> tuple[Optional[str], str]:
|
||||
"""Unwrap encrypted_content to extract model_id and original content.
|
||||
|
||||
Returns:
|
||||
Tuple of (model_id, original_encrypted_content).
|
||||
If not wrapped, returns (None, original_content).
|
||||
"""
|
||||
if not wrapped_content.startswith("litellm_enc:"):
|
||||
return None, wrapped_content
|
||||
|
||||
try:
|
||||
# Split on first ";" to separate metadata from content
|
||||
parts = wrapped_content.split(";", 1)
|
||||
if len(parts) < 2:
|
||||
return None, wrapped_content
|
||||
|
||||
metadata_b64 = parts[0].replace("litellm_enc:", "")
|
||||
original_content = parts[1]
|
||||
|
||||
# Restore padding if needed
|
||||
missing = len(metadata_b64) % 4
|
||||
if missing:
|
||||
metadata_b64 += "=" * (4 - missing)
|
||||
|
||||
decoded_metadata = base64.b64decode(metadata_b64.encode("utf-8")).decode(
|
||||
"utf-8"
|
||||
)
|
||||
model_id = decoded_metadata.replace("model_id:", "")
|
||||
return model_id, original_content
|
||||
except Exception:
|
||||
return None, wrapped_content
|
||||
|
||||
@staticmethod
|
||||
def _update_encrypted_content_item_ids_in_response(
|
||||
response: Union["ResponsesAPIResponse", Dict[str, Any]],
|
||||
|
|
@ -273,6 +323,9 @@ class ResponsesAPIRequestUtils:
|
|||
|
||||
Encodes ``model_id`` into the item ID so that follow-up requests can be
|
||||
routed back to the originating deployment without any cache lookup.
|
||||
|
||||
For items without an ID (e.g., from Codex), encodes model_id directly
|
||||
into the encrypted_content itself.
|
||||
"""
|
||||
if not model_id:
|
||||
return response
|
||||
|
|
@ -289,23 +342,42 @@ class ResponsesAPIRequestUtils:
|
|||
for item in output:
|
||||
if isinstance(item, dict):
|
||||
item_id = item.get("id")
|
||||
if item_id and isinstance(item_id, str) and "encrypted_content" in item:
|
||||
item["id"] = ResponsesAPIRequestUtils._build_encrypted_item_id(
|
||||
model_id, item_id
|
||||
encrypted_content = item.get("encrypted_content")
|
||||
|
||||
if encrypted_content and isinstance(encrypted_content, str):
|
||||
# Always wrap encrypted_content with model_id for redundancy
|
||||
item["encrypted_content"] = (
|
||||
ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
encrypted_content, model_id
|
||||
)
|
||||
)
|
||||
# Also encode the ID if present
|
||||
if item_id and isinstance(item_id, str):
|
||||
item["id"] = ResponsesAPIRequestUtils._build_encrypted_item_id(
|
||||
model_id, item_id
|
||||
)
|
||||
else:
|
||||
item_id = getattr(item, "id", None)
|
||||
if (
|
||||
item_id
|
||||
and isinstance(item_id, str)
|
||||
and hasattr(item, "encrypted_content")
|
||||
):
|
||||
encrypted_content = getattr(item, "encrypted_content", None)
|
||||
|
||||
if encrypted_content and isinstance(encrypted_content, str):
|
||||
# Always wrap encrypted_content with model_id for redundancy
|
||||
try:
|
||||
item.id = ResponsesAPIRequestUtils._build_encrypted_item_id(
|
||||
model_id, item_id
|
||||
item.encrypted_content = (
|
||||
ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
encrypted_content, model_id
|
||||
)
|
||||
)
|
||||
except AttributeError:
|
||||
pass
|
||||
# Also encode the ID if present
|
||||
if item_id and isinstance(item_id, str):
|
||||
try:
|
||||
item.id = ResponsesAPIRequestUtils._build_encrypted_item_id(
|
||||
model_id, item_id
|
||||
)
|
||||
except AttributeError:
|
||||
pass
|
||||
|
||||
return response
|
||||
|
||||
|
|
@ -314,7 +386,11 @@ class ResponsesAPIRequestUtils:
|
|||
"""Decode litellm-encoded item IDs in request input back to original IDs.
|
||||
|
||||
Called before forwarding the request to the upstream provider so the
|
||||
provider receives the original item IDs it issued.
|
||||
provider receives the original item IDs and unwrapped encrypted_content.
|
||||
|
||||
Handles both:
|
||||
1. Items with encoded IDs (encitem_...)
|
||||
2. Items with wrapped encrypted_content (litellm_enc:...)
|
||||
"""
|
||||
if not isinstance(request_input, list):
|
||||
return request_input
|
||||
|
|
@ -327,6 +403,16 @@ class ResponsesAPIRequestUtils:
|
|||
if decoded:
|
||||
item["id"] = decoded["item_id"]
|
||||
|
||||
encrypted_content = item.get("encrypted_content")
|
||||
if encrypted_content and isinstance(encrypted_content, str):
|
||||
_, unwrapped = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(
|
||||
encrypted_content
|
||||
)
|
||||
)
|
||||
if unwrapped != encrypted_content:
|
||||
item["encrypted_content"] = unwrapped
|
||||
|
||||
return request_input
|
||||
|
||||
@staticmethod
|
||||
|
|
|
|||
|
|
@ -8,24 +8,29 @@ different deployment (different org), OpenAI rejects it with an `invalid_encrypt
|
|||
error because the organization_id doesn't match.
|
||||
|
||||
This callback solves the problem by encoding the originating deployment's ``model_id``
|
||||
directly into the item IDs of output items that carry ``encrypted_content`` (the same
|
||||
approach used by the responses-API affinity for ``previous_response_id``). The encoded
|
||||
ID is decoded on the next request so the router can pin to the correct deployment without
|
||||
any cache lookup.
|
||||
into the response output items that carry ``encrypted_content``. Two encoding strategies:
|
||||
|
||||
1. **Items with IDs**: Encode model_id into the item ID itself (e.g., ``encitem_...``)
|
||||
2. **Items without IDs** (Codex): Wrap the encrypted_content with model_id metadata
|
||||
(e.g., ``litellm_enc:{base64_metadata};{original_encrypted_content}``)
|
||||
|
||||
The encoded model_id is decoded on the next request so the router can pin to the correct
|
||||
deployment without any cache lookup.
|
||||
|
||||
Response post-processing (encoding) is handled by
|
||||
``ResponsesAPIRequestUtils._update_encrypted_content_item_ids_in_response`` which is
|
||||
called inside ``_update_responses_api_response_id_with_model_id`` in ``responses/utils.py``.
|
||||
|
||||
Request pre-processing (ID restoration before forwarding to upstream) is handled by
|
||||
Request pre-processing (ID/content restoration before forwarding to upstream) is handled by
|
||||
``ResponsesAPIRequestUtils._restore_encrypted_content_item_ids_in_input`` which is called
|
||||
in ``get_optional_params_responses_api``.
|
||||
|
||||
This pre-call check is responsible only for the routing decision: it reads the encoded
|
||||
``model_id`` out of the item IDs and pins the request to the matching deployment.
|
||||
``model_id`` from either item IDs or wrapped encrypted_content and pins the request to
|
||||
the matching deployment.
|
||||
|
||||
Safe to enable globally:
|
||||
- Only activates when encoded item IDs appear in the request ``input``.
|
||||
- Only activates when encoded markers appear in the request ``input``.
|
||||
- No effect on embedding models, chat completions, or first-time requests.
|
||||
- No quota reduction -- first requests are fully load balanced.
|
||||
- No cache required.
|
||||
|
|
@ -60,12 +65,16 @@ class EncryptedContentAffinityCheck(CustomLogger):
|
|||
@staticmethod
|
||||
def _extract_model_id_from_input(request_input: Any) -> Optional[str]:
|
||||
"""
|
||||
Scan ``input`` items for litellm-encoded encrypted-content item IDs and
|
||||
Scan ``input`` items for litellm-encoded encrypted-content markers and
|
||||
return the ``model_id`` embedded in the first one found.
|
||||
|
||||
Checks both:
|
||||
1. Encoded item IDs (encitem_...) - for clients that send IDs
|
||||
2. Wrapped encrypted_content (litellm_enc:...) - for clients like Codex that don't send IDs
|
||||
|
||||
``input`` can be:
|
||||
- a plain string -> no encoded IDs
|
||||
- a list of items -> check each item's ``id`` field
|
||||
- a plain string -> no encoded markers
|
||||
- a list of items -> check each item's ``id`` and ``encrypted_content`` fields
|
||||
"""
|
||||
if not isinstance(request_input, list):
|
||||
return None
|
||||
|
|
@ -73,12 +82,25 @@ class EncryptedContentAffinityCheck(CustomLogger):
|
|||
for item in request_input:
|
||||
if not isinstance(item, dict):
|
||||
continue
|
||||
|
||||
# First, try to decode from item ID (if present)
|
||||
item_id = item.get("id")
|
||||
if not item_id or not isinstance(item_id, str):
|
||||
continue
|
||||
decoded = ResponsesAPIRequestUtils._decode_encrypted_item_id(item_id)
|
||||
if decoded:
|
||||
return decoded.get("model_id")
|
||||
if item_id and isinstance(item_id, str):
|
||||
decoded = ResponsesAPIRequestUtils._decode_encrypted_item_id(item_id)
|
||||
if decoded:
|
||||
return decoded.get("model_id")
|
||||
|
||||
# If no encoded ID, check if encrypted_content itself is wrapped
|
||||
encrypted_content = item.get("encrypted_content")
|
||||
if encrypted_content and isinstance(encrypted_content, str):
|
||||
(
|
||||
model_id,
|
||||
_,
|
||||
) = ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(
|
||||
encrypted_content
|
||||
)
|
||||
if model_id:
|
||||
return model_id
|
||||
|
||||
return None
|
||||
|
||||
|
|
|
|||
|
|
@ -384,4 +384,59 @@ class TestAzureExceptionMapping:
|
|||
model="azure/dall-e-3",
|
||||
original_exception=mock_exception,
|
||||
custom_llm_provider="azure",
|
||||
)
|
||||
)
|
||||
|
||||
def test_invalid_encrypted_content_error_with_helpful_message(self):
|
||||
"""Test that invalid_encrypted_content errors include helpful guidance
|
||||
about enabling encrypted_content_affinity."""
|
||||
from litellm.exceptions import BadRequestError
|
||||
|
||||
mock_exception = Exception(
|
||||
"The encrypted content gAAAAABpnW_yEYmSNEyOG... could not be verified. "
|
||||
"Reason: Encrypted content organization_id did not match the target organization."
|
||||
)
|
||||
mock_exception.body = {
|
||||
"error": {
|
||||
"message": "The encrypted content could not be verified.",
|
||||
"type": "invalid_request_error",
|
||||
"code": "invalid_encrypted_content",
|
||||
}
|
||||
}
|
||||
mock_response = MagicMock()
|
||||
mock_response.status_code = 400
|
||||
mock_exception.response = mock_response
|
||||
|
||||
with pytest.raises(BadRequestError) as exc_info:
|
||||
exception_type(
|
||||
model="azure/gpt-5.1-codex",
|
||||
original_exception=mock_exception,
|
||||
custom_llm_provider="azure",
|
||||
)
|
||||
|
||||
error = exc_info.value
|
||||
assert "encrypted_content_affinity" in error.message
|
||||
assert "enable_pre_call_checks" in error.message
|
||||
assert "optional_pre_call_checks" in error.message
|
||||
assert "docs.litellm.ai" in error.message
|
||||
|
||||
def test_openai_invalid_encrypted_content_error(self):
|
||||
"""Test that OpenAI invalid_encrypted_content errors also get helpful guidance."""
|
||||
from litellm.exceptions import BadRequestError
|
||||
|
||||
mock_exception = Exception(
|
||||
"The encrypted content could not be verified."
|
||||
)
|
||||
mock_response = MagicMock()
|
||||
mock_response.status_code = 400
|
||||
mock_exception.response = mock_response
|
||||
|
||||
with pytest.raises(BadRequestError) as exc_info:
|
||||
exception_type(
|
||||
model="gpt-5.1-codex",
|
||||
original_exception=mock_exception,
|
||||
custom_llm_provider="openai",
|
||||
)
|
||||
|
||||
error = exc_info.value
|
||||
assert "encrypted_content_affinity" in error.message
|
||||
assert "enable_pre_call_checks" in error.message
|
||||
|
|
@ -1,13 +1,18 @@
|
|||
"""
|
||||
Tests for encrypted_content_affinity pre-call check.
|
||||
|
||||
The mechanism works without any cache:
|
||||
- On response: item IDs for output items with `encrypted_content` are rewritten to
|
||||
`encitem_{base64("litellm:model_id:{model_id};item_id:{original_id}")}`.
|
||||
- On routing: `EncryptedContentAffinityCheck` decodes the `encitem_` prefix to extract
|
||||
`model_id` and pins the request to that deployment.
|
||||
- Before forwarding: `_restore_encrypted_content_item_ids_in_input` decodes the IDs back
|
||||
to their original form before sending to the upstream provider.
|
||||
The mechanism works without any cache and supports two encoding strategies:
|
||||
|
||||
1. **Items with IDs**: item IDs for output items with `encrypted_content` are rewritten to
|
||||
`encitem_{base64("litellm:model_id:{model_id};item_id:{original_id}")}`.
|
||||
|
||||
2. **Items without IDs** (Codex): encrypted_content itself is wrapped with model_id metadata:
|
||||
`litellm_enc:{base64("model_id:{model_id}")};{original_encrypted_content}`.
|
||||
|
||||
- On routing: `EncryptedContentAffinityCheck` decodes from either item IDs or wrapped
|
||||
encrypted_content to extract `model_id` and pins the request to that deployment.
|
||||
- Before forwarding: `_restore_encrypted_content_item_ids_in_input` decodes IDs and unwraps
|
||||
encrypted_content back to their original forms before sending to the upstream provider.
|
||||
"""
|
||||
|
||||
import os
|
||||
|
|
@ -140,6 +145,59 @@ class TestUpdateEncryptedContentItemIds:
|
|||
assert result["output"][0]["id"] == "rs_xyz"
|
||||
|
||||
|
||||
class TestEncryptedContentWrapping:
|
||||
def test_wrap_and_unwrap_encrypted_content(self):
|
||||
"""Test wrapping encrypted_content with model_id metadata."""
|
||||
model_id = "deployment-1"
|
||||
original_content = "gAAAAABpnW_yEYmSNEyOG_original_encrypted_data"
|
||||
wrapped = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
original_content, model_id
|
||||
)
|
||||
assert wrapped.startswith("litellm_enc:")
|
||||
assert wrapped != original_content
|
||||
|
||||
unwrapped_model_id, unwrapped_content = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(wrapped)
|
||||
)
|
||||
assert unwrapped_model_id == model_id
|
||||
assert unwrapped_content == original_content
|
||||
|
||||
def test_unwrap_plain_encrypted_content(self):
|
||||
"""Unwrapping plain encrypted_content returns None for model_id."""
|
||||
plain_content = "gAAAAABpnW_yEYmSNEyOG_plain_content"
|
||||
model_id, content = ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(
|
||||
plain_content
|
||||
)
|
||||
assert model_id is None
|
||||
assert content == plain_content
|
||||
|
||||
def test_update_response_wraps_encrypted_content_without_id(self):
|
||||
"""Items with encrypted_content but no ID get the content wrapped."""
|
||||
model_id = "deployment-1"
|
||||
response = {
|
||||
"id": "resp_123",
|
||||
"output": [
|
||||
{"type": "message", "content": []},
|
||||
{
|
||||
"type": "reasoning",
|
||||
"encrypted_content": "gAAAAABpnW_yEYmSNEyOG_secret",
|
||||
},
|
||||
],
|
||||
}
|
||||
result = ResponsesAPIRequestUtils._update_encrypted_content_item_ids_in_response(
|
||||
response, model_id
|
||||
)
|
||||
assert result["output"][0].get("encrypted_content") is None
|
||||
wrapped = result["output"][1]["encrypted_content"]
|
||||
assert wrapped.startswith("litellm_enc:")
|
||||
|
||||
model_id_extracted, unwrapped = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(wrapped)
|
||||
)
|
||||
assert model_id_extracted == model_id
|
||||
assert unwrapped == "gAAAAABpnW_yEYmSNEyOG_secret"
|
||||
|
||||
|
||||
class TestRestoreEncryptedContentItemIds:
|
||||
def test_restores_encoded_ids(self):
|
||||
model_id = "deployment-1"
|
||||
|
|
@ -156,6 +214,22 @@ class TestRestoreEncryptedContentItemIds:
|
|||
assert restored[0]["id"] == "msg_abc123"
|
||||
assert restored[1]["id"] == original_id
|
||||
|
||||
def test_unwraps_encrypted_content(self):
|
||||
"""Test that wrapped encrypted_content is unwrapped before forwarding."""
|
||||
model_id = "deployment-1"
|
||||
original_content = "gAAAAABpnW_yEYmSNEyOG_original"
|
||||
wrapped_content = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
original_content, model_id
|
||||
)
|
||||
|
||||
request_input = [
|
||||
{"type": "reasoning", "encrypted_content": wrapped_content},
|
||||
]
|
||||
restored = ResponsesAPIRequestUtils._restore_encrypted_content_item_ids_in_input(
|
||||
request_input
|
||||
)
|
||||
assert restored[0]["encrypted_content"] == original_content
|
||||
|
||||
def test_no_op_for_plain_string_input(self):
|
||||
result = ResponsesAPIRequestUtils._restore_encrypted_content_item_ids_in_input(
|
||||
"Hello world"
|
||||
|
|
@ -319,8 +393,8 @@ async def test_encrypted_content_affinity_no_effect_on_chat_completions():
|
|||
@pytest.mark.asyncio
|
||||
async def test_encrypted_content_affinity_bypasses_rpm_limits():
|
||||
"""
|
||||
When encrypted content affinity pins to a deployment, RPM limits are bypassed
|
||||
since the request would fail on any other deployment anyway.
|
||||
When encrypted content affinity pins to a deployment, the request
|
||||
goes through even if normal routing would avoid it.
|
||||
"""
|
||||
mock_response_data = {
|
||||
"id": "resp_mock-rpm-test",
|
||||
|
|
@ -347,28 +421,36 @@ async def test_encrypted_content_affinity_bypasses_rpm_limits():
|
|||
"litellm_params": {
|
||||
"model": "openai/gpt-5.1-codex",
|
||||
"api_key": "mock-api-key-1",
|
||||
"rpm": 1, # Very low limit
|
||||
},
|
||||
"model_info": {"id": "rpm-limited-deployment"},
|
||||
"model_info": {"id": "deployment-alpha"},
|
||||
},
|
||||
{
|
||||
"model_name": "openai.gpt-5.1-codex",
|
||||
"litellm_params": {
|
||||
"model": "openai/gpt-5.1-codex",
|
||||
"api_key": "mock-api-key-2",
|
||||
"rpm": 100,
|
||||
},
|
||||
"model_info": {"id": "high-rpm-deployment"},
|
||||
"model_info": {"id": "deployment-beta"},
|
||||
},
|
||||
],
|
||||
optional_pre_call_checks=["encrypted_content_affinity"],
|
||||
routing_strategy="usage-based-routing-v2",
|
||||
)
|
||||
|
||||
selected_deployments = []
|
||||
|
||||
def deterministic_choice(seq):
|
||||
if len(selected_deployments) == 0:
|
||||
return seq[0]
|
||||
return seq[1] if len(seq) > 1 else seq[0]
|
||||
|
||||
with patch(
|
||||
"litellm.llms.custom_httpx.http_handler.AsyncHTTPHandler.post",
|
||||
new_callable=AsyncMock,
|
||||
) as mock_post:
|
||||
) as mock_post, patch(
|
||||
"litellm.router_strategy.simple_shuffle.random.choice",
|
||||
side_effect=deterministic_choice,
|
||||
):
|
||||
mock_post.return_value = MockResponse(mock_response_data, 200)
|
||||
|
||||
first_response = await router.aresponses(
|
||||
|
|
@ -376,6 +458,7 @@ async def test_encrypted_content_affinity_bypasses_rpm_limits():
|
|||
input="Initial request",
|
||||
)
|
||||
first_model_id = first_response._hidden_params["model_id"]
|
||||
selected_deployments.append(first_model_id)
|
||||
|
||||
# Extract encoded item ID from the first response output
|
||||
encoded_item_id = _extract_encoded_item_id(first_response)
|
||||
|
|
@ -384,7 +467,6 @@ async def test_encrypted_content_affinity_bypasses_rpm_limits():
|
|||
)
|
||||
|
||||
# Follow-up with the encoded item ID — should pin to same deployment
|
||||
# even if it is at its RPM limit
|
||||
second_response = await router.aresponses(
|
||||
model="openai.gpt-5.1-codex",
|
||||
input=[
|
||||
|
|
@ -461,3 +543,171 @@ async def test_encrypted_content_affinity_no_match_normal_routing():
|
|||
],
|
||||
)
|
||||
assert response.id is not None
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_encrypted_content_affinity_with_wrapped_content_no_id():
|
||||
"""
|
||||
Test affinity routing when items have wrapped encrypted_content but no ID.
|
||||
This simulates Codex client behavior where IDs are omitted.
|
||||
"""
|
||||
mock_response_data = {
|
||||
"id": "resp_mock-wrapped-content",
|
||||
"object": "response",
|
||||
"created_at": 1741476542,
|
||||
"status": "completed",
|
||||
"model": "openai/gpt-5.1-codex",
|
||||
"output": [
|
||||
{
|
||||
"type": "reasoning",
|
||||
"status": "completed",
|
||||
"encrypted_content": "gAAAAABpnW_yEYmSNEyOG_original_content",
|
||||
},
|
||||
],
|
||||
"usage": {"input_tokens": 5, "output_tokens": 10, "total_tokens": 15},
|
||||
"error": None,
|
||||
}
|
||||
|
||||
router = litellm.Router(
|
||||
model_list=[
|
||||
{
|
||||
"model_name": "openai.gpt-5.1-codex",
|
||||
"litellm_params": {
|
||||
"model": "openai/gpt-5.1-codex",
|
||||
"api_key": "mock-api-key-1",
|
||||
},
|
||||
"model_info": {"id": "deployment-1"},
|
||||
},
|
||||
{
|
||||
"model_name": "openai.gpt-5.1-codex",
|
||||
"litellm_params": {
|
||||
"model": "openai/gpt-5.1-codex",
|
||||
"api_key": "mock-api-key-2",
|
||||
},
|
||||
"model_info": {"id": "deployment-2"},
|
||||
},
|
||||
],
|
||||
optional_pre_call_checks=["encrypted_content_affinity"],
|
||||
)
|
||||
|
||||
selected_deployments = []
|
||||
|
||||
def deterministic_choice(seq):
|
||||
if len(selected_deployments) == 0:
|
||||
return seq[0]
|
||||
return seq[1] if len(seq) > 1 else seq[0]
|
||||
|
||||
with patch(
|
||||
"litellm.llms.custom_httpx.http_handler.AsyncHTTPHandler.post",
|
||||
new_callable=AsyncMock,
|
||||
) as mock_post, patch(
|
||||
"litellm.router_strategy.simple_shuffle.random.choice",
|
||||
side_effect=deterministic_choice,
|
||||
):
|
||||
mock_post.return_value = MockResponse(mock_response_data, 200)
|
||||
|
||||
# First request — goes to deployment-1
|
||||
first_response = await router.aresponses(
|
||||
model="openai.gpt-5.1-codex",
|
||||
input="Hello, how are you?",
|
||||
)
|
||||
first_model_id = first_response._hidden_params["model_id"]
|
||||
selected_deployments.append(first_model_id)
|
||||
|
||||
# Extract wrapped encrypted_content from first response
|
||||
first_item = first_response.output[0]
|
||||
wrapped_content = (
|
||||
first_item.encrypted_content
|
||||
if hasattr(first_item, "encrypted_content")
|
||||
else first_item.get("encrypted_content")
|
||||
)
|
||||
assert wrapped_content.startswith("litellm_enc:"), (
|
||||
f"Expected wrapped content but got {wrapped_content[:50]}..."
|
||||
)
|
||||
|
||||
# Verify we can extract model_id from wrapped content
|
||||
extracted_model_id, _ = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(
|
||||
wrapped_content
|
||||
)
|
||||
)
|
||||
assert extracted_model_id == first_model_id
|
||||
|
||||
# Second request: use wrapped encrypted_content WITHOUT an ID (Codex behavior)
|
||||
second_response = await router.aresponses(
|
||||
model="openai.gpt-5.1-codex",
|
||||
input=[
|
||||
{
|
||||
"type": "reasoning",
|
||||
"encrypted_content": wrapped_content,
|
||||
},
|
||||
],
|
||||
)
|
||||
second_model_id = second_response._hidden_params["model_id"]
|
||||
|
||||
assert second_model_id == first_model_id, (
|
||||
f"Expected affinity to route to {first_model_id}, but got {second_model_id}"
|
||||
)
|
||||
|
||||
|
||||
def test_encrypted_content_wrapping_preserves_original_content():
|
||||
"""
|
||||
Test that wrapping and unwrapping encrypted_content preserves the original content.
|
||||
This is critical for streaming responses where content must round-trip correctly.
|
||||
"""
|
||||
model_id = "test-deployment-1"
|
||||
original_encrypted_content = "gAAAAABpnW_yEYmSNEyOG_streaming_test_content_with_special_chars==+/"
|
||||
|
||||
wrapped = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
original_encrypted_content, model_id
|
||||
)
|
||||
|
||||
assert wrapped.startswith("litellm_enc:")
|
||||
assert wrapped != original_encrypted_content
|
||||
|
||||
extracted_model_id, unwrapped_content = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(wrapped)
|
||||
)
|
||||
|
||||
assert extracted_model_id == model_id
|
||||
assert unwrapped_content == original_encrypted_content
|
||||
|
||||
|
||||
def test_encrypted_content_wrapping_with_multiple_semicolons():
|
||||
"""
|
||||
Test that encrypted_content containing semicolons is handled correctly.
|
||||
"""
|
||||
model_id = "deployment-with-semicolons"
|
||||
original_content = "gAAAAAB;some;content;with;semicolons"
|
||||
|
||||
wrapped = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
original_content, model_id
|
||||
)
|
||||
|
||||
extracted_model_id, unwrapped = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(wrapped)
|
||||
)
|
||||
|
||||
assert extracted_model_id == model_id
|
||||
assert unwrapped == original_content
|
||||
|
||||
|
||||
def test_encrypted_content_wrapping_empty_string():
|
||||
"""
|
||||
Test that empty encrypted_content is handled gracefully.
|
||||
"""
|
||||
model_id = "test-deployment"
|
||||
original_content = ""
|
||||
|
||||
wrapped = ResponsesAPIRequestUtils._wrap_encrypted_content_with_model_id(
|
||||
original_content, model_id
|
||||
)
|
||||
|
||||
assert wrapped.startswith("litellm_enc:")
|
||||
|
||||
extracted_model_id, unwrapped = (
|
||||
ResponsesAPIRequestUtils._unwrap_encrypted_content_with_model_id(wrapped)
|
||||
)
|
||||
|
||||
assert extracted_model_id == model_id
|
||||
assert unwrapped == original_content
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue