diff --git a/tests/e2e/batches/batch_cleanup.py b/tests/e2e/batches/batch_cleanup.py index f1142a60782..f8fc1bbf01f 100644 --- a/tests/e2e/batches/batch_cleanup.py +++ b/tests/e2e/batches/batch_cleanup.py @@ -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) diff --git a/tests/e2e/batches/batch_client.py b/tests/e2e/batches/batch_client.py index 8745140a818..b02b09e557b 100644 --- a/tests/e2e/batches/batch_client.py +++ b/tests/e2e/batches/batch_client.py @@ -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}", diff --git a/tests/e2e/batches/capabilities.py b/tests/e2e/batches/capabilities.py index d510426dee2..c03b481f060 100644 --- a/tests/e2e/batches/capabilities.py +++ b/tests/e2e/batches/capabilities.py @@ -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") diff --git a/tests/e2e/batches/test_batches_e2e.py b/tests/e2e/batches/test_batches_e2e.py index 8da2deb4010..54bcb80ed24 100644 --- a/tests/e2e/batches/test_batches_e2e.py +++ b/tests/e2e/batches/test_batches_e2e.py @@ -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, diff --git a/tests/e2e/batches/test_managed_files_enforcement_e2e.py b/tests/e2e/batches/test_managed_files_enforcement_e2e.py index 2f5d0588aca..43a2488289c 100644 --- a/tests/e2e/batches/test_managed_files_enforcement_e2e.py +++ b/tests/e2e/batches/test_managed_files_enforcement_e2e.py @@ -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: diff --git a/tests/e2e/mcp/conftest.py b/tests/e2e/mcp/conftest.py index e6094ab95ea..91f55a528b5 100644 --- a/tests/e2e/mcp/conftest.py +++ b/tests/e2e/mcp/conftest.py @@ -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: ... diff --git a/tests/e2e/mcp/datadog_mcp.py b/tests/e2e/mcp/datadog_mcp.py index 352b4446cfd..63194af221e 100644 --- a/tests/e2e/mcp/datadog_mcp.py +++ b/tests/e2e/mcp/datadog_mcp.py @@ -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, diff --git a/tests/e2e/mcp/mcp_client.py b/tests/e2e/mcp/mcp_client.py index 56f7fffba29..e373a8e31ae 100644 --- a/tests/e2e/mcp/mcp_client.py +++ b/tests/e2e/mcp/mcp_client.py @@ -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, diff --git a/tests/e2e/mcp/oauth_chat_client.py b/tests/e2e/mcp/oauth_chat_client.py index 0c5c6106259..68e838d8745 100644 --- a/tests/e2e/mcp/oauth_chat_client.py +++ b/tests/e2e/mcp/oauth_chat_client.py @@ -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 diff --git a/tests/e2e/mcp/oauth_gateway.py b/tests/e2e/mcp/oauth_gateway.py index 029b0135900..2abe8262b6e 100644 --- a/tests/e2e/mcp/oauth_gateway.py +++ b/tests/e2e/mcp/oauth_gateway.py @@ -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" diff --git a/tests/e2e/mcp/test_mcp_access_group_e2e.py b/tests/e2e/mcp/test_mcp_access_group_e2e.py index f72b75fd43d..744acb78aea 100644 --- a/tests/e2e/mcp/test_mcp_access_group_e2e.py +++ b/tests/e2e/mcp/test_mcp_access_group_e2e.py @@ -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: diff --git a/tests/e2e/mcp/test_mcp_chat_completion_oauth_e2e.py b/tests/e2e/mcp/test_mcp_chat_completion_oauth_e2e.py index 01e94f7b86f..9d7e4713963 100644 --- a/tests/e2e/mcp/test_mcp_chat_completion_oauth_e2e.py +++ b/tests/e2e/mcp/test_mcp_chat_completion_oauth_e2e.py @@ -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: diff --git a/tests/e2e/mcp/test_mcp_datadog_e2e.py b/tests/e2e/mcp/test_mcp_datadog_e2e.py index 031fbf6d936..e383e9ab766 100644 --- a/tests/e2e/mcp/test_mcp_datadog_e2e.py +++ b/tests/e2e/mcp/test_mcp_datadog_e2e.py @@ -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, diff --git a/tests/e2e/mcp/test_mcp_guardrail_e2e.py b/tests/e2e/mcp/test_mcp_guardrail_e2e.py index 92c632cb316..d34bd80e564 100644 --- a/tests/e2e/mcp/test_mcp_guardrail_e2e.py +++ b/tests/e2e/mcp/test_mcp_guardrail_e2e.py @@ -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: diff --git a/tests/e2e/mcp/test_mcp_key_access_e2e.py b/tests/e2e/mcp/test_mcp_key_access_e2e.py index c00d67bc9cf..9f5e04ccb9b 100644 --- a/tests/e2e/mcp/test_mcp_key_access_e2e.py +++ b/tests/e2e/mcp/test_mcp_key_access_e2e.py @@ -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, diff --git a/tests/e2e/mcp/test_mcp_oauth_happy_path_e2e.py b/tests/e2e/mcp/test_mcp_oauth_happy_path_e2e.py index 81470c21d51..b57c941cb18 100644 --- a/tests/e2e/mcp/test_mcp_oauth_happy_path_e2e.py +++ b/tests/e2e/mcp/test_mcp_oauth_happy_path_e2e.py @@ -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( diff --git a/tests/e2e/mcp/test_mcp_toolset_enforcement_e2e.py b/tests/e2e/mcp/test_mcp_toolset_enforcement_e2e.py index 6b901145eb1..cff72f3c95d 100644 --- a/tests/e2e/mcp/test_mcp_toolset_enforcement_e2e.py +++ b/tests/e2e/mcp/test_mcp_toolset_enforcement_e2e.py @@ -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) diff --git a/tests/e2e/router/reliability_support.py b/tests/e2e/router/reliability_support.py index 04fa30a0d14..03e585273c0 100644 --- a/tests/e2e/router/reliability_support.py +++ b/tests/e2e/router/reliability_support.py @@ -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, diff --git a/tests/e2e/router/test_auto_router_regressions_e2e.py b/tests/e2e/router/test_auto_router_regressions_e2e.py index 374badcf5fc..898b164e713 100644 --- a/tests/e2e/router/test_auto_router_regressions_e2e.py +++ b/tests/e2e/router/test_auto_router_regressions_e2e.py @@ -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: diff --git a/tests/e2e/router/test_complexity_router_e2e.py b/tests/e2e/router/test_complexity_router_e2e.py index e8508c963b8..c8f20102223 100644 --- a/tests/e2e/router/test_complexity_router_e2e.py +++ b/tests/e2e/router/test_complexity_router_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_cache_e2e.py b/tests/e2e/router/test_reliability_cache_e2e.py index f7a2f2ffeb7..5452853a57e 100644 --- a/tests/e2e/router/test_reliability_cache_e2e.py +++ b/tests/e2e/router/test_reliability_cache_e2e.py @@ -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", ), diff --git a/tests/e2e/router/test_reliability_cancel_on_disconnect_e2e.py b/tests/e2e/router/test_reliability_cancel_on_disconnect_e2e.py index 06174e97d20..e39f1ddcb90 100644 --- a/tests/e2e/router/test_reliability_cancel_on_disconnect_e2e.py +++ b/tests/e2e/router/test_reliability_cancel_on_disconnect_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_cooldowns_e2e.py b/tests/e2e/router/test_reliability_cooldowns_e2e.py index 2456bfb5f85..bdb02256976 100644 --- a/tests/e2e/router/test_reliability_cooldowns_e2e.py +++ b/tests/e2e/router/test_reliability_cooldowns_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_fallbacks_e2e.py b/tests/e2e/router/test_reliability_fallbacks_e2e.py index 54c11b163d0..61115e00d07 100644 --- a/tests/e2e/router/test_reliability_fallbacks_e2e.py +++ b/tests/e2e/router/test_reliability_fallbacks_e2e.py @@ -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) diff --git a/tests/e2e/router/test_reliability_memory_e2e.py b/tests/e2e/router/test_reliability_memory_e2e.py index 77d3a68cae5..9ae6edb3eb2 100644 --- a/tests/e2e/router/test_reliability_memory_e2e.py +++ b/tests/e2e/router/test_reliability_memory_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_prompt_caching_e2e.py b/tests/e2e/router/test_reliability_prompt_caching_e2e.py index 667b398cd16..7bf23bf967f 100644 --- a/tests/e2e/router/test_reliability_prompt_caching_e2e.py +++ b/tests/e2e/router/test_reliability_prompt_caching_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_retries_e2e.py b/tests/e2e/router/test_reliability_retries_e2e.py index a90efa52b5f..c6ae8c457e4 100644 --- a/tests/e2e/router/test_reliability_retries_e2e.py +++ b/tests/e2e/router/test_reliability_retries_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_routing_strategies_e2e.py b/tests/e2e/router/test_reliability_routing_strategies_e2e.py index 2abc2ee5f54..ee0334f989d 100644 --- a/tests/e2e/router/test_reliability_routing_strategies_e2e.py +++ b/tests/e2e/router/test_reliability_routing_strategies_e2e.py @@ -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: diff --git a/tests/e2e/router/test_reliability_timeouts_e2e.py b/tests/e2e/router/test_reliability_timeouts_e2e.py index f24d5139e66..926d8539d6d 100644 --- a/tests/e2e/router/test_reliability_timeouts_e2e.py +++ b/tests/e2e/router/test_reliability_timeouts_e2e.py @@ -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: