test(e2e): read datadog log delivery back from the real datadog api (#33604)

* test(e2e): read datadog log delivery back from the real datadog api

* test(e2e): compare datadog-read cost with math.isclose, not bit-equality

The response_cost now round-trips through DataDog's attribute indexing
pipeline, whose float serialization is not guaranteed to preserve the
exact bit pattern the proxy shipped. rel_tol=1e-9 (equal to 9 significant
digits) still fails on any real cost discrepancy while tolerating
representation drift. Addresses the Greptile P2 on this PR.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* test(e2e): widen the duplicate-settle window to 30s for real DataDog

Against the local sink one poll interval (5s) after the first hit was
enough to catch a same-call duplicate, because both events arrived in the
same flush batch. Against real DataDog, ingestion jitter can make one
call's two events searchable tens of seconds apart, so a 5s settle could
let the LIT-4447 duplicate slip past the exactly-one assertion. The reader
now keeps re-reading for DD_SETTLE_SECONDS (default 30s, env-overridable
via E2E_DD_SETTLE_SECONDS) after the first event appears, returning early
only when a duplicate is already visible - more waiting cannot clear it.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
yucheng-berri 2026-07-16 16:06:35 -07:00 • committed by GitHub
parent c6778b79c3
commit ab6d7578ce
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 188 additions and 183 deletions

View file

@ -1,41 +1,5 @@
# local setup to run e2e tests
configs:
dd_sink_script:
content: |
# Minimal DataDog logs-intake sink for the logging suite: records every
# POST (gunzipping the compressed batches the integration sends) and
# replays them as JSON on GET /requests so tests can assert delivery.
import gzip, json
from http.server import BaseHTTPRequestHandler, HTTPServer
REQUESTS = []
class Handler(BaseHTTPRequestHandler):
def do_POST(self):
body = self.rfile.read(int(self.headers.get("Content-Length", 0)))
if self.headers.get("Content-Encoding") == "gzip":
body = gzip.decompress(body)
REQUESTS.append({"path": self.path, "body": body.decode("utf-8", "replace")})
self.send_response(202)
self.end_headers()
self.wfile.write(b"{}")
def do_GET(self):
self.send_response(200)
if self.path == "/health":
self.send_header("Content-Type", "text/plain")
self.end_headers()
self.wfile.write(b"ok")
return
self.send_header("Content-Type", "application/json")
self.end_headers()
self.wfile.write(json.dumps({"requests": REQUESTS}).encode())
def log_message(self, *args):
pass
HTTPServer(("0.0.0.0", 8080), Handler).serve_forever()
litellm_config:
content: |
general_settings:
@ -129,15 +93,16 @@ services:
condition: service_healthy
jaeger:
condition: service_healthy
dd-sink:
condition: service_healthy
env_file: .env
environment:
LITELLM_MASTER_KEY: sk-1234
STORE_MODEL_IN_DB: "True"
DD_API_KEY: local-sink-noauth
DD_SITE: datadoghq.com
DD_BASE_URL: http://dd-sink:8080
# Real DataDog delivery (no local sink): the key comes from the
# environment - the cluster's secret manager injects it, locally
# tests/e2e/.env provides it. Tests read delivery back via the DataDog
# Logs Search API (DD_APP_KEY, test-side only - see logging/datadog_reader.py).
DD_API_KEY: ${DD_API_KEY:-}
DD_SITE: ${DD_SITE:-datadoghq.com}
LITELLM_OTEL_V2: "true"
PHOENIX_COLLECTOR_HTTP_ENDPOINT: http://jaeger:4318/v1/traces
PHOENIX_API_KEY: local-jaeger-noauth
@ -198,19 +163,3 @@ services:
interval: 3s
timeout: 3s
retries: 20
# throwaway DataDog logs-intake sink (records POSTs, replays on GET /requests;
# see E2E_DD_SINK_URL)
dd-sink:
image: python:3.12-alpine
command: ["python", "/sink.py"]
configs:
- source: dd_sink_script
target: /sink.py
ports:
- "9915:8080"
healthcheck:
test: ["CMD", "wget", "-qO-", "http://127.0.0.1:8080/health"]
interval: 3s
timeout: 3s
retries: 20

View file

@ -32,9 +32,19 @@ CHEAP_OPENAI_MODEL = os.environ.get("E2E_CHEAP_OPENAI_MODEL", "gpt-5.5")
# read exported spans back through it.
OTEL_QUERY_URL = os.environ.get("E2E_OTEL_QUERY_URL", "http://localhost:16686").rstrip("/")
# Query URL of the compose stack's DataDog logs-intake sink (the `dd-sink`
# service records every intake POST and replays them on GET /requests).
DD_SINK_URL = os.environ.get("E2E_DD_SINK_URL", "http://localhost:9915").rstrip("/")
# Real-DataDog read-back (no local sink - destination fakes cannot be deployed
# on the cluster): the proxy delivers with DD_API_KEY as in production, and the
# tests read ingested events back through the DataDog Logs Search API, which
# additionally needs an application key. On the cluster the secret manager
# injects both; locally tests/e2e/.env provides them.
DD_SITE = os.environ.get("DD_SITE", "datadoghq.com").strip()
DD_API_KEY = os.environ.get("DD_API_KEY", "").strip()
DD_APP_KEY = os.environ.get("DD_APP_KEY", "").strip()
# After the first event is searchable, keep watching this long for a late
# duplicate before the exactly-one assertion: real-DataDog ingestion jitter can
# make one call's two events searchable tens of seconds apart, and a duplicate
# that surfaces late IS the bug (LIT-4447), so one poll interval is not enough.
DD_SETTLE_SECONDS = float(os.environ.get("E2E_DD_SETTLE_SECONDS", "30"))
# Writes on the proxy are eventually consistent (e.g. spend rows flush on
# proxy_batch_write_at, ~60s). Read-backs poll to this deadline, never sleep-once.

View file

@ -11,7 +11,7 @@ import os
import pytest
from logging_client import LangfuseCreds, LoggingClient, build_logging_client, load_langfuse_creds
from datadog_sink import DdSinkReader, build_dd_sink_reader
from datadog_reader import DdLogsReader, build_dd_logs_reader
from otel_client import OtelReader, build_otel_reader
@ -37,9 +37,10 @@ def otel_reader() -> OtelReader:
@pytest.fixture(scope="session")
def dd_sink() -> DdSinkReader:
"""Read-back client for the compose stack's DataDog logs-intake sink."""
return build_dd_sink_reader()
def dd_logs() -> DdLogsReader:
"""Read-back client for the real DataDog Logs Search API (keys from the
secret manager on the cluster, tests/e2e/.env locally)."""
return build_dd_logs_reader()
@pytest.fixture

View file

@ -0,0 +1,140 @@
"""Read-back for the DataDog logging tests against the real DataDog Logs
Search API.
Delivery is judged on what DataDog itself ingested: the proxy ships logs with
DD_API_KEY exactly as in production (no base-URL override, no local sink), and
the tests search the ingested events back with POST /api/v2/logs/events/search,
authenticated with the same DD_API_KEY plus a DD_APP_KEY application key. On
the cluster the secret manager injects both keys; locally tests/e2e/.env
provides them. Missing keys or a failed search call are hard failures, 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, Field
from e2e_config import (
DD_API_KEY,
DD_APP_KEY,
DD_SETTLE_SECONDS,
DD_SITE,
POLL_INTERVAL,
POLL_TIMEOUT,
)
from e2e_http import URL, Headers, Success, post
class _DdAuthHeaders(Headers):
api_key: str = Field(serialization_alias="DD-API-KEY")
app_key: str = Field(serialization_alias="DD-APPLICATION-KEY")
class _SearchFilter(BaseModel):
query: str
#: Wide enough to cover a full suite run plus DataDog's ingestion lag;
#: markers are unique per test, so a wide window cannot match foreign events.
from_: str = Field(default="now-30m", serialization_alias="from")
to: str = "now"
class _SearchPage(BaseModel):
limit: int = 100
class _SearchRequest(BaseModel):
filter: _SearchFilter
page: _SearchPage = _SearchPage()
sort: str = "timestamp"
class DdLogEvent(BaseModel):
"""One ingested log event as the search API returns it: the indexed
envelope (service/status/tags) plus ``attributes`` - DataDog's parse of the
JSON message the integration shipped, i.e. the StandardLoggingPayload
fields."""
model_config = ConfigDict(extra="ignore")
service: str | None = None
status: str | None = None
tags: list[str] = []
attributes: dict[str, object] = {}
class _SearchEvent(BaseModel):
model_config = ConfigDict(extra="ignore")
attributes: DdLogEvent
class _SearchResponse(BaseModel):
model_config = ConfigDict(extra="ignore")
data: list[_SearchEvent] = []
@dataclass(frozen=True, slots=True)
class DdLogsReader:
site: str
api_key: str
app_key: str
def events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""Every ingested event matching the marker (full-text, exact phrase).
More than one hit for one call IS the duplicate-delivery bug, so this
never collapses to a single event."""
result = post(
URL(f"https://api.{self.site}/api/v2/logs/events/search"),
headers=_DdAuthHeaders(api_key=self.api_key, app_key=self.app_key),
json=_SearchRequest(filter=_SearchFilter(query=f'"{marker}"')),
response_type=_SearchResponse,
timeout=30.0,
)
match result:
case Success(data=page):
return [event.attributes for event in page.data]
case failure:
pytest.fail(f"DataDog Logs Search API at api.{self.site} failed: {failure}")
def poll_events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""Poll until at least one matching event is searchable (the callback
flushes in periodic batches and DataDog ingestion adds seconds of lag),
then keep re-reading for DD_SETTLE_SECONDS so a late duplicate cannot
hide from the exactly-one assertion - real-DataDog jitter can surface
one call's two events tens of seconds apart. 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:
return self._settled_events_for_marker(marker, events)
time.sleep(POLL_INTERVAL)
return self.events_for_marker(marker)
def _settled_events_for_marker(
self, marker: str, events: list[DdLogEvent]
) -> list[DdLogEvent]:
"""Re-read at every poll interval until the settle window closes; a
duplicate ends the watch early because more waiting cannot clear it."""
settle_deadline = time.monotonic() + DD_SETTLE_SECONDS
while time.monotonic() < settle_deadline:
time.sleep(POLL_INTERVAL)
events = self.events_for_marker(marker)
if len(events) > 1:
return events
return events
def build_dd_logs_reader() -> DdLogsReader:
if not DD_API_KEY or not DD_APP_KEY:
pytest.fail(
"DD_API_KEY and DD_APP_KEY must be set: the DataDog tests deliver to and "
"read back from the real DataDog API (on the cluster the secret manager "
"injects them; locally set them in tests/e2e/.env)"
)
return DdLogsReader(site=DD_SITE, api_key=DD_API_KEY, app_key=DD_APP_KEY)

View file

@ -1,103 +0,0 @@
"""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)

View file

@ -3,24 +3,27 @@
Covers logging.datadog.success.exports_metric: one successful call on each
route must reach the DataDog logs intake as EXACTLY ONE log event whose
message (the StandardLoggingPayload) carries the model, the token counts, and
the response cost. Delivery is judged on what the intake actually received:
the compose stack's dd-sink service records every batch the datadog callback
ships (DD_BASE_URL override) and the tests read it back, so a dropped event, a
duplicated event, or a payload missing the cost all fail here.
the response cost. Delivery is judged on what DataDog itself ingested: the
proxy ships with DD_API_KEY exactly as in production, and the tests search the
events back through the DataDog Logs Search API (DD_APP_KEY, keys from the
secret manager on the cluster), so a dropped event, a duplicated event, or a
payload missing the cost all fail here.
Both halves of the contract are asserted: the recorded state (the proxy
reports the DataDogLogger callback active via /health/readiness/details) and
the enforced behavior (the event at the intake, with the cost cross-checked
exactly against the x-litellm-response-cost header of the very response the
caller received).
against the x-litellm-response-cost header of the very response the caller
received).
"""
from __future__ import annotations
import math
import pytest
from pydantic import BaseModel, ConfigDict
from datadog_sink import DdLogEvent, DdSinkReader
from datadog_reader import DdLogEvent, DdLogsReader
from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import NoBody, StreamingResponse
from lifecycle import ResourceManager
@ -71,10 +74,12 @@ def _assert_exactly_one_event(
"for the currently known /v1/messages instance)"
)
event = events[0]
assert event.ddsource == "litellm", f"event ddsource must be litellm, got {event.ddsource!r}"
assert "source:litellm" in event.tags, (
f"the ingested event must carry the litellm source (shipped as ddsource), got tags {event.tags!r}"
)
assert event.status == "info", f"success events ship at status info, got {event.status!r}"
payload = _DdMessagePayload.model_validate_json(event.message)
payload = _DdMessagePayload.model_validate(event.attributes)
assert payload.status == "success", f"payload status must be success, got {payload.status!r}"
assert payload.model_group == model_group, (
f"payload model_group must be {model_group!r}, got {payload.model_group!r}"
@ -86,7 +91,10 @@ def _assert_exactly_one_event(
assert outcome.response_cost is not None and outcome.response_cost > 0, (
f"the response must report x-litellm-response-cost, got {outcome.response_cost!r}"
)
assert abs(payload.response_cost - outcome.response_cost) < 1e-12, (
# Relative tolerance, not bit-equality: the cost round-trips through
# DataDog's attribute indexing, whose float serialization may drift in the
# last bits; 9 significant digits still catches any real cost discrepancy.
assert math.isclose(payload.response_cost, outcome.response_cost, rel_tol=1e-9), (
f"payload response_cost {payload.response_cost} must equal the response header "
f"cost {outcome.response_cost}"
)
@ -95,7 +103,7 @@ def _assert_exactly_one_event(
class TestDataDogLogDelivery:
@pytest.mark.covers("logging.datadog.success.exports_metric", exercised_on=["chat_completions"])
def test_chat_completions_emits_one_log_event(
self, client: LoggingClient, dd_sink: DdSinkReader, resources: ResourceManager
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
"""One successful non-streaming /chat/completions call must reach the
DataDog logs intake as exactly one log event whose payload carries the
@ -110,14 +118,14 @@ class TestDataDogLogDelivery:
client,
lambda: client.chat_raw(key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", max_tokens=16),
)
events = dd_sink.poll_events_for_marker(marker)
events = dd_logs.poll_events_for_marker(marker)
_assert_exactly_one_event(
events, model_group=CHEAP_ANTHROPIC_MODEL, call_type="acompletion", outcome=outcome
)
@pytest.mark.covers("logging.datadog.success.exports_metric", exercised_on=["messages"])
def test_messages_emits_one_log_event(
self, client: LoggingClient, dd_sink: DdSinkReader, resources: ResourceManager
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
"""One successful non-streaming /v1/messages call must reach the
DataDog logs intake as exactly one log event whose payload carries the
@ -134,14 +142,14 @@ class TestDataDogLogDelivery:
client,
lambda: client.messages_raw(key, CHEAP_ANTHROPIC_MODEL, f"reply with one word {marker}", max_tokens=16),
)
events = dd_sink.poll_events_for_marker(marker)
events = dd_logs.poll_events_for_marker(marker)
_assert_exactly_one_event(
events, model_group=CHEAP_ANTHROPIC_MODEL, call_type="anthropic_messages", outcome=outcome
)
@pytest.mark.covers("logging.datadog.success.exports_metric", exercised_on=["responses"])
def test_responses_emits_one_log_event(
self, client: LoggingClient, dd_sink: DdSinkReader, resources: ResourceManager
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
"""One successful non-streaming /v1/responses call must reach the
DataDog logs intake as exactly one log event whose payload carries the
@ -156,7 +164,7 @@ class TestDataDogLogDelivery:
client,
lambda: client.responses_raw(key, CHEAP_OPENAI_MODEL, f"reply with one word {marker}"),
)
events = dd_sink.poll_events_for_marker(marker)
events = dd_logs.poll_events_for_marker(marker)
_assert_exactly_one_event(
events, model_group=CHEAP_OPENAI_MODEL, call_type="aresponses", outcome=outcome
)