diff --git a/litellm/integrations/datadog/datadog.py b/litellm/integrations/datadog/datadog.py index 21e1d562224..08bf5034f0e 100644 --- a/litellm/integrations/datadog/datadog.py +++ b/litellm/integrations/datadog/datadog.py @@ -27,6 +27,12 @@ import litellm from litellm._logging import verbose_logger from litellm._uuid import uuid from litellm.integrations.custom_batch_logger import CustomBatchLogger +from litellm.integrations.datadog.datadog_handler import ( + get_datadog_hostname, + get_datadog_service, + get_datadog_source, + get_datadog_tags, +) from litellm.llms.custom_httpx.http_handler import ( _get_httpx_client, get_async_httpx_client, @@ -67,23 +73,23 @@ class DataDogLogger( Optional environment variables (DataDog Agent): `LITELLM_DD_AGENT_HOST` - hostname or IP of DataDog agent, example = `"localhost"` `LITELLM_DD_AGENT_PORT` - port of DataDog agent (default: 10518 for logs) - + Note: We use LITELLM_DD_AGENT_HOST instead of DD_AGENT_HOST to avoid conflicts with ddtrace which automatically sets DD_AGENT_HOST for APM tracing. """ try: verbose_logger.debug("Datadog: in init datadog logger") - + ######################################################### # Handle datadog_params set as litellm.datadog_params ######################################################### dict_datadog_params = self._get_datadog_params() kwargs.update(dict_datadog_params) - + self.async_client = get_async_httpx_client( llm_provider=httpxSpecialProvider.LoggingCallback ) - + # Configure DataDog endpoint (Agent or Direct API) # Use LITELLM_DD_AGENT_HOST to avoid conflicts with ddtrace's DD_AGENT_HOST dd_agent_host = os.getenv("LITELLM_DD_AGENT_HOST") @@ -91,7 +97,7 @@ class DataDogLogger( self._configure_dd_agent(dd_agent_host=dd_agent_host) else: self._configure_dd_direct_api() - + # Optional override for testing self._apply_dd_base_url_override() self.sync_client = _get_httpx_client() @@ -118,17 +124,21 @@ class DataDogLogger( dict_datadog_params = litellm.datadog_params.model_dump() elif isinstance(litellm.datadog_params, Dict): # only allow params that are of DatadogInitParams - dict_datadog_params = DatadogInitParams(**litellm.datadog_params).model_dump() + dict_datadog_params = DatadogInitParams( + **litellm.datadog_params + ).model_dump() return dict_datadog_params def _configure_dd_agent(self, dd_agent_host: str) -> None: """ Configure DataDog Agent for log forwarding - + Args: dd_agent_host: Hostname or IP of DataDog agent """ - dd_agent_port = os.getenv("LITELLM_DD_AGENT_PORT", "10518") # default port for logs + dd_agent_port = os.getenv( + "LITELLM_DD_AGENT_PORT", "10518" + ) # default port for logs self.intake_url = f"http://{dd_agent_host}:{dd_agent_port}/api/v2/logs" self.DD_API_KEY = os.getenv("DD_API_KEY") # Optional when using agent verbose_logger.debug(f"Datadog: Using DD Agent at {self.intake_url}") @@ -136,7 +146,7 @@ class DataDogLogger( def _configure_dd_direct_api(self) -> None: """ Configure direct DataDog API connection - + Raises: Exception: If required environment variables are not set """ @@ -144,11 +154,9 @@ class DataDogLogger( raise Exception("DD_API_KEY is not set, set 'DD_API_KEY=<>") if os.getenv("DD_SITE", None) is None: raise Exception("DD_SITE is not set in .env, set 'DD_SITE=<>") - + self.DD_API_KEY = os.getenv("DD_API_KEY") - self.intake_url = ( - f"https://http-intake.logs.{os.getenv('DD_SITE')}/api/v2/logs" - ) + self.intake_url = f"https://http-intake.logs.{os.getenv('DD_SITE')}/api/v2/logs" def _apply_dd_base_url_override(self) -> None: """ @@ -270,7 +278,7 @@ class DataDogLogger( # Add API key if available (required for direct API, optional for agent) if self.DD_API_KEY: headers["DD-API-KEY"] = self.DD_API_KEY - + response = self.sync_client.post( url=self.intake_url, json=dd_payload, # type: ignore @@ -318,16 +326,15 @@ class DataDogLogger( status: DataDogStatus, ) -> DatadogPayload: from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + json_payload = safe_dumps(standard_logging_object) verbose_logger.debug("Datadog: Logger - Logging payload = %s", json_payload) dd_payload = DatadogPayload( - ddsource=self._get_datadog_source(), - ddtags=self._get_datadog_tags( - standard_logging_object=standard_logging_object - ), - hostname=self._get_datadog_hostname(), + ddsource=get_datadog_source(), + ddtags=get_datadog_tags(standard_logging_object=standard_logging_object), + hostname=get_datadog_hostname(), message=json_payload, - service=self._get_datadog_service(), + service=get_datadog_service(), status=status, ) return dd_payload @@ -384,18 +391,19 @@ class DataDogLogger( import gzip from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + compressed_data = gzip.compress(safe_dumps(data).encode("utf-8")) - + # Build headers headers = { "Content-Encoding": "gzip", "Content-Type": "application/json", } - + # Add API key if available (required for direct API, optional for agent) if self.DD_API_KEY: headers["DD-API-KEY"] = self.DD_API_KEY - + response = await self.async_client.post( url=self.intake_url, data=compressed_data, # type: ignore @@ -421,13 +429,14 @@ class DataDogLogger( _payload_dict = payload.model_dump() _payload_dict.update(event_metadata or {}) from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + _dd_message_str = safe_dumps(_payload_dict) _dd_payload = DatadogPayload( - ddsource=self._get_datadog_source(), - ddtags=self._get_datadog_tags(), - hostname=self._get_datadog_hostname(), + ddsource=get_datadog_source(), + ddtags=get_datadog_tags(), + hostname=get_datadog_hostname(), message=_dd_message_str, - service=self._get_datadog_service(), + service=get_datadog_service(), status=DataDogStatus.WARN, ) @@ -462,13 +471,14 @@ class DataDogLogger( _payload_dict.update(event_metadata or {}) from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + _dd_message_str = safe_dumps(_payload_dict) _dd_payload = DatadogPayload( - ddsource=self._get_datadog_source(), - ddtags=self._get_datadog_tags(), - hostname=self._get_datadog_hostname(), + ddsource=get_datadog_source(), + ddtags=get_datadog_tags(), + hostname=get_datadog_hostname(), message=_dd_message_str, - service=self._get_datadog_service(), + service=get_datadog_service(), status=DataDogStatus.INFO, ) @@ -530,7 +540,6 @@ class DataDogLogger( else: clean_metadata[key] = value - # Build the initial payload payload = { "id": id, @@ -550,69 +559,21 @@ class DataDogLogger( } from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + json_payload = safe_dumps(payload) verbose_logger.debug("Datadog: Logger - Logging payload = %s", json_payload) dd_payload = DatadogPayload( - ddsource=self._get_datadog_source(), - ddtags=self._get_datadog_tags(), - hostname=self._get_datadog_hostname(), + ddsource=get_datadog_source(), + ddtags=get_datadog_tags(), + hostname=get_datadog_hostname(), message=json_payload, - service=self._get_datadog_service(), + service=get_datadog_service(), status=DataDogStatus.INFO, ) return dd_payload - @staticmethod - def _get_datadog_tags( - standard_logging_object: Optional[StandardLoggingPayload] = None, - ) -> str: - """ - Get the datadog tags for the request - - DD tags need to be as follows: - - tags: ["user_handle:dog@gmail.com", "app_version:1.0.0"] - """ - base_tags = { - "env": os.getenv("DD_ENV", "unknown"), - "service": os.getenv("DD_SERVICE", "litellm"), - "version": os.getenv("DD_VERSION", "unknown"), - "HOSTNAME": DataDogLogger._get_datadog_hostname(), - "POD_NAME": os.getenv("POD_NAME", "unknown"), - } - - tags = [f"{k}:{v}" for k, v in base_tags.items()] - - if standard_logging_object: - _request_tags: List[str] = ( - standard_logging_object.get("request_tags", []) or [] - ) - request_tags = [f"request_tag:{tag}" for tag in _request_tags] - tags.extend(request_tags) - - return ",".join(tags) - - @staticmethod - def _get_datadog_source(): - return os.getenv("DD_SOURCE", "litellm") - - @staticmethod - def _get_datadog_service(): - return os.getenv("DD_SERVICE", "litellm-server") - - @staticmethod - def _get_datadog_hostname(): - return os.getenv("HOSTNAME", "") - - @staticmethod - def _get_datadog_env(): - return os.getenv("DD_ENV", "unknown") - - @staticmethod - def _get_datadog_pod_name(): - return os.getenv("POD_NAME", "unknown") - async def async_health_check(self) -> IntegrationHealthCheckStatus: """ Check if the service is healthy @@ -651,4 +612,4 @@ class DataDogLogger( start_time_utc: Optional[datetimeObj], end_time_utc: Optional[datetimeObj], ) -> Optional[dict]: - pass \ No newline at end of file + pass diff --git a/litellm/integrations/datadog/datadog_handler.py b/litellm/integrations/datadog/datadog_handler.py new file mode 100644 index 00000000000..26fab77759e --- /dev/null +++ b/litellm/integrations/datadog/datadog_handler.py @@ -0,0 +1,50 @@ +"""Shared helpers for Datadog integrations.""" + +from __future__ import annotations + +import os +from typing import List, Optional + +from litellm.types.utils import StandardLoggingPayload + + +def get_datadog_source() -> str: + return os.getenv("DD_SOURCE", "litellm") + + +def get_datadog_service() -> str: + return os.getenv("DD_SERVICE", "litellm-server") + + +def get_datadog_hostname() -> str: + return os.getenv("HOSTNAME", "") + + +def get_datadog_env() -> str: + return os.getenv("DD_ENV", "unknown") + + +def get_datadog_pod_name() -> str: + return os.getenv("POD_NAME", "unknown") + + +def get_datadog_tags( + standard_logging_object: Optional[StandardLoggingPayload] = None, +) -> str: + """Build Datadog tags string used by multiple integrations.""" + + base_tags = { + "env": get_datadog_env(), + "service": get_datadog_service(), + "version": os.getenv("DD_VERSION", "unknown"), + "HOSTNAME": get_datadog_hostname(), + "POD_NAME": get_datadog_pod_name(), + } + + tags: List[str] = [f"{k}:{v}" for k, v in base_tags.items()] + + if standard_logging_object: + request_tags = standard_logging_object.get("request_tags", []) or [] + tags.extend(f"request_tag:{tag}" for tag in request_tags) + + return ",".join(tags) diff --git a/litellm/integrations/datadog/datadog_llm_obs.py b/litellm/integrations/datadog/datadog_llm_obs.py index b44762d0af8..65ed8a795c0 100644 --- a/litellm/integrations/datadog/datadog_llm_obs.py +++ b/litellm/integrations/datadog/datadog_llm_obs.py @@ -18,7 +18,10 @@ import httpx import litellm from litellm._logging import verbose_logger from litellm.integrations.custom_batch_logger import CustomBatchLogger -from litellm.integrations.datadog.datadog import DataDogLogger +from litellm.integrations.datadog.datadog_handler import ( + get_datadog_service, + get_datadog_tags, +) from litellm.litellm_core_utils.dd_tracing import tracer from litellm.litellm_core_utils.prompt_templates.common_utils import ( handle_any_messages_to_chat_completion_str_messages_conversion, @@ -36,7 +39,7 @@ from litellm.types.utils import ( ) -class DataDogLLMObsLogger(DataDogLogger, CustomBatchLogger): +class DataDogLLMObsLogger(CustomBatchLogger): def __init__(self, **kwargs): try: verbose_logger.debug("DataDogLLMObs: Initializing logger") @@ -142,8 +145,8 @@ class DataDogLLMObsLogger(DataDogLogger, CustomBatchLogger): "data": DDIntakePayload( type="span", attributes=DDSpanAttributes( - ml_app=self._get_datadog_service(), - tags=[self._get_datadog_tags()], + ml_app=get_datadog_service(), + tags=[get_datadog_tags()], spans=self.log_queue, ), ), @@ -243,9 +246,7 @@ class DataDogLLMObsLogger(DataDogLogger, CustomBatchLogger): duration=int((end_time - start_time).total_seconds() * 1e9), metrics=metrics, status="error" if error_info else "ok", - tags=[ - self._get_datadog_tags(standard_logging_object=standard_logging_payload) - ], + tags=[get_datadog_tags(standard_logging_object=standard_logging_payload)], ) apm_trace_id = self._get_apm_trace_id() diff --git a/tests/logging_callback_tests/test_datadog.py b/tests/logging_callback_tests/test_datadog.py index c877f34ac03..7fa4dfe7222 100644 --- a/tests/logging_callback_tests/test_datadog.py +++ b/tests/logging_callback_tests/test_datadog.py @@ -2,6 +2,14 @@ import io import os import sys +from litellm.integrations.datadog.datadog_handler import ( + get_datadog_source, + get_datadog_service, + get_datadog_env, + get_datadog_pod_name, + get_datadog_hostname, + get_datadog_tags, +) sys.path.insert(0, os.path.abspath("../..")) @@ -452,16 +460,16 @@ def test_datadog_static_methods(): """Test the static helper methods in DataDogLogger class""" # Test with default environment variables - assert DataDogLogger._get_datadog_source() == "litellm" - assert DataDogLogger._get_datadog_service() == "litellm-server" - assert DataDogLogger._get_datadog_hostname() is not None - assert DataDogLogger._get_datadog_env() == "unknown" - assert DataDogLogger._get_datadog_pod_name() == "unknown" + assert get_datadog_source() == "litellm" + assert get_datadog_service() == "litellm-server" + assert get_datadog_hostname() is not None + assert get_datadog_env() == "unknown" + assert get_datadog_pod_name() == "unknown" # Test tags format with default values assert ( - "env:unknown,service:litellm,version:unknown,HOSTNAME:" - in DataDogLogger._get_datadog_tags() + "env:unknown,service:litellm-server,version:unknown,HOSTNAME:" + in get_datadog_tags() ) # Test with custom environment variables @@ -475,31 +483,31 @@ def test_datadog_static_methods(): } with patch.dict(os.environ, test_env): - assert DataDogLogger._get_datadog_source() == "custom-source" + assert get_datadog_source() == "custom-source" print( - "DataDogLogger._get_datadog_source()", DataDogLogger._get_datadog_source() + "DataDogLogger._get_datadog_source()", get_datadog_source() ) - assert DataDogLogger._get_datadog_service() == "custom-service" + assert get_datadog_service() == "custom-service" print( - "DataDogLogger._get_datadog_service()", DataDogLogger._get_datadog_service() + "DataDogLogger._get_datadog_service()", get_datadog_service() ) - assert DataDogLogger._get_datadog_hostname() == "test-host" + assert get_datadog_hostname() == "test-host" print( "DataDogLogger._get_datadog_hostname()", - DataDogLogger._get_datadog_hostname(), + get_datadog_hostname(), ) - assert DataDogLogger._get_datadog_env() == "production" - print("DataDogLogger._get_datadog_env()", DataDogLogger._get_datadog_env()) - assert DataDogLogger._get_datadog_pod_name() == "pod-123" + assert get_datadog_env() == "production" + print("DataDogLogger._get_datadog_env()", get_datadog_env()) + assert get_datadog_pod_name() == "pod-123" print( "DataDogLogger._get_datadog_pod_name()", - DataDogLogger._get_datadog_pod_name(), + get_datadog_pod_name(), ) # Test tags format with custom values expected_custom_tags = "env:production,service:custom-service,version:1.0.0,HOSTNAME:test-host,POD_NAME:pod-123" - print("DataDogLogger._get_datadog_tags()", DataDogLogger._get_datadog_tags()) - assert DataDogLogger._get_datadog_tags() == expected_custom_tags + print("DataDogLogger._get_datadog_tags()", get_datadog_tags()) + assert get_datadog_tags() == expected_custom_tags @pytest.mark.asyncio @@ -539,7 +547,7 @@ async def test_datadog_non_serializable_messages(): def test_get_datadog_tags(): """Test the _get_datadog_tags static method with various inputs""" # Test with no standard_logging_object and default env vars - base_tags = DataDogLogger._get_datadog_tags() + base_tags = get_datadog_tags() assert "env:" in base_tags assert "service:" in base_tags assert "version:" in base_tags @@ -555,7 +563,7 @@ def test_get_datadog_tags(): "POD_NAME": "pod-123", } with patch.dict(os.environ, test_env): - custom_tags = DataDogLogger._get_datadog_tags() + custom_tags = get_datadog_tags() assert "env:production" in custom_tags assert "service:custom-service" in custom_tags assert "version:1.0.0" in custom_tags @@ -566,18 +574,18 @@ def test_get_datadog_tags(): standard_logging_obj = create_standard_logging_payload() standard_logging_obj["request_tags"] = ["tag1", "tag2"] - tags_with_request = DataDogLogger._get_datadog_tags(standard_logging_obj) + tags_with_request = get_datadog_tags(standard_logging_obj) assert "request_tag:tag1" in tags_with_request assert "request_tag:tag2" in tags_with_request # Test with empty request_tags standard_logging_obj["request_tags"] = [] - tags_empty_request = DataDogLogger._get_datadog_tags(standard_logging_obj) + tags_empty_request = get_datadog_tags(standard_logging_obj) assert "request_tag:" not in tags_empty_request # Test with None request_tags standard_logging_obj["request_tags"] = None - tags_none_request = DataDogLogger._get_datadog_tags(standard_logging_obj) + tags_none_request = get_datadog_tags(standard_logging_obj) assert "request_tag:" not in tags_none_request diff --git a/tests/test_litellm/litellm_core_utils/test_litellm_logging.py b/tests/test_litellm/litellm_core_utils/test_litellm_logging.py index 08d71d4bcd7..9b150fd89f4 100644 --- a/tests/test_litellm/litellm_core_utils/test_litellm_logging.py +++ b/tests/test_litellm/litellm_core_utils/test_litellm_logging.py @@ -161,6 +161,40 @@ def test_logging_prevent_double_logging(logging_obj): assert logging_obj.should_run_logging(event_type="async_failure") == True +@pytest.mark.asyncio +async def test_datadog_logger_not_shadowed_by_llm_obs(monkeypatch): + """Ensure DataDog logger instantiates even when LLM Obs logger already cached.""" + + # Ensure required env vars exist for Datadog loggers + monkeypatch.setenv("DD_API_KEY", "test") + monkeypatch.setenv("DD_SITE", "us5.datadoghq.com") + + from litellm.litellm_core_utils import litellm_logging as logging_module + from litellm.integrations.datadog.datadog import DataDogLogger + from litellm.integrations.datadog.datadog_llm_obs import DataDogLLMObsLogger + + logging_module._in_memory_loggers.clear() + + try: + # Cache an LLM Obs logger first to mirror callbacks=["datadog_llm_observability", ...] + obs_logger = DataDogLLMObsLogger() + logging_module._in_memory_loggers.append(obs_logger) + + datadog_logger = logging_module._init_custom_logger_compatible_class( + logging_integration="datadog", + internal_usage_cache=None, + llm_router=None, + custom_logger_init_args={}, + ) + + # Regression check: we expect a distinct DataDogLogger, not the LLM Obs logger + assert type(datadog_logger) is DataDogLogger + assert any(isinstance(cb, DataDogLLMObsLogger) for cb in logging_module._in_memory_loggers) + assert any(type(cb) is DataDogLogger for cb in logging_module._in_memory_loggers) + finally: + logging_module._in_memory_loggers.clear() + + @pytest.mark.asyncio async def test_logging_result_for_bridge_calls(logging_obj): """