test(e2e): tolerate the readiness 503 from a transient db blip in the callback-config probes

This commit is contained in:
Yucheng Zhu 2026-08-28 18:35:30 -07:00
parent 3e036dcdca
commit 72353b3a2c
6 changed files with 41 additions and 73 deletions

View file

@ -47,6 +47,4 @@ def dd_logs() -> DdLogsReader:
def datadog_creds() -> None:
"""Require Datadog shipping credentials. Hard-fail when absent; never skip."""
if not (os.getenv("DD_API_KEY") and os.getenv("DD_SITE")):
pytest.fail(
"Datadog e2e requires DD_API_KEY and DD_SITE; missing credentials is a hard failure, not a skip"
)
pytest.fail("Datadog e2e requires DD_API_KEY and DD_SITE; missing credentials is a hard failure, not a skip")

View file

@ -480,12 +480,8 @@ class LoggingClient:
stream=True if stream else None,
)
if stream:
return self.proxy.transport.stream(
"/v1/messages", headers=self.proxy.transport.bearer(key), json=body
)
return self.proxy.transport.send(
"/v1/messages", headers=self.proxy.transport.bearer(key), json=body
)
return self.proxy.transport.stream("/v1/messages", headers=self.proxy.transport.bearer(key), json=body)
return self.proxy.transport.send("/v1/messages", headers=self.proxy.transport.bearer(key), json=body)
def responses_raw(
self, key: str, model: str, text: str, *, max_output_tokens: int = 64, stream: bool = False
@ -499,12 +495,8 @@ class LoggingClient:
model=model, input=text, max_output_tokens=max_output_tokens, stream=True if stream else None
)
if stream:
return self.proxy.transport.stream(
"/v1/responses", headers=self.proxy.transport.bearer(key), json=body
)
return self.proxy.transport.send(
"/v1/responses", headers=self.proxy.transport.bearer(key), json=body
)
return self.proxy.transport.stream("/v1/responses", headers=self.proxy.transport.bearer(key), json=body)
return self.proxy.transport.send("/v1/responses", headers=self.proxy.transport.bearer(key), json=body)
def scrape_metrics(self) -> str:
return self.proxy.probe("/metrics", params=NoBody()).body
@ -530,9 +522,7 @@ class LoggingClient:
return False
return True
rows = self.proxy.poll_logs_for_key(
key, min_rows=1, predicate=lambda rs: any(_matches(r) for r in rs)
)
rows = self.proxy.poll_logs_for_key(key, min_rows=1, predicate=lambda rs: any(_matches(r) for r in rs))
for row in rows:
if _matches(row):
return row
@ -593,9 +583,7 @@ class LoggingClient:
deadline = time.monotonic() + POLL_TIMEOUT
last: LangfuseObservation | None = None
while time.monotonic() < deadline:
last = self.find_langfuse_observation(
creds, key_alias=key_alias, prompt_marker=prompt_marker
)
last = self.find_langfuse_observation(creds, key_alias=key_alias, prompt_marker=prompt_marker)
if last is not None:
cost = observation_spend(last)
if not require_positive_cost or (cost is not None and cost > 0):
@ -611,9 +599,7 @@ class LoggingClient:
prompt_marker: str,
) -> list[LangfuseObservation]:
"""Generation plus any sibling/child observations (guardrail spans, etc.)."""
gen = self.poll_langfuse_observation(
creds, key_alias=key_alias, prompt_marker=prompt_marker
)
gen = self.poll_langfuse_observation(creds, key_alias=key_alias, prompt_marker=prompt_marker)
if gen is None or not gen.trace_id:
return [] if gen is None else [gen]
return self.list_langfuse_observations(creds, trace_id=gen.trace_id) or [gen]
@ -636,3 +622,15 @@ def first_ok(client: LoggingClient, send: Callable[[], StreamingResponse]) -> St
def build_logging_client(proxy: ProxyClient) -> LoggingClient:
return LoggingClient(proxy=proxy)
def readiness_details_body(client: LoggingClient) -> str:
"""/health/readiness/details, tolerating the 503 it serves while the
ephemeral stack's DB leg blips: the recorded state the logging suites check
here is the callback list, which the body carries either way."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
db_blip = result.status_code == 503 and '"db":"disconnected"' in result.body
assert result.status_code == 200 or db_blip, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
return result.body

View file

@ -26,9 +26,8 @@ from pydantic import BaseModel, ConfigDict
from datadog_reader import DdLogEvent, DdLogsReader
from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import NoBody
from lifecycle import ResourceManager
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok, readiness_details_body
from models import LiteLLMParamsBody
pytestmark = pytest.mark.e2e
@ -55,13 +54,10 @@ def _assert_datadog_configured(client: LoggingClient) -> None:
"""Recorded state: the proxy reports the DataDog callback among its active
callbacks, so a missing destination config fails here, before any
delivery-based assertion can time out confusingly."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
assert result.status_code == 200, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
assert DD_LOGGER_NAME in result.body, (
body = readiness_details_body(client)
assert DD_LOGGER_NAME in body, (
f"the proxy must report the {DD_LOGGER_NAME} callback active "
f"(callbacks + DD_* env in the compose config); got: {result.body[:400]}"
f"(callbacks + DD_* env in the compose config); got: {body[:400]}"
)

View file

@ -22,10 +22,9 @@ import math
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_http import NoBody
from gcs_reader import GcsLogReader, build_gcs_reader, utc_now
from lifecycle import ResourceManager
from logging_client import LoggingClient, completion_response_id, first_ok
from logging_client import LoggingClient, completion_response_id, first_ok, readiness_details_body
pytestmark = pytest.mark.e2e
@ -43,14 +42,11 @@ def _assert_gcs_configured(client: LoggingClient) -> None:
active callbacks, so a missing destination config (or a missing enterprise
license - gcs_bucket refuses to initialize without one) fails here, before
any delivery-based assertion can time out confusingly."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
assert result.status_code == 200, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
assert GCS_LOGGER_NAME in result.body, (
body = readiness_details_body(client)
assert GCS_LOGGER_NAME in body, (
f"the proxy must report the {GCS_LOGGER_NAME} callback active "
f"(litellm_settings.callbacks: ['gcs_bucket'] + GCS_BUCKET_NAME env + enterprise license); "
f"got: {result.body[:400]}"
f"got: {body[:400]}"
)

View file

@ -23,9 +23,8 @@ import pytest
from pydantic import BaseModel, ConfigDict, ValidationError
from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import NoBody
from lifecycle import ResourceManager
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok, readiness_details_body
from models import LiteLLMParamsBody
from otel_client import JaegerSpan, JaegerTrace, OtelReader
@ -48,11 +47,7 @@ def _assert_otel_destination_configured(client: LoggingClient) -> None:
"""Recorded state: the proxy reports the OTEL v2 logger among its active
callbacks, so a missing/failed destination config fails here, before any
traffic-based assertion can time out confusingly."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
assert result.status_code == 200, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
details = _ReadinessDetails.model_validate_json(result.body)
details = _ReadinessDetails.model_validate_json(readiness_details_body(client))
assert OTEL_V2_LOGGER_NAME in details.success_callbacks, (
f"the proxy must report the {OTEL_V2_LOGGER_NAME} callback active "
f"(LITELLM_OTEL_V2 + arize_phoenix preset in the compose config); got: {details.success_callbacks}"
@ -164,17 +159,14 @@ def served_genai_spans(trace: JaegerTrace, genai_span: str) -> list[JaegerSpan]:
these tests fail whenever the upstream 429s, 529s, or hands back a stale
credential on the first try."""
return [
span
for span in trace.spans
if span.operation_name == genai_span and _tag(span, ERROR_STATUS_TAG) != "ERROR"
span for span in trace.spans if span.operation_name == genai_span and _tag(span, ERROR_STATUS_TAG) != "ERROR"
]
def one_served_genai_span(trace: JaegerTrace, genai_span: str) -> JaegerSpan:
served = served_genai_spans(trace, genai_span)
assert len(served) == 1, (
f"a streamed call must produce exactly ONE served gen-AI span, got {len(served)}; "
f"spans: {trace.span_names()}"
f"a streamed call must produce exactly ONE served gen-AI span, got {len(served)}; spans: {trace.span_names()}"
)
return served[0]
@ -190,8 +182,7 @@ def _assert_real_ttft(hits: list[JaegerTrace], *, genai_span: str) -> None:
"(nothing tagged with its call id was found)"
)
assert len(hits) == 1, (
f"expected exactly ONE trace for the call, got {len(hits)}: "
f"{[(t.trace_id, t.span_names()) for t in hits]}"
f"expected exactly ONE trace for the call, got {len(hits)}: {[(t.trace_id, t.span_names()) for t in hits]}"
)
trace = hits[0]
span = one_served_genai_span(trace, genai_span)
@ -280,9 +271,7 @@ def _assert_error_span_contract(span: JaegerSpan) -> None:
"the span status description must carry the same untruncated message as error.message"
)
stack = _tag(span, "litellm.provider.error.stack_trace")
assert isinstance(stack, str) and stack, (
"the error span must carry a non-empty litellm.provider.error.stack_trace"
)
assert isinstance(stack, str) and stack, "the error span must carry a non-empty litellm.provider.error.stack_trace"
class TestOtelTraceCompleteness:
@ -313,9 +302,7 @@ class TestOtelTraceCompleteness:
resources.defer(lambda: client.delete_key(key))
marker = unique_marker()
outcome = first_ok(
client, lambda: client.chat_raw(key, MODEL, f"reply with one word {marker}", max_tokens=16)
)
outcome = first_ok(client, lambda: client.chat_raw(key, MODEL, f"reply with one word {marker}", max_tokens=16))
assert outcome.call_id is not None, "success response must carry x-litellm-call-id"
hits = otel_reader.poll_traces_for_call(
@ -520,9 +507,7 @@ class TestOtelTraceCompleteness:
route = "/v1/responses"
_assert_otel_destination_configured(client)
key = client.key_with_alias(
f"otel-stream-responses-{unique_marker()}", models=[CHEAP_OPENAI_MODEL]
)
key = client.key_with_alias(f"otel-stream-responses-{unique_marker()}", models=[CHEAP_OPENAI_MODEL])
resources.defer(lambda: client.delete_key(key))
marker = unique_marker()
@ -660,9 +645,7 @@ class TestOtelTraceCompleteness:
route = "/v1/responses"
_assert_otel_destination_configured(client)
key = client.key_with_alias(
f"otel-ttft-responses-{unique_marker()}", models=[CHEAP_OPENAI_MODEL]
)
key = client.key_with_alias(f"otel-ttft-responses-{unique_marker()}", models=[CHEAP_OPENAI_MODEL])
resources.defer(lambda: client.delete_key(key))
marker = unique_marker()

View file

@ -27,13 +27,13 @@ import time
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_http import NoBody
from lifecycle import ResourceManager
from logging_client import (
INVALID_UPSTREAM_API_KEY,
LoggingClient,
completion_response_id,
first_ok,
readiness_details_body,
)
from models import LiteLLMParamsBody
from s3_reader import S3LogReader, build_s3_reader
@ -53,14 +53,11 @@ def _assert_s3_configured(client: LoggingClient) -> None:
"""Recorded state: the proxy reports the s3_v2 callback among its active
callbacks, so a missing destination config fails here, before any
delivery-based assertion can time out confusingly."""
result = client.proxy.probe("/health/readiness/details", params=NoBody())
assert result.status_code == 200, (
f"/health/readiness/details must answer 200, got {result.status_code}: {result.body[:300]}"
)
assert S3_LOGGER_NAME in result.body, (
body = readiness_details_body(client)
assert S3_LOGGER_NAME in body, (
f"the proxy must report the {S3_LOGGER_NAME} callback active "
f"(litellm_settings.callbacks: ['s3_v2'] + s3_callback_params in the proxy config); "
f"got: {result.body[:400]}"
f"got: {body[:400]}"
)