Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/native_compaction.py: 53%
43 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1import asyncio
2from collections.abc import Awaitable, Mapping
3from contextvars import Context
4from types import MappingProxyType
5from typing import Final, Literal, TypeVar
7from fastapi import Request
8from pydantic import TypeAdapter, ValidationError
9from starlette.types import ASGIApp
11from litellm.exceptions import BadRequestError
12from litellm.litellm_core_utils.initialize_dynamic_callback_params import (
13 inherit_message_logging_privacy,
14 initialize_standard_callback_dynamic_params,
15)
16from litellm.llms.custom_httpx.asgi_handler import get_async_asgi_client
17from litellm.proxy.litellm_pre_call_utils import UNTRUSTED_REQUEST_HEADER_CONTROL_FIELDS
18from litellm.router_strategy.complexity_router.context_compaction import (
19 compaction_executor,
20 native_compaction_call,
21)
23_ResultT: Final = TypeVar("_ResultT")
24_ASGI_APP: Final = TypeAdapter[ASGIApp](ASGIApp)
25_ROOT_PATH: Final = TypeAdapter(str)
26_JSON_OBJECT: Final = TypeAdapter(Mapping[str, object])
27_REMOVED_HEADERS: Final = frozenset(
28 (
29 b"content-length",
30 b"x-litellm-call-id",
31 b"x-litellm-num-retries",
32 b"x-litellm-timeout",
33 b"x-litellm-stream-timeout",
34 )
35)
38async def with_proxy_compaction_executor(call: Awaitable[_ResultT], request: Request) -> _ResultT:
39 async def execute(
40 protocol: Literal["chat", "messages"], payload: Mapping[str, object], parent_model: str | None = None
41 ) -> Mapping[str, object]:
42 logging_disabled: Final = initialize_standard_callback_dynamic_params().get("turn_off_message_logging") is True
44 async def dispatch() -> Mapping[str, object]:
45 scope: Final = _JSON_OBJECT.validate_python(request.scope)
46 root_path: Final = _ROOT_PATH.validate_python(scope.get("root_path", ""))
47 path: Final = "/v1/chat/completions" if protocol == "chat" else "/v1/messages"
48 url: Final = str(request.url.replace(path=root_path.rstrip("/") + path, query="", fragment=""))
49 headers: Final = tuple(
50 (name, value)
51 for name, value in request.headers.raw
52 if name.lower() not in _REMOVED_HEADERS
53 and not (logging_disabled and name.decode("latin-1").lower() in UNTRUSTED_REQUEST_HEADER_CONTROL_FIELDS)
54 )
55 with (
56 native_compaction_call(parent_model, str(payload["model"])),
57 inherit_message_logging_privacy(logging_disabled),
58 ):
59 with get_async_asgi_client(
60 app=_ASGI_APP.validate_python(scope["app"]),
61 root_path=root_path,
62 client=request.client,
63 ) as client:
64 async with client.stream(
65 "POST", url, headers=headers, json=_JSON_OBJECT.validate_python(payload)
66 ) as response:
67 if not response.is_success:
68 raise BadRequestError(
69 message=f"Native compaction child request failed (HTTP {response.status_code})",
70 model="context_compaction",
71 llm_provider="",
72 )
73 body: Final = await response.aread()
74 try:
75 return MappingProxyType(_JSON_OBJECT.validate_json(body))
76 except ValidationError:
77 raise BadRequestError(
78 message="Native compaction child returned an invalid JSON object",
79 model="context_compaction",
80 llm_provider="",
81 ) from None
83 task: Final = Context().run(asyncio.create_task, dispatch())
84 return await task
86 token: Final = compaction_executor.set(execute)
87 try:
88 return await call
89 finally:
90 compaction_executor.reset(token)