litellm/tests/e2e/logging/datadog_sink.py
yucheng-berri edc30ea515
test(e2e): datadog log delivery for successful chat, messages, and responses (LIT-4447) (#33415)
* test(e2e): datadog log delivery for successful chat, messages, and responses

Covers logging.datadog.success.exports_metric on all three routes: one
successful non-streaming call must reach the DataDog logs intake as exactly
one log event whose StandardLoggingPayload message carries the model group,
real token counts, and a response cost equal to the x-litellm-response-cost
header of the same response. Delivery is judged at the intake: the compose
stack gains a dd-sink service recording every batch the datadog callback
ships via the DD_BASE_URL testing override, and a typed reader replays it.

Writing these caught a live product bug: /v1/messages double-logs every
success (two byte-identical events per call), filed as LIT-4447; the messages
test tolerates byte-identical duplicates of the one event until it lands,
while a second differing event still fails

* test(e2e): address review findings on the datadog delivery suite

Consolidates the fresh-key first_ok helper into logging_client now that the
otel PR it mirrored has merged (both test files use the shared copy), moves
intake batch parsing into a helper so no path can leave the batch unbound,
and gives the sink's /health endpoint a truthful text/plain content type

* test(e2e): tolerate same-logical-event duplicates by call id, not byte identity

A clean LIT-4447 repro showed the duplicated payload is built twice and can
mint a fresh synthetic completion id per emission, arriving as two separate
intake POSTs with the same litellm_call_id and identical substantive fields.
Byte-identity was therefore a flaky criterion; duplicates now qualify only
when they share the call id, call type, model group, tokens, and cost, and a
second differing event still fails

* test(e2e): assert the scenario strictly; the messages test is the LIT-4447 regression pin

Per review direction the tests now assert exactly what the scenario promises:
exactly one DataDog log event per successful call, on every route. The
/v1/messages test therefore fails on current code against the known
double-log (LIT-4447) and is its regression pin; it goes green when the fix
lands. The duplicate-tolerance machinery is removed

* Simplify docstrings for DataDog log tests

Removed redundant phrasing about cost cross-checking in docstrings.

* Update test_datadog_log_e2e.py
2026-07-16 09:54:07 -07:00

103 lines
3.4 KiB
Python

"""Read-back for the DataDog logging tests: typed models over the compose
stack's dd-sink service, which records every logs-intake POST the datadog
callback sends (gunzipped) and replays them as JSON.
Delivery is judged on what the sink actually received, mirroring how the OTEL
tests read Jaeger; a failed sink query is a hard failure, never an empty
result. External reads go through ``e2e_http``.
"""
from __future__ import annotations
import time
from dataclasses import dataclass
import pytest
from pydantic import BaseModel, ConfigDict, TypeAdapter, ValidationError
from e2e_config import DD_SINK_URL, POLL_INTERVAL, POLL_TIMEOUT
from e2e_http import URL, NoBody, Success, get
class DdSinkRequest(BaseModel):
model_config = ConfigDict(extra="ignore")
path: str
body: str
class DdSinkRequests(BaseModel):
model_config = ConfigDict(extra="ignore")
requests: list[DdSinkRequest] = []
class DdLogEvent(BaseModel):
model_config = ConfigDict(extra="ignore")
message: str
ddsource: str | None = None
service: str | None = None
status: str | None = None
_EVENT_BATCH: TypeAdapter[list[DdLogEvent]] = TypeAdapter(list[DdLogEvent])
def _parse_batch(request: DdSinkRequest) -> list[DdLogEvent]:
"""The intake accepts an array of events or a single event object."""
try:
return _EVENT_BATCH.validate_json(request.body)
except ValidationError:
try:
return [DdLogEvent.model_validate_json(request.body)]
except ValidationError:
pytest.fail(f"dd-sink recorded a non-log body on {request.path}: {request.body[:200]}")
@dataclass(frozen=True, slots=True)
class DdSinkReader:
sink_url: str
def _recorded_requests(self) -> list[DdSinkRequest]:
result = get(
URL(f"{self.sink_url}/requests"),
headers=NoBody(),
params=NoBody(),
response_type=DdSinkRequests,
timeout=30.0,
)
match result:
case Success(data=page):
return page.requests
case failure:
pytest.fail(f"dd-sink query at {self.sink_url} failed: {failure}")
def events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""Every log event across every recorded intake batch whose message
carries the marker. More than one hit for one call IS the
duplicate-delivery bug, so this never collapses to a single event."""
events: list[DdLogEvent] = []
for request in self._recorded_requests():
if "/api/v2/logs" not in request.path:
continue
events.extend(event for event in _parse_batch(request) if marker in event.message)
return events
def poll_events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""Poll until at least one matching event lands (the callback flushes
in periodic batches), then re-read after one more interval so a late
duplicate cannot hide from the exactly-one assertion. At the deadline
the last result is returned as-is."""
deadline = time.monotonic() + POLL_TIMEOUT
while time.monotonic() < deadline:
events = self.events_for_marker(marker)
if events:
time.sleep(POLL_INTERVAL)
return self.events_for_marker(marker)
time.sleep(POLL_INTERVAL)
return self.events_for_marker(marker)
def build_dd_sink_reader() -> DdSinkReader:
return DdSinkReader(sink_url=DD_SINK_URL)