mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
Response cache reads and writes open cache.get llm_response and cache.set llm_response phase spans with their Redis spans nested underneath, on the Python path and on the native Rust path, and deployment selection runs inside a route {model_group} phase so the cooldown, usage and model-id reads the router issues nest under it before chat {model}. The autorouter classifier call nests under that route phase as well and carries its typed internal origin on litellm.request.purpose, so it is told apart from the provider attempt. Service spans are named {service}.{verb} {target} from a low-cardinality key family the producer declares (llm_response, auth_objects, spend_counters, router_cooldowns, claude_code_session_router_binding, rate_limits, pod_lock, budget_reset, ...) instead of the raw method or a per-request pipeline length; a pipeline flush is targeted by the one family its ops share or by mixed with the sorted families on litellm.redis.families, a batch op keeps the family it was declared under whichever pipeline or standalone read settles it, and the ambient family labels Redis spans only, never the DB write-back a task spawned inside that context performs later. The raw method stays on litellm.service.call_type and on the Prometheus and Datadog labels. Caller attribution is carried across asyncio task boundaries on a ContextVar so forwarder-only chains no longer surface, the raw cache key is dropped from Redis span metadata, pipeline op counts land as an integer attribute, every call_type the Redis cache layer emits maps to a verb, and a scan over litellm/ and enterprise/ fails when a Redis producer, batch reservation included, declares no key family.
A V2 logger built for a key or team logging entry while the operator's V2 logger is already registered keeps only the exporters its own preset contributed, whether or not the operator holds credentials for that backend, so every chat span no longer reaches the operator's collector twice. A span the success callback has to open itself, with no pre-call carrier, starts at the provider handoff (api_call_start_time) instead of the logging object's creation.
Co-authored-by: yassin <yassin@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
428 lines
17 KiB
Python
428 lines
17 KiB
Python
import asyncio
|
|
from collections.abc import Callable, Coroutine
|
|
from datetime import datetime, timedelta
|
|
from typing import TYPE_CHECKING, Any, Final, Protocol
|
|
|
|
import litellm
|
|
from litellm._internal_context import current_service_target
|
|
from litellm._logging import verbose_logger
|
|
|
|
from .integrations.custom_logger import CustomLogger
|
|
from .integrations.datadog.datadog import DataDogLogger
|
|
from .integrations.opentelemetry import OpenTelemetry
|
|
from .integrations.prometheus_services import PrometheusServicesLogger
|
|
from .types.services import ServiceLoggerPayload, ServiceTypes
|
|
|
|
if TYPE_CHECKING:
|
|
from opentelemetry.trace import Span as _Span
|
|
|
|
from litellm.proxy._types import UserAPIKeyAuth
|
|
|
|
Span = _Span | Any
|
|
OTELClass = OpenTelemetry
|
|
else:
|
|
Span = Any
|
|
OTELClass = Any
|
|
UserAPIKeyAuth = Any
|
|
|
|
|
|
class _ServiceSpanLogger(Protocol):
|
|
"""The OTel logger surface this module drives: the two service-span hooks it calls."""
|
|
|
|
async def async_service_success_hook(
|
|
self,
|
|
payload: ServiceLoggerPayload,
|
|
parent_otel_span: Span | None = None,
|
|
start_time: datetime | float | None = None,
|
|
end_time: datetime | float | None = None,
|
|
event_metadata: dict | None = None,
|
|
) -> None: ...
|
|
|
|
async def async_service_failure_hook(
|
|
self,
|
|
payload: ServiceLoggerPayload,
|
|
error: str | None = "",
|
|
parent_otel_span: Span | None = None,
|
|
start_time: datetime | float | None = None,
|
|
end_time: datetime | float | None = None,
|
|
event_metadata: dict | None = None,
|
|
) -> None: ...
|
|
|
|
|
|
def _get_otel_v2_class() -> type[_ServiceSpanLogger] | None:
|
|
"""Return the ``OpenTelemetryV2`` class, or ``None`` if the OTel SDK is absent.
|
|
|
|
Imported lazily: ``litellm.integrations.otel.logger`` imports the OpenTelemetry
|
|
SDK at module scope, so importing it eagerly would break installs without the
|
|
SDK. The V2 logger only exists when ``LITELLM_OTEL_V2`` is enabled (which
|
|
requires the SDK), so a failed import simply means "no V2 logger in play".
|
|
"""
|
|
try:
|
|
from litellm.integrations.otel.logger import OpenTelemetryV2
|
|
|
|
return OpenTelemetryV2
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
class ServiceLogging(CustomLogger):
|
|
"""
|
|
Separate class used for monitoring health of litellm-adjacent services (redis/postgres).
|
|
"""
|
|
|
|
def __init__(self, mock_testing: bool = False) -> None:
|
|
self.mock_testing = mock_testing
|
|
self.mock_testing_sync_success_hook = 0
|
|
self.mock_testing_async_success_hook = 0
|
|
self.mock_testing_sync_failure_hook = 0
|
|
self.mock_testing_async_failure_hook = 0
|
|
if "prometheus_system" in litellm.service_callback:
|
|
self.prometheusServicesLogger = PrometheusServicesLogger()
|
|
|
|
def _resolve_otel_service_logger(self, callback: object) -> _ServiceSpanLogger | None:
|
|
"""Resolve the OTel logger (legacy or V2) to emit a service span on.
|
|
|
|
Returns the logger instance whose ``async_service_*_hook`` should fire for
|
|
this ``callback``, or ``None`` when ``callback`` is not an OTel callback.
|
|
|
|
The V2 ``OpenTelemetryV2`` logger is a plain ``CustomLogger`` and is NOT a
|
|
subclass of the legacy ``OpenTelemetry``, so the legacy ``isinstance``
|
|
check alone misses it — which is why redis/postgres service spans never
|
|
showed up under ``LITELLM_OTEL_V2``. Match both the legacy and V2 types,
|
|
whether the callback is the logger instance itself or the ``"otel"`` string
|
|
(which routes to the proxy's registered ``open_telemetry_logger``).
|
|
"""
|
|
otel_v2_cls: Final = _get_otel_v2_class()
|
|
|
|
def _as_otel_logger(obj: object) -> _ServiceSpanLogger | None:
|
|
if isinstance(obj, OpenTelemetry):
|
|
return obj
|
|
if otel_v2_cls is not None and isinstance(obj, otel_v2_cls):
|
|
return obj
|
|
return None
|
|
|
|
resolved_callback: Final = _as_otel_logger(callback)
|
|
if resolved_callback is not None:
|
|
return resolved_callback
|
|
if callback == "otel":
|
|
from litellm.proxy.proxy_server import open_telemetry_logger
|
|
|
|
if open_telemetry_logger is not None:
|
|
return _as_otel_logger(open_telemetry_logger)
|
|
return None
|
|
|
|
@staticmethod
|
|
def _sync_dispatch_loop() -> asyncio.AbstractEventLoop | None:
|
|
"""The event loop a blocking caller can dispatch on, or ``None`` if it has none."""
|
|
try:
|
|
loop: Final = asyncio.get_event_loop()
|
|
except RuntimeError:
|
|
return None
|
|
return None if loop.is_closed() else loop
|
|
|
|
@staticmethod
|
|
async def _emit_guarded(hook: Callable[[], Coroutine[object, object, None]]) -> None:
|
|
"""Emit one service event, absorbing anything the callbacks raise.
|
|
|
|
Monitoring must not break the call it monitors. Sync callers are the ones that
|
|
swallow their own service failures (a Redis batch read returns an empty dict),
|
|
so an exception from a misconfigured callback would replace a Redis outage with
|
|
a callback error and skip the caller's fallback handling.
|
|
"""
|
|
try:
|
|
await hook()
|
|
except Exception as e:
|
|
verbose_logger.exception("Error emitting service event - %s", e)
|
|
|
|
@staticmethod
|
|
def _dispatch_from_sync(hook: Callable[[], Coroutine[object, object, None]]) -> None:
|
|
"""Run an async service hook from a blocking caller, whatever event loop it holds.
|
|
|
|
Takes a factory rather than a coroutine so the hook is built on the path that
|
|
runs it, and only ever once.
|
|
"""
|
|
loop: Final = ServiceLogging._sync_dispatch_loop()
|
|
try:
|
|
if loop is None:
|
|
asyncio.run(ServiceLogging._emit_guarded(hook))
|
|
elif loop.is_running():
|
|
loop.create_task(ServiceLogging._emit_guarded(hook))
|
|
else:
|
|
loop.run_until_complete(ServiceLogging._emit_guarded(hook))
|
|
except Exception as e:
|
|
verbose_logger.exception("Error dispatching service event - %s", e)
|
|
|
|
def service_success_hook(
|
|
self,
|
|
service: ServiceTypes,
|
|
duration: float,
|
|
call_type: str,
|
|
parent_otel_span: Span | None = None,
|
|
start_time: datetime | float | None = None,
|
|
end_time: float | datetime | None = None,
|
|
caller: str | None = None,
|
|
):
|
|
"""
|
|
Handles both sync and async monitoring by checking for existing event loop.
|
|
"""
|
|
|
|
if self.mock_testing:
|
|
self.mock_testing_sync_success_hook += 1
|
|
|
|
self._dispatch_from_sync(
|
|
lambda: self.async_service_success_hook(
|
|
service=service,
|
|
duration=duration,
|
|
call_type=call_type,
|
|
caller=caller,
|
|
parent_otel_span=parent_otel_span,
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
)
|
|
)
|
|
|
|
def service_failure_hook(
|
|
self,
|
|
service: ServiceTypes,
|
|
duration: float,
|
|
error: Exception,
|
|
call_type: str,
|
|
parent_otel_span: Span | None = None,
|
|
start_time: datetime | float | None = None,
|
|
end_time: float | datetime | None = None,
|
|
caller: str | None = None,
|
|
):
|
|
"""
|
|
Handles both sync and async monitoring by checking for existing event loop.
|
|
"""
|
|
if self.mock_testing:
|
|
self.mock_testing_sync_failure_hook += 1
|
|
|
|
self._dispatch_from_sync(
|
|
lambda: self.async_service_failure_hook(
|
|
service=service,
|
|
duration=duration,
|
|
error=error,
|
|
call_type=call_type,
|
|
caller=caller,
|
|
parent_otel_span=parent_otel_span,
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
)
|
|
)
|
|
|
|
async def async_service_success_hook(
|
|
self,
|
|
service: ServiceTypes,
|
|
call_type: str,
|
|
duration: float,
|
|
parent_otel_span: Span | None = None,
|
|
start_time: datetime | float | None = None,
|
|
end_time: datetime | float | None = None,
|
|
event_metadata: dict | None = None,
|
|
caller: str | None = None,
|
|
):
|
|
"""
|
|
- For counting if the redis, postgres call is successful
|
|
"""
|
|
if self.mock_testing:
|
|
self.mock_testing_async_success_hook += 1
|
|
|
|
payload: Final = ServiceLoggerPayload(
|
|
is_error=False,
|
|
error=None,
|
|
service=service,
|
|
duration=duration,
|
|
call_type=call_type,
|
|
caller=caller,
|
|
target=current_service_target() if service == ServiceTypes.REDIS else None,
|
|
event_metadata=event_metadata,
|
|
)
|
|
|
|
# OTel loggers already fired this event. ``service_callback`` can hold more
|
|
# than one reference that resolves to the *same* logger — the ``"otel"``
|
|
# string AND the registered instance both map to ``open_telemetry_logger``
|
|
# (the V2 logger self-registers its instance even when the string is
|
|
# present, unlike V1). Without this guard each such reference emits its own
|
|
# span, so a single DB call shows up as duplicate ``postgres ...`` spans.
|
|
emitted_otel_logger_ids: Final[set] = set()
|
|
for callback in litellm.service_callback:
|
|
if callback == "prometheus_system":
|
|
await self.init_prometheus_services_logger_if_none()
|
|
await self.prometheusServicesLogger.async_service_success_hook(payload=payload)
|
|
elif callback == "datadog" or isinstance(callback, DataDogLogger):
|
|
await self.init_datadog_logger_if_none()
|
|
await self.dd_logger.async_service_success_hook(
|
|
payload=payload,
|
|
parent_otel_span=parent_otel_span,
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
event_metadata=event_metadata,
|
|
)
|
|
else:
|
|
_otel_logger_to_use = self._resolve_otel_service_logger(callback)
|
|
# No ``parent_otel_span is not None`` gate: a background service
|
|
# call (no request on the stack) has no parent, and dropping it
|
|
# here is what hid those calls from traces entirely. The OTel
|
|
# logger decides what to do with a missing parent — legacy V1
|
|
# no-ops, V2 emits a root span (and skips metrics-only pings).
|
|
if _otel_logger_to_use is not None and id(_otel_logger_to_use) not in emitted_otel_logger_ids:
|
|
emitted_otel_logger_ids.add(id(_otel_logger_to_use))
|
|
await _otel_logger_to_use.async_service_success_hook(
|
|
payload=payload,
|
|
parent_otel_span=parent_otel_span,
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
event_metadata=event_metadata,
|
|
)
|
|
|
|
async def init_prometheus_services_logger_if_none(self):
|
|
"""
|
|
initializes prometheusServicesLogger if it is None or no attribute exists on ServiceLogging Object
|
|
|
|
"""
|
|
if not hasattr(self, "prometheusServicesLogger"):
|
|
self.prometheusServicesLogger = PrometheusServicesLogger()
|
|
elif self.prometheusServicesLogger is None:
|
|
self.prometheusServicesLogger = self.prometheusServicesLogger()
|
|
|
|
async def init_datadog_logger_if_none(self):
|
|
"""
|
|
initializes dd_logger if it is None or no attribute exists on ServiceLogging Object
|
|
|
|
"""
|
|
from litellm.integrations.datadog.datadog import DataDogLogger
|
|
|
|
if not hasattr(self, "dd_logger"):
|
|
self.dd_logger: DataDogLogger = DataDogLogger()
|
|
|
|
async def init_otel_logger_if_none(self):
|
|
"""
|
|
initializes otel_logger if it is None or no attribute exists on ServiceLogging Object
|
|
|
|
"""
|
|
from litellm.proxy.proxy_server import open_telemetry_logger
|
|
|
|
if not hasattr(self, "otel_logger"):
|
|
if open_telemetry_logger is not None and isinstance(open_telemetry_logger, OpenTelemetry):
|
|
self.otel_logger: OpenTelemetry = open_telemetry_logger
|
|
else:
|
|
verbose_logger.warning(
|
|
"ServiceLogger: open_telemetry_logger is None or not an instance of OpenTelemetry"
|
|
)
|
|
|
|
async def async_service_failure_hook(
|
|
self,
|
|
service: ServiceTypes,
|
|
duration: float,
|
|
error: str | Exception,
|
|
call_type: str,
|
|
parent_otel_span: Span | None = None,
|
|
start_time: datetime | float | None = None,
|
|
end_time: float | datetime | None = None,
|
|
event_metadata: dict | None = None,
|
|
caller: str | None = None,
|
|
):
|
|
"""
|
|
- For counting if the redis, postgres call is unsuccessful
|
|
"""
|
|
if self.mock_testing:
|
|
self.mock_testing_async_failure_hook += 1
|
|
|
|
error_message = ""
|
|
if isinstance(error, Exception):
|
|
error_message = str(error)
|
|
elif isinstance(error, str):
|
|
error_message = error
|
|
|
|
payload: Final = ServiceLoggerPayload(
|
|
is_error=True,
|
|
error=error_message,
|
|
service=service,
|
|
duration=duration,
|
|
call_type=call_type,
|
|
caller=caller,
|
|
target=current_service_target() if service == ServiceTypes.REDIS else None,
|
|
event_metadata=event_metadata,
|
|
)
|
|
|
|
# Dedupe OTel loggers per event — see ``async_service_success_hook`` for why
|
|
# the same logger can be referenced twice in ``service_callback``.
|
|
emitted_otel_logger_ids: Final[set] = set()
|
|
for callback in litellm.service_callback:
|
|
if callback == "prometheus_system":
|
|
await self.init_prometheus_services_logger_if_none()
|
|
await self.prometheusServicesLogger.async_service_failure_hook(
|
|
payload=payload,
|
|
error=error,
|
|
)
|
|
elif callback == "datadog" or isinstance(callback, DataDogLogger):
|
|
await self.init_datadog_logger_if_none()
|
|
await self.dd_logger.async_service_failure_hook(
|
|
payload=payload,
|
|
error=error_message,
|
|
parent_otel_span=parent_otel_span,
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
event_metadata=event_metadata,
|
|
)
|
|
else:
|
|
_otel_logger_to_use = self._resolve_otel_service_logger(callback)
|
|
|
|
if not isinstance(error, str):
|
|
error = str(error)
|
|
|
|
# See the success hook: no parent gate, so background failures
|
|
# are traced too. V1 no-ops without a parent; V2 emits a root.
|
|
if _otel_logger_to_use is not None and id(_otel_logger_to_use) not in emitted_otel_logger_ids:
|
|
emitted_otel_logger_ids.add(id(_otel_logger_to_use))
|
|
await _otel_logger_to_use.async_service_failure_hook(
|
|
payload=payload,
|
|
error=error,
|
|
parent_otel_span=parent_otel_span,
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
event_metadata=event_metadata,
|
|
)
|
|
|
|
async def async_post_call_failure_hook(
|
|
self,
|
|
request_data: dict,
|
|
original_exception: Exception,
|
|
user_api_key_dict: UserAPIKeyAuth,
|
|
traceback_str: str | None = None,
|
|
):
|
|
"""
|
|
Hook to track failed litellm-service calls
|
|
"""
|
|
return await super().async_post_call_failure_hook(
|
|
request_data,
|
|
original_exception,
|
|
user_api_key_dict,
|
|
)
|
|
|
|
async def async_log_success_event(self, kwargs, response_obj, start_time, end_time):
|
|
"""
|
|
Hook to track latency for litellm proxy llm api calls
|
|
"""
|
|
try:
|
|
_duration = end_time - start_time
|
|
if isinstance(_duration, timedelta):
|
|
_duration = _duration.total_seconds()
|
|
elif isinstance(_duration, float):
|
|
pass
|
|
else:
|
|
raise Exception(
|
|
f"Duration={_duration} is not a float or timedelta object. type={type(_duration)}"
|
|
) # invalid _duration value
|
|
# Batch polling callbacks (check_batch_cost) don't include call_type in kwargs.
|
|
# Use .get() to avoid KeyError.
|
|
await self.async_service_success_hook(
|
|
service=ServiceTypes.LITELLM,
|
|
duration=_duration,
|
|
call_type=kwargs.get("call_type", "unknown"),
|
|
start_time=start_time,
|
|
end_time=end_time,
|
|
)
|
|
except Exception as e:
|
|
raise e
|