test(e2e): tag router, batches and mcp tests with Subject metadata and record client steps (#44964)

* test(e2e): add enum values, auto-discovering label gates and secret hiding for e2e metadata

* test(e2e): tag router, batches and mcp tests with Subject metadata and record client steps

* test(e2e): leave the batches cleanup harness unit tests untagged

* docs(e2e): name every markerless harness test file that carries no Subject

* test(e2e): keep the step discovery comprehensions to one for clause

* test(e2e): name every driven model on the vllm batch, prompt caching and complexity router subjects
This commit is contained in:
ryan-crabbe-berri 2026-10-07 10:32:58 -07:00 • committed by GitHub
parent 77fc3315e5
commit c943d650f4
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
29 changed files with 773 additions and 36 deletions

View file

@ -8,6 +8,7 @@ from typing import Final, Protocol
from batch_client import BatchObject, FileDeleteResponse
from capabilities import is_cloud_storage_id, is_managed_id
from e2e_http import NetworkError, RateLimitedError, Result, Success, UnknownApiError
from e2e_metadata import STEP_FRAMES, step
from pydantic import BaseModel
CLEANUP_DELAYS: Final = (1.0, 2.0, 4.0)
@ -52,6 +53,7 @@ def _require_cleanup_success[R: BaseModel](result: Result[R], operation: str) ->
raise AssertionError(f"{operation} failed: {result.kind}")
@step("Clean up the uploaded file")
def cleanup_file(client: BatchCleanupClient, file_id: str, *, key: str, provider: str | None = None) -> None:
delete: Final[Callable[[], Result[FileDeleteResponse]]] = (
(lambda: client.delete_file_as_admin(file_id, provider=provider))
@ -65,7 +67,7 @@ def cleanup_file(client: BatchCleanupClient, file_id: str, *, key: str, provider
warnings.warn(
f"Left file {file_id} in place: LiteLLM refused to delete it while a batch still references it",
UserWarning,
stacklevel=2,
stacklevel=2 + STEP_FRAMES,
)
return
deleted: Final = _require_cleanup_success(result, f"Delete file {file_id}")
@ -74,6 +76,7 @@ def cleanup_file(client: BatchCleanupClient, file_id: str, *, key: str, provider
), f"Delete file {file_id} did not confirm deletion"
@step("Cancel the batch if it is still running")
def cleanup_batch(
client: BatchCleanupClient,
batch_id: str,
@ -137,7 +140,7 @@ def cleanup_batch(
warnings.warn(
f"Left batch {batch_id} cancelling after {BATCH_CANCEL_TIMEOUT_SECONDS}s for the provider to finish",
UserWarning,
stacklevel=2,
stacklevel=2 + STEP_FRAMES,
)
return
wait(BATCH_CANCEL_POLL_SECONDS)

View file

@ -17,6 +17,7 @@ from typing import Final, Literal
from pydantic import BaseModel, Field
from e2e_metadata import step
from proxy_client import ProxyClient
from e2e_http import (
FileUploadForm,
@ -136,12 +137,15 @@ def is_result_access_denied[R: BaseModel](result: Result[R]) -> bool:
class BatchClient:
proxy: ProxyClient
@step("Add a batch deployment named {model_name} that calls {litellm_params.model}")
def create_model(self, model_name: str, litellm_params: LiteLLMParamsBody) -> str:
return self.proxy.create_model(model_name, litellm_params, mode="batch")
@step("Delete the batch deployment")
def delete_model(self, model_id: str) -> None:
self.proxy.delete_model(model_id)
@step("Upload a batch input file to /v1/files")
def upload_file(
self,
*,
@ -161,6 +165,7 @@ class BatchClient:
response_type=FileObject,
)
@step("Retrieve the uploaded file")
def retrieve_file(
self, file_id: str, *, key: str, provider: str | None = None
) -> Result[FileObject]:
@ -171,6 +176,7 @@ class BatchClient:
response_type=FileObject,
)
@step("List the files the key can see from /v1/files")
def list_files(self, *, key: str, provider: str | None = None) -> Result[FileList]:
return self.proxy.transport.get(
_files_path(provider),
@ -179,6 +185,7 @@ class BatchClient:
response_type=FileList,
)
@step("Create a batch of {body.endpoint} requests from the uploaded file")
def create_batch(
self, *, body: BatchCreateBody, key: str, provider: str | None = None
) -> StreamingResponse:
@ -188,6 +195,7 @@ class BatchClient:
json=body,
)
@step("Retrieve the batch")
def retrieve_batch(
self, batch_id: str, *, key: str, provider: str | None = None
) -> Result[BatchObject]:
@ -198,6 +206,7 @@ class BatchClient:
response_type=BatchObject,
)
@step("Cancel the batch")
def cancel_batch(
self, batch_id: str, *, key: str, provider: str | None = None
) -> Result[BatchObject]:
@ -208,6 +217,7 @@ class BatchClient:
response_type=BatchObject,
)
@step("List the batches the key can see from /v1/batches")
def list_batches(
self,
*,
@ -223,6 +233,7 @@ class BatchClient:
response_type=BatchList,
)
@step("Delete the uploaded file")
def delete_file(
self, file_id: str, *, key: str, provider: str | None = None
) -> Result[FileDeleteResponse]:
@ -233,6 +244,7 @@ class BatchClient:
response_type=FileDeleteResponse,
)
@step("Delete the uploaded file as the proxy admin")
def delete_file_as_admin(self, file_id: str, *, provider: str | None = None) -> Result[FileDeleteResponse]:
return self.proxy.transport.delete(
f"{_files_path(provider)}/{file_id}",

View file

@ -7,7 +7,11 @@ import os
from dataclasses import dataclass
from typing import Final, Literal
import pytest
from e2e_config import provider_edge_base, unique_marker
from e2e_metadata import Domain, Mode, Route, Subject, meta
from e2e_metadata import Provider as MetaProvider
from models import LiteLLMParamsBody
_BATCH_RUN = unique_marker()
@ -18,6 +22,9 @@ def batch_model_name(base: str) -> str:
OPENAI_BATCH_BACKEND: Final = "gpt-4o-mini"
AZURE_BATCH_BACKEND: Final = "gpt-5.4-mini-batch"
VERTEX_BATCH_BACKEND: Final = "gemini-2.5-flash"
BEDROCK_BATCH_BACKEND: Final = "bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0"
def openai_batch_params() -> LiteLLMParamsBody:
@ -65,14 +72,14 @@ class Provider:
return openai_batch_params()
case "azure":
return LiteLLMParamsBody(
model="azure/gpt-5.4-mini-batch",
model=f"azure/{AZURE_BATCH_BACKEND}",
api_base="os.environ/AZURE_API_BASE",
api_key="os.environ/AZURE_API_KEY",
api_version="2025-04-01-preview",
)
case "vertex_ai":
return LiteLLMParamsBody(
model="vertex_ai/gemini-2.5-flash",
model=f"vertex_ai/{VERTEX_BATCH_BACKEND}",
vertex_project="os.environ/VERTEXAI_PROJECT",
vertex_location="us-central1",
vertex_credentials="os.environ/VERTEXAI_CREDENTIALS",
@ -81,7 +88,7 @@ class Provider:
)
case "bedrock":
return LiteLLMParamsBody(
model="bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0",
model=BEDROCK_BATCH_BACKEND,
aws_access_key_id="os.environ/AWS_ACCESS_KEY_ID",
aws_secret_access_key="os.environ/AWS_SECRET_ACCESS_KEY",
aws_region_name="os.environ/AWS_REGION",
@ -132,21 +139,21 @@ PROVIDERS: tuple[Provider, ...] = (
Provider(
"azure",
batch_model_name("azure-batch"),
"gpt-5.4-mini-batch",
AZURE_BATCH_BACKEND,
can_cancel=True,
can_list=True,
),
Provider(
"vertex_ai",
batch_model_name("vertex-batch"),
"gemini-2.5-flash",
VERTEX_BATCH_BACKEND,
can_cancel=True,
can_list=True,
),
Provider(
"bedrock",
batch_model_name("bedrock-batch"),
"bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0",
BEDROCK_BATCH_BACKEND,
can_cancel=True,
can_list=True,
),
@ -181,6 +188,29 @@ CAPABILITIES: tuple[Capability, ...] = tuple(
)
def lifecycle_meta(cap: Capability) -> pytest.MarkDecorator:
return meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider(cap.provider),),
models=(cap.raw_model,),
mode=Mode.BATCH,
)
)
def file_content_meta(provider: Provider) -> pytest.MarkDecorator:
return meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider(provider.name),),
models=(provider.raw_model,),
)
)
def raw_id_matches_provider(provider: str, batch_id: str) -> bool:
if provider in ("openai", "azure"):
return batch_id.startswith("batch")

View file

@ -44,17 +44,22 @@ from capabilities import (
OPENAI_BATCH_BACKEND,
OPENAI_BATCH_MODEL,
PROVIDERS,
VERTEX_BATCH_BACKEND,
Capability,
Provider,
batch_model_name,
coverage_cells_for_lifecycle,
decoded_model_from_id,
file_content_meta,
is_managed_id,
lifecycle_meta,
matches_id_shape,
openai_batch_params,
raw_id_matches_provider,
)
from e2e_config import MASTER_KEY, PROXY_BASE_URL, unique_marker
from e2e_metadata import Domain, Mode, Route, Subject, meta
from e2e_metadata import Provider as MetaProvider
from e2e_http import (
FileUploadForm,
Result,
@ -249,7 +254,7 @@ def assert_batch_object(batch: BatchObject) -> None:
pytest.param(
cap,
id=cap.id,
marks=pytest.mark.covers(*coverage_cells_for_lifecycle(cap)),
marks=(pytest.mark.covers(*coverage_cells_for_lifecycle(cap)), lifecycle_meta(cap)),
)
for cap in CAPABILITIES
],
@ -350,6 +355,15 @@ def test_batch_lifecycle(
@pytest.mark.covers("llm.batches.openai.key_model_access_denied.nonstream.works")
@meta(
Subject(
domain=Domain.PROXY_AUTH,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_batch_key_model_access_denied(
client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -389,6 +403,14 @@ def test_batch_key_model_access_denied(
"llm.files.openai.upload.nonstream.works",
"llm.files.openai.delete.nonstream.works",
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
)
)
def test_file_upload_and_delete_outputs(
client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -433,6 +455,15 @@ def unattributed_rows(rows: list[SpendLogRow]) -> list[SpendLogRow]:
"once the fetch is bounded."
)
)
@meta(
Subject(
domain=Domain.SPEND_BUDGETS,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_rate_limited_batch_create_leaves_no_unattributed_spend_row(
client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -471,7 +502,7 @@ def test_rate_limited_batch_create_leaves_no_unattributed_spend_row(
file = unwrap(
client.upload_file(
content=render_jsonl("gpt-4o-mini"),
content=render_jsonl(OPENAI_BATCH_BACKEND),
form=FileUploadForm(purpose="batch"),
model=OPENAI_BATCH_MODEL,
key=key,
@ -520,6 +551,14 @@ class TestBatchFileContent:
"llm.files.openai.content.nonstream.works",
exercised_on=["files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
)
)
def test_file_content_matches_upload(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -558,8 +597,9 @@ class TestBatchFileContent:
pytest.param(
p,
id=p.name,
marks=pytest.mark.covers(
FILE_CONTENT_CELLS[p.name], exercised_on=["files"]
marks=(
pytest.mark.covers(FILE_CONTENT_CELLS[p.name], exercised_on=["files"]),
file_content_meta(p),
),
)
for p in PROVIDERS
@ -632,6 +672,14 @@ class TestOpenAIFiles:
"marker when LIT-4820 is fixed; do not relax the assertion to make it pass."
)
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
)
)
def test_uploaded_file_appears_in_list(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -662,6 +710,7 @@ class TestOpenAIFiles:
"llm.files.openai.list_isolation.nonstream.works",
exercised_on=["files"],
)
@meta(Subject(domain=Domain.LLM_TRANSLATION, route=Route.FILES, providers=(MetaProvider.OPENAI,)))
def test_list_page_cursors_address_only_the_callers_own_files(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -697,6 +746,14 @@ class TestOpenAIFiles:
"llm.files.openai.retrieve.nonstream.works",
exercised_on=["files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
)
)
def test_retrieve_round_trips_metadata(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -760,6 +817,15 @@ class TestBatchRateLimitErrorMapping:
"quota_management.ratelimit.batch_rpm.blocks_over_limit",
exercised_on=["batches"],
)
@meta(
Subject(
domain=Domain.SPEND_BUDGETS,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_batch_create_over_rpm_returns_mapped_429(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -773,7 +839,7 @@ class TestBatchRateLimitErrorMapping:
file = unwrap(
client.upload_file(
content=_multi_request_jsonl("gpt-4o-mini", BATCH_RL_REQUEST_LINES),
content=_multi_request_jsonl(OPENAI_BATCH_BACKEND, BATCH_RL_REQUEST_LINES),
form=FileUploadForm(purpose="batch"),
model=OPENAI_BATCH_MODEL,
key=key,
@ -826,7 +892,7 @@ class TestBatchEnqueuedTokenLimit:
) -> FileObject:
file = unwrap(
client.upload_file(
content=_multi_request_jsonl("gpt-4o-mini", BATCH_RL_REQUEST_LINES),
content=_multi_request_jsonl(OPENAI_BATCH_BACKEND, BATCH_RL_REQUEST_LINES),
form=FileUploadForm(purpose="batch"),
model=OPENAI_BATCH_MODEL,
key=key,
@ -859,6 +925,15 @@ class TestBatchEnqueuedTokenLimit:
"quota_management.ratelimit.batch_enqueued_tokens.accepts_over_rpm",
exercised_on=["batches"],
)
@meta(
Subject(
domain=Domain.SPEND_BUDGETS,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_enqueued_allowance_accepts_batch_over_key_rpm(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -890,6 +965,15 @@ class TestBatchEnqueuedTokenLimit:
"quota_management.ratelimit.batch_enqueued_tokens.refunds_on_cancel",
exercised_on=["batches"],
)
@meta(
Subject(
domain=Domain.SPEND_BUDGETS,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_exhausted_allowance_blocks_until_cancel_refunds(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -985,6 +1069,15 @@ class TestBedrockBatchAssumeRole:
"llm.files.bedrock.upload.nonstream.works",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.BEDROCK,),
models=(ASSUME_ROLE_RAW_MODEL,),
mode=Mode.BATCH,
)
)
def test_unified_batch_create_with_assume_role(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -1052,6 +1145,14 @@ class TestBedrockBatchSplitS3Credentials:
"llm.files.bedrock.split_s3_credentials.nonstream.works",
exercised_on=["files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider.BEDROCK,),
models=(ASSUME_ROLE_RAW_MODEL,),
)
)
def test_file_lifecycle_signs_s3_with_s3_credentials(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -1123,6 +1224,15 @@ class TestBedrockBatchGovCloud:
"llm.files.bedrock.govcloud_partition.nonstream.works",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.BEDROCK,),
models=(GOVCLOUD_RAW_MODEL,),
mode=Mode.BATCH,
)
)
def test_unified_file_upload_and_batch_create_in_govcloud(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -1193,6 +1303,14 @@ class TestGeminiFiles:
"llm.files.gemini.upload.nonstream.works",
exercised_on=["files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
providers=(MetaProvider.GEMINI,),
models=(GEMINI_FILES_RAW_MODEL,),
)
)
def test_gemini_file_upload(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -1227,7 +1345,7 @@ def _vllm_params(api_base: str, api_key: str | None, model_id: str) -> LiteLLMPa
)
HOSTED_VLLM_DEFAULT_MODEL = "Qwen/Qwen2.5-0.5B-Instruct"
HOSTED_VLLM_MODEL: Final = (os.environ.get("HOSTED_VLLM_MODEL") or "Qwen/Qwen2.5-0.5B-Instruct").strip()
HOSTED_VLLM_BAD_LINE_CUSTOM_ID = "req-bad"
@ -1236,9 +1354,8 @@ def _hosted_vllm_deployment(client: BatchClient, resources: ResourceManager) ->
if api_base is None:
pytest.skip("set HOSTED_VLLM_API_BASE (the live vLLM server this deployment targets)")
api_key = (os.environ.get("HOSTED_VLLM_API_KEY") or "").strip() or None
model_id = (os.environ.get("HOSTED_VLLM_MODEL") or HOSTED_VLLM_DEFAULT_MODEL).strip()
proxy_name = batch_model_name("hosted-vllm-batch")
model_row_id = client.create_model(proxy_name, _vllm_params(api_base, api_key, model_id))
model_row_id = client.create_model(proxy_name, _vllm_params(api_base, api_key, HOSTED_VLLM_MODEL))
resources.defer(lambda: client.delete_model(model_row_id))
return proxy_name
@ -1290,6 +1407,15 @@ class TestHostedVllmBatch:
"llm.files.hosted_vllm.upload.nonstream.works",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.HOSTED_VLLM,),
models=(HOSTED_VLLM_MODEL,),
mode=Mode.BATCH,
)
)
def test_batch_runs_to_completion_with_a_downloadable_output(
self, client: BatchClient, resources: ResourceManager, upload_route: str
) -> None:
@ -1337,6 +1463,15 @@ class TestHostedVllmBatch:
)
@pytest.mark.covers("llm.batches.hosted_vllm.basic.nonstream.works", exercised_on=["batches", "files"])
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.HOSTED_VLLM,),
models=(HOSTED_VLLM_MODEL,),
mode=Mode.BATCH,
)
)
def test_failing_line_lands_in_the_error_file_not_the_batch_status(
self, client: BatchClient, resources: ResourceManager
) -> None:
@ -1417,6 +1552,7 @@ class TestBatchFailurePaths:
"llm.batches.openai.malformed_jsonl.nonstream.works",
exercised_on=["files"],
)
@meta(Subject(domain=Domain.LLM_TRANSLATION, route=Route.FILES))
def test_malformed_jsonl_upload_rejected(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -1439,13 +1575,22 @@ class TestBatchFailurePaths:
"llm.batches.openai.cancel_terminal.nonstream.works",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_endpoint_mismatch_fails_batch_and_cancel_conflicts(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
key = resources.key()
file = unwrap(
client.upload_file(
content=_mismatched_endpoint_jsonl("gpt-4o-mini"),
content=_mismatched_endpoint_jsonl(OPENAI_BATCH_BACKEND),
form=FileUploadForm(purpose="batch"),
model=OPENAI_BATCH_MODEL,
key=key,
@ -1495,6 +1640,15 @@ class TestBatchFailurePaths:
"llm.batches.openai.foreign_file_id.nonstream.works",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.AZURE,),
models=(AZURE_BATCH_RAW_MODEL,),
mode=Mode.BATCH,
)
)
def test_foreign_encoded_file_id_routes_by_file_model(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -1544,6 +1698,15 @@ class TestBatchSecondHop:
"llm.batches.openai.second_hop.nonstream.works",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.LITELLM_PROXY, MetaProvider.OPENAI),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_unified_create_and_retrieve_via_chained_gateway(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -1561,7 +1724,7 @@ class TestBatchSecondHop:
file = unwrap(
client.upload_file(
content=render_jsonl("gpt-4o-mini"),
content=render_jsonl(OPENAI_BATCH_BACKEND),
form=FileUploadForm(purpose="batch", target_model_names=hop_name),
key=key,
)
@ -1680,13 +1843,22 @@ class TestBatchTerminalState:
"llm.batches.openai.terminal_state.nonstream.cost_logged",
exercised_on=["batches", "files"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_completed_batch_downloads_output_and_books_cost(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
key = resources.key()
file = unwrap(
client.upload_file(
content=render_jsonl("gpt-4o-mini"),
content=render_jsonl(OPENAI_BATCH_BACKEND),
form=FileUploadForm(purpose="batch"),
model=OPENAI_BATCH_MODEL,
key=key,
@ -1786,6 +1958,15 @@ class TestVertexNativePassthrough:
"llm.batches.vertex.native_passthrough.nonstream.works",
exercised_on=["files", "batches"],
)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.BATCHES,
providers=(MetaProvider.VERTEX_AI,),
models=(VERTEX_BATCH_BACKEND,),
mode=Mode.BATCH,
)
)
def test_native_jsonl_round_trips_untouched_and_starts_a_batch(
self, client: BatchClient, resources: ResourceManager, batch_deployments: None
) -> None:
@ -1848,6 +2029,7 @@ class TestVertexNativePassthrough:
),
],
)
@meta(Subject(domain=Domain.LLM_TRANSLATION, route=Route.FILES))
def test_passthrough_upload_is_rejected_outside_a_native_vertex_batch(
self,
content: bytes,

View file

@ -22,9 +22,10 @@ import pytest
from batch_client import BatchClient, FileObject
from batch_cleanup import cleanup_file
from capabilities import batch_model_name, is_managed_id, openai_batch_params
from capabilities import OPENAI_BATCH_BACKEND, batch_model_name, is_managed_id, openai_batch_params
from e2e_config import unique_marker
from e2e_http import FileUploadForm, Result, UnknownApiError, unwrap
from e2e_metadata import Domain, Provider, Route, Subject, meta
from lifecycle import ResourceManager
pytestmark = [pytest.mark.e2e, pytest.mark.managed_files]
@ -64,6 +65,12 @@ def managed_model(client: BatchClient) -> Iterator[str]:
@pytest.mark.covers(UPLOAD_ROW)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
)
)
def test_upload_without_target_model_names_rejected(
client: BatchClient, scoped_key: str, managed_model: str
) -> None:
@ -76,6 +83,12 @@ def test_upload_without_target_model_names_rejected(
@pytest.mark.covers(UPLOAD_ROW)
@meta(
Subject(
domain=Domain.LLM_TRANSLATION,
route=Route.FILES,
)
)
def test_upload_with_model_param_rejected(
client: BatchClient, scoped_key: str, managed_model: str
) -> None:
@ -89,12 +102,26 @@ def test_upload_with_model_param_rejected(
@pytest.mark.covers(ISOLATION_ROW)
@meta(
Subject(
domain=Domain.PROXY_AUTH,
route=Route.FILES,
)
)
def test_raw_provider_file_id_rejected(client: BatchClient, scoped_key: str) -> None:
result = client.retrieve_file("file-e2e-raw-provider-id", key=scoped_key)
expect_api_error(result, 400, "Raw provider file ids cannot be used")
@pytest.mark.covers(ISOLATION_ROW)
@meta(
Subject(
domain=Domain.PROXY_AUTH,
route=Route.FILES,
providers=(Provider.OPENAI,),
models=(OPENAI_BATCH_BACKEND,),
)
)
def test_cross_user_managed_id_denied_owner_allowed(
client: BatchClient, resources: ResourceManager, managed_model: str
) -> None:

View file

@ -15,14 +15,11 @@ from typing import Protocol, cast
import pytest
from datadog_mcp import DdLogsReader
from mcp_client import McpClient, build_client
from proxy_client import ProxyClient
class DdLogsReader(Protocol):
def poll_events_for_marker(self, marker: str) -> list[object]: ...
class _DdLogsReaderBuilder(Protocol):
def __call__(self) -> DdLogsReader: ...

View file

@ -4,14 +4,20 @@ from __future__ import annotations
import os
from collections.abc import Sequence
from typing import Protocol
from e2e_config import datadog_mcp_url, unique_marker
from e2e_metadata import step
from lifecycle import ResourceManager
from mcp_client import McpClient
SEARCH_LOGS_TOOL = "search_datadog_logs"
class DdLogsReader(Protocol):
def poll_events_for_marker(self, marker: str) -> list[object]: ...
def _dd_api_key() -> str:
return os.environ.get("DD_API_KEY", "").strip()
@ -31,6 +37,7 @@ def assert_dd_mcp_creds() -> None:
)
@step("Register the Datadog remote MCP server with its credentials from the environment")
def register_datadog_mcp(
client: McpClient,
resources: ResourceManager,

View file

@ -21,6 +21,7 @@ from pydantic import BaseModel, ConfigDict, Field, RootModel
from e2e_config import settle_propagation
from e2e_http import Headers, NoBody, Result, Success, UnknownApiError, unwrap
from e2e_metadata import step
from models import KeyGenerateBody, McpServerListResponse, McpServerRow, ObjectPermission
from proxy_client import ProxyClient
@ -153,6 +154,7 @@ class McpCallToolResponse(BaseModel):
class McpClient:
proxy: ProxyClient
@step("Register the MCP server {server_name} with the alias {alias}")
def register_server(
self,
*,
@ -183,6 +185,7 @@ class McpClient:
)
).server_id
@step("Delete the MCP server")
def delete_server(self, server_id: str) -> None:
_ = self.proxy.transport.delete(
f"/v1/mcp/server/{server_id}",
@ -191,6 +194,7 @@ class McpClient:
response_type=NoBody,
)
@step("List the MCP servers from /v1/mcp/server")
def registered_servers(self) -> list[McpServerRow]:
return unwrap(
self.proxy.transport.get(
@ -201,6 +205,7 @@ class McpClient:
)
).root
@step("List the MCP servers the key can see from /v1/mcp/server")
def list_servers(self, key: str) -> Result[McpServerListResponse]:
return self.proxy.transport.get(
"/v1/mcp/server",
@ -209,6 +214,7 @@ class McpClient:
response_type=McpServerListResponse,
)
@step("Check the health of the MCP servers the key can see from /v1/mcp/server/health")
def server_health(self, key: str, server_ids: list[str] | None = None) -> Result[McpHealthResponse]:
return self.proxy.transport.get(
"/v1/mcp/server/health",
@ -217,6 +223,7 @@ class McpClient:
response_type=McpHealthResponse,
)
@step("Wait for every proxy replica to list the MCP server in /v1/mcp/server")
def await_registered(self, server_id: str) -> McpServerRow:
"""Wait for every configured replica to list the server and return its row."""
registered = self.proxy.read_body_back_everywhere(
@ -228,6 +235,7 @@ class McpClient:
row for response in registered.values() for row in response.root if row.server_id == server_id
)
@step("Generate a virtual key for the user {user_id}")
def generate_key(
self,
*,
@ -254,6 +262,7 @@ class McpClient:
)
)
@step("List the MCP tools the key can see from /mcp-rest/tools/list")
def list_tools(self, key: str) -> Result[McpToolsListResponse]:
return self.proxy.transport.get(
"/mcp-rest/tools/list",
@ -262,6 +271,7 @@ class McpClient:
response_type=McpToolsListResponse,
)
@step('Wait for /mcp-rest/tools/list to show the MCP server\'s tool matching "{needle}"')
def await_tool(self, key: str, server_id: str, needle: str) -> str:
"""Poll tools/list until `server_id` serves a tool matching `needle`, and
return its fully-qualified name. Fails at poll_timeout.
@ -287,6 +297,7 @@ class McpClient:
)
time.sleep(self.proxy.poll_interval)
@step("Wait for /mcp-rest/tools/list to show the key exactly the expected tools on the MCP server")
def await_tools(self, key: str, server_id: str, *, expected: frozenset[str]) -> frozenset[str]:
"""Poll tools/list until `server_id`'s tools as `key` sees them are exactly
`expected`, and return the last listing either way, so the caller's equality
@ -301,6 +312,7 @@ class McpClient:
return unwrap(result).tool_names_for_server(server_id)
time.sleep(self.proxy.poll_interval)
@step("Call the MCP tool {name} through /mcp-rest/tools/call")
def await_call_tool(
self,
key: str,
@ -329,6 +341,7 @@ class McpClient:
)
time.sleep(self.proxy.poll_interval)
@step("Call the MCP tool {name} through /mcp-rest/tools/call and wait for a 403")
def await_call_tool_denied(
self,
key: str,
@ -355,6 +368,7 @@ class McpClient:
)
time.sleep(self.proxy.poll_interval)
@step('Create the guardrail {name} that blocks MCP tool calls containing "{blocked_keyword}"')
def register_mcp_content_filter(self, *, name: str, blocked_keyword: str) -> str:
"""Register a default-on content-filter guardrail that runs on the MCP
tool-call hook (pre_mcp_call) and blocks a single keyword. The keyword is
@ -378,6 +392,7 @@ class McpClient:
settle_propagation(time.monotonic())
return guardrail_id
@step("Delete the guardrail")
def delete_guardrail(self, guardrail_id: str) -> None:
_ = self.proxy.transport.delete(
f"/guardrails/{guardrail_id}",
@ -386,6 +401,7 @@ class McpClient:
response_type=NoBody,
)
@step("Call the MCP tool {name} through /mcp-rest/tools/call with {arguments}")
def call_tool(
self,
key: str,

View file

@ -26,6 +26,7 @@ import httpx2
import pytest
from e2e_config import PROXY_BASE_URL, REQUEST_TIMEOUT
from e2e_http import AuthHeaders, NoBody, unwrap
from e2e_metadata import step
from idp import Identity
from mcp import ClientSession
from mcp.client.auth import OAuthClientProvider
@ -318,6 +319,7 @@ async def _list_and_call(
class ChatMcpClient:
proxy: ProxyClient
@step("Register the MCP server with the alias {body.alias}")
def create_server(self, body: McpServerCreateBody) -> McpServerInfo:
return unwrap(
self.proxy.transport.post(
@ -328,6 +330,7 @@ class ChatMcpClient:
)
)
@step("Read the MCP server back from /v1/mcp/server")
def server_info(self, server_id: str) -> McpServerInfo:
return unwrap(
self.proxy.transport.get(
@ -338,6 +341,7 @@ class ChatMcpClient:
)
)
@step("Delete the MCP server")
def delete_server(self, server_id: str) -> None:
_ = self.proxy.transport.delete(
f"/v1/mcp/server/{server_id}",
@ -346,6 +350,7 @@ class ChatMcpClient:
response_type=NoBody,
)
@step("Sign the key's user in to the MCP server {alias} through the OAuth consent flow")
def seed_user_token(self, alias: str, key: str, storage_state_path: str) -> tuple[str, ...]:
"""Drive the interactive authorize dance for `key`'s user so the gateway
stores their upstream token, retried to the shared deadline since the
@ -367,6 +372,7 @@ class ChatMcpClient:
f"last error: {last_error!r}"
)
@step("List the tools on the MCP server {alias} and call {tool} over the MCP protocol")
def list_and_call(
self,
alias: str,
@ -394,6 +400,7 @@ class ChatMcpClient:
)
)
@step("List the users with a stored OAuth token for the MCP server")
def server_user_credentials(self, server_id: str) -> tuple[McpServerUserCredentialRow, ...]:
return unwrap(
self.proxy.transport.get(
@ -404,6 +411,7 @@ class ChatMcpClient:
)
).root
@step("Revoke the user's stored OAuth token for the MCP server")
def revoke_user_token(self, server_id: str, headers: AuthHeaders) -> None:
_ = unwrap(
self.proxy.transport.delete(
@ -414,6 +422,7 @@ class ChatMcpClient:
)
)
@step("Send a /chat/completions request to {body.model} with an MCP server attached as a tool")
def chat_with_mcp(self, headers: AuthHeaders, body: ChatBody) -> ChatResponse:
"""POST /chat/completions carrying the LiteLLM key in `headers` (either
ingress form) with an MCP server attached in `body.tools`. The gateway

View file

@ -21,6 +21,7 @@ from typing import Final
import psycopg
from e2e_config import INHERITED_ENV_PREFIXES, available_port
from e2e_http import NoBody
from e2e_metadata import step
from idp import Keycloak, stop_process_group
from proxy_client import ProxyClient, build_proxy_client
from psycopg.rows import class_row
@ -37,6 +38,7 @@ class CredentialRow:
credential_b64: str = field(repr=False)
@step("Read the user's stored OAuth credential for the MCP server from the database and decrypt it")
def stored_oauth(user_id: str, server_id: str) -> StoredOAuth:
"""Read the encrypted credential because management APIs omit the plaintext token."""
from litellm.proxy.common_utils.encrypt_decrypt_utils import decrypt_value_helper
@ -108,6 +110,7 @@ class OAuthGateway:
_log_path: Path
_child: subprocess.Popen[bytes] | None = field(default=None, init=False, repr=False)
@step("Start the separate LiteLLM proxy and wait for /health/liveliness")
def start(self) -> None:
with self._log_path.open("ab") as log:
self._child = subprocess.Popen(
@ -126,11 +129,13 @@ class OAuthGateway:
time.sleep(0.5)
raise AssertionError("owned OAuth gateway did not become ready")
@step("Stop the separate LiteLLM proxy")
def stop(self) -> None:
if self._child is not None:
stop_process_group(self._child)
assert self._child.poll() is not None, "old gateway process is still alive"
@step("Restart the separate LiteLLM proxy process so its in-memory caches start empty")
def restart(self) -> None:
assert self._child is not None
previous: Final = self._child.pid
@ -139,6 +144,7 @@ class OAuthGateway:
assert self._child.pid != previous, "gateway restart did not create a new process"
@step("Start a separate LiteLLM proxy from source with JWT auth against Keycloak")
def owned_gateway(idp: Keycloak, directory: Path, cleanup: ExitStack) -> OAuthGateway:
for name in ("DATABASE_URL", "LITELLM_LICENSE", "LITELLM_SALT_KEY", "LITELLM_MASTER_KEY"):
assert os.environ.get(name), f"{name} is required for the owned OAuth gateway"

View file

@ -16,6 +16,7 @@ import pytest
from datadog_mcp import SEARCH_LOGS_TOOL, register_datadog_mcp
from e2e_config import unique_marker
from e2e_http import unwrap
from e2e_metadata import Domain, Route, Subject, meta
from lifecycle import ResourceManager
from mcp_client import McpClient
@ -24,6 +25,12 @@ pytestmark = pytest.mark.e2e
class TestMcpAccessGroupToolSelection:
@pytest.mark.covers("mcp.list_tools.api_key.access_group_scoped")
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
def test_access_group_scopes_tool_selection(
self, client: McpClient, resources: ResourceManager
) -> None:

View file

@ -30,6 +30,7 @@ import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, LINEAR_MCP_URL, LINEAR_STORAGE_STATE, unique_marker
from e2e_http import AuthHeaders
from e2e_metadata import Capability, Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from models import ChatBody, ChatMessage, KeyGenerateBody, McpChatTool, McpServerCreateBody, ObjectPermission
from proxy_client import ProxyClient
@ -70,6 +71,16 @@ class TestMcpChatCompletionOauth:
@pytest.mark.covers("mcp.list_tools.oauth.succeeds")
@pytest.mark.covers("mcp.call_tool.oauth.succeeds")
@meta(
Subject(
domain=Domain.MCP,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
capabilities=(Capability.FUNCTION_CALLING,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completion_uses_linear_with_x_litellm_api_key_header(
self, chat_client: ChatMcpClient, resources: ResourceManager
) -> None:
@ -134,6 +145,16 @@ class TestMcpChatCompletionOauth:
@pytest.mark.covers("mcp.list_tools.oauth.succeeds")
@pytest.mark.covers("mcp.call_tool.oauth.succeeds")
@meta(
Subject(
domain=Domain.MCP,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
capabilities=(Capability.FUNCTION_CALLING,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completion_uses_linear_with_authorization_bearer_header(
self, chat_client: ChatMcpClient, resources: ResourceManager
) -> None:

View file

@ -13,10 +13,10 @@ from __future__ import annotations
import pytest
from conftest import DdLogsReader
from datadog_mcp import SEARCH_LOGS_TOOL, assert_dd_mcp_creds, register_datadog_mcp
from datadog_mcp import SEARCH_LOGS_TOOL, DdLogsReader, assert_dd_mcp_creds, register_datadog_mcp
from e2e_config import CHEAP_ANTHROPIC_MODEL, DD_SEARCH_FROM, unique_marker
from e2e_http import NoBody, unwrap
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from mcp_client import McpClient
from models import ChatBody, ChatMessage
@ -50,6 +50,15 @@ def _seed_completion(proxy: ProxyClient, *, key: str, marker: str) -> None:
class TestDatadogMcpRoundTrip:
@pytest.mark.covers("mcp.list_tools.api_key.succeeds", "mcp.call_tool.api_key.succeeds")
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_search_logs_finds_seeded_completion(
self,
client: McpClient,

View file

@ -23,6 +23,7 @@ import pytest
from datadog_mcp import SEARCH_LOGS_TOOL, assert_dd_mcp_creds, register_datadog_mcp
from e2e_config import DD_SEARCH_FROM, unique_marker
from e2e_http import Result, Success, UnknownApiError
from e2e_metadata import Domain, Route, Subject, meta
from lifecycle import ResourceManager
from mcp_client import McpCallToolResponse, McpClient, McpToolArguments
@ -82,6 +83,12 @@ class TestMcpToolCallGuardrail:
"guardrail.litellm_content_filter.pre_mcp_call.blocks",
exercised_on=["mcp_operations"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MCP,
)
)
def test_content_filter_blocks_banned_keyword_in_tool_args(
self, client: McpClient, resources: ResourceManager
) -> None:

View file

@ -18,6 +18,7 @@ from typing import Final
from datadog_mcp import SEARCH_LOGS_TOOL, register_datadog_mcp
from e2e_config import DD_SEARCH_FROM, unique_marker
from e2e_http import unwrap
from e2e_metadata import Domain, Route, Subject, meta
from lifecycle import ResourceManager
from mcp_client import McpClient
from models import KeyGenerateBody, ObjectPermission
@ -33,6 +34,12 @@ def _key(client: McpClient, resources: ResourceManager, *, mcp_servers: list[str
class TestMcpKeyGrantByAlias:
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
def test_alias_grant_persists_verbatim_and_lists_tools(
self,
client: McpClient,
@ -62,6 +69,12 @@ class TestMcpKeyGrantByAlias:
class TestMcpKeyWithoutAccessIsDenied:
@pytest.mark.covers("mcp.list_tools.api_key.denied_without_permission")
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
def test_list_tools_denied_without_permission(
self,
client: McpClient,
@ -82,6 +95,12 @@ class TestMcpKeyWithoutAccessIsDenied:
)
@pytest.mark.covers("mcp.call_tool.api_key.denied_without_permission")
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
def test_call_tool_denied_without_permission(
self,
client: McpClient,
@ -113,6 +132,12 @@ class TestMcpKeyWithoutAccessIsDenied:
class TestMcpHealthVisibility:
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
def test_route_restricted_health_matches_server_grants(
self,
client: McpClient,

View file

@ -18,6 +18,7 @@ from typing import Final, Literal
import pytest
from e2e_config import LINEAR_MCP_URL, LINEAR_READONLY_TOOL, LINEAR_STORAGE_STATE, unique_marker
from e2e_http import AuthHeaders, NoBody, get_external, unwrap
from e2e_metadata import Domain, Route, Subject, meta
from idp import Identity, Keycloak
from lifecycle import ResourceManager
from models import (
@ -88,6 +89,12 @@ class TestMcpOauthHappyPath:
@pytest.mark.covers("mcp.list_tools.oauth.succeeds")
@pytest.mark.covers("mcp.call_tool.oauth.succeeds")
@pytest.mark.covers("mcp.call_tool.oauth.persists_across_processes")
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
@pytest.mark.parametrize("route", ("aggregate_sso", "explicit_header_jwt"))
@pytest.mark.parametrize("observed", (False, True), ids=("direct", "observed"))
def test_consent_list_call_and_cold_restart(

View file

@ -18,6 +18,7 @@ import pytest
from datadog_mcp import SEARCH_LOGS_TOOL, register_datadog_mcp
from e2e_config import unique_marker
from e2e_http import unwrap
from e2e_metadata import Domain, Route, Subject, meta
from lifecycle import ResourceManager
from mcp_client import McpClient
from models import ToolsetCreateBody, ToolsetTool
@ -60,6 +61,12 @@ def _wire_prefix(wire_name: str, tool_name: str, catalog: frozenset[str]) -> str
class TestMcpToolsetEnforcement:
@pytest.mark.covers("mcp.list_tools.api_key.toolset_scoped")
@meta(
Subject(
domain=Domain.MCP,
route=Route.MCP,
)
)
def test_key_granted_a_toolset_lists_exactly_its_tools(self, client: McpClient, resources: ResourceManager) -> None:
server_id: Final = register_datadog_mcp(client, resources, allowed_tools=None)
client.await_registered(server_id)

View file

@ -23,6 +23,7 @@ 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 e2e_metadata import step
from models import (
CacheControl,
ChatMessage,
@ -79,6 +80,7 @@ def cached_system_turn(marker: str) -> ChatMessage:
return ChatMessage(role="system", content=[TextContentPart(text=filler, cache_control=CacheControl())])
@step(f"Add a deployment named {{name}} that calls {REAL_MODEL} at an unreachable address")
def create_bad_base_deployment(proxy: ProxyClient, name: str) -> str:
"""Register a deployment pointing at an unreachable base, so every call to it
fails with a real connection error the fallback can reroute around."""
@ -87,6 +89,7 @@ def create_bad_base_deployment(proxy: ProxyClient, name: str) -> str:
)
@step(f"Add a deployment named {{name}} that calls {REAL_MODEL} at an unreachable address and is never benched")
def create_never_benched_refusing_deployment(proxy: ProxyClient, name: str) -> str:
return proxy.create_model(
name,
@ -94,6 +97,7 @@ def create_never_benched_refusing_deployment(proxy: ProxyClient, name: str) -> s
)
@step(f"Add a deployment named {{name}} that calls {REAL_MODEL} with a 1ms timeout")
def create_timeout_deployment(proxy: ProxyClient, name: str) -> str:
"""Register a deployment with a 1ms deadline the real backend always exceeds."""
return proxy.create_model(
@ -101,12 +105,14 @@ def create_timeout_deployment(proxy: ProxyClient, name: str) -> str:
)
@step(f"Add a deployment named {{name}} that calls the small-context model {SMALL_CONTEXT_MODEL}")
def create_small_context_deployment(proxy: ProxyClient, name: str) -> str:
"""Register a deployment on the smallest-context model OpenAI still serves, so an
oversized prompt earns a real context-window refusal from the provider."""
return proxy.create_model(name, LiteLLMParamsBody(model=SMALL_CONTEXT_MODEL, api_key=REAL_KEY))
@step(f"Add a deployment named {{name}} that calls {AZURE_MODEL} behind Azure's content filter")
def create_content_filtered_deployment(proxy: ProxyClient, name: str) -> str:
"""Register the Azure OpenAI deployment whose content filter refuses
CONTENT_POLICY_PROMPT with a real policy-violation 400 (the one live trigger
@ -124,6 +130,10 @@ def create_content_filtered_deployment(proxy: ProxyClient, name: str) -> str:
)
@step(
f"Add a deployment named {{name}} that calls {AZURE_MODEL} and is benched for {{cooldown_time}}s"
" on its first failure"
)
def create_azure_benched_on_first_failure_deployment(proxy: ProxyClient, name: str, cooldown_time: float) -> str:
"""The live Azure OpenAI deployment holding all of the group's shuffle weight,
benched on its first failure of any class, with the client's own retries off."""
@ -144,6 +154,7 @@ def create_azure_benched_on_first_failure_deployment(proxy: ProxyClient, name: s
)
@step(f"Add a deployment named {{name}} that calls {CACHING_MODEL} with prompt caching")
def create_caching_deployment(proxy: ProxyClient, name: str) -> str:
"""Register the Anthropic deployment whose prompt cache the affinity check pins to."""
return proxy.create_model(name, LiteLLMParamsBody(model=CACHING_MODEL, api_key=CACHING_KEY, weight=1))
@ -165,6 +176,7 @@ def _register_benched_on_first_failure(
)
@step(f"Add a deployment named {{name}} that calls {REAL_MODEL} with a 1ms timeout and is benched on its first timeout")
def create_always_timing_out_deployment(proxy: ProxyClient, name: str, cooldown_time: float | None = None) -> str:
"""A 1ms deadline the real backend always exceeds, benched on its first Timeout."""
return _register_benched_on_first_failure(
@ -176,6 +188,7 @@ def create_always_timing_out_deployment(proxy: ProxyClient, name: str, cooldown_
)
@step(f"Add a deployment named {{name}} that calls {REAL_MODEL} with an invalid key and is benched on its first 401")
def create_always_unauthorized_deployment(proxy: ProxyClient, name: str, cooldown_time: float | None = None) -> str:
"""A key the real backend rejects with a 401, benched on its first AuthenticationError."""
return _register_benched_on_first_failure(
@ -204,6 +217,10 @@ def _nested_proxy_params(upstream_group: str, upstream_key: str, cooldown_time:
)
@step(
"Add a deployment named {name} that fronts {upstream_group} on this proxy, so it always gets a 500"
" and is benched on the first one"
)
def create_always_5xx_deployment(
proxy: ProxyClient, name: str, upstream_group: str, upstream_key: str, cooldown_time: float | None = None
) -> str:
@ -217,6 +234,10 @@ def create_always_5xx_deployment(
)
@step(
"Add a deployment named {name} that fronts {upstream_group} on this proxy with a key out of rpm,"
" so it always gets a 429 and is benched on the first one"
)
def create_always_rate_limited_deployment(
proxy: ProxyClient, name: str, upstream_group: str, upstream_key: str, cooldown_time: float | None = None
) -> str:
@ -227,6 +248,7 @@ def create_always_rate_limited_deployment(
)
@step(f"Use up the rpm-limited key's one allowed request with a /chat/completions call to {CHEAP_OPENAI_MODEL}")
def spend_only_request_of(proxy: ProxyClient, spent_key: str) -> None:
"""Uses up the one request an rpm_limit=1 key allows. The proxy's rate limiter
opens the key's 60s window on this call, so it goes right before the calls that
@ -239,6 +261,7 @@ def spend_only_request_of(proxy: ProxyClient, spent_key: str) -> None:
)
@step(f"Add a deployment named {{name}} that calls {SMALL_CONTEXT_MODEL} and takes all of its group's traffic")
def create_always_picked_small_context_deployment(proxy: ProxyClient, name: str) -> str:
"""The always-picked half of a retry pair on the smallest-context model OpenAI
still serves: it holds all of the model group's shuffle weight, so an oversized
@ -253,12 +276,14 @@ def create_always_picked_small_context_deployment(proxy: ProxyClient, name: str)
)
@step(f"Add a deployment named {{name}} for {REAL_MODEL} that answers with a canned reply")
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"))
@step(f"Add a zero-weight backup deployment named {{name}} that calls {REAL_MODEL}")
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,
@ -273,6 +298,7 @@ def create_zero_weight_backup_deployment(proxy: ProxyClient, name: str) -> str:
)
@step("Send a /chat/completions request to {model} with a full message history and stream set to {stream}")
def chat_turns_override(
proxy: ProxyClient,
key: str,
@ -300,6 +326,7 @@ def chat_turns_override(
)
@step("Send a /chat/completions request to {model} with stream set to {stream}")
def chat_override(
proxy: ProxyClient,
key: str,
@ -322,6 +349,7 @@ def chat_override(
)
@step('Open a streaming /chat/completions request to {model} with the prompt "{content}" and leave it in flight')
def open_chat_stream(
proxy: ProxyClient,
key: str,

View file

@ -49,6 +49,7 @@ import pytest
from pydantic import BaseModel, ConfigDict, Field
from e2e_config import unique_marker
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from e2e_http import AnthropicHeaders, AuthHeaders, UnauthorizedError, unwrap
from lifecycle import ResourceManager
from models import (
@ -341,6 +342,15 @@ def credentialed_alias(proxy: ProxyClient, router_stack: ExitStack) -> Credentia
class TestTagSplitRouting:
@pytest.mark.covers("reliability.routing.tagged_marker.request_tag_selects_marker")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_body_tagged_chat_routes_through_the_marker_to_its_tier(
self, proxy: ProxyClient, resources: ResourceManager, plain_first_split: TagSplitDeployment
) -> None:
@ -355,6 +365,15 @@ class TestTagSplitRouting:
_assert_served_only_by(rows, CHEAP_SERVED | {plain_first_split.tier}, "body-tagged chat on the shared name")
@pytest.mark.covers("reliability.routing.tagged_marker.untagged_request_served_by_plain_deployment")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(PLAIN_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_untagged_chat_is_always_served_by_the_plain_deployment(
self, proxy: ProxyClient, resources: ResourceManager, plain_first_split: TagSplitDeployment
) -> None:
@ -370,6 +389,15 @@ class TestTagSplitRouting:
_assert_served_only_by(rows, PLAIN_SERVED | {plain_first_split.shared}, "untagged chat on the shared name")
@pytest.mark.covers("reliability.routing.tagged_marker.untagged_request_served_by_plain_deployment")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(PLAIN_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_untagged_messages_is_served_by_the_plain_deployment(
self, proxy: ProxyClient, resources: ResourceManager, plain_first_split: TagSplitDeployment
) -> None:
@ -387,6 +415,15 @@ class TestTagSplitRouting:
class TestUntaggedTierDeployments:
@pytest.mark.covers("reliability.routing.tagged_marker.header_tag_selects_marker")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_header_tagged_messages_routes_through_the_marker_to_an_untagged_tier(
self, proxy: ProxyClient, resources: ResourceManager, marker_first_split: TagSplitDeployment
) -> None:
@ -413,6 +450,15 @@ class TestUntaggedTierDeployments:
)
@pytest.mark.covers("reliability.routing.tagged_marker.untagged_tier_deployments_still_served")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_body_tagged_chat_reaches_the_untagged_tier_after_marker_rewrite(
self, proxy: ProxyClient, resources: ResourceManager, marker_first_split: TagSplitDeployment
) -> None:
@ -431,6 +477,12 @@ class TestUntaggedTierDeployments:
_assert_served_only_by(rows, CHEAP_SERVED | {marker_first_split.tier}, "body-tagged chat with untagged tier")
@pytest.mark.covers("reliability.routing.tagged_marker.tag_semantics_stay_strict")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.CHAT_COMPLETIONS,
)
)
def test_tagged_call_straight_at_an_untagged_deployment_stays_denied(
self, proxy: ProxyClient, resources: ResourceManager, marker_first_split: TagSplitDeployment
) -> None:
@ -449,6 +501,15 @@ class TestUntaggedTierDeployments:
class TestResponsesApiTagRouting:
@pytest.mark.covers("reliability.routing.tagged_marker.responses_input_routes_through_marker")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.RESPONSES,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_header_tagged_responses_with_string_input_routes_to_the_tier(
self, proxy: ProxyClient, resources: ResourceManager, plain_first_split: TagSplitDeployment
) -> None:
@ -471,6 +532,15 @@ class TestResponsesApiTagRouting:
)
@pytest.mark.covers("reliability.routing.tagged_marker.responses_input_routes_through_marker")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.RESPONSES,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_body_tagged_responses_with_list_input_routes_to_the_tier(
self, proxy: ProxyClient, resources: ResourceManager, plain_first_split: TagSplitDeployment
) -> None:
@ -497,6 +567,15 @@ class TestResponsesApiTagRouting:
_assert_served_only_by(rows, CHEAP_SERVED | {plain_first_split.tier}, "body-tagged /v1/responses list input")
@pytest.mark.covers("reliability.routing.tagged_marker.untagged_request_served_by_plain_deployment")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.RESPONSES,
providers=(Provider.ANTHROPIC,),
models=(PLAIN_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_untagged_responses_is_served_by_the_plain_deployment(
self, proxy: ProxyClient, resources: ResourceManager, plain_first_split: TagSplitDeployment
) -> None:
@ -524,6 +603,14 @@ class TestResponsesApiTagRouting:
class TestStrategyAliasPricing:
@pytest.mark.covers("reliability.routing.strategy_alias.custom_pricing_ignored")
@meta(
Subject(
domain=Domain.SPEND_BUDGETS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_zero_priced_alias_still_logs_spend_at_the_tier_rate(
self, proxy: ProxyClient, resources: ResourceManager, zero_priced_alias: ZeroPricedAlias
) -> None:
@ -546,6 +633,14 @@ class TestStrategyAliasPricing:
class TestComplexityHeuristicScope:
@pytest.mark.covers("reliability.routing.complexity_heuristic.scores_current_ask_only")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_trivial_ask_behind_keyword_heavy_system_prompt_stays_on_the_cheap_tier(
self, proxy: ProxyClient, resources: ResourceManager, heuristic_split: HeuristicSplit
) -> None:
@ -573,6 +668,15 @@ class TestComplexityHeuristicScope:
class TestSemanticAutoRouterResponses:
@pytest.mark.covers("reliability.routing.semantic_auto_router.responses_input_routed")
@meta(
Subject(
domain=Domain.ROUTING,
route=Route.RESPONSES,
providers=(Provider.ANTHROPIC, Provider.OPENAI,),
models=(CHEAP_MODEL, EMBEDDING_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_responses_input_reaches_the_semantic_auto_router(
self, proxy: ProxyClient, resources: ResourceManager, semantic_auto_router: SemanticAutoRouter
) -> None:
@ -617,6 +721,14 @@ class TestSemanticAutoRouterResponses:
class TestAliasParamForwarding:
@pytest.mark.covers("reliability.routing.tagged_marker.alias_connection_params_stay_with_tier")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_alias_api_key_never_overrides_the_tier_credential(
self, proxy: ProxyClient, resources: ResourceManager, credentialed_alias: CredentialedAlias
) -> None:

View file

@ -19,9 +19,12 @@ anthropic proves the classifier ran and openai proves it silently fell back - th
exact failure before the fix.
"""
from typing import Final
import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from e2e_http import unwrap
from models import ChatBody, ChatMessage
@ -33,9 +36,11 @@ LEXICALLY_SIMPLE_HARD_PROMPT = "Should I pay off my mortgage early or invest the
# SIMPLE tier backend; served only when the classifier silently falls back to heuristic.
# Spend logs may store the alias (gpt-5.5) or the provider-prefixed form depending on
# how the deployment is registered (compose vs /model/new).
HEURISTIC_TIER_MODELS = frozenset({"openai/gpt-5.5", "gpt-5.5"})
HEURISTIC_TIER_BACKEND: Final = "openai/gpt-5.5"
HEURISTIC_TIER_MODELS = frozenset({HEURISTIC_TIER_BACKEND, "gpt-5.5"})
# MEDIUM/COMPLEX/REASONING tier backend; served only when the LLM classifier runs.
LLM_TIER_MODELS = frozenset({"anthropic/claude-haiku-4-5", "claude-haiku-4-5"})
LLM_TIER_BACKEND: Final = "anthropic/claude-haiku-4-5"
LLM_TIER_MODELS = frozenset({LLM_TIER_BACKEND, "claude-haiku-4-5"})
@pytest.mark.usefixtures("_ensure_complexity_smart_router")
@ -45,6 +50,14 @@ class TestComplexityRouterLlmClassifier:
"(e.g. Is P equal to NP?); re-enable when classifier tier quality is fixed"
)
@pytest.mark.covers("reliability.routing.complexity_llm_classifier.routes_by_llm_tier")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI, Provider.ANTHROPIC),
models=(HEURISTIC_TIER_BACKEND, LLM_TIER_BACKEND),
mode=Mode.NONSTREAM,
)
)
def test_llm_classifier_runs_and_routes_by_semantic_tier(
self, client: ComplexityRouterClient, complexity_key: str
) -> None:

View file

@ -18,6 +18,7 @@ from e2e_config import (
REQUEST_TIMEOUT,
unique_marker,
)
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import ChatBody, ChatMessage, ChatResponse, LiteLLMParamsBody
from provider_edge import ProviderRequestObservation, observed_provider_edge
@ -25,6 +26,8 @@ from pydantic import BaseModel, JsonValue
pytestmark = [pytest.mark.e2e, pytest.mark.replayable]
CACHE_MODEL: Final = "openai/gpt-5.6"
class _CacheChatBody(ChatBody):
ttl: int = 600
@ -38,6 +41,14 @@ class _CachedAnswer(BaseModel):
class TestReliabilityCache:
@pytest.mark.covers("reliability.cache.exact.returns_cached")
@meta(
Subject(
domain=Domain.CACHING,
providers=(Provider.OPENAI,),
models=(CACHE_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_exact_cache_returns_cached(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -57,7 +68,7 @@ class TestReliabilityCache:
model_id: Final = client.proxy.create_model(
model,
LiteLLMParamsBody(
model="openai/gpt-5.6",
model=CACHE_MODEL,
api_key="os.environ/OPENAI_API_KEY",
api_base=f"{edge.api_base('openai')}/v1",
),

View file

@ -29,9 +29,11 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import unique_marker
from e2e_http import AbandonedRequest, StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import ChatMessage, ReliabilityChatBody, RouterSettingsOverride
from reliability_support import (
AZURE_MODEL,
REPLICA_PROPAGATION_SECONDS,
chat_override,
create_azure_benched_on_first_failure_deployment,
@ -102,6 +104,14 @@ def _hang_up_mid_answer(client: ComplexityRouterClient, key: str, group: str) ->
class TestReliabilityCancelOnDisconnect:
@pytest.mark.covers("reliability.cooldown.client_disconnect.stays_healthy")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.AZURE,),
models=(AZURE_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_client_hanging_up_never_benches_the_deployment(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -71,10 +71,12 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import KeyGenerateBody, RouterSettingsOverride
from reliability_support import (
COOLDOWN_SECONDS,
REAL_MODEL,
REPLICA_PROPAGATION_SECONDS,
chat_override,
create_always_5xx_deployment,
@ -212,6 +214,14 @@ def _assert_trips_then_recovers(
class TestReliabilityCooldowns:
@pytest.mark.covers("reliability.cooldown.5xx.trips_then_recovers")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_5xx_trips_cooldown_then_recovers(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -230,6 +240,14 @@ 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")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_sibling_replica_serves_backup_within_redis_read_interval(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -272,6 +290,14 @@ class TestReliabilityCooldowns:
)
@pytest.mark.covers("reliability.cooldown.429.trips_then_recovers")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL, REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_429_trips_cooldown_then_recovers(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -292,6 +318,14 @@ class TestReliabilityCooldowns:
_assert_trips_then_recovers(client, scoped_key, group, failing, backup, failure_status=429)
@pytest.mark.covers("reliability.cooldown.auth.trips_then_recovers")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_auth_failure_trips_cooldown_then_recovers(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -304,6 +338,14 @@ class TestReliabilityCooldowns:
_assert_trips_then_recovers(client, scoped_key, group, failing, backup, failure_status=401)
@pytest.mark.covers("reliability.cooldown.timeout.trips_then_recovers")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_timeout_trips_cooldown_then_recovers(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -31,10 +31,14 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import unique_marker
from e2e_http import StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import RouterSettingsOverride
from reliability_support import (
AZURE_MODEL,
CONTENT_POLICY_PROMPT,
REAL_MODEL,
SMALL_CONTEXT_MODEL,
azure_prompt_filter_skipped,
chat_override,
completion_tokens_of,
@ -50,6 +54,8 @@ from reliability_support import (
pytestmark = pytest.mark.e2e
FALLBACK_MODEL: Final = "gpt-5.5"
def _assert_served_by_fallback(resp: StreamingResponse) -> None:
assert resp.status_code == 200, f"expected 200 after fallback, got {resp.status_code}: {resp.body[:300]}"
@ -95,6 +101,14 @@ def _filter_verdict(resp: StreamingResponse) -> str:
class TestReliabilityFallbacks:
@pytest.mark.covers("reliability.fallback.5xx.routes_to_fallback")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(FALLBACK_MODEL, REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_5xx_routes_to_fallback(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -107,11 +121,19 @@ class TestReliabilityFallbacks:
scoped_key,
primary,
f"say hi {unique_marker()}",
override=RouterSettingsOverride(fallbacks=[{primary: ["gpt-5.5"]}]),
override=RouterSettingsOverride(fallbacks=[{primary: [FALLBACK_MODEL]}]),
)
_assert_served_by_fallback(resp)
@pytest.mark.covers("reliability.fallback.timeout.routes_to_fallback")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(FALLBACK_MODEL, REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_timeout_routes_to_fallback(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -124,11 +146,19 @@ class TestReliabilityFallbacks:
scoped_key,
primary,
f"say hi {unique_marker()}",
override=RouterSettingsOverride(fallbacks=[{primary: ["gpt-5.5"]}]),
override=RouterSettingsOverride(fallbacks=[{primary: [FALLBACK_MODEL]}]),
)
_assert_served_by_fallback(resp)
@pytest.mark.covers("reliability.fallback.context_window.routes_to_fallback")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(FALLBACK_MODEL, SMALL_CONTEXT_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_context_window_routes_to_fallback(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -141,11 +171,19 @@ class TestReliabilityFallbacks:
scoped_key,
primary,
oversized_prompt(unique_marker()),
override=RouterSettingsOverride(context_window_fallbacks=[{primary: ["gpt-5.5"]}]),
override=RouterSettingsOverride(context_window_fallbacks=[{primary: [FALLBACK_MODEL]}]),
)
_assert_served_by_fallback(resp)
@pytest.mark.covers("reliability.fallback.content_policy.routes_to_fallback")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.AZURE, Provider.OPENAI,),
models=(AZURE_MODEL, FALLBACK_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_content_policy_routes_to_fallback(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -167,7 +205,7 @@ class TestReliabilityFallbacks:
scoped_key,
primary,
f"{CONTENT_POLICY_PROMPT} {unique_marker()}",
override=RouterSettingsOverride(content_policy_fallbacks=[{primary: ["gpt-5.5"]}]),
override=RouterSettingsOverride(content_policy_fallbacks=[{primary: [FALLBACK_MODEL]}]),
)
)
_assert_served_by_fallback(resp)

View file

@ -80,11 +80,12 @@ from e2e_config import (
PROXY_REPLICA_URLS,
unique_marker,
)
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from memory_readings import RssCapture, RssReading, WorkerKey, read_rss_everywhere
from models import ChatMessage, RouterSettingsOverride, SpendLogRow
from proxy_client import ProxyClient
from reliability_support import chat_override, create_never_benched_refusing_deployment
from reliability_support import REAL_MODEL, chat_override, create_never_benched_refusing_deployment
pytestmark = [pytest.mark.e2e, pytest.mark.quiet_stack]
@ -222,6 +223,11 @@ def _stored_request_kb(proxy: ProxyClient, call: FailedCall) -> float:
class TestReliabilityMemory:
@pytest.mark.covers("reliability.perf.idle_memory.under_slo")
@meta(
Subject(
domain=Domain.DEPLOY_OPS,
)
)
def test_workers_idle_under_rss_budget_before_traffic(self, idle_rss: RssCapture) -> None:
assert not idle_rss.failures, (
f"{len(idle_rss.failures)} replica(s) gave no RSS reading when the session started, so their idle "
@ -242,6 +248,14 @@ class TestReliabilityMemory:
)
@pytest.mark.covers("reliability.perf.memory.under_slo")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_failing_requests_do_not_grow_rss_or_stored_request(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -23,9 +23,11 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import unique_marker
from e2e_metadata import Capability, Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import ChatMessage, LiteLLMParamsBody, ModelInfoBody, ModelNewBody
from reliability_support import (
CACHING_MODEL,
REAL_KEY,
REAL_MODEL,
cached_system_turn,
@ -42,6 +44,15 @@ FOLLOW_UPS = 3
class TestReliabilityPromptCachingAffinity:
@pytest.mark.covers("reliability.cache.prompt_caching_model_select.returns_cached")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.ANTHROPIC, Provider.OPENAI),
models=(CACHING_MODEL, REAL_MODEL),
capabilities=(Capability.PROMPT_CACHING,),
mode=Mode.NONSTREAM,
)
)
def test_cached_conversation_stays_on_deployment_holding_its_cache(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -30,9 +30,12 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import KeyGenerateBody, RouterSettingsOverride
from reliability_support import (
REAL_MODEL,
SMALL_CONTEXT_MODEL,
chat_override,
completion_tokens_of,
content_of,
@ -84,6 +87,14 @@ def _retry_once(client: ComplexityRouterClient, key: str, group: str) -> Streami
class TestReliabilityRetries:
@pytest.mark.covers("reliability.retry.timeout.succeeds_within_retries")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_timeout_on_first_deployment_succeeds_on_retry(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -96,6 +107,14 @@ class TestReliabilityRetries:
_assert_served_after_retry(_retry_once(client, scoped_key, group))
@pytest.mark.covers("reliability.retry.5xx.succeeds_within_retries")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_5xx_on_first_deployment_succeeds_on_retry(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -112,6 +131,14 @@ class TestReliabilityRetries:
_assert_served_after_retry(_retry_once(client, scoped_key, group))
@pytest.mark.covers("reliability.retry.429.succeeds_within_retries")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL, REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_429_on_first_deployment_succeeds_on_retry(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -130,6 +157,14 @@ class TestReliabilityRetries:
_assert_served_after_retry(_retry_once(client, scoped_key, group))
@pytest.mark.covers("reliability.retry.auth.succeeds_within_retries")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_auth_failure_on_first_deployment_succeeds_on_retry(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -142,6 +177,14 @@ class TestReliabilityRetries:
_assert_served_after_retry(_retry_once(client, scoped_key, group))
@pytest.mark.covers("reliability.retry.context_window.succeeds_within_retries")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL, SMALL_CONTEXT_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_context_window_refusal_on_first_deployment_succeeds_on_retry(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -63,6 +63,7 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import unique_marker
from e2e_http import StreamChunk, StreamHead, StreamStep, StreamTruncation
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from models import LiteLLMParamsBody, ModelInfoBody, ModelNewBody, RouterSettingsOverride, RoutingStrategy
from reliability_support import REAL_KEY, REAL_MODEL, chat_override, model_id_of, open_chat_stream
@ -173,6 +174,14 @@ def _assert_shuffle_control_lands_on(client: ComplexityRouterClient, key: str, g
class TestReliabilityRoutingStrategies:
@pytest.mark.covers("reliability.routing.simple_shuffle.picks_healthy_deployment")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_simple_shuffle_honors_weights(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -191,6 +200,14 @@ class TestReliabilityRoutingStrategies:
)
@pytest.mark.covers("reliability.routing.cost_based.picks_lowest_cost")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_cost_based_picks_cheapest_deployment(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -206,6 +223,14 @@ class TestReliabilityRoutingStrategies:
_assert_shuffle_control_lands_on(client, scoped_key, group, pricey)
@pytest.mark.covers("reliability.routing.usage_based.picks_under_tpm")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_usage_based_picks_deployment_with_tpm_headroom(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -223,6 +248,14 @@ class TestReliabilityRoutingStrategies:
"so latency-based has no signal to route on"
)
@pytest.mark.covers("reliability.routing.latency_based.picks_lowest_latency")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_latency_based_routes_around_deployment_that_times_out(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -253,6 +286,13 @@ class TestReliabilityRoutingStrategies:
"so least-busy has no signal to route on"
)
@pytest.mark.covers("reliability.routing.least_busy.picks_lowest_traffic")
@meta(
Subject(
domain=Domain.ROUTING,
providers=(Provider.OPENAI,),
models=(REAL_MODEL,),
)
)
def test_least_busy_avoids_deployment_with_request_in_flight(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -13,14 +13,16 @@ import pytest
from complexity_router_client import ComplexityRouterClient
from e2e_config import unique_marker
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from reliability_support import chat_override, create_timeout_deployment
from reliability_support import REAL_MODEL, chat_override, create_timeout_deployment
pytestmark = pytest.mark.e2e
class TestReliabilityTimeouts:
@pytest.mark.covers("reliability.timeout.request_timeout.exceeds_deadline")
@meta(Subject(domain=Domain.ROUTING, providers=(Provider.OPENAI,), models=(REAL_MODEL,), mode=Mode.NONSTREAM))
def test_request_timeout_exceeds_deadline(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -35,6 +37,7 @@ class TestReliabilityTimeouts:
assert "timeout" in resp.body.lower(), f"the 408 body should name the timeout, got: {resp.body[:300]}"
@pytest.mark.covers("reliability.timeout.stream_timeout.exceeds_deadline")
@meta(Subject(domain=Domain.ROUTING, providers=(Provider.OPENAI,), models=(REAL_MODEL,), mode=Mode.STREAM))
def test_stream_timeout_exceeds_deadline(
self, client: ComplexityRouterClient, resources: ResourceManager, scoped_key: str
) -> None: