mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-07 02:59:05 +00:00
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>
This commit is contained in:
parent
061d60b04c
commit
2156d6c3c7
4 changed files with 83 additions and 90 deletions
|
|
@ -324,3 +324,5 @@ other.<area>.<case>.<assertion>
|
|||
- spin up a local proxy by running the litellm proxy locally (`litellm --config <your-e2e-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
|
||||
|
|
|
|||
|
|
@ -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"}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue