diff --git a/docs/my-website/docs/observability/datadog.md b/docs/my-website/docs/observability/datadog.md index 9385b0020c..e83cfcbafe 100644 --- a/docs/my-website/docs/observability/datadog.md +++ b/docs/my-website/docs/observability/datadog.md @@ -7,6 +7,7 @@ import TabItem from '@theme/TabItem'; LiteLLM Supports logging to the following Datdog Integrations: - `datadog` [Datadog Logs](https://docs.datadoghq.com/logs/) - `datadog_llm_observability` [Datadog LLM Observability](https://www.datadoghq.com/product/llm-observability/) +- `datadog_metrics` [Datadog Custom Metrics](#datadog-custom-metrics) - `datadog_cost_management` [Datadog Cloud Cost Management](#datadog-cloud-cost-management) - `ddtrace-run` [Datadog Tracing](#datadog-tracing) @@ -168,6 +169,65 @@ On the Datadog LLM Observability page, you should see that both input messages a +## Datadog Custom Metrics + +| Feature | Details | +|---------|---------| +| **What is logged** | Latency metrics, request counts by status code | +| **Events** | Success + Failure | +| **Product Link** | [Datadog Metrics](https://docs.datadoghq.com/metrics/) | + +Publishes the following metrics to Datadog via the `/api/v2/series` endpoint: + +| Metric | Type | Description | +|--------|------|-------------| +| `litellm.request.total_latency` | Gauge | End-to-end request latency (seconds) | +| `litellm.llm_api.latency` | Gauge | Time spent waiting for the LLM provider response (seconds) | +| `litellm.llm_api.request_count` | Count | Request count, tagged with status code | + +Using `total_latency` and `llm_api.latency`, you can derive **internal latency** = `total_latency - llm_api.latency`. + +All metrics include the following tags: `env`, `service`, `version`, `HOSTNAME`, `POD_NAME`, `provider`, `model_name`, `model_group`, `team`, `status_code`. + +**Step 1**: Create a `config.yaml` file + +```yaml +model_list: + - model_name: gpt-3.5-turbo + litellm_params: + model: gpt-3.5-turbo +litellm_settings: + success_callback: ["datadog_metrics"] + failure_callback: ["datadog_metrics"] +``` + +**Step 2**: Set required env variables + +```shell +DD_API_KEY="your-api-key" +DD_SITE="us5.datadoghq.com" # your datadog site +``` + +**Step 3**: Start the proxy and make a test request + +```shell +litellm --config config.yaml +``` + +```shell +curl --location 'http://0.0.0.0:4000/chat/completions' \ + --header 'Content-Type: application/json' \ + --header 'Authorization: Bearer sk-1234' \ + --data '{ + "model": "gpt-3.5-turbo", + "messages": [{"role": "user", "content": "hello"}] +}' +``` + +**Step 4**: View metrics in Datadog Metrics Explorer + +Navigate to **Metrics > Explorer** in Datadog and search for `litellm.request.total_latency`, `litellm.llm_api.latency`, or `litellm.llm_api.request_count`. + ## Datadog Cloud Cost Management | Feature | Details | diff --git a/litellm/__init__.py b/litellm/__init__.py index 50fa0e7675..3c61aca3b8 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -105,6 +105,7 @@ _custom_logger_compatible_callbacks_literal = Literal[ "prometheus", "otel", "datadog", + "datadog_metrics", "datadog_llm_observability", "galileo", "braintrust", diff --git a/litellm/integrations/callback_configs.json b/litellm/integrations/callback_configs.json index 6a003b8c49..c2b0c4ddce 100644 --- a/litellm/integrations/callback_configs.json +++ b/litellm/integrations/callback_configs.json @@ -83,6 +83,27 @@ }, "description": "Datadog Logging Integration" }, + { + "id": "datadog_metrics", + "displayName": "Datadog Metrics", + "logo": "datadog.png", + "supports_key_team_logging": false, + "dynamic_params": { + "dd_api_key": { + "type": "password", + "ui_name": "API Key", + "description": "Datadog API key for authentication", + "required": true + }, + "dd_site": { + "type": "text", + "ui_name": "Site", + "description": "Datadog site URL (e.g., us5.datadoghq.com)", + "required": true + } + }, + "description": "Datadog Custom Metrics Integration" + }, { "id": "datadog_cost_management", "displayName": "Datadog Cost Management", @@ -434,4 +455,4 @@ }, "description": "SQS Queue (AWS) Logging Integration" } -] \ No newline at end of file +] diff --git a/litellm/integrations/datadog/datadog_metrics.py b/litellm/integrations/datadog/datadog_metrics.py new file mode 100644 index 0000000000..fcf40701e2 --- /dev/null +++ b/litellm/integrations/datadog/datadog_metrics.py @@ -0,0 +1,286 @@ +import asyncio +import gzip +import os +import time +from datetime import datetime +from typing import List, Optional, Union + +from litellm._logging import verbose_logger +from litellm.integrations.custom_batch_logger import CustomBatchLogger +from litellm.integrations.datadog.datadog_handler import ( + get_datadog_env, + get_datadog_hostname, + get_datadog_pod_name, + get_datadog_service, +) +from litellm.litellm_core_utils.safe_json_dumps import safe_dumps +from litellm.llms.custom_httpx.http_handler import ( + get_async_httpx_client, + httpxSpecialProvider, +) +from litellm.types.integrations.base_health_check import IntegrationHealthCheckStatus +from litellm.types.integrations.datadog_metrics import ( + DatadogMetricPoint, + DatadogMetricSeries, + DatadogMetricsPayload, +) +from litellm.types.utils import StandardLoggingPayload + + +class DatadogMetricsLogger(CustomBatchLogger): + def __init__(self, start_periodic_flush: bool = True, **kwargs): + self.dd_api_key = os.getenv("DD_API_KEY") + self.dd_app_key = os.getenv("DD_APP_KEY") + self.dd_site = os.getenv("DD_SITE", "datadoghq.com") + + if not self.dd_api_key: + verbose_logger.warning( + "Datadog Metrics: DD_API_KEY is required. Integration will not work." + ) + + self.upload_url = f"https://api.{self.dd_site}/api/v2/series" + + self.async_client = get_async_httpx_client( + llm_provider=httpxSpecialProvider.LoggingCallback + ) + + # Initialize lock + self.flush_lock = asyncio.Lock() + + # Only set flush_lock if not already provided by caller + if "flush_lock" not in kwargs: + kwargs["flush_lock"] = self.flush_lock + + # Send metrics more quickly to datadog (every 5 seconds) + if "flush_interval" not in kwargs: + kwargs["flush_interval"] = 5 + + super().__init__(**kwargs) + + # Start periodic flush task only if instructed + if start_periodic_flush: + asyncio.create_task(self.periodic_flush()) + + def _extract_tags( + self, + log: StandardLoggingPayload, + status_code: Optional[Union[str, int]] = None, + ) -> List[str]: + """ + Builds the list of tags for a Datadog metric point + """ + # Base tags + tags = [ + f"env:{get_datadog_env()}", + f"service:{get_datadog_service()}", + f"version:{os.getenv('DD_VERSION', 'unknown')}", + f"HOSTNAME:{get_datadog_hostname()}", + f"POD_NAME:{get_datadog_pod_name()}", + ] + + # Add metric-specific tags + if provider := log.get("custom_llm_provider"): + tags.append(f"provider:{provider}") + + if model := log.get("model"): + tags.append(f"model_name:{model}") + + if model_group := log.get("model_group"): + tags.append(f"model_group:{model_group}") + + if status_code is not None: + tags.append(f"status_code:{status_code}") + + # Extract team tag + metadata = log.get("metadata", {}) or {} + team_tag = ( + metadata.get("user_api_key_team_alias") + or metadata.get("team_alias") # type: ignore + or metadata.get("user_api_key_team_id") + or metadata.get("team_id") # type: ignore + ) + + if team_tag: + tags.append(f"team:{team_tag}") + + return tags + + def _add_metrics_from_log( + self, + log: StandardLoggingPayload, + kwargs: dict, + status_code: Union[str, int] = "200", + ): + """ + Extracts latencies and appends Datadog metric series to the queue + """ + tags = self._extract_tags(log, status_code=status_code) + + # We record metrics with the end_time as the timestamp for the point + end_time_dt = kwargs.get("end_time") or datetime.now() + timestamp = int(end_time_dt.timestamp()) + + # 1. Total Request Latency Metric (End to End) + start_time_dt = kwargs.get("start_time") + if start_time_dt and end_time_dt: + total_duration = (end_time_dt - start_time_dt).total_seconds() + series_total_latency: DatadogMetricSeries = { + "metric": "litellm.request.total_latency", + "type": 3, # gauge + "points": [{"timestamp": timestamp, "value": total_duration}], + "tags": tags, + } + self.log_queue.append(series_total_latency) + + # 2. LLM API Latency Metric (Provider alone) + api_call_start_time = kwargs.get("api_call_start_time") + if api_call_start_time and end_time_dt: + llm_api_duration = (end_time_dt - api_call_start_time).total_seconds() + series_llm_latency: DatadogMetricSeries = { + "metric": "litellm.llm_api.latency", + "type": 3, # gauge + "points": [{"timestamp": timestamp, "value": llm_api_duration}], + "tags": tags, + } + self.log_queue.append(series_llm_latency) + + # 3. Request Count / Status Code + series_count: DatadogMetricSeries = { + "metric": "litellm.llm_api.request_count", + "type": 1, # count + "points": [{"timestamp": timestamp, "value": 1.0}], + "tags": tags, + "interval": self.flush_interval, + } + self.log_queue.append(series_count) + + async def async_log_success_event(self, kwargs, response_obj, start_time, end_time): + try: + standard_logging_object: Optional[StandardLoggingPayload] = kwargs.get( + "standard_logging_object", None + ) + + if standard_logging_object is None: + return + + self._add_metrics_from_log( + log=standard_logging_object, kwargs=kwargs, status_code="200" + ) + + if len(self.log_queue) >= self.batch_size: + await self.flush_queue() + + except Exception as e: + verbose_logger.exception( + f"Datadog Metrics: Error in async_log_success_event: {str(e)}" + ) + + async def async_log_failure_event(self, kwargs, response_obj, start_time, end_time): + try: + standard_logging_object: Optional[StandardLoggingPayload] = kwargs.get( + "standard_logging_object", None + ) + + if standard_logging_object is None: + return + + # Extract status code from error information + status_code = "500" # default + error_information = ( + standard_logging_object.get("error_information", {}) or {} + ) + error_code = error_information.get("error_code") # type: ignore + if error_code is not None: + status_code = str(error_code) + + self._add_metrics_from_log( + log=standard_logging_object, kwargs=kwargs, status_code=status_code + ) + + if len(self.log_queue) >= self.batch_size: + await self.flush_queue() + + except Exception as e: + verbose_logger.exception( + f"Datadog Metrics: Error in async_log_failure_event: {str(e)}" + ) + + async def async_send_batch(self): + if not self.log_queue: + return + + batch = self.log_queue.copy() + payload_data: DatadogMetricsPayload = {"series": batch} + + try: + await self._upload_to_datadog(payload_data) + except Exception as e: + verbose_logger.exception( + f"Datadog Metrics: Error in async_send_batch: {str(e)}" + ) + raise + + async def _upload_to_datadog(self, payload: DatadogMetricsPayload): + if not self.dd_api_key: + return + + headers = { + "Content-Type": "application/json", + "DD-API-KEY": self.dd_api_key, + } + + if self.dd_app_key: + headers["DD-APPLICATION-KEY"] = self.dd_app_key + + json_data = safe_dumps(payload) + compressed_data = gzip.compress(json_data.encode("utf-8")) + headers["Content-Encoding"] = "gzip" + + response = await self.async_client.post( + self.upload_url, content=compressed_data, headers=headers # type: ignore + ) + + response.raise_for_status() + + verbose_logger.debug( + f"Datadog Metrics: Uploaded {len(payload['series'])} metric points. Status: {response.status_code}" + ) + + async def async_health_check(self) -> IntegrationHealthCheckStatus: + """ + Check if the service is healthy + """ + try: + # Send a test metric point to Datadog + test_metric_point: DatadogMetricPoint = { + "timestamp": int(time.time()), + "value": 1.0, + } + test_metric_series: DatadogMetricSeries = { + "metric": "litellm.health_check", + "type": 3, # Gauge + "points": [test_metric_point], + "tags": ["env:health_check"], + } + + payload_data: DatadogMetricsPayload = {"series": [test_metric_series]} + + await self._upload_to_datadog(payload_data) + + return IntegrationHealthCheckStatus( + status="healthy", + error_message=None, + ) + except Exception as e: + return IntegrationHealthCheckStatus( + status="unhealthy", + error_message=str(e), + ) + + async def get_request_response_payload( + self, + request_id: str, + start_time_utc: Optional[datetime], + end_time_utc: Optional[datetime], + ) -> Optional[dict]: + pass diff --git a/litellm/litellm_core_utils/custom_logger_registry.py b/litellm/litellm_core_utils/custom_logger_registry.py index fc73701ea9..2d483f7861 100644 --- a/litellm/litellm_core_utils/custom_logger_registry.py +++ b/litellm/litellm_core_utils/custom_logger_registry.py @@ -20,6 +20,7 @@ from litellm.integrations.braintrust_logging import BraintrustLogger from litellm.integrations.cloudzero.cloudzero import CloudZeroLogger from litellm.integrations.datadog.datadog import DataDogLogger from litellm.integrations.datadog.datadog_llm_obs import DataDogLLMObsLogger +from litellm.integrations.datadog.datadog_metrics import DatadogMetricsLogger from litellm.integrations.deepeval import DeepEvalLogger from litellm.integrations.dotprompt import DotpromptManager from litellm.integrations.focus.focus_logger import FocusLogger @@ -66,6 +67,7 @@ class CustomLoggerRegistry: "prometheus": PrometheusLogger, "datadog": DataDogLogger, "datadog_llm_observability": DataDogLLMObsLogger, + "datadog_metrics": DatadogMetricsLogger, "gcs_bucket": GCSBucketLogger, "opik": OpikLogger, "argilla": ArgillaLogger, diff --git a/litellm/litellm_core_utils/litellm_logging.py b/litellm/litellm_core_utils/litellm_logging.py index e450b233c7..5e5a6cea1b 100644 --- a/litellm/litellm_core_utils/litellm_logging.py +++ b/litellm/litellm_core_utils/litellm_logging.py @@ -133,6 +133,7 @@ from ..integrations.azure_sentinel.azure_sentinel import AzureSentinelLogger from ..integrations.azure_storage.azure_storage import AzureBlobStorageLogger from ..integrations.custom_prompt_management import CustomPromptManagement from ..integrations.datadog.datadog import DataDogLogger +from ..integrations.datadog.datadog_metrics import DatadogMetricsLogger from ..integrations.datadog.datadog_llm_obs import DataDogLLMObsLogger from ..integrations.dotprompt import DotpromptManager from ..integrations.dynamodb import DyanmoDBLogger @@ -3661,6 +3662,14 @@ def _init_custom_logger_compatible_class( # noqa: PLR0915 _datadog_logger = DataDogLogger() _in_memory_loggers.append(_datadog_logger) return _datadog_logger # type: ignore + elif logging_integration == "datadog_metrics": + for callback in _in_memory_loggers: + if isinstance(callback, DatadogMetricsLogger): + return callback # type: ignore + + _datadog_metrics_logger = DatadogMetricsLogger() + _in_memory_loggers.append(_datadog_metrics_logger) + return _datadog_metrics_logger # type: ignore elif logging_integration == "datadog_llm_observability": _datadog_llm_obs_logger = DataDogLLMObsLogger() _in_memory_loggers.append(_datadog_llm_obs_logger) @@ -4268,6 +4277,10 @@ def get_custom_logger_compatible_class( # noqa: PLR0915 for callback in _in_memory_loggers: if isinstance(callback, DataDogLogger): return callback + elif logging_integration == "datadog_metrics": + for callback in _in_memory_loggers: + if isinstance(callback, DatadogMetricsLogger): + return callback elif logging_integration == "datadog_llm_observability": for callback in _in_memory_loggers: if isinstance(callback, DataDogLLMObsLogger): diff --git a/litellm/proxy/common_utils/callback_utils.py b/litellm/proxy/common_utils/callback_utils.py index 62ca6dc2ae..9ecae363ed 100644 --- a/litellm/proxy/common_utils/callback_utils.py +++ b/litellm/proxy/common_utils/callback_utils.py @@ -390,9 +390,7 @@ def get_logging_caching_headers(request_data: Dict) -> Optional[Dict]: ) if "applied_policies" in _metadata: - headers["x-litellm-applied-policies"] = ",".join( - _metadata["applied_policies"] - ) + headers["x-litellm-applied-policies"] = ",".join(_metadata["applied_policies"]) if "policy_sources" in _metadata: sources = _metadata["policy_sources"] @@ -449,9 +447,7 @@ def add_policy_to_applied_policies_header( request_data["metadata"] = _metadata -def add_policy_sources_to_metadata( - request_data: Dict, policy_sources: Dict[str, str] -): +def add_policy_sources_to_metadata(request_data: Dict, policy_sources: Dict[str, str]): """ Store policy match reasons in metadata for x-litellm-policy-sources header. diff --git a/litellm/proxy/health_endpoints/_health_endpoints.py b/litellm/proxy/health_endpoints/_health_endpoints.py index 95b1836d8a..5a4a8eae38 100644 --- a/litellm/proxy/health_endpoints/_health_endpoints.py +++ b/litellm/proxy/health_endpoints/_health_endpoints.py @@ -230,6 +230,7 @@ async def health_services_endpoint( # noqa: PLR0915 "custom_callback_api", "langsmith", "datadog", + "datadog_metrics", "datadog_llm_observability", "generic_api", "arize", @@ -284,6 +285,30 @@ async def health_services_endpoint( # noqa: PLR0915 else "Datadog is healthy" ), } + elif service == "datadog_metrics": + from litellm.integrations.datadog.datadog_metrics import ( + DatadogMetricsLogger, + ) + from litellm.litellm_core_utils.litellm_logging import ( + get_custom_logger_compatible_class, + ) + + datadog_metrics_logger = get_custom_logger_compatible_class( + "datadog_metrics" + ) + if datadog_metrics_logger is None: + datadog_metrics_logger = DatadogMetricsLogger( + start_periodic_flush=False + ) + response = await datadog_metrics_logger.async_health_check() + return { + "status": response["status"], + "message": ( + response["error_message"] + if response["status"] == "unhealthy" + else "Datadog Metrics is healthy" + ), + } elif service == "arize": from litellm.integrations.arize.arize import ArizeLogger diff --git a/litellm/types/integrations/datadog_metrics.py b/litellm/types/integrations/datadog_metrics.py new file mode 100644 index 0000000000..4c980cdee6 --- /dev/null +++ b/litellm/types/integrations/datadog_metrics.py @@ -0,0 +1,20 @@ +from typing import List, Optional + +from typing_extensions import TypedDict + + +class DatadogMetricPoint(TypedDict): + timestamp: int # Unix epoch seconds + value: float # The metric value + + +class DatadogMetricSeries(TypedDict, total=False): + metric: str + type: int # 0=unspecified, 1=count, 2=rate, 3=gauge + points: List[DatadogMetricPoint] + tags: List[str] + interval: Optional[int] # Required for count (type=1) and rate (type=2) metrics + + +class DatadogMetricsPayload(TypedDict): + series: List[DatadogMetricSeries] diff --git a/tests/test_litellm/integrations/datadog/test_datadog_metrics.py b/tests/test_litellm/integrations/datadog/test_datadog_metrics.py new file mode 100644 index 0000000000..757c558c29 --- /dev/null +++ b/tests/test_litellm/integrations/datadog/test_datadog_metrics.py @@ -0,0 +1,273 @@ +import os +import time +from datetime import datetime, timedelta +from unittest.mock import AsyncMock + +import pytest +from httpx import Request, Response + +from litellm.integrations.datadog.datadog_metrics import DatadogMetricsLogger +from litellm.types.utils import StandardLoggingPayload + + +@pytest.fixture +def clean_env(): + """Set test env vars and restore originals after test.""" + keys = ["DD_API_KEY", "DD_APP_KEY", "DD_SITE", "DD_ENV", "DD_SERVICE", "DD_VERSION"] + originals = {k: os.environ.get(k) for k in keys} + + os.environ["DD_API_KEY"] = "test_api_key" + os.environ["DD_APP_KEY"] = "test_app_key" + os.environ["DD_SITE"] = "test.datadoghq.com" + os.environ["DD_ENV"] = "test-env" + os.environ["DD_SERVICE"] = "test-service" + os.environ["DD_VERSION"] = "1.0.0" + + yield + + for k, v in originals.items(): + if v is not None: + os.environ[k] = v + elif k in os.environ: + del os.environ[k] + + +@pytest.mark.asyncio +async def test_init(clean_env): + """Test initialization sets up clients and url correctly.""" + logger = DatadogMetricsLogger(start_periodic_flush=False) + assert logger.upload_url == "https://api.test.datadoghq.com/api/v2/series" + + +@pytest.mark.asyncio +async def test_extract_tags(clean_env): + """Test tag extraction from a StandardLoggingPayload.""" + logger = DatadogMetricsLogger(start_periodic_flush=False) + + payload = StandardLoggingPayload( + custom_llm_provider="openai", + model="gpt-4o", + model_group="gpt-4", + metadata={"user_api_key_team_alias": "test-team"}, + ) + + tags = logger._extract_tags(log=payload, status_code="200") + + assert "env:test-env" in tags + assert "service:test-service" in tags + assert "version:1.0.0" in tags + assert "provider:openai" in tags + assert "model_name:gpt-4o" in tags + assert "model_group:gpt-4" in tags + assert "status_code:200" in tags + assert "team:test-team" in tags + + +@pytest.mark.asyncio +async def test_extract_tags_no_team(clean_env): + """Test tag extraction when no team info is present.""" + logger = DatadogMetricsLogger(start_periodic_flush=False) + + payload = StandardLoggingPayload( + custom_llm_provider="anthropic", + model="claude-3-sonnet", + ) + + tags = logger._extract_tags(log=payload, status_code="500") + + assert "provider:anthropic" in tags + assert "model_name:claude-3-sonnet" in tags + assert "status_code:500" in tags + assert not any(tag.startswith("team:") for tag in tags) + + +@pytest.mark.asyncio +async def test_add_metrics_from_log(clean_env): + """Test that _add_metrics_from_log appends the correct metric series to the queue.""" + logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False) + + now = datetime.now() + start_time = now - timedelta(seconds=2) + api_call_start_time = now - timedelta(seconds=1) + + payload = StandardLoggingPayload( + custom_llm_provider="openai", + model="gpt-4o", + ) + + kwargs = { + "start_time": start_time, + "api_call_start_time": api_call_start_time, + "end_time": now, + } + + logger._add_metrics_from_log(log=payload, kwargs=kwargs, status_code="200") + + # Should have 3 series: total_latency, llm_api_latency, request_count + assert len(logger.log_queue) == 3 + + metrics = {s["metric"]: s for s in logger.log_queue} + + # Total latency ~2s + total = metrics["litellm.request.total_latency"] + assert total["type"] == 3 # gauge + assert abs(total["points"][0]["value"] - 2.0) < 0.1 + + # LLM API latency ~1s + llm = metrics["litellm.llm_api.latency"] + assert llm["type"] == 3 # gauge + assert abs(llm["points"][0]["value"] - 1.0) < 0.1 + + # Request count + count = metrics["litellm.llm_api.request_count"] + assert count["type"] == 1 # count + assert count["points"][0]["value"] == 1.0 + assert "status_code:200" in count["tags"] + + +@pytest.mark.asyncio +async def test_async_log_success_event(clean_env): + """Test that success events are added to the queue.""" + logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False) + + now = datetime.now() + start_time = now - timedelta(seconds=1) + + await logger.async_log_success_event( + kwargs={ + "standard_logging_object": StandardLoggingPayload( + custom_llm_provider="openai", + model="gpt-4o", + ), + "start_time": start_time, + "end_time": now, + }, + response_obj=None, + start_time=start_time, + end_time=now, + ) + + # At least request_count and total_latency + assert len(logger.log_queue) >= 2 + + +@pytest.mark.asyncio +async def test_async_log_success_event_no_standard_logging_object(clean_env): + """Test that events without standard_logging_object are skipped.""" + logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False) + + await logger.async_log_success_event( + kwargs={}, + response_obj=None, + start_time=datetime.now(), + end_time=datetime.now(), + ) + + assert len(logger.log_queue) == 0 + + +@pytest.mark.asyncio +async def test_async_log_failure_event_extracts_status_code(clean_env): + """Test that failure events extract the error status code.""" + logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False) + + now = datetime.now() + start_time = now - timedelta(seconds=1) + + await logger.async_log_failure_event( + kwargs={ + "standard_logging_object": StandardLoggingPayload( + custom_llm_provider="openai", + model="gpt-4o", + error_information={"error_code": "429"}, + ), + "start_time": start_time, + "end_time": now, + }, + response_obj=None, + start_time=start_time, + end_time=now, + ) + + count_series = next( + (s for s in logger.log_queue if s["metric"] == "litellm.llm_api.request_count"), + None, + ) + assert count_series is not None + assert "status_code:429" in count_series["tags"] + + +@pytest.mark.asyncio +async def test_async_log_failure_event_default_status_code(clean_env): + """Test that failure events default to 500 when no error_code is present.""" + logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False) + + now = datetime.now() + + await logger.async_log_failure_event( + kwargs={ + "standard_logging_object": StandardLoggingPayload( + custom_llm_provider="openai", + model="gpt-4o", + ), + "start_time": now, + "end_time": now, + }, + response_obj=None, + start_time=now, + end_time=now, + ) + + count_series = next( + (s for s in logger.log_queue if s["metric"] == "litellm.llm_api.request_count"), + None, + ) + assert count_series is not None + assert "status_code:500" in count_series["tags"] + + +@pytest.mark.asyncio +async def test_async_send_batch(clean_env): + """Test that async_send_batch uploads metrics to Datadog.""" + logger = DatadogMetricsLogger(start_periodic_flush=False) + logger.async_client = AsyncMock() + mock_request = Request("POST", "https://api.test.datadoghq.com/api/v2/series") + logger.async_client.post.return_value = Response( + 202, json={"status": "ok"}, request=mock_request + ) + + # Manually add a metric series to the queue + logger.log_queue = [ + { + "metric": "litellm.request.total_latency", + "type": 3, + "points": [{"timestamp": int(time.time()), "value": 1.5}], + "tags": ["env:test"], + } + ] + + await logger.async_send_batch() + + assert logger.async_client.post.called + call_args = logger.async_client.post.call_args + assert call_args[0][0] == "https://api.test.datadoghq.com/api/v2/series" + + # Verify gzip + JSON payload + import gzip + import json + + compressed = call_args[1]["content"] + payload = json.loads(gzip.decompress(compressed).decode("utf-8")) + assert len(payload["series"]) == 1 + assert payload["series"][0]["metric"] == "litellm.request.total_latency" + + +@pytest.mark.asyncio +async def test_async_send_batch_empty_queue(clean_env): + """Test that async_send_batch does nothing when queue is empty.""" + logger = DatadogMetricsLogger(start_periodic_flush=False) + logger.async_client = AsyncMock() + + await logger.async_send_batch() + + assert not logger.async_client.post.called