test(load): drive /v1/messages alongside /chat/completions in the Redis chaos test

The Anthropic Messages route reaches the same Redis touchpoints and cost-tracking
callback through its own request path, so a failure-path regression there would not
surface from chat completions alone. Each simulated user now picks one endpoint round
robin and stays on it, and the per-endpoint split is asserted and reported so a run
that silently drove only one route fails instead of passing.

Co-Authored-By: Claude Code <noreply@anthropic.com>
This commit is contained in:
Kerry Lu 2026-09-11 12:28:02 -07:00
parent 1699f2d6dc
commit 1270ecb781
6 changed files with 135 additions and 17 deletions

View file

@ -18,7 +18,7 @@ Each subdirectory under `tests/e2e/` is one suite, scoped to an endpoint family
- `logging/` - logging-integration delivery (datadog and friends)
- `security/` - secret handling and log-leak protection
- `router/` - routing and reliability behavior (fallbacks, cooldowns)
- `load/` - performance-category tests, kept OUT of the main suite: throughput/load SLO tests are a different testing category from functional e2e (variance-driven, historically flaky) and live outside this suite until re-implemented as their own pipeline (LIT-5163); do not add a live load test that runs in the default collection. What lives here: the weekly session-anomaly test (`test_weekly_session_anomaly_e2e.py`, Claude Code-shaped multi-turn sessions against real providers with ceilings on error rate, cache read/write, turn time, and spend; marked `weekly` and deselected unless `E2E_WEEKLY_ANOMALY` is set, driven by `.github/workflows/weekly_load_anomaly.yml`), the Redis chaos test (`test_redis_chaos_e2e.py`, locust load against mock deployments with `CLIENT PAUSE ALL` on the proxy's Redis mid-run to simulate it being down outright, asserting zero failed requests and budgeting p50/p90/p99 latency, RSS, and CPU-per-request as ratios against the same run's healthy phase; needs a proxy booted from `gateway/redis_chaos_ci_config.yml` on the same host with `E2E_PROXY_PID` set, marked `redis_chaos`, deselected unless `E2E_REDIS_CHAOS` is set, driven weekly by `.github/workflows/test-e2e-redis-chaos.yml`), and markerless harness unit tests for the locust, process-usage, and session-anomaly aggregation logic
- `load/` - performance-category tests, kept OUT of the main suite: throughput/load SLO tests are a different testing category from functional e2e (variance-driven, historically flaky) and live outside this suite until re-implemented as their own pipeline (LIT-5163); do not add a live load test that runs in the default collection. What lives here: the weekly session-anomaly test (`test_weekly_session_anomaly_e2e.py`, Claude Code-shaped multi-turn sessions against real providers with ceilings on error rate, cache read/write, turn time, and spend; marked `weekly` and deselected unless `E2E_WEEKLY_ANOMALY` is set, driven by `.github/workflows/weekly_load_anomaly.yml`), the Redis chaos test (`test_redis_chaos_e2e.py`, locust load against mock deployments split round robin over `/chat/completions` and `/v1/messages`, one endpoint per simulated user, with `CLIENT PAUSE ALL` on the proxy's Redis mid-run to simulate it being down outright, asserting zero failed requests on every endpoint and budgeting p50/p90/p99 latency, RSS, and CPU-per-request as ratios against the same run's healthy phase; needs a proxy booted from `gateway/redis_chaos_ci_config.yml` on the same host with `E2E_PROXY_PID` set, marked `redis_chaos`, deselected unless `E2E_REDIS_CHAOS` is set, driven weekly by `.github/workflows/test-e2e-redis-chaos.yml`), and markerless harness unit tests for the locust, process-usage, and session-anomaly aggregation logic
- `other/` - the holding-pen suite for the `other.*` registry cluster with no home of its own yet: the master-key auth gate and the process-lifecycle health probes (liveness, public readiness, authenticated readiness diagnostics). Promote a cluster out once it is large/stable enough for its own suite
- `gateway/` - proxy configuration only (`litellm-config.yml`); no tests
- `claude_code/` - the Claude Code compatibility matrix: drives the real `claude` CLI (and HTTP probes) against a proxy for each feature x provider cell, reporting tagged-union outcomes via the `compat_result` fixture; ships its own driver/builder/publisher plus `_*_unit_tests/` trees. The HTTP probes ride the shared transport (`ProxyClient.count_tokens` / `ProxyClient.messages`); the CLI-driving path stays bespoke

View file

@ -31,7 +31,7 @@
- {id: reliability.cache.exact.returns_cached, module: reliability, tier: P1, behavior: cache, variant: exact, assertions: [returns_cached], exercised_on: [chat_completions, messages, embeddings], source: "litellm/caching/caching.py", rationale: "Response cache returns cached on exact match"}
- {id: reliability.cache.prompt_caching_model_select.returns_cached, module: reliability, tier: P1, behavior: cache, variant: prompt_caching_model_select, assertions: [returns_cached], exercised_on: [chat_completions], source: "router_utils/prompt_caching_cache.py", rationale: "Selects model supporting prompt caching for cacheable prefix"}
- {id: reliability.circuit_breaker.redis.trips_then_recovers, module: reliability, tier: P0, behavior: circuit_breaker, variant: redis, assertions: [trips_then_recovers], exercised_on: [chat_completions, messages, embeddings], source: "litellm/caching/redis_cache.py:99", rationale: "Redis breaker CLOSED->OPEN->HALF_OPEN; guards all cache/rate-limit ops"}
- {id: reliability.circuit_breaker.redis_timeout.stays_responsive, module: reliability, tier: P1, behavior: circuit_breaker, variant: redis_timeout, assertions: [stays_responsive], exercised_on: [chat_completions], source: "litellm/proxy/hooks/proxy_track_cost_callback.py:386", fail_before_fix: proven, rationale: "Under locust load with every request retrying through failing mock deployments, pausing Redis writes mid-run trips the breaker and every request still succeeds, with latency, RSS, and CPU reported as p50/p90/p99 against the pre-pause baseline; on v1.100.0 the failed-tracking alert body doubled per request until the worker OOMed (LIT-6780)"}
- {id: reliability.circuit_breaker.redis_timeout.stays_responsive, module: reliability, tier: P1, behavior: circuit_breaker, variant: redis_timeout, assertions: [stays_responsive], exercised_on: [chat_completions, messages], source: "litellm/proxy/hooks/proxy_track_cost_callback.py:386", fail_before_fix: proven, rationale: "Under locust load split round robin over /chat/completions and /v1/messages with every request retrying through failing mock deployments, holding Redis in CLIENT PAUSE ALL for the phase trips the breaker and every request still succeeds, with latency, RSS, and CPU reported as p50/p90/p99 against the pre-pause baseline; on v1.100.0 the failed-tracking alert body doubled per request until the worker OOMed (LIT-6780)"}
- {id: reliability.timeout.request_timeout.exceeds_deadline, module: reliability, tier: P1, behavior: timeout, variant: request_timeout, assertions: [exceeds_deadline], exercised_on: [chat_completions, messages], source: "litellm/router.py:545-551", rationale: "Per-request timeout raises Timeout"}
- {id: reliability.timeout.stream_timeout.exceeds_deadline, module: reliability, tier: P1, behavior: timeout, variant: stream_timeout, assertions: [exceeds_deadline], exercised_on: [chat_completions], source: "litellm/router.py:551", rationale: "Streaming chunk-delivery timeout"}
- {id: reliability.perf.throughput.under_slo, module: reliability, tier: P1, behavior: perf, variant: throughput, assertions: [under_slo], exercised_on: [chat_completions, messages], source: grammar, rationale: "Throughput SLO under load"}

View file

@ -5,6 +5,7 @@ import os
import subprocess
import sys
import tempfile
from collections.abc import Sequence
from dataclasses import dataclass
from itertools import accumulate
from pathlib import Path
@ -19,6 +20,7 @@ _MAX_REPORTED_ERRORS = 5
class LocustStatEntry(BaseModel):
name: str
num_requests: int
num_failures: int
start_time: float
@ -36,6 +38,16 @@ class LoadError:
occurrences: int
@dataclass(frozen=True, slots=True)
class EndpointLoad:
"""One route's share of a phase, so a run that silently drove only one of them is visible."""
name: str
requests: int
failures: int
p50_seconds: float
@dataclass(frozen=True, slots=True)
class LoadResult:
requests: int
@ -44,6 +56,7 @@ class LoadResult:
p50_seconds: float
p90_seconds: float
p99_seconds: float
endpoints: tuple[EndpointLoad, ...]
errors: tuple[LoadError, ...]
generator_warnings: tuple[str, ...]
@ -65,8 +78,15 @@ class LoadResult:
def latency_summary(self) -> str:
return f"p50 {self.p50_seconds:.3f}s, p90 {self.p90_seconds:.3f}s, p99 {self.p99_seconds:.3f}s"
def endpoint_summary(self) -> str:
return ", ".join(
f"{endpoint.name} {endpoint.requests} requests, {endpoint.failures} failures, "
f"p50 {endpoint.p50_seconds:.3f}s"
for endpoint in self.endpoints
)
def percentile_seconds(entries: list[LocustStatEntry], fraction: float) -> float:
def percentile_seconds(entries: Sequence[LocustStatEntry], fraction: float) -> float:
"""The response time at `fraction` of the merged histograms, in seconds.
Locust buckets response times by millisecond, so this reads the first bucket whose
@ -82,13 +102,29 @@ def percentile_seconds(entries: list[LocustStatEntry], fraction: float) -> float
return next(milliseconds for (milliseconds, _), seen in zip(samples, running) if seen >= rank) / 1000.0
def per_endpoint(entries: Sequence[LocustStatEntry]) -> tuple[EndpointLoad, ...]:
"""Each locust request name's own totals, in the order the names first appear."""
names: Final = tuple(dict.fromkeys(entry.name for entry in entries))
grouped: Final = ((name, tuple(entry for entry in entries if entry.name == name)) for name in names)
return tuple(
EndpointLoad(
name=name,
requests=sum(entry.num_requests for entry in group),
failures=sum(entry.num_failures for entry in group),
p50_seconds=percentile_seconds(group, 0.5),
)
for name, group in grouped
)
def aggregate_stats(
entries: list[LocustStatEntry],
entries: Sequence[LocustStatEntry],
errors: tuple[LoadError, ...],
generator_warnings: tuple[str, ...],
) -> LoadResult:
requests = sum(entry.num_requests for entry in entries)
failures = sum(entry.num_failures for entry in entries)
endpoints = per_endpoint(entries)
if not entries or requests == 0:
return LoadResult(
requests=requests,
@ -97,6 +133,7 @@ def aggregate_stats(
p50_seconds=0.0,
p90_seconds=0.0,
p99_seconds=0.0,
endpoints=endpoints,
errors=errors,
generator_warnings=generator_warnings,
)
@ -108,6 +145,7 @@ def aggregate_stats(
p50_seconds=percentile_seconds(entries, 0.5),
p90_seconds=percentile_seconds(entries, 0.9),
p99_seconds=percentile_seconds(entries, 0.99),
endpoints=endpoints,
errors=errors,
generator_warnings=generator_warnings,
)
@ -142,19 +180,21 @@ def read_generator_warnings(stderr: str) -> tuple[str, ...]:
return tuple(dict.fromkeys(saturated))
def run_chat_load(
def run_gateway_load(
*,
base_url: str,
api_keys: tuple[str, ...],
model: str,
endpoints: tuple[str, ...],
users: int,
spawn_rate: float,
duration_seconds: float,
) -> LoadResult:
"""Drive /chat/completions from headless locust and aggregate what it reported.
"""Drive `endpoints` from headless locust and aggregate what it reported.
Each simulated user picks one of `api_keys`, so auth and budget lookups spread over a
pool of virtual keys instead of keeping one key's cache entry permanently warm.
pool of virtual keys instead of keeping one key's cache entry permanently warm, and one
of `endpoints` round robin, so the run covers every route the caller asked for.
"""
with tempfile.TemporaryDirectory(prefix="e2e-load-") as report_dir:
csv_prefix = Path(report_dir) / _CSV_PREFIX
@ -180,7 +220,12 @@ def run_chat_load(
"--exit-code-on-error",
"0",
],
env={**os.environ, "LOAD_API_KEYS": ",".join(api_keys), "LOAD_MODEL": model},
env={
**os.environ,
"LOAD_API_KEYS": ",".join(api_keys),
"LOAD_MODEL": model,
"LOAD_ENDPOINTS": ",".join(endpoints),
},
capture_output=True,
text=True,
timeout=duration_seconds + 120,

View file

@ -3,16 +3,22 @@ from __future__ import annotations
import os
import random
import uuid
from itertools import cycle
from typing import Final
from locust import FastHttpUser, constant, task
_MODEL: Final = os.environ["LOAD_MODEL"]
_API_KEYS: Final = tuple(os.environ["LOAD_API_KEYS"].split(","))
_NEXT_ENDPOINT: Final = cycle(os.environ["LOAD_ENDPOINTS"].split(","))
def _payload() -> dict[str, object]:
"""A prompt no other request sent, so the response cache never answers for the deployment."""
"""A prompt no other request sent, so the response cache never answers for the deployment.
Both endpoints take the same body: /v1/messages requires max_tokens, which /chat/completions
also accepts, so one payload serves the whole round robin.
"""
return {
"model": _MODEL,
"messages": [{"role": "user", "content": f"load test ping {uuid.uuid4().hex}"}],
@ -20,17 +26,24 @@ def _payload() -> dict[str, object]:
}
class ChatUser(FastHttpUser):
class GatewayUser(FastHttpUser):
"""One simulated user, pinned to one endpoint for its lifetime.
Endpoints are handed out round robin as users spawn, so a run spreads evenly over them
while each user's traffic stays on a single route, the way a real client behaves.
"""
wait_time = constant(0)
def on_start(self) -> None:
self.headers = {"Authorization": f"Bearer {random.choice(_API_KEYS)}"}
self.endpoint = next(_NEXT_ENDPOINT)
@task
def chat(self) -> None:
def call(self) -> None:
self.client.post( # pyright: ignore[reportUnknownMemberType] # locust FastHttpSession.post types json/**kwargs as Any
"/chat/completions",
self.endpoint,
json=_payload(),
headers=self.headers,
name="/chat/completions",
name=self.endpoint,
)

View file

@ -1,6 +1,7 @@
from __future__ import annotations
from pathlib import Path
from typing import Final
from locust_load import (
LoadError,
@ -18,12 +19,14 @@ _FAILURES_HEADER = "Method,Name,Error,Occurrences,First Seen,Last Seen\n"
def _entry(
*,
num_requests: int,
name: str = "/chat/completions",
num_failures: int = 0,
start_time: float = 1000.0,
last_request_timestamp: float = 1010.0,
response_times: dict[int, int] | None = None,
) -> LocustStatEntry:
return LocustStatEntry(
name=name,
num_requests=num_requests,
num_failures=num_failures,
start_time=start_time,
@ -44,6 +47,7 @@ def _result(
p50_seconds=0.05,
p90_seconds=0.08,
p99_seconds=0.1,
endpoints=(),
errors=errors,
generator_warnings=generator_warnings,
)
@ -126,6 +130,49 @@ class TestAggregate:
assert result.requests == 0
assert result.requests_per_second == 0.0
assert result.failure_ratio == 1.0
assert result.endpoints == ()
class TestPerEndpoint:
def test_each_route_keeps_its_own_requests_failures_and_median(self) -> None:
entries: Final = (
_entry(name="/chat/completions", num_requests=100, response_times={20: 100}),
_entry(name="/v1/messages", num_requests=40, num_failures=3, response_times={900: 40}),
)
result: Final = aggregate_stats(entries, (), ())
assert tuple((one.name, one.requests, one.failures, one.p50_seconds) for one in result.endpoints) == (
("/chat/completions", 100, 0, 0.02),
("/v1/messages", 40, 3, 0.9),
)
def test_several_stats_entries_for_one_route_fold_into_a_single_row(self) -> None:
entries: Final = (
_entry(name="/v1/messages", num_requests=10, response_times={30: 10}),
_entry(name="/v1/messages", num_requests=30, num_failures=1, response_times={30: 30}),
)
result: Final = aggregate_stats(entries, (), ())
assert tuple((one.name, one.requests, one.failures) for one in result.endpoints) == (("/v1/messages", 40, 1),)
def test_a_route_that_never_ran_is_absent_so_a_one_sided_run_cannot_pass_unnoticed(self) -> None:
result: Final = aggregate_stats((_entry(name="/chat/completions", num_requests=10),), (), ())
assert tuple(one.name for one in result.endpoints) == ("/chat/completions",)
def test_the_summary_names_every_route_with_its_counts(self) -> None:
entries: Final = (
_entry(name="/chat/completions", num_requests=2, response_times={20: 2}),
_entry(name="/v1/messages", num_requests=1, num_failures=1, response_times={500: 1}),
)
result: Final = aggregate_stats(entries, (), ())
assert result.endpoint_summary() == (
"/chat/completions 2 requests, 0 failures, p50 0.020s, /v1/messages 1 requests, 1 failures, p50 0.500s"
)
class TestErrorBreakdown:

View file

@ -11,6 +11,11 @@ retries on the failing pair (a 500 is retryable, so retries keep re-picking insi
order) and the router's order-based fallback then re-targets order 2. Every request is expected
to succeed, and each one carries retry breadcrumbs into cost tracking.
Traffic is split round robin between /chat/completions and /v1/messages, one endpoint per
simulated user: the Redis touchpoints and the cost-tracking callback are shared by both, but
the Anthropic Messages route reaches them through its own request path, so a regression that
only shows up there would not surface from chat completions alone.
Phase A is a baseline with Redis healthy; phase B holds Redis in CLIENT PAUSE ALL for the
length of the phase, simulating Redis being down outright rather than merely slow to write.
Every touchpoint times out: the auth cache read falls back to Postgres, the response cache
@ -40,7 +45,7 @@ from e2e_config import PROXY_BASE_URL, unique_marker
from e2e_http import NoBody
from lifecycle import ResourceManager
from load_client import LoadClient
from locust_load import LoadResult, run_chat_load
from locust_load import LoadResult, run_gateway_load
from models import KeyGenerateBody, LiteLLMParamsBody
from phase_budget import Budget, violations
from proxy_client import ProxyClient
@ -55,6 +60,7 @@ SERVING_DEPLOYMENTS: Final = 1
FAILING_ORDER: Final = 1
SERVING_ORDER: Final = 2
KEY_POOL_SIZE: Final = 8
LOAD_ENDPOINTS: Final = ("/chat/completions", "/v1/messages")
LOCUST_USERS: Final = 50
LOCUST_SPAWN_RATE: Final = 50.0
BASELINE_SECONDS: Final = 60.0
@ -119,7 +125,8 @@ class Phase:
f"{self.name}: {self.load.requests} requests, {self.load.failures} failures, "
f"{self.load.requests_per_second:.0f} rps, {self.load.latency_summary()}; {self.usage.summary()}; "
f"{self.cpu_seconds_per_request * 1000:.1f} ms CPU per request; "
f"{self.timeouts_per_request:.2f} Redis timeouts per request"
f"{self.timeouts_per_request:.2f} Redis timeouts per request; "
f"by endpoint: {self.load.endpoint_summary()}"
)
@ -241,10 +248,11 @@ def _generate_key_pool(proxy: ProxyClient, resources: ResourceManager) -> tuple[
def _drive(keys: tuple[str, ...], seconds: float) -> LoadResult:
return run_chat_load(
return run_gateway_load(
base_url=PROXY_BASE_URL,
api_keys=keys,
model=MODEL_GROUP,
endpoints=LOAD_ENDPOINTS,
users=LOCUST_USERS,
spawn_rate=LOCUST_SPAWN_RATE,
duration_seconds=seconds,
@ -310,7 +318,7 @@ def _chaos_budgets(baseline: Phase, chaos: Phase) -> tuple[Budget, ...]:
class TestRedisChaos:
@pytest.mark.covers(
"reliability.circuit_breaker.redis_timeout.stays_responsive",
exercised_on=("chat_completions",),
exercised_on=("chat_completions", "messages"),
)
def test_load_survives_redis_being_down(
self,
@ -356,6 +364,11 @@ class TestRedisChaos:
assert phase.load.requests > 0, (
f"{phase.name} drove no traffic at all, so it proved nothing: {phase.load.diagnosis()}. {report}"
)
assert frozenset(endpoint.name for endpoint in phase.load.endpoints) == frozenset(LOAD_ENDPOINTS), (
f"{phase.name} drove {tuple(endpoint.name for endpoint in phase.load.endpoints)} rather than every "
f"endpoint in {LOAD_ENDPOINTS}; the round robin hands one endpoint to each simulated user, so a "
f"missing one means a route never ran and its request path was never exercised. {report}"
)
assert phase.load.failures == 0, (
f"{phase.name} had {phase.load.failures} of {phase.load.requests} requests fail. Every request "
f"must succeed: the failing deployments sit at order {FAILING_ORDER} and the serving one at order "