diff --git a/tests/code_coverage_tests/test_e2e_metadata.py b/tests/code_coverage_tests/test_e2e_metadata.py index 5989bc993dc..57d04262969 100644 --- a/tests/code_coverage_tests/test_e2e_metadata.py +++ b/tests/code_coverage_tests/test_e2e_metadata.py @@ -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] = [] diff --git a/tests/e2e/e2e_config.py b/tests/e2e/e2e_config.py index 8a332aa8f24..902c161e425 100644 --- a/tests/e2e/e2e_config.py +++ b/tests/e2e/e2e_config.py @@ -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. diff --git a/tests/e2e/e2e_http.py b/tests/e2e/e2e_http.py index e5d50d05c87..94e135535ce 100644 --- a/tests/e2e/e2e_http.py +++ b/tests/e2e/e2e_http.py @@ -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, *, diff --git a/tests/e2e/guardrails/guardrails_client.py b/tests/e2e/guardrails/guardrails_client.py index b193896bf6c..08124a03313 100644 --- a/tests/e2e/guardrails/guardrails_client.py +++ b/tests/e2e/guardrails/guardrails_client.py @@ -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 diff --git a/tests/e2e/load/locust_load.py b/tests/e2e/load/locust_load.py index 40f9f333db5..6d0aa8e1159 100644 --- a/tests/e2e/load/locust_load.py +++ b/tests/e2e/load/locust_load.py @@ -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, diff --git a/tests/e2e/load/session_anomaly.py b/tests/e2e/load/session_anomaly.py index c29b833635d..b68482e965f 100644 --- a/tests/e2e/load/session_anomaly.py +++ b/tests/e2e/load/session_anomaly.py @@ -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, diff --git a/tests/e2e/logging/logging_client.py b/tests/e2e/logging/logging_client.py index d522c01c054..c4d09800c79 100644 --- a/tests/e2e/logging/logging_client.py +++ b/tests/e2e/logging/logging_client.py @@ -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 diff --git a/tests/e2e/otel_client.py b/tests/e2e/otel_client.py index cf8fac42e88..5f93c677d89 100644 --- a/tests/e2e/otel_client.py +++ b/tests/e2e/otel_client.py @@ -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, *, diff --git a/tests/e2e/provider_edge.py b/tests/e2e/provider_edge.py index 3680375b6af..67b8d8e2980 100644 --- a/tests/e2e/provider_edge.py +++ b/tests/e2e/provider_edge.py @@ -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, diff --git a/tests/e2e/transport.py b/tests/e2e/transport.py index 87aad0d08de..0e37094f4b4 100644 --- a/tests/e2e/transport.py +++ b/tests/e2e/transport.py @@ -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)