From 2156d6c3c796cbdda3b3f753d0d323a2261c1c7b Mon Sep 17 00:00:00 2001 From: "devin-ai-integration[bot]" <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Tue, 6 Oct 2026 02:12:57 +0000 Subject: [PATCH] test(e2e): assert the sibling-replica cooldown through the router (#44706) * test(e2e): assert the sibling-replica cooldown through the router * test(e2e): warm the cooldown reads concurrently so every pod's read lands just before the trip * test(e2e): send the trip right behind the warm so every pod's cooldown read is pinned to it * test(e2e): warm every pod with a canned-answer group and trip only after every warm call answered * test(e2e): trim the sibling cell's module docstring to what the design needs --------- Co-authored-by: mateo-berri <277851410+mateo-berri@users.noreply.github.com> --- tests/e2e/AGENTS.md | 2 + tests/e2e/coverage_registry/reliability.yaml | 2 +- tests/e2e/router/reliability_support.py | 38 ++--- .../router/test_reliability_cooldowns_e2e.py | 131 ++++++++++-------- 4 files changed, 83 insertions(+), 90 deletions(-) diff --git a/tests/e2e/AGENTS.md b/tests/e2e/AGENTS.md index 3fd13439fd8..e3fe1421869 100644 --- a/tests/e2e/AGENTS.md +++ b/tests/e2e/AGENTS.md @@ -324,3 +324,5 @@ other... - spin up a local proxy by running the litellm proxy locally (`litellm --config .yml --port 4000`; see CONTRIBUTING.md), make sure all tests pass. if a test fails due to an internally found issue, let users know to create a linear ticket for it. - do not use xfail markers, tests should be written in a form that the end user expects it to pass + +- a cell addresses the stack through its front door (`PROXY_BASE_URL`) the way a customer does, never a gateway pod by address, and proves a cross-replica property with N independent calls through that door, naming the miss odds (2^-N at two pods) in its docstring. `PROXY_REPLICA_URLS` is for read-backs that are per pod by nature (polling until a management write has converged on every replica, `/metrics`, RSS), never for steering the scenario a cell asserts on at a chosen pod diff --git a/tests/e2e/coverage_registry/reliability.yaml b/tests/e2e/coverage_registry/reliability.yaml index 5040f5f4dcf..7382cc215f2 100644 --- a/tests/e2e/coverage_registry/reliability.yaml +++ b/tests/e2e/coverage_registry/reliability.yaml @@ -9,7 +9,7 @@ - {id: reliability.retry.auth.succeeds_within_retries, module: reliability, tier: P1, behavior: retry, variant: auth, assertions: [succeeds_within_retries], exercised_on: [chat_completions], source: "get_retry_from_policy.py:42", rationale: "Transient auth glitch retry"} - {id: reliability.retry.context_window.succeeds_within_retries, module: reliability, tier: P1, behavior: retry, variant: context_window, assertions: [succeeds_within_retries], exercised_on: [chat_completions], source: "get_retry_from_policy.py:51", fail_before_fix: proven, rationale: "A context-window 400 under BadRequestErrorRetries retries onto a sibling deployment in the same model group, instead of coming straight back as the 400 the deployment that just refused it returned"} - {id: reliability.cooldown.5xx.trips_then_recovers, module: reliability, tier: P0, behavior: cooldown, variant: "5xx", assertions: [trips_then_recovers], exercised_on: [chat_completions], source: "cooldown_handlers.py:40", rationale: "Deployment cools after repeated 5xx, recovers after cooldown_time"} -- {id: reliability.cooldown.sibling_replica.serves_backup_within_read_interval, module: reliability, tier: P1, behavior: cooldown, variant: sibling_replica, assertions: [serves_backup_within_read_interval], exercised_on: [chat_completions], source: "cooldown_cache.py:44", fail_before_fix: proven, rationale: "A bench taken on one gateway reaches a sibling that already holds the key's read timer within the 1s Redis read interval plus margin, so its next call lands on the backup"} +- {id: reliability.cooldown.sibling_replica.serves_backup_within_read_interval, module: reliability, tier: P1, behavior: cooldown, variant: sibling_replica, assertions: [serves_backup_within_read_interval], exercised_on: [chat_completions], source: "cooldown_cache.py:44", fail_before_fix: proven, rationale: "A bench taken on one gateway reaches every sibling within the 1s Redis read interval plus margin, proven through the front door with no pod addresses: 10 warm calls start every pod's read timer on the key (a pod left unwarmed has odds 2^-9), one call trips, and all 10 probes after the wait land on the backup, a stale pod being missed with odds 2^-10 (0.1%)"} - {id: reliability.cooldown.429.trips_then_recovers, module: reliability, tier: P0, behavior: cooldown, variant: "429", assertions: [trips_then_recovers], exercised_on: [chat_completions], source: "cooldown_handlers.py:69", rationale: "Cools on 429, avoids hammering exhausted provider"} - {id: reliability.cooldown.auth.trips_then_recovers, module: reliability, tier: P1, behavior: cooldown, variant: auth, assertions: [trips_then_recovers], exercised_on: [chat_completions], source: "cooldown_handlers.py:74", rationale: "Cools on 401 auth error"} - {id: reliability.cooldown.timeout.trips_then_recovers, module: reliability, tier: P1, behavior: cooldown, variant: timeout, assertions: [trips_then_recovers], exercised_on: [chat_completions], source: "cooldown_handlers.py:77", rationale: "Cools on 408 timeout"} diff --git a/tests/e2e/router/reliability_support.py b/tests/e2e/router/reliability_support.py index f600531e663..04fa30a0d14 100644 --- a/tests/e2e/router/reliability_support.py +++ b/tests/e2e/router/reliability_support.py @@ -23,7 +23,6 @@ from pydantic import BaseModel, ValidationError from proxy_client import ProxyClient from e2e_config import CHEAP_OPENAI_MODEL, PROXY_BASE_URL, unique_marker from e2e_http import NetworkError, StreamHead, StreamingResponse -from transport import Transport from models import ( CacheControl, ChatMessage, @@ -254,6 +253,12 @@ def create_always_picked_small_context_deployment(proxy: ProxyClient, name: str) ) +def create_canned_deployment(proxy: ProxyClient, name: str) -> str: + """A deployment that answers from a canned reply, so a call to it goes through the + router's deployment pick like any other but never reaches a provider.""" + return proxy.create_model(name, LiteLLMParamsBody(model=REAL_MODEL, mock_response="ok")) + + def create_zero_weight_backup_deployment(proxy: ProxyClient, name: str) -> str: """The other half of a retry pair: healthy, but weight 0, so the weighted shuffle never opens on it. It is reachable only once its sibling is out of the running, @@ -280,36 +285,9 @@ def chat_turns_override( ) -> StreamingResponse: """POST /chat/completions with an optional per-request router_settings_override, returning the raw outcome so tests read status, body, and reliability headers.""" - return chat_turns_override_via( - proxy.transport, key, model, turns, override=override, stream=stream, cache=cache, max_tokens=max_tokens - ) - - -def chat_override_via( - transport: Transport, - key: str, - model: str, - content: str, - override: RouterSettingsOverride | None = None, -) -> StreamingResponse: - """`chat_override` aimed at one replica's transport (from `proxy.replicas`) instead of - the client's default, for cells that must know which gateway took the call.""" - return chat_turns_override_via(transport, key, model, [ChatMessage(role="user", content=content)], override=override) - - -def chat_turns_override_via( - transport: Transport, - key: str, - model: str, - turns: Sequence[ChatMessage], - override: RouterSettingsOverride | None = None, - stream: bool = False, - cache: dict[str, bool] | None = {"no-cache": True}, - max_tokens: int = 512, -) -> StreamingResponse: - return transport.send( + return proxy.transport.send( "/chat/completions", - headers=transport.bearer(key), + headers=proxy.transport.bearer(key), json=ReliabilityChatBody( model=model, messages=turns, diff --git a/tests/e2e/router/test_reliability_cooldowns_e2e.py b/tests/e2e/router/test_reliability_cooldowns_e2e.py index ce3bce32880..2456bfb5f85 100644 --- a/tests/e2e/router/test_reliability_cooldowns_e2e.py +++ b/tests/e2e/router/test_reliability_cooldowns_e2e.py @@ -23,16 +23,34 @@ is the recovery, since a benched deployment is one the router will try again, not one it forgot. Its deadline counts from the last failure a stale replica caused, because every failure re-arms the cooldown. -The sibling cell is the one that asserts the speed. It addresses two gateways -from PROXY_REPLICA_URLS by name, warms the second with a healthy call so its -router has already read the failing deployment's cooldown key from Redis and -started the read interval on it, trips the deployment through the first, waits -the interval plus a margin, and then sends the second replica exactly one call, -which has to come back from the backup. One call, because a poll that reached -the failing deployment through the second replica would bench it there too and -hide whether the first replica's bench ever travelled. A stack addressed only -through its load balancer cannot pin which replica takes a call, so the cell is -skipped at collection unless LITELLM_PROXY_REPLICA_URLS names at least two. +The sibling cell is the one that asserts the speed, and it sends every call +through the stack's front door the way a customer does, never to a gateway pod +by address: the litellm-e2e-pr gate fronts two pods with an nginx ingress that +picks the pod per connection, so each call is an independent draw over the two +routers. The cell registers a warm group whose deployment answers from a +canned reply, sends it COOLDOWN_WARM_CALLS calls at once, and waits for every +answer before it trips. A router reads the cooldown keys when it picks the +deployment, at the start of a call, so every pod's read of the failing +deployment's key, and the read interval that starts with it, is over before +the trip is sent; a pod the warm never reached (odds 2^(1-COOLDOWN_WARM_CALLS) +at two pods) would read Redis on its first touch of the key and pass even +under a regressed interval. The warm answers from a canned reply rather than a +live model because ten live answers would spread the reads across their +latencies, and tripping before they land would let a warm call reach a pod +after the bench and hand it a fresh read of the bench itself. The trip is one +call, retries off, that surfaces the deployment's own 500; after the interval +plus a margin, SIBLING_PROBES calls each have to come back from the backup, +since a pod that has not seen the bench answers 500 to any probe it takes, and +no probe reaching it has odds 2^-SIBLING_PROBES, 0.1% at ten. Other workers' +traffic can refresh a stale pod's read anywhere within a regressed interval of +the trip, so under the full suite a regression is caught on the runs whose +first probe reaches the stale pod before that read, while the per-file run, +which nothing else shares, catches every regression wider than the span from +the warm to the stale pod's first probe, about six seconds at two pods on the +two-process rig that proved the cell at ten. The bench is +SIBLING_COOLDOWN_SECONDS rather than COOLDOWN_SECONDS because the cell never +waits for the recovery and its probes, ten live calls to the backup, have to +land before the bench can lapse. The failures are the same real ones the retry tests use: a 1ms deadline and a bogus key on the real backend, and this proxy standing in as the upstream for @@ -45,11 +63,13 @@ from __future__ import annotations import time from collections.abc import Iterator +from concurrent.futures import ThreadPoolExecutor from dataclasses import dataclass +from typing import Final import pytest from complexity_router_client import ComplexityRouterClient -from e2e_config import CHEAP_OPENAI_MODEL, PROXY_REPLICA_URLS, unique_marker +from e2e_config import CHEAP_OPENAI_MODEL, unique_marker from e2e_http import StreamingResponse from lifecycle import ResourceManager from models import KeyGenerateBody, RouterSettingsOverride @@ -57,17 +77,16 @@ from reliability_support import ( COOLDOWN_SECONDS, REPLICA_PROPAGATION_SECONDS, chat_override, - chat_override_via, create_always_5xx_deployment, create_always_rate_limited_deployment, create_always_timing_out_deployment, create_always_unauthorized_deployment, create_bad_base_deployment, + create_canned_deployment, create_zero_weight_backup_deployment, model_id_of, spend_only_request_of, ) -from transport import Transport pytestmark = pytest.mark.e2e @@ -76,6 +95,9 @@ PROPAGATION_POLL_SECONDS = 0.25 BENCH_MARGIN_SECONDS = 4.0 COOLDOWN_REDIS_READ_INTERVAL_SECONDS = 1.0 SIBLING_READ_MARGIN_SECONDS = 1.0 +COOLDOWN_WARM_CALLS = 10 +SIBLING_PROBES = 10 +SIBLING_COOLDOWN_SECONDS = 120.0 def _call_without_retries(client: ComplexityRouterClient, key: str, group: str) -> StreamingResponse: @@ -84,29 +106,19 @@ def _call_without_retries(client: ComplexityRouterClient, key: str, group: str) ) -def _call_replica_without_retries(transport: Transport, key: str, group: str) -> StreamingResponse: - return chat_override_via( - transport, key, group, f"say hi {unique_marker()}", override=RouterSettingsOverride(num_retries=0) - ) +def _warm_call(client: ComplexityRouterClient, key: str, warm_group: str) -> StreamingResponse: + return chat_override(client.proxy, key, warm_group, f"say hi {unique_marker()}") -@dataclass(frozen=True, slots=True) -class _Replica: - url: str - transport: Transport - - -def _two_replicas(client: ComplexityRouterClient) -> tuple[_Replica, _Replica]: - first, second, *_ = (_Replica(url, transport) for url, transport in client.proxy.replicas.items()) - return first, second - - -def _warm_cooldown_reads(replica: _Replica, key: str) -> None: - warmed = chat_override_via(replica.transport, key, CHEAP_OPENAI_MODEL, f"say hi {unique_marker()}") - assert warmed.status_code == 200, ( - f"{replica.url} should have answered a healthy {CHEAP_OPENAI_MODEL} call before the trip, got " - f"{warmed.status_code}: {warmed.body[:300]}" - ) +def _warm_every_pod(client: ComplexityRouterClient, key: str, warm_group: str) -> None: + with ThreadPoolExecutor(max_workers=COOLDOWN_WARM_CALLS) as pool: + warm: Final = tuple(pool.submit(_warm_call, client, key, warm_group) for _ in range(COOLDOWN_WARM_CALLS)) + answered: Final = tuple(future.result() for future in warm) + for call, resp in enumerate(answered, start=1): + assert resp.status_code == 200, ( + f"warm call {call} of {COOLDOWN_WARM_CALLS} to {warm_group} should have answered 200, " + f"got {resp.status_code}: {resp.body[:300]}" + ) def _assert_served_by_backup(resp: StreamingResponse, backup: str, when: str) -> None: @@ -218,45 +230,46 @@ class TestReliabilityCooldowns: _assert_trips_then_recovers(client, scoped_key, group, failing, backup, failure_status=500) @pytest.mark.covers("reliability.cooldown.sibling_replica.serves_backup_within_read_interval") - @pytest.mark.skipif( - len(PROXY_REPLICA_URLS) < 2, - reason=( - "this cell trips a deployment through one gateway and reads the bench from another, so " - f"LITELLM_PROXY_REPLICA_URLS has to name at least two distinct gateways, got {PROXY_REPLICA_URLS}" - ), - ) def test_sibling_replica_serves_backup_within_redis_read_interval( self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str ) -> None: - tripping, sibling = _two_replicas(client) - - upstream = f"reliability-cooldown-sibling-upstream-{unique_marker()}" - upstream_id = create_bad_base_deployment(client.proxy, upstream) + upstream: Final = f"reliability-cooldown-sibling-upstream-{unique_marker()}" + upstream_id: Final = create_bad_base_deployment(client.proxy, upstream) resources.defer(lambda: client.proxy.delete_model(upstream_id)) - group = f"reliability-cooldown-sibling-{unique_marker()}" - failing = create_always_5xx_deployment( - client.proxy, group, upstream, scoped_key, cooldown_time=COOLDOWN_SECONDS + group: Final = f"reliability-cooldown-sibling-{unique_marker()}" + failing: Final = create_always_5xx_deployment( + client.proxy, group, upstream, scoped_key, cooldown_time=SIBLING_COOLDOWN_SECONDS ) resources.defer(lambda: client.proxy.delete_model(failing)) - backup = create_zero_weight_backup_deployment(client.proxy, group) + backup: Final = create_zero_weight_backup_deployment(client.proxy, group) resources.defer(lambda: client.proxy.delete_model(backup)) - _warm_cooldown_reads(sibling, scoped_key) + warm_group: Final = f"reliability-cooldown-sibling-warm-{unique_marker()}" + warm_id: Final = create_canned_deployment(client.proxy, warm_group) + resources.defer(lambda: client.proxy.delete_model(warm_id)) - tripped = _call_replica_without_retries(tripping.transport, scoped_key, group) + _warm_every_pod(client, scoped_key, warm_group) + tripped: Final = _call_without_retries(client, scoped_key, group) assert tripped.status_code == 500, ( - f"the first call through {tripping.url} should have surfaced the deployment's own 500, got " - f"{tripped.status_code}: {tripped.body[:300]}" + f"the first call should have surfaced the deployment's own 500, got {tripped.status_code}: " + f"{tripped.body[:300]}" ) - tripped_at = time.monotonic() + tripped_at: Final = time.monotonic() time.sleep(COOLDOWN_REDIS_READ_INTERVAL_SECONDS + SIBLING_READ_MARGIN_SECONDS) - _assert_served_by_backup( - _call_replica_without_retries(sibling.transport, scoped_key, group), - backup, - f"{time.monotonic() - tripped_at:.1f}s after {tripping.url} benched {failing}, on {sibling.url}", - ) + bench_lapses_at: Final = tripped_at + SIBLING_COOLDOWN_SECONDS - BENCH_MARGIN_SECONDS + for probe in range(1, SIBLING_PROBES + 1): + assert time.monotonic() < bench_lapses_at, ( + f"probe {probe} of {SIBLING_PROBES} would start after the {SIBLING_COOLDOWN_SECONDS:.0f}s bench can " + "lapse, so the earlier probes answered too slowly for this run to say anything about the read interval" + ) + _assert_served_by_backup( + _call_without_retries(client, scoped_key, group), + backup, + f"probe {probe} of {SIBLING_PROBES}, {time.monotonic() - tripped_at:.1f}s after the trip benched " + f"{failing},", + ) @pytest.mark.covers("reliability.cooldown.429.trips_then_recovers") def test_429_trips_cooldown_then_recovers(