mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
fix(anthropic adapter): correct sync streaming, surface polyfill failures, decouple sync path from proxy router
- translate_completion_output_params_streaming: add is_async flag so the sync handler returns Iterator[bytes] instead of an unusable AsyncIterator. Async callers keep the existing behavior via the default is_async=True. - _run_polyfill_if_enabled: when the polyfill crashes and the spec requested non-compact edits (e.g. clear_tool_uses_20250919), raise an AnthropicContextManagementError instead of silently returning None so those edits are not dropped without an error surface. The compaction-block-slicing safety net remains for compact-only specs. - anthropic_messages_handler (sync): stop auto-attaching the proxy llm_router. run_async_function bridges to a new thread's event loop; reusing the proxy's loop-bound httpx clients there causes 'Event loop is closed' errors. The summary editor falls back to litellm.acompletion when llm_router is None. Co-authored-by: Yassin Kortam <yassin@berri.ai>
This commit is contained in:
parent
108798b78f
commit
d50ea4325a
2 changed files with 102 additions and 24 deletions
|
|
@ -4,6 +4,7 @@ from typing import (
|
|||
AsyncIterator,
|
||||
Coroutine,
|
||||
Dict,
|
||||
Iterator,
|
||||
List,
|
||||
Optional,
|
||||
Tuple,
|
||||
|
|
@ -121,19 +122,69 @@ def _polyfill_will_run(
|
|||
skip only applies when the dispatcher will actually invoke
|
||||
``apply_compact_20260112`` (which has its own compaction-block slicing).
|
||||
"""
|
||||
if not context_management_spec:
|
||||
return False
|
||||
|
||||
effective_drop_params = (
|
||||
drop_params if drop_params is not None else litellm.drop_params
|
||||
edits = _normalize_spec_edits(
|
||||
context_management_spec=context_management_spec,
|
||||
drop_params=drop_params,
|
||||
)
|
||||
if effective_drop_params:
|
||||
if edits is None:
|
||||
return False
|
||||
|
||||
from litellm.llms.anthropic.experimental_pass_through.context_management.constants import (
|
||||
COMPACT_EDIT_TYPE,
|
||||
)
|
||||
|
||||
return any(
|
||||
isinstance(edit, dict) and edit.get("type") == COMPACT_EDIT_TYPE
|
||||
for edit in edits
|
||||
)
|
||||
|
||||
|
||||
def _spec_has_non_compact_edits(
|
||||
*,
|
||||
context_management_spec: Any,
|
||||
drop_params: Optional[bool],
|
||||
) -> bool:
|
||||
"""Return True when the spec includes edits other than ``compact_20260112``.
|
||||
|
||||
Used to decide whether a polyfill failure can be silently swallowed
|
||||
(compact-only specs have a safe compaction-block slicing fallback) or
|
||||
must be surfaced (other editors like ``clear_tool_uses_20250919`` have
|
||||
no slice-only fallback and would otherwise be dropped without notice).
|
||||
"""
|
||||
edits = _normalize_spec_edits(
|
||||
context_management_spec=context_management_spec,
|
||||
drop_params=drop_params,
|
||||
)
|
||||
if edits is None:
|
||||
return False
|
||||
|
||||
from litellm.llms.anthropic.experimental_pass_through.context_management.constants import (
|
||||
COMPACT_EDIT_TYPE,
|
||||
)
|
||||
|
||||
return any(
|
||||
isinstance(edit, dict)
|
||||
and isinstance(edit.get("type"), str)
|
||||
and edit.get("type") != COMPACT_EDIT_TYPE
|
||||
for edit in edits
|
||||
)
|
||||
|
||||
|
||||
def _normalize_spec_edits(
|
||||
*,
|
||||
context_management_spec: Any,
|
||||
drop_params: Optional[bool],
|
||||
) -> Optional[List[Dict[str, Any]]]:
|
||||
"""Return the normalized ``edits`` list, or ``None`` if the polyfill won't run."""
|
||||
if not context_management_spec:
|
||||
return None
|
||||
|
||||
effective_drop_params = (
|
||||
drop_params if drop_params is not None else litellm.drop_params
|
||||
)
|
||||
if effective_drop_params:
|
||||
return None
|
||||
|
||||
spec = context_management_spec
|
||||
if isinstance(spec, list):
|
||||
try:
|
||||
|
|
@ -141,16 +192,12 @@ def _polyfill_will_run(
|
|||
|
||||
spec = AnthropicConfig.map_openai_context_management_to_anthropic(spec)
|
||||
except Exception:
|
||||
return False
|
||||
return None
|
||||
|
||||
edits = spec.get("edits") if isinstance(spec, dict) else None
|
||||
if not isinstance(edits, list):
|
||||
return False
|
||||
|
||||
return any(
|
||||
isinstance(edit, dict) and edit.get("type") == COMPACT_EDIT_TYPE
|
||||
for edit in edits
|
||||
)
|
||||
return None
|
||||
return edits
|
||||
|
||||
|
||||
async def _run_polyfill_if_enabled(
|
||||
|
|
@ -198,6 +245,21 @@ async def _run_polyfill_if_enabled(
|
|||
verbose_logger.exception(
|
||||
"context_management polyfill: skipping edits due to error: %s", e
|
||||
)
|
||||
# Best-effort swallow is only safe for compact-only specs, where the
|
||||
# caller's compaction-block-slicing safety net produces a correct
|
||||
# (if degraded) result. When the spec also requested non-compact
|
||||
# edits (e.g. ``clear_tool_uses_20250919``), the safety net does
|
||||
# NOT re-run those editors, so silently returning ``None`` would
|
||||
# drop them with no error surface. Raise instead so the endpoint
|
||||
# emits an Anthropic-format error.
|
||||
if _spec_has_non_compact_edits(
|
||||
context_management_spec=context_management_spec,
|
||||
drop_params=drop_params,
|
||||
):
|
||||
raise AnthropicContextManagementError(
|
||||
status_code=500,
|
||||
message=f"context_management polyfill failed: {e}",
|
||||
) from e
|
||||
return None
|
||||
|
||||
|
||||
|
|
@ -532,6 +594,7 @@ class LiteLLMMessagesToCompletionTransformationHandler:
|
|||
model=model,
|
||||
tool_name_mapping=tool_name_mapping,
|
||||
polyfill_result=polyfill_result,
|
||||
is_async=True,
|
||||
)
|
||||
)
|
||||
if transformed_stream is not None:
|
||||
|
|
@ -567,7 +630,7 @@ class LiteLLMMessagesToCompletionTransformationHandler:
|
|||
**kwargs,
|
||||
) -> Union[
|
||||
AnthropicMessagesResponse,
|
||||
AsyncIterator[Any],
|
||||
Iterator[bytes],
|
||||
Coroutine[Any, Any, Union[AnthropicMessagesResponse, AsyncIterator[Any]]],
|
||||
]:
|
||||
"""Handle non-Anthropic models using the adapter."""
|
||||
|
|
@ -597,14 +660,18 @@ class LiteLLMMessagesToCompletionTransformationHandler:
|
|||
# bridge to it via ``run_async_function``.
|
||||
context_management = kwargs.pop("context_management", None)
|
||||
drop_params: Optional[bool] = kwargs.get("drop_params", None)
|
||||
# Deliberately do NOT auto-attach the proxy ``llm_router`` here:
|
||||
# ``run_async_function`` spawns a new event loop in a worker thread
|
||||
# to bridge to the async dispatcher, but the proxy router's httpx
|
||||
# ``AsyncClient`` instances are bound to the proxy's main event loop.
|
||||
# Reusing them from the new thread's loop violates httpx's single-loop
|
||||
# invariant and can raise ``RuntimeError: Event loop is closed`` or
|
||||
# produce stalled connections. The summary editor falls back to
|
||||
# ``litellm.acompletion`` (which creates a fresh client per call) when
|
||||
# ``llm_router`` is ``None``, which is safe to call from the bridged
|
||||
# loop. The async ``async_anthropic_messages_handler`` path is
|
||||
# unaffected because it ``await``s within the original event loop.
|
||||
litellm_router = kwargs.pop("litellm_router", None)
|
||||
if litellm_router is None:
|
||||
try:
|
||||
from litellm.proxy.proxy_server import llm_router as _proxy_router
|
||||
|
||||
litellm_router = _proxy_router
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
polyfill_result = run_async_function(
|
||||
_prepare_context_managed_request,
|
||||
|
|
@ -655,6 +722,7 @@ class LiteLLMMessagesToCompletionTransformationHandler:
|
|||
model=model,
|
||||
tool_name_mapping=tool_name_mapping,
|
||||
polyfill_result=polyfill_result,
|
||||
is_async=False,
|
||||
)
|
||||
)
|
||||
if transformed_stream is not None:
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ from typing import (
|
|||
Any,
|
||||
AsyncIterator,
|
||||
Dict,
|
||||
Iterator,
|
||||
List,
|
||||
Literal,
|
||||
Optional,
|
||||
|
|
@ -225,7 +226,8 @@ class AnthropicAdapter:
|
|||
model: str,
|
||||
tool_name_mapping: Optional[Dict[str, str]] = None,
|
||||
polyfill_result: Optional[PolyfillResult] = None,
|
||||
) -> Union[AsyncIterator[bytes], None]:
|
||||
is_async: bool = True,
|
||||
) -> Union[AsyncIterator[bytes], Iterator[bytes], None]:
|
||||
"""
|
||||
Translate OpenAI streaming response to Anthropic format.
|
||||
|
||||
|
|
@ -234,6 +236,12 @@ class AnthropicAdapter:
|
|||
model: The model name
|
||||
tool_name_mapping: Optional mapping of truncated tool names to original names.
|
||||
polyfill_result: PolyfillResult from context_management polyfill.
|
||||
is_async: When ``True`` (default, for back-compat with existing
|
||||
async callers) returns an ``AsyncIterator[bytes]``. When
|
||||
``False`` returns a sync ``Iterator[bytes]`` so sync callers
|
||||
(e.g. ``litellm.anthropic.messages.create(stream=True)`` via
|
||||
the sync handler) don't get back an async iterator they
|
||||
can't iterate without an event loop.
|
||||
"""
|
||||
applied_edits = (
|
||||
polyfill_result.applied_edits_for_response() if polyfill_result else None
|
||||
|
|
@ -252,8 +260,10 @@ class AnthropicAdapter:
|
|||
compaction_block=compaction_block,
|
||||
iterations_usage=iterations_usage,
|
||||
)
|
||||
# Return the SSE-wrapped version for proper event formatting
|
||||
return anthropic_wrapper.async_anthropic_sse_wrapper()
|
||||
# Return the SSE-wrapped version for proper event formatting.
|
||||
if is_async:
|
||||
return anthropic_wrapper.async_anthropic_sse_wrapper()
|
||||
return anthropic_wrapper.anthropic_sse_wrapper()
|
||||
|
||||
|
||||
class LiteLLMAnthropicMessagesAdapter:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue