diff --git a/litellm/integrations/datadog/datadog_metrics.py b/litellm/integrations/datadog/datadog_metrics.py index 5dda336dc94..2eee53c0915 100644 --- a/litellm/integrations/datadog/datadog_metrics.py +++ b/litellm/integrations/datadog/datadog_metrics.py @@ -2,8 +2,10 @@ import asyncio import gzip import os import time +import zlib +from collections.abc import Sequence from datetime import datetime -from typing import Final +from typing import Final, Literal from litellm._logging import verbose_logger from litellm.integrations.custom_batch_logger import CustomBatchLogger @@ -21,6 +23,8 @@ from litellm.llms.custom_httpx.http_handler import ( ) from litellm.types.integrations.base_health_check import IntegrationHealthCheckStatus from litellm.types.integrations.datadog_metrics import ( + DatadogDistributionPayload, + DatadogDistributionSeries, DatadogMetricPoint, DatadogMetricSeries, DatadogMetricsPayload, @@ -38,6 +42,7 @@ class DatadogMetricsLogger(CustomBatchLogger): 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.distribution_upload_url = f"https://api.{self.dd_site}/api/v1/distribution_points" self.async_client = get_async_httpx_client(llm_provider=httpxSpecialProvider.LoggingCallback) @@ -71,10 +76,12 @@ class DatadogMetricsLogger(CustomBatchLogger): 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()}", ] + if hostname := get_datadog_hostname(): + tags.append(f"HOSTNAME:{hostname}") + # Add metric-specific tags if provider := log.get("custom_llm_provider"): tags.append(f"provider:{provider}") @@ -102,6 +109,25 @@ class DatadogMetricsLogger(CustomBatchLogger): return tags + def _add_latency_metric(self, metric: str, seconds: float, timestamp: int, tags: Sequence[str]) -> None: + """ + Queues a latency sample as a gauge (legacy metric name) and as a distribution + (`.distribution`) so Datadog computes percentiles over every request + """ + gauge: Final[DatadogMetricSeries] = { + "metric": metric, + "type": 3, # gauge + "points": [{"timestamp": timestamp, "value": seconds}], + "tags": tags, + } + distribution: Final[DatadogDistributionSeries] = { + "metric": f"{metric}.distribution", + "type": "distribution", + "points": ((timestamp, (seconds,)),), + "tags": tuple(tags), + } + self.log_queue.extend((gauge, distribution)) + def _add_metrics_from_log( self, log: StandardLoggingPayload, @@ -121,43 +147,22 @@ class DatadogMetricsLogger(CustomBatchLogger): start_time_dt: Final = kwargs.get("start_time") if start_time_dt and end_time_dt: total_duration: Final = (end_time_dt - start_time_dt).total_seconds() - series_total_latency: Final[DatadogMetricSeries] = { - "metric": "litellm.request.total_latency", - "type": 3, # gauge - "points": [{"timestamp": timestamp, "value": total_duration}], - "tags": tags, - } - self.log_queue.append(series_total_latency) + self._add_latency_metric("litellm.request.total_latency", total_duration, timestamp, tags) # 2. LLM API Latency Metric (Provider alone) api_call_start_time: Final = kwargs.get("api_call_start_time") if api_call_start_time and end_time_dt: llm_api_duration: Final = (end_time_dt - api_call_start_time).total_seconds() - series_llm_latency: Final[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) + self._add_latency_metric("litellm.llm_api.latency", llm_api_duration, timestamp, tags) # 3. LiteLLM Overhead Latency Metric (total - llm_api time) hidden_params: Final = log.get("hidden_params", {}) or {} litellm_overhead_time_ms: Final = hidden_params.get("litellm_overhead_time_ms") if litellm_overhead_time_ms is not None: overhead_tags: Final = self._extract_tags(log) # no status_code on latency metric - series_overhead: Final[DatadogMetricSeries] = { - "metric": "litellm.overhead.latency", - "type": 3, # gauge - "points": [ - { - "timestamp": timestamp, - "value": litellm_overhead_time_ms / 1000, # convert ms → seconds - } - ], - "tags": overhead_tags, - } - self.log_queue.append(series_overhead) + self._add_latency_metric( + "litellm.overhead.latency", litellm_overhead_time_ms / 1000, timestamp, overhead_tags + ) # 4. Request Count / Status Code series_count: Final[DatadogMetricSeries] = { @@ -210,42 +215,55 @@ class DatadogMetricsLogger(CustomBatchLogger): if not self.log_queue: return - batch: Final = self.log_queue.copy() - payload_data: Final[DatadogMetricsPayload] = {"series": batch} + batch: Final[tuple[DatadogMetricSeries | DatadogDistributionSeries, ...]] = tuple(self.log_queue) + series: Final[tuple[DatadogMetricSeries, ...]] = tuple(s for s in batch if s["type"] != "distribution") + distributions: Final[tuple[DatadogDistributionSeries, ...]] = tuple( + s for s in batch if s["type"] == "distribution" + ) try: - await self._upload_to_datadog(payload_data) + if series: + await self._upload_to_datadog({"series": series}) except Exception as e: verbose_logger.exception("Datadog Metrics: Error in async_send_batch: %s", e) raise + try: + if distributions: + await self._upload_distributions_to_datadog({"series": distributions}) + except Exception as e: + self.log_queue[: len(batch)] = distributions + verbose_logger.exception("Datadog Metrics: Error in async_send_batch: %s", e) + raise + async def _upload_to_datadog(self, payload: DatadogMetricsPayload): + await self._post_compressed(self.upload_url, safe_dumps(payload), encoding="gzip") + + async def _upload_distributions_to_datadog(self, payload: DatadogDistributionPayload): + # /api/v1/distribution_points only accepts deflate-compressed bodies + await self._post_compressed(self.distribution_upload_url, safe_dumps(payload), encoding="deflate") + + async def _post_compressed(self, url: str, json_data: str, encoding: Literal["gzip", "deflate"]) -> None: if not self.dd_api_key: return headers: Final = { "Content-Type": "application/json", + "Content-Encoding": encoding, "DD-API-KEY": self.dd_api_key, } if self.dd_app_key: headers["DD-APPLICATION-KEY"] = self.dd_app_key - json_data: Final = safe_dumps(payload) - compressed_data: Final = gzip.compress(json_data.encode("utf-8")) - headers["Content-Encoding"] = "gzip" + raw: Final = json_data.encode("utf-8") + compressed_data: Final = gzip.compress(raw) if encoding == "gzip" else zlib.compress(raw) - response: Final = await self.async_client.post( - self.upload_url, - content=compressed_data, - headers=headers, - ) + response: Final = await self.async_client.post(url, content=compressed_data, headers=headers) response.raise_for_status() - verbose_logger.debug( - "Datadog Metrics: Uploaded %s metric points. Status: %s", len(payload["series"]), response.status_code - ) + verbose_logger.debug("Datadog Metrics: Uploaded metrics to %s. Status: %s", url, response.status_code) async def async_health_check(self) -> IntegrationHealthCheckStatus: """ diff --git a/litellm/types/integrations/datadog_metrics.py b/litellm/types/integrations/datadog_metrics.py index c294fdd2522..44801ec8b0c 100644 --- a/litellm/types/integrations/datadog_metrics.py +++ b/litellm/types/integrations/datadog_metrics.py @@ -1,4 +1,7 @@ -from typing_extensions import TypedDict +from collections.abc import Sequence +from typing import Literal + +from typing_extensions import NotRequired, ReadOnly, TypedDict class DatadogMetricPoint(TypedDict): @@ -6,13 +9,27 @@ class DatadogMetricPoint(TypedDict): 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: int | None # Required for count (type=1) and rate (type=2) metrics +class DatadogMetricSeries(TypedDict): + metric: ReadOnly[str] + type: ReadOnly[Literal[0, 1, 2, 3]] # 0=unspecified, 1=count, 2=rate, 3=gauge + points: ReadOnly[Sequence[DatadogMetricPoint]] + tags: ReadOnly[Sequence[str]] + interval: ReadOnly[NotRequired[int]] # Required for count (type=1) and rate (type=2) metrics class DatadogMetricsPayload(TypedDict): - series: list[DatadogMetricSeries] + series: ReadOnly[Sequence[DatadogMetricSeries]] + + +DatadogDistributionPoint = tuple[int, tuple[float, ...]] + + +class DatadogDistributionSeries(TypedDict): + metric: ReadOnly[str] + type: ReadOnly[Literal["distribution"]] + points: ReadOnly[tuple[DatadogDistributionPoint, ...]] + tags: ReadOnly[tuple[str, ...]] + + +class DatadogDistributionPayload(TypedDict): + series: ReadOnly[tuple[DatadogDistributionSeries, ...]] diff --git a/tests/test_litellm/integrations/datadog/test_datadog_metrics.py b/tests/test_litellm/integrations/datadog/test_datadog_metrics.py index eade92d6672..ab8b2c49321 100644 --- a/tests/test_litellm/integrations/datadog/test_datadog_metrics.py +++ b/tests/test_litellm/integrations/datadog/test_datadog_metrics.py @@ -1,3 +1,5 @@ +import gzip +import json import time from datetime import datetime, timedelta from unittest.mock import AsyncMock @@ -125,9 +127,9 @@ async def test_add_metrics_from_log(clean_env): logger._add_metrics_from_log(log=payload, kwargs=kwargs, status_code="200") - # Should have 3 series: total_latency, llm_api_latency, request_count + # total_latency and llm_api_latency (each as gauge + distribution) plus request_count # (no overhead metric because payload has no hidden_params litellm_overhead_time_ms) - assert len(logger.log_queue) == 3 + assert len(logger.log_queue) == 5 metrics = {s["metric"]: s for s in logger.log_queue} @@ -148,6 +150,65 @@ async def test_add_metrics_from_log(clean_env): assert "status_code:200" in count["tags"] +@pytest.mark.asyncio +async def test_latency_metrics_also_emitted_as_distributions(clean_env): + """Each latency sample is queued as a distribution point so Datadog can compute per-request percentiles.""" + logger = DatadogMetricsLogger(batch_size=100, start_periodic_flush=False) + + now = datetime.now() + payload = StandardLoggingPayload( + custom_llm_provider="openai", + model="gpt-4o", + hidden_params={"litellm_overhead_time_ms": 250}, + ) + kwargs = { + "start_time": now - timedelta(seconds=2), + "api_call_start_time": now - timedelta(seconds=1), + "end_time": now, + } + + logger._add_metrics_from_log(log=payload, kwargs=kwargs, status_code="200") + + distributions = {s["metric"]: s for s in logger.log_queue if s["type"] == "distribution"} + assert set(distributions) == { + "litellm.request.total_latency.distribution", + "litellm.llm_api.latency.distribution", + "litellm.overhead.latency.distribution", + } + + expected_seconds = { + "litellm.request.total_latency.distribution": 2.0, + "litellm.llm_api.latency.distribution": 1.0, + "litellm.overhead.latency.distribution": 0.25, + } + for metric, seconds in expected_seconds.items(): + ((timestamp, values),) = distributions[metric]["points"] + assert timestamp == int(now.timestamp()) + assert len(values) == 1 + assert abs(values[0] - seconds) < 0.1 + assert "provider:openai" in distributions[metric]["tags"] + assert "model_name:gpt-4o" in distributions[metric]["tags"] + + assert "status_code:200" in distributions["litellm.request.total_latency.distribution"]["tags"] + assert not any( + tag.startswith("status_code:") for tag in distributions["litellm.overhead.latency.distribution"]["tags"] + ) + + +@pytest.mark.asyncio +async def test_extract_tags_omits_hostname_when_unset(clean_env, monkeypatch: pytest.MonkeyPatch): + """An unset HOSTNAME must not produce an empty `HOSTNAME:` tag.""" + monkeypatch.delenv("HOSTNAME", raising=False) + logger = DatadogMetricsLogger(start_periodic_flush=False) + + tags = logger._extract_tags(log=StandardLoggingPayload(model="gpt-4o")) + + assert not any(tag.startswith("HOSTNAME") for tag in tags) + + monkeypatch.setenv("HOSTNAME", "pod-abc") + assert "HOSTNAME:pod-abc" in logger._extract_tags(log=StandardLoggingPayload(model="gpt-4o")) + + @pytest.mark.asyncio async def test_overhead_latency_metric_emitted(clean_env): """Test that litellm.overhead.latency is emitted when hidden_params contains litellm_overhead_time_ms.""" @@ -176,9 +237,9 @@ async def test_overhead_latency_metric_emitted(clean_env): metrics = {s["metric"]: s for s in logger.log_queue} # Overhead metric must be present - assert ( - "litellm.overhead.latency" in metrics - ), f"Expected 'litellm.overhead.latency' in emitted metrics, got: {list(metrics.keys())}" + assert "litellm.overhead.latency" in metrics, ( + f"Expected 'litellm.overhead.latency' in emitted metrics, got: {list(metrics.keys())}" + ) overhead = metrics["litellm.overhead.latency"] assert overhead["type"] == 3 # gauge # 250 ms → 0.25 s @@ -321,9 +382,7 @@ async def test_async_send_batch(clean_env): 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 - ) + 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 = [ @@ -351,6 +410,139 @@ async def test_async_send_batch(clean_env): assert payload["series"][0]["metric"] == "litellm.request.total_latency" +@pytest.mark.asyncio +async def test_async_send_batch_routes_distributions_to_v1_endpoint(clean_env): + """Distribution series go to /api/v1/distribution_points (deflate), gauges/counts stay on /api/v2/series.""" + import gzip + import json + import zlib + + logger = DatadogMetricsLogger(start_periodic_flush=False) + logger.async_client = AsyncMock() + logger.async_client.post.return_value = Response( + 202, json={"status": "ok"}, request=Request("POST", "https://api.test.datadoghq.com") + ) + + timestamp = int(time.time()) + logger.log_queue = [ + { + "metric": "litellm.request.total_latency", + "type": 3, + "points": [{"timestamp": timestamp, "value": 1.5}], + "tags": ["env:test"], + }, + { + "metric": "litellm.request.total_latency.distribution", + "type": "distribution", + "points": ((timestamp, (1.5,)),), + "tags": ("env:test",), + }, + ] + + await logger.async_send_batch() + + calls = {call.args[0]: call.kwargs for call in logger.async_client.post.call_args_list} + assert set(calls) == { + "https://api.test.datadoghq.com/api/v2/series", + "https://api.test.datadoghq.com/api/v1/distribution_points", + } + + series_call = calls["https://api.test.datadoghq.com/api/v2/series"] + assert series_call["headers"]["Content-Encoding"] == "gzip" + series_payload = json.loads(gzip.decompress(series_call["content"])) + assert [s["metric"] for s in series_payload["series"]] == ["litellm.request.total_latency"] + + distribution_call = calls["https://api.test.datadoghq.com/api/v1/distribution_points"] + assert distribution_call["headers"]["Content-Encoding"] == "deflate" + assert distribution_call["headers"]["DD-API-KEY"] == "test_api_key" + distribution_payload = json.loads(zlib.decompress(distribution_call["content"])) + assert distribution_payload == { + "series": [ + { + "metric": "litellm.request.total_latency.distribution", + "type": "distribution", + "points": [[timestamp, [1.5]]], + "tags": ["env:test"], + } + ] + } + + +@pytest.mark.asyncio +async def test_async_send_batch_skips_v2_when_only_distributions_queued(clean_env): + logger = DatadogMetricsLogger(start_periodic_flush=False) + logger.async_client = AsyncMock() + logger.async_client.post.return_value = Response( + 202, json={"status": "ok"}, request=Request("POST", "https://api.test.datadoghq.com") + ) + logger.log_queue = [ + { + "metric": "litellm.llm_api.latency.distribution", + "type": "distribution", + "points": ((int(time.time()), (0.4,)),), + "tags": ("env:test",), + } + ] + + await logger.async_send_batch() + + assert [call.args[0] for call in logger.async_client.post.call_args_list] == [ + "https://api.test.datadoghq.com/api/v1/distribution_points" + ] + + +@pytest.mark.asyncio +async def test_flush_retries_only_distributions_after_v1_failure(clean_env): + logger = DatadogMetricsLogger(start_periodic_flush=False) + logger.async_client = AsyncMock() + v2_url = "https://api.test.datadoghq.com/api/v2/series" + v1_url = "https://api.test.datadoghq.com/api/v1/distribution_points" + responses = { + v2_url: Response(202, json={"status": "ok"}, request=Request("POST", v2_url)), + v1_url: Response(503, json={"errors": ["down"]}, request=Request("POST", v1_url)), + } + timestamp = int(time.time()) + late_count = { + "metric": "litellm.llm_api.request_count", + "type": 1, + "points": [{"timestamp": timestamp, "value": 1}], + "tags": ["env:test"], + } + + def post(url, **_): + if url == v1_url and late_count not in logger.log_queue: + logger.log_queue.append(late_count) + return responses[url] + + logger.async_client.post.side_effect = post + + gauge = { + "metric": "litellm.request.total_latency", + "type": 3, + "points": [{"timestamp": timestamp, "value": 1.5}], + "tags": ["env:test"], + } + distribution = { + "metric": "litellm.request.total_latency.distribution", + "type": "distribution", + "points": ((timestamp, (1.5,)),), + "tags": ("env:test",), + } + logger.log_queue = [gauge, distribution] + + await logger.flush_queue() + + assert logger.log_queue == [distribution, late_count] + + responses[v1_url] = Response(202, json={"status": "ok"}, request=Request("POST", v1_url)) + await logger.flush_queue() + + assert logger.log_queue == [] + assert [call.args[0] for call in logger.async_client.post.call_args_list] == [v2_url, v1_url, v2_url, v1_url] + second_v2 = json.loads(gzip.decompress(logger.async_client.post.call_args_list[2].kwargs["content"])) + assert [s["metric"] for s in second_v2["series"]] == ["litellm.llm_api.request_count"] + + @pytest.mark.asyncio async def test_async_send_batch_empty_queue(clean_env): """Test that async_send_batch does nothing when queue is empty."""