This commit is contained in:
devin-ai-integration[bot] 2026-09-23 14:46:05 +00:00 • committed by GitHub
commit 10bc081c5b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 285 additions and 58 deletions

View file

@ -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
(`<metric>.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:
"""

View file

@ -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, ...]]

View file

@ -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."""