mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-08 03:08:45 +00:00
test(e2e): record steps for raw transport calls and poll helpers (#45150)
Add @step labels to the HttpTransport methods, the poll and wait helpers and the boot helpers that did real IO without recording a step, so a test that reaches the proxy through them no longer reports an empty or gappy step timeline in the JUnit report.
This commit is contained in:
parent
048a1500df
commit
679f7e636e
10 changed files with 43 additions and 0 deletions
|
|
@ -17,6 +17,7 @@ import re
|
|||
import string
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
import warnings
|
||||
from collections import Counter
|
||||
from collections.abc import Callable, Generator, Iterator, Mapping
|
||||
|
|
@ -44,6 +45,7 @@ from e2e_metadata import (
|
|||
environment_secrets,
|
||||
meta,
|
||||
step,
|
||||
step_properties,
|
||||
subject_properties,
|
||||
)
|
||||
from junit_properties import package_from_nodeid, result_properties, source_from_item
|
||||
|
|
@ -319,6 +321,15 @@ class TestStepRecording:
|
|||
delete_team()
|
||||
assert [Path(warning.filename).name for warning in caught] == [Path(__file__).name]
|
||||
|
||||
def test_a_harness_wait_the_test_calls_directly_is_a_step_in_its_report(self) -> None:
|
||||
"""A test that only waits through a bare harness helper, never a typed
|
||||
client, still has that wait in its JUnit story. The stamp is old enough
|
||||
that the helper returns without sleeping."""
|
||||
from e2e_config import PROPAGATION_TIMEOUT, settle_propagation
|
||||
|
||||
settle_propagation(written_at=time.monotonic() - PROPAGATION_TIMEOUT)
|
||||
assert step_properties() == (("step", "Wait for the last control-plane write to reach every proxy replica"),)
|
||||
|
||||
|
||||
class _KeyBody(BaseModel):
|
||||
models: list[str] = []
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ from pathlib import Path
|
|||
from typing import Final
|
||||
|
||||
from dotenv import load_dotenv
|
||||
from e2e_metadata import step
|
||||
from fixture_mode import deterministic_marker, parse_fixture_mode, registration_owner
|
||||
from provider_edge import provider_edge_api_base
|
||||
from pydantic import TypeAdapter
|
||||
|
|
@ -308,6 +309,7 @@ def available_port() -> int:
|
|||
return TypeAdapter(tuple[str, int]).validate_python(listener.getsockname())[1]
|
||||
|
||||
|
||||
@step("Wait for the last control-plane write to reach every proxy replica")
|
||||
def settle_propagation(written_at: float) -> None:
|
||||
"""Block until PROPAGATION_TIMEOUT has elapsed since `written_at`, a
|
||||
`time.monotonic()` stamp taken the moment a control-plane write returned.
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ from typing import Final, Generic, Literal, NewType, Protocol, TypeVar, cast
|
|||
|
||||
import pytest
|
||||
import requests
|
||||
from e2e_metadata import step
|
||||
from pydantic import BaseModel, ConfigDict, Field
|
||||
|
||||
URL = NewType("URL", str)
|
||||
|
|
@ -473,6 +474,7 @@ def get[R: BaseModel](
|
|||
return classify(resp, response_type)
|
||||
|
||||
|
||||
@step("GET the external URL {url}")
|
||||
def get_external[R: BaseModel](
|
||||
url: str,
|
||||
*,
|
||||
|
|
|
|||
|
|
@ -607,6 +607,7 @@ def build_client(proxy: ProxyClient) -> GuardrailsClient:
|
|||
return GuardrailsClient(proxy=proxy)
|
||||
|
||||
|
||||
@step("Retry the call until the guardrail {guardrail_name} is applied")
|
||||
def poll_until_guardrail_applied(
|
||||
call: Callable[[], StreamingResponse],
|
||||
guardrail_name: str,
|
||||
|
|
@ -630,6 +631,7 @@ def poll_until_guardrail_applied(
|
|||
return result
|
||||
|
||||
|
||||
@step("Retry the call until a guardrail blocks it")
|
||||
def poll_until_blocked[R: BaseModel](call: Callable[[], Result[R]]) -> Result[R]:
|
||||
"""Retry a call that a guardrail should reject until it is, returning the last result.
|
||||
|
||||
|
|
@ -657,6 +659,7 @@ def poll_until_blocked[R: BaseModel](call: Callable[[], Result[R]]) -> Result[R]
|
|||
_TRANSIENT_STREAM_STATUSES = frozenset({-1, 401, 429})
|
||||
|
||||
|
||||
@step("Retry the streamed call until a guardrail blocks it")
|
||||
def poll_until_blocked_stream(call: Callable[[], StreamingResponse]) -> StreamingResponse:
|
||||
"""poll_until_blocked for raw/streamed sends, which return a StreamingResponse
|
||||
instead of a Result: retry while the call still succeeds (the data-plane worker
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ from itertools import accumulate
|
|||
from pathlib import Path
|
||||
from typing import Final
|
||||
|
||||
from e2e_metadata import step
|
||||
from pydantic import BaseModel, TypeAdapter
|
||||
|
||||
_LOCUSTFILE = Path(__file__).with_name("locustfile.py")
|
||||
|
|
@ -180,6 +181,7 @@ def read_generator_warnings(stderr: str) -> tuple[str, ...]:
|
|||
return tuple(dict.fromkeys(saturated))
|
||||
|
||||
|
||||
@step("Drive {users} locust users at {endpoints} for {duration_seconds}s")
|
||||
def run_gateway_load(
|
||||
*,
|
||||
base_url: str,
|
||||
|
|
|
|||
|
|
@ -9,6 +9,7 @@ from pydantic import BaseModel
|
|||
|
||||
from e2e_config import unique_marker
|
||||
from e2e_http import Result, Success
|
||||
from e2e_metadata import step
|
||||
from models import CacheControl, RichMessage, TextBlock
|
||||
from transport import Transport
|
||||
|
||||
|
|
@ -230,6 +231,7 @@ def run_session(
|
|||
)
|
||||
|
||||
|
||||
@step("Run {sessions} concurrent sessions of {turns_per_session} turns against {model}")
|
||||
def run_concurrent_sessions(
|
||||
transport: Transport,
|
||||
key: str,
|
||||
|
|
@ -248,6 +250,7 @@ def run_concurrent_sessions(
|
|||
return tuple(turn for future in futures for turn in future.result())
|
||||
|
||||
|
||||
@step("Poll the key's spend until it holds steady for {settle_seconds}s")
|
||||
def settled_spend(
|
||||
read_spend: Callable[[], float],
|
||||
poll_interval: float,
|
||||
|
|
|
|||
|
|
@ -704,6 +704,7 @@ class LoggingClient:
|
|||
return self.list_langfuse_observations(creds, trace_id=gen.trace_id) or [gen]
|
||||
|
||||
|
||||
@step("Retry the call until the fresh key stops answering 401")
|
||||
def first_ok(client: LoggingClient, send: Callable[[], StreamingResponse]) -> StreamingResponse:
|
||||
"""First successful call on a fresh key. A fresh key may briefly 401 until
|
||||
the data plane's auth cache picks it up, so retry on 401 to a deadline; a
|
||||
|
|
|
|||
|
|
@ -32,6 +32,7 @@ from pydantic import BaseModel, ConfigDict, Field
|
|||
|
||||
from e2e_config import OTEL_QUERY_URL, POLL_INTERVAL, POLL_TIMEOUT
|
||||
from e2e_http import URL, NetworkError, NoBody, Result, Success, get
|
||||
from e2e_metadata import step
|
||||
|
||||
#: OTEL resource service.name the proxy exports under (OTEL_SERVICE_NAME default).
|
||||
JAEGER_SERVICE = "litellm"
|
||||
|
|
@ -231,6 +232,7 @@ class OtelReader:
|
|||
case failure:
|
||||
pytest.fail(f"Jaeger query API at {self.query_url} failed: {failure}")
|
||||
|
||||
@step("Poll Jaeger for the traces of call {call_id}")
|
||||
def poll_traces_for_call(
|
||||
self,
|
||||
*,
|
||||
|
|
|
|||
|
|
@ -66,6 +66,7 @@ from e2e_http import (
|
|||
StreamTruncation,
|
||||
forward_stream,
|
||||
)
|
||||
from e2e_metadata import step
|
||||
from fixture_bundle import (
|
||||
BundleRecorder,
|
||||
Interaction,
|
||||
|
|
@ -1176,6 +1177,7 @@ class RunningEdge:
|
|||
self.server.server_close()
|
||||
|
||||
|
||||
@step("Start a provider edge server")
|
||||
def start_provider_edge(
|
||||
backend: EdgeBackend,
|
||||
*,
|
||||
|
|
@ -1335,6 +1337,7 @@ def _shared_cache_edge(bind_host: str, advertise_host: str, forward_timeout: flo
|
|||
).edge
|
||||
|
||||
|
||||
@step("Run an observed provider edge server")
|
||||
@contextmanager
|
||||
def observed_provider_edge(
|
||||
observation: ProviderRequestObservation,
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ from e2e_http import (
|
|||
StreamHead,
|
||||
StreamingResponse,
|
||||
)
|
||||
from e2e_metadata import step
|
||||
from pydantic import BaseModel
|
||||
|
||||
|
||||
|
|
@ -132,6 +133,7 @@ class HttpTransport:
|
|||
def master(self) -> AuthHeaders:
|
||||
return self.bearer(self.master_key)
|
||||
|
||||
@step("POST {path}")
|
||||
def post[R: BaseModel](
|
||||
self,
|
||||
path: str,
|
||||
|
|
@ -151,6 +153,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout if timeout is None else timeout,
|
||||
)
|
||||
|
||||
@step("GET {path}")
|
||||
def get[R: BaseModel](
|
||||
self,
|
||||
path: str,
|
||||
|
|
@ -170,6 +173,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout if timeout is None else timeout,
|
||||
)
|
||||
|
||||
@step("DELETE {path}")
|
||||
def delete[R: BaseModel](
|
||||
self,
|
||||
path: str,
|
||||
|
|
@ -188,6 +192,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout,
|
||||
)
|
||||
|
||||
@step("PATCH {path}")
|
||||
def patch[R: BaseModel](
|
||||
self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]
|
||||
) -> Result[R]:
|
||||
|
|
@ -199,6 +204,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout,
|
||||
)
|
||||
|
||||
@step("PUT {path}")
|
||||
def put[R: BaseModel](self, path: str, *, headers: BaseModel, json: BaseModel, response_type: type[R]) -> Result[R]:
|
||||
return e2e_http.put(
|
||||
self._url(path),
|
||||
|
|
@ -208,12 +214,15 @@ class HttpTransport:
|
|||
timeout=self.request_timeout,
|
||||
)
|
||||
|
||||
@step("Stream a POST to {path}")
|
||||
def stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamingResponse:
|
||||
return e2e_http.stream(self._url(path), headers=headers, json=json, timeout=self.request_timeout)
|
||||
|
||||
@step("Open a stream to {path}")
|
||||
def open_stream(self, path: str, *, headers: BaseModel, json: BaseModel) -> StreamHead | NetworkError:
|
||||
return e2e_http.open_stream(self._url(path), headers=headers, json=json, timeout=self.request_timeout)
|
||||
|
||||
@step("Stream binary from {path}")
|
||||
def stream_binary(
|
||||
self,
|
||||
path: str,
|
||||
|
|
@ -230,6 +239,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout,
|
||||
)
|
||||
|
||||
@step("Send a request to {path}")
|
||||
def send(
|
||||
self,
|
||||
path: str,
|
||||
|
|
@ -248,11 +258,13 @@ class HttpTransport:
|
|||
timeout=self.request_timeout,
|
||||
)
|
||||
|
||||
@step("Abandon the request to {path} after {after}s")
|
||||
def abandon(
|
||||
self, path: str, *, headers: BaseModel, json: BaseModel, after: float
|
||||
) -> AbandonedRequest | StreamingResponse:
|
||||
return e2e_http.abandon(self._url(path), headers=headers, json=json, after=after)
|
||||
|
||||
@step("Probe {path}")
|
||||
def probe(self, path: str, *, params: BaseModel, headers: BaseModel | None = None) -> ProbeResult:
|
||||
return e2e_http.probe(
|
||||
self._url(path),
|
||||
|
|
@ -261,6 +273,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout,
|
||||
)
|
||||
|
||||
@step("Upload {filename} to {path}")
|
||||
def upload[R: BaseModel](
|
||||
self,
|
||||
path: str,
|
||||
|
|
@ -288,6 +301,7 @@ class HttpTransport:
|
|||
timeout=self.request_timeout if timeout is None else timeout,
|
||||
)
|
||||
|
||||
@step("Download {path}")
|
||||
def download(self, path: str, *, headers: BaseModel) -> StreamingResponse:
|
||||
return e2e_http.download(self._url(path), headers=headers, timeout=self.request_timeout)
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue