test(e2e): tag guardrails and logging tests with Subject metadata and record client steps (#44963)

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

* test(e2e): tag guardrails and logging tests with Subject metadata and record client steps

* test(e2e): leave the guardrails and logging harness unit tests untagged

* test(e2e): let the inner create_model step name the guardrail backend deployment

* 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): declare the default guardrail backend model on the tests that drive it
This commit is contained in:
ryan-crabbe-berri 2026-10-07 10:32:25 -07:00 • committed by GitHub
parent d7cdc88c66
commit 79209b92a1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
31 changed files with 821 additions and 30 deletions

View file

@ -11,6 +11,7 @@ from typing import Final, Literal
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT, SLOW_PROVIDER_TIMEOUT_SECONDS, settle_propagation, unique_marker
from e2e_http import NoBody, Result, StreamingResponse, Success, unwrap
from e2e_metadata import step
from lifecycle import ResourceManager
from models import (
AnthropicMessagesBody,
@ -35,7 +36,9 @@ from models import (
VideoCreateResponse,
)
from proxy_client import ProxyClient
from pydantic import BaseModel
from pydantic import BaseModel, Field
GUARDRAIL_BACKEND: Final = "gemini/gemini-2.5-flash"
GuardrailMode = Literal["pre_call", "post_call", "during_call", "logging_only"]
PiiEntity = Literal["EMAIL_ADDRESS", "PHONE_NUMBER", "PERSON", "CREDIT_CARD", "US_SSN"]
@ -62,14 +65,14 @@ class BedrockGuardrailParamsBody(GuardrailParamsBase):
guardrail: Literal["bedrock"] = "bedrock"
guardrailIdentifier: str
guardrailVersion: str
aws_access_key_id: str | None = None
aws_secret_access_key: str | None = None
aws_access_key_id: str | None = Field(default=None, repr=False)
aws_secret_access_key: str | None = Field(default=None, repr=False)
aws_region_name: str | None = None
class OpenAIModerationParamsBody(GuardrailParamsBase):
guardrail: Literal["openai_moderation"] = "openai_moderation"
api_key: str | None = None
api_key: str | None = Field(default=None, repr=False)
model: str | None = None
@ -193,6 +196,7 @@ class _ResponsesGuardrailBody(BaseModel):
class GuardrailsClient:
proxy: ProxyClient
@step("Register the content filter guardrail {name} that blocks prompts containing {blocked_keyword}")
def create_content_filter_guardrail(self, name: str, blocked_keyword: str, *, default_on: bool = True) -> str:
return self.register(
name,
@ -203,6 +207,7 @@ class GuardrailsClient:
),
)
@step("Register the Bedrock guardrail {name}")
def create_bedrock_guardrail(
self,
name: str,
@ -235,7 +240,7 @@ class GuardrailsClient:
resources: ResourceManager,
prefix: str = "e2e-guard-backend",
*,
backend: str = "gemini/gemini-2.5-flash",
backend: str = GUARDRAIL_BACKEND,
api_key: str = "os.environ/GEMINI_API_KEY",
) -> str:
"""Register a chat deployment for a guardrail test to run against
@ -250,6 +255,7 @@ class GuardrailsClient:
resources.defer(lambda: self.proxy.delete_model(model_id))
return model_name
@step("Register the {params.guardrail} guardrail {name} with mode {params.mode}")
def register(self, name: str, params: GuardrailParamsBody) -> str:
"""Register any guardrail via POST /guardrails and return its id, once every
replica can be expected to serve it. New built-ins register with
@ -274,6 +280,7 @@ class GuardrailsClient:
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}",
@ -282,6 +289,7 @@ class GuardrailsClient:
response_type=NoBody,
)
@step("Create the guardrail policy {body.policy_name} that adds the guardrails {body.guardrails_add}")
def create_policy(self, body: PolicyCreateBody) -> str:
"""Create a policy via POST /policies and return its name once every replica
can be expected to serve it (policies reach the data plane on the periodic
@ -297,6 +305,7 @@ class GuardrailsClient:
settle_propagation(time.monotonic())
return created.policy_name
@step("Delete every version of the guardrail policy {policy_name}")
def delete_policy(self, policy_name: str) -> None:
_ = self.proxy.transport.delete(
f"/policies/name/{policy_name}/all-versions",
@ -305,6 +314,7 @@ class GuardrailsClient:
response_type=NoBody,
)
@step("Attach the guardrail policy {policy_name} to requests tagged {tags}")
def attach_policy_to_tags(self, policy_name: str, tags: list[str]) -> str:
attachment_id = unwrap(
self.proxy.transport.post(
@ -317,6 +327,7 @@ class GuardrailsClient:
settle_propagation(time.monotonic())
return attachment_id
@step("Delete the guardrail policy attachment")
def delete_policy_attachment(self, attachment_id: str) -> None:
_ = self.proxy.transport.delete(
f"/policies/attachments/{attachment_id}",
@ -325,6 +336,7 @@ class GuardrailsClient:
response_type=NoBody,
)
@step("Create the team {alias} opted out of global guardrails and wait until /team/info returns it")
def create_team_opted_out_of_global_guardrails(self, alias: str) -> str:
team_id = unwrap(
self.proxy.transport.post(
@ -340,6 +352,7 @@ class GuardrailsClient:
self._await_team(team_id)
return team_id
@step("Delete the team")
def delete_team(self, team_id: str) -> None:
_ = self.proxy.transport.post(
"/team/delete",
@ -348,9 +361,11 @@ class GuardrailsClient:
response_type=NoBody,
)
@step("Generate a virtual key in the team")
def create_key_in_team(self, team_id: str) -> str:
return self.proxy.generate_key(KeyGenerateBody(team_id=team_id, user_id="e2e-guardrails-user"))
@step("Generate a virtual key with the guardrails {guardrails}")
def create_key_with_guardrails(self, resources: ResourceManager, guardrails: list[str]) -> str:
key = self.proxy.generate_key(
KeyGenerateBody(user_id="e2e-guardrails-user", metadata=KeyMetadata(guardrails=guardrails))
@ -358,6 +373,7 @@ class GuardrailsClient:
resources.defer(lambda: self.proxy.delete_key(key))
return key
@step("Send a /v1/videos request to {model}")
def create_video(self, key: str, model: str, prompt: str) -> Result[VideoCreateResponse]:
return self.proxy.transport.post(
"/v1/videos",
@ -366,6 +382,7 @@ class GuardrailsClient:
response_type=VideoCreateResponse,
)
@step("Send a /v1/images/edits request to {model}")
def edit_image(self, key: str, model: str, prompt: str, image: bytes) -> Result[ImageGenerationResponse]:
return self.proxy.transport.upload(
"/v1/images/edits",
@ -379,6 +396,7 @@ class GuardrailsClient:
timeout=SLOW_PROVIDER_TIMEOUT_SECONDS,
)
@step("Send a /chat/completions request to {model}")
def chat(
self,
key: str,
@ -407,6 +425,7 @@ class GuardrailsClient:
),
)
@step("Send a /chat/completions request to {model}")
def chat_raw(
self,
key: str,
@ -438,6 +457,7 @@ class GuardrailsClient:
),
)
@step("Send a streaming /chat/completions request to {model}")
def chat_stream_raw(
self,
key: str,
@ -462,6 +482,7 @@ class GuardrailsClient:
),
)
@step("Send a /v1/messages request to {model}")
def messages(
self,
key: str,
@ -481,6 +502,7 @@ class GuardrailsClient:
),
)
@step("Send a /v1/messages request to {model}")
def messages_raw(
self,
key: str,
@ -501,6 +523,7 @@ class GuardrailsClient:
),
)
@step("Send a streaming /v1/messages request to {model}")
def messages_stream_raw(
self,
key: str,
@ -521,6 +544,7 @@ class GuardrailsClient:
),
)
@step("Send a /v1/responses request to {model}")
def responses(
self,
key: str,
@ -535,6 +559,7 @@ class GuardrailsClient:
json=_ResponsesGuardrailBody(model=model, input=text, guardrails=guardrails),
)
@step("Send a streaming /v1/responses request to {model}")
def responses_stream_raw(
self,
key: str,
@ -553,6 +578,7 @@ class GuardrailsClient:
stream=True,
)
@step("Apply the guardrail {name} to a piece of text with /guardrails/apply_guardrail")
def apply_guardrail(self, key: str, *, name: str, text: str) -> Result[ApplyGuardrailResponse]:
return self.proxy.transport.post(
"/guardrails/apply_guardrail",

View file

@ -11,6 +11,7 @@ import pytest
from e2e_config import MASTER_KEY, unique_marker
from e2e_http import Success, UnauthorizedError, UnknownApiError
from e2e_metadata import Domain, Route, Subject, meta
from guardrails_client import GuardrailsClient
from lifecycle import ResourceManager
@ -23,6 +24,12 @@ class TestApplyGuardrailEndpoint:
"guardrail.litellm_content_filter.apply_endpoint.allows",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.GUARDRAILS,
)
)
def test_apply_guardrail_blocks_banned_and_allows_clean(
self, client: GuardrailsClient, resources: ResourceManager
) -> None:

View file

@ -22,6 +22,7 @@ from typing import Final
import pytest
from e2e_config import unique_marker
from e2e_http import StreamingResponse, UnknownApiError
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from guardrails_client import (
BedrockGuardrailParamsBody,
GuardrailsClient,
@ -59,6 +60,15 @@ class TestBedrockGuardrail:
"guardrail.bedrock.pre_call.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_bedrock_pre_call_blocks_harmful_prompt(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -96,6 +106,14 @@ class TestBedrockGuardrail:
"guardrail.bedrock.post_call.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_bedrock_post_call_blocks_denied_model_output(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -138,6 +156,15 @@ class TestBedrockGuardrail:
pytest.fail(f"bedrock post_call guardrail did not block denied model output; got {result}")
@pytest.mark.covers("guardrail.bedrock.pre_call.blocks", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_bedrock_pre_call_blocks_on_messages(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -149,6 +176,15 @@ class TestBedrockGuardrail:
_assert_policy_block(result, "/v1/messages")
@pytest.mark.covers("guardrail.bedrock.pre_call.blocks", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.RESPONSES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_bedrock_pre_call_blocks_on_responses(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -160,6 +196,14 @@ class TestBedrockGuardrail:
_assert_policy_block(result, "/v1/responses")
@pytest.mark.covers("guardrail.bedrock.post_call.blocks", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_bedrock_post_call_blocks_denied_streamed_output_and_passes_clean_streams(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -20,7 +20,8 @@ import pytest
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT, unique_marker
from e2e_http import unwrap
from guardrails_client import BlockCodeExecutionParamsBody, GuardrailsClient
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import GUARDRAIL_BACKEND, BlockCodeExecutionParamsBody, GuardrailsClient
from lifecycle import ResourceManager
from models import ChatResponse
@ -45,6 +46,14 @@ class TestBlockCodeExecutionGuardrail:
"guardrail.block_code_execution.pre_call.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(GUARDRAIL_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_blocks_execution_request_but_allows_explanation(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -11,6 +11,7 @@ from __future__ import annotations
import pytest
from e2e_config import unique_marker
from e2e_http import UnknownApiError, ValidationError
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import GuardrailsClient
pytestmark = pytest.mark.e2e
@ -28,6 +29,14 @@ MODEL = "gemini-2.5-flash"
"guardrail.dispatch.pre_call.rejects_unknown_name",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_request_naming_an_unknown_guardrail_fails_closed(client: GuardrailsClient, scoped_key: str) -> None:
result = client.chat(scoped_key, MODEL, "say hi", guardrails=[f"e2e-no-such-guardrail-{unique_marker()}"])

View file

@ -9,6 +9,7 @@ import pytest
from e2e_config import unique_marker
from e2e_http import unwrap
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import (
BlockedWordBody,
ContentFilterParamsBody,
@ -19,6 +20,8 @@ from models import ChatResponse, GuardrailInformationEntry
pytestmark = pytest.mark.e2e
BACKEND_MODEL: Final = "openai/gpt-4.1-mini"
GUARDRAIL_PROPAGATION_DEADLINE_SECONDS: Final = 40.0
GUARDRAIL_PROPAGATION_POLL_INTERVAL_SECONDS: Final = 5.0
@ -60,6 +63,14 @@ class TestGuardrailInformationResponse:
"guardrail.litellm_content_filter.pre_call.returns_guardrail_information",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.OPENAI,),
models=(BACKEND_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_flag_returns_guardrail_information_for_the_guardrail_that_ran(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -68,7 +79,7 @@ class TestGuardrailInformationResponse:
model = client.create_backend_model(
resources,
prefix="e2e-guardrail-info-backend",
backend="openai/gpt-4.1-mini",
backend=BACKEND_MODEL,
api_key="os.environ/OPENAI_API_KEY",
)
deadline = time.monotonic() + GUARDRAIL_PROPAGATION_DEADLINE_SECONDS
@ -91,6 +102,14 @@ class TestGuardrailInformationResponse:
)
time.sleep(GUARDRAIL_PROPAGATION_POLL_INTERVAL_SECONDS)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.OPENAI,),
models=(BACKEND_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_without_flag_response_has_no_guardrail_information(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -99,7 +118,7 @@ class TestGuardrailInformationResponse:
model = client.create_backend_model(
resources,
prefix="e2e-guardrail-info-backend",
backend="openai/gpt-4.1-mini",
backend=BACKEND_MODEL,
api_key="os.environ/OPENAI_API_KEY",
)

View file

@ -6,6 +6,7 @@ from typing import Final
import pytest
from e2e_config import unique_marker
from e2e_http import Success, UnknownApiError
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from guardrails_client import GuardrailsClient, poll_until_blocked
from lifecycle import ResourceManager
from models import LiteLLMParamsBody
@ -41,6 +42,15 @@ class TestKeyAttachedGuardrailOnImageEdits:
"guardrail.litellm_content_filter.pre_call.blocks_image_edit",
exercised_on=["images_edits"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.IMAGES,
providers=(Provider.GEMINI, Provider.OPENAI,),
models=(CHAT_MODEL, IMAGE_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_key_attached_content_filter_blocks_banned_image_edit_prompt(
self, client: GuardrailsClient, resources: ResourceManager
) -> None:

View file

@ -3,6 +3,7 @@ from __future__ import annotations
import pytest
from e2e_config import unique_marker
from e2e_http import Success, UnknownApiError
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import GuardrailsClient, poll_until_blocked
from lifecycle import ResourceManager
from models import LiteLLMParamsBody
@ -38,6 +39,14 @@ class TestKeyAttachedGuardrailOnVideos:
"guardrail.litellm_content_filter.pre_call.blocks_video",
exercised_on=["videos"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI, Provider.VERTEX_AI,),
models=(CHAT_MODEL, VIDEO_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_key_attached_content_filter_blocks_banned_video_prompt(
self, client: GuardrailsClient, resources: ResourceManager
) -> None:

View file

@ -7,15 +7,22 @@ body that names moderation; a refine-wrapper bypass must also be blocked.
from __future__ import annotations
from typing import Final
import pytest
from e2e_config import unique_marker
from e2e_http import Result, UnknownApiError
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from guardrails_client import GuardrailsClient, OpenAIModerationParamsBody
from lifecycle import ResourceManager
from models import AnthropicMessagesResponse, ChatResponse
pytestmark = pytest.mark.e2e
GEMINI_BACKEND: Final = "gemini/gemini-2.5-flash"
ANTHROPIC_BACKEND: Final = "anthropic/claude-haiku-4-5"
OPENAI_BACKEND: Final = "openai/gpt-4o-mini"
CATEGORY_PROMPTS: tuple[tuple[str, str], ...] = (
(
"violence",
@ -79,6 +86,15 @@ class TestOpenAIModerationCategoryMatrix:
"guardrail.openai_moderations.pre_call.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(GEMINI_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_chat_blocks_category(
self,
client: GuardrailsClient,
@ -89,7 +105,7 @@ class TestOpenAIModerationCategoryMatrix:
client,
resources,
prefix="e2e-mod-cat-chat",
backend="gemini/gemini-2.5-flash",
backend=GEMINI_BACKEND,
api_key="os.environ/GEMINI_API_KEY",
)
for category, prompt in CATEGORY_PROMPTS:
@ -99,6 +115,15 @@ class TestOpenAIModerationCategoryMatrix:
"guardrail.openai_moderations.pre_call.blocks",
exercised_on=["messages"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(ANTHROPIC_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_messages_blocks_category(
self,
client: GuardrailsClient,
@ -109,7 +134,7 @@ class TestOpenAIModerationCategoryMatrix:
client,
resources,
prefix="e2e-mod-cat-msg",
backend="anthropic/claude-haiku-4-5",
backend=ANTHROPIC_BACKEND,
api_key="os.environ/ANTHROPIC_API_KEY",
)
for category, prompt in CATEGORY_PROMPTS:
@ -119,6 +144,15 @@ class TestOpenAIModerationCategoryMatrix:
"guardrail.openai_moderations.pre_call.blocks",
exercised_on=["responses"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(OPENAI_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_responses_blocks_category(
self,
client: GuardrailsClient,
@ -129,7 +163,7 @@ class TestOpenAIModerationCategoryMatrix:
client,
resources,
prefix="e2e-mod-cat-resp",
backend="openai/gpt-4o-mini",
backend=OPENAI_BACKEND,
api_key="os.environ/OPENAI_API_KEY",
)
for category, prompt in CATEGORY_PROMPTS:

View file

@ -18,7 +18,9 @@ import pytest
from e2e_config import unique_marker
from e2e_http import UnknownApiError, unwrap
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from guardrails_client import (
GUARDRAIL_BACKEND,
GuardrailsClient,
OpenAIModerationParamsBody,
poll_until_blocked,
@ -37,6 +39,15 @@ class TestOpenAIModerationGuardrail:
"guardrail.openai_moderations.pre_call.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(GUARDRAIL_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_moderation_blocks_flagged_input(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -76,6 +87,15 @@ class TestOpenAIModerationGuardrail:
"guardrail.openai_moderations.pre_call.blocks",
exercised_on=["messages"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.GEMINI,),
models=(GUARDRAIL_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_moderation_blocks_flagged_input_on_messages(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -16,6 +16,7 @@ from __future__ import annotations
import pytest
from e2e_config import CHEAP_OPENAI_MODEL, unique_marker
from e2e_http import StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import (
GuardrailsClient,
PolicyConditionBody,
@ -73,6 +74,14 @@ def _setup_child_policy_attached_to_tag(
class TestPolicyInheritedGuardrail:
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.OPENAI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_child_condition_miss_still_applies_inherited_parent_guardrail(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -101,6 +110,14 @@ class TestPolicyInheritedGuardrail:
f"the child's own guardrail must not run when its condition fails; got {outcome.headers}"
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.OPENAI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_child_condition_match_applies_child_and_inherited_parent_guardrails(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -40,6 +40,7 @@ from pydantic import BaseModel, JsonValue, TypeAdapter
from e2e_config import unique_marker
from e2e_http import Result, StreamingResponse, Success
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from guardrails_client import GuardrailMode, GuardrailsClient, PiiAction, PiiEntity, PresidioParamsBody
from lifecycle import ResourceManager
from models import (
@ -248,6 +249,15 @@ class TestPresidioPreCallMasking:
"guardrail.presidio.pre_call.masks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_pre_call_masks_pii_on_chat_completions(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -267,6 +277,15 @@ class TestPresidioPreCallMasking:
"guardrail.presidio.pre_call.masks",
exercised_on=["messages"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_pre_call_masks_pii_on_messages(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -312,6 +331,14 @@ class TestPresidioPostCallMasking:
"guardrail.presidio.post_call.masks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_post_call_masks_pii_in_model_output(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -393,6 +420,15 @@ class TestPresidioCreditCardOutputMasking:
"guardrail.presidio.post_call.masks_generated_output",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_ui_default_scope_masks_a_card_number_the_model_generates_on_chat_completions(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -421,6 +457,15 @@ class TestPresidioCreditCardOutputMasking:
"guardrail.presidio.post_call.masks_generated_output",
exercised_on=["chat_completions_stream"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_ui_default_scope_masks_a_card_number_the_model_generates_on_streaming_chat_completions(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -453,6 +498,15 @@ class TestPresidioCreditCardOutputMasking:
"guardrail.presidio.post_call.masks_generated_output",
exercised_on=["anthropic_messages_stream"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_ui_default_scope_masks_a_card_number_the_model_generates_on_streaming_anthropic_messages(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -571,6 +625,15 @@ class TestPresidioSpendLogStoresMaskedOutput:
)
@pytest.mark.covers(_CELL, exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_spend_log_stores_masked_output_on_chat_completions(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -585,6 +648,15 @@ class TestPresidioSpendLogStoresMaskedOutput:
)
@pytest.mark.covers(_CELL, exercised_on=["chat_completions_stream"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_spend_log_stores_masked_output_on_streaming_chat_completions(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -599,6 +671,15 @@ class TestPresidioSpendLogStoresMaskedOutput:
)
@pytest.mark.covers(_CELL, exercised_on=["messages"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_spend_log_stores_masked_output_on_anthropic_messages(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -613,6 +694,15 @@ class TestPresidioSpendLogStoresMaskedOutput:
)
@pytest.mark.covers(_CELL, exercised_on=["anthropic_messages_stream"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.MESSAGES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_spend_log_stores_masked_output_on_streaming_anthropic_messages(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -627,6 +717,15 @@ class TestPresidioSpendLogStoresMaskedOutput:
)
@pytest.mark.covers(_CELL, exercised_on=["responses"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.RESPONSES,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_spend_log_stores_masked_output_on_responses(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -673,6 +772,14 @@ class TestPresidioSpendLogRecord:
"guardrail.presidio.pre_call.logs_masked_entities",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_masking_run_is_recorded_on_the_spend_log(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -7,12 +7,15 @@ from typing import Final
import pytest
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT, unique_marker
from e2e_http import StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from guardrails_client import CustomCodeParamsBody, GuardrailsClient
from lifecycle import ResourceManager
from pydantic import BaseModel, TypeAdapter
pytestmark = pytest.mark.e2e
BACKEND_MODEL: Final = "openai/gpt-4.1-mini"
DENIAL: Final = "This model is not currently available. Please contact support if you think this is a mistake."
CUSTOM_CODE: Final = f'''
@ -109,12 +112,21 @@ class TestResponsesPreCallBlock:
return name
@pytest.mark.covers("guardrail.custom_code.pre_call.blocks", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(BACKEND_MODEL,),
mode=Mode.STREAM,
)
)
def test_stream_block_is_sse_with_completed_assistant_message(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
name: Final = self._register_block(client, resources)
model: Final = client.create_backend_model(
resources, prefix="e2e-responses-block", backend="openai/gpt-4.1-mini", api_key="os.environ/OPENAI_API_KEY"
resources, prefix="e2e-responses-block", backend=BACKEND_MODEL, api_key="os.environ/OPENAI_API_KEY"
)
result: Final = _poll_for_block(
@ -137,12 +149,21 @@ class TestResponsesPreCallBlock:
_assert_blocked_response(completed[0].response)
@pytest.mark.covers("guardrail.custom_code.pre_call.blocks", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(BACKEND_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_non_stream_block_is_schema_valid_json(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
name: Final = self._register_block(client, resources)
model: Final = client.create_backend_model(
resources, prefix="e2e-responses-block", backend="openai/gpt-4.1-mini", api_key="os.environ/OPENAI_API_KEY"
resources, prefix="e2e-responses-block", backend=BACKEND_MODEL, api_key="os.environ/OPENAI_API_KEY"
)
result: Final = _poll_for_block(lambda: client.responses(scoped_key, model, "say hi", guardrails=[name]))

View file

@ -22,6 +22,7 @@ import os
import pytest
from e2e_config import unique_marker
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import (
BedrockGuardrailParamsBody,
GuardrailsClient,
@ -39,6 +40,14 @@ class TestBedrockDuringCallStreaming:
"guardrail.bedrock.during.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_during_call_blocks_stream_before_first_chunk(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -14,6 +14,7 @@ import pytest
from e2e_config import unique_marker
from e2e_http import UnknownApiError, unwrap
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from guardrails_client import GuardrailsClient
from lifecycle import ResourceManager
@ -59,6 +60,14 @@ class TestTeamDisableGlobalGuardrail:
"guardrail.litellm_content_filter.pre_call.blocks",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_global_guardrail_blocks_key_without_team_opt_out(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -72,6 +81,14 @@ class TestTeamDisableGlobalGuardrail:
"guardrail.litellm_content_filter.pre_call.allows",
exercised_on=["chat_completions"],
)
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_team_with_disable_flag_bypasses_global_guardrail(
self, client: GuardrailsClient, resources: ResourceManager
) -> None:

View file

@ -25,6 +25,7 @@ import pytest
from e2e_config import unique_marker
from e2e_http import StreamingResponse, UnknownApiError
from e2e_metadata import Capability, Domain, Mode, Provider, Subject, meta
from guardrails_client import (
GuardrailsClient,
ToolPermissionParamsBody,
@ -101,6 +102,15 @@ def _tool_call_names(response: ChatResponse) -> tuple[str, ...]:
class TestToolPermissionPreCall:
@pytest.mark.covers("guardrail.tool_permission.pre_call.blocks", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
capabilities=(Capability.FUNCTION_CALLING,),
mode=Mode.NONSTREAM,
)
)
def test_pre_call_blocks_tool_outside_the_allow_list(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:
@ -135,6 +145,15 @@ class TestToolPermissionPreCall:
pytest.fail(f"tool_permission let a tool outside the allow-list through; got {result}")
@pytest.mark.covers("guardrail.tool_permission.pre_call.allows", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.GUARDRAILS,
providers=(Provider.GEMINI,),
models=(MODEL,),
capabilities=(Capability.FUNCTION_CALLING,),
mode=Mode.NONSTREAM,
)
)
def test_pre_call_allows_permitted_tool(
self, client: GuardrailsClient, resources: ResourceManager, scoped_key: str
) -> None:

View file

@ -32,6 +32,7 @@ from e2e_config import (
POLL_TIMEOUT,
)
from e2e_http import URL, Headers, StreamingResponse, send
from e2e_metadata import step
type SearchCall = Callable[[str, float], StreamingResponse]
@ -115,6 +116,7 @@ class DdLogsReader:
sleep: Callable[[float], None] = field(default=time.sleep, repr=False)
jitter: Callable[[], float] = field(default=random.random, repr=False)
@step("Search DataDog for logs carrying the marker {marker}")
def events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""Every ingested event whose attributes carry the marker. DataDog
consumes the shipped JSON message into ``attributes`` and leaves the
@ -124,6 +126,7 @@ class DdLogsReader:
it)."""
return self.events_for_query(f"*:*{marker}*")
@step("Search DataDog for logs matching {query}")
def events_for_query(self, query: str) -> list[DdLogEvent]:
"""Every ingested event the search query matches (failure payloads
carry no prompt to mark, so failure scenarios query indexed attributes
@ -156,10 +159,12 @@ class DdLogsReader:
timeout=timeout,
)
@step("Wait for DataDog to ingest logs carrying the marker {marker}, then watch for duplicates")
def poll_events_for_marker(self, marker: str) -> list[DdLogEvent]:
"""``poll_events_for_query`` over the every-attribute marker scan."""
return self.poll_events_for_query(f"*:*{marker}*")
@step("Wait for DataDog to ingest logs matching {query}, then watch for duplicates")
def poll_events_for_query(self, query: str) -> list[DdLogEvent]:
"""Poll until at least one matching event is searchable (the callback
flushes in periodic batches and DataDog ingestion adds seconds of lag),

View file

@ -30,6 +30,7 @@ from pydantic import BaseModel, ConfigDict, Field
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT
from e2e_http import URL, Headers, probe
from e2e_metadata import step
_GCS_API = "https://storage.googleapis.com"
#: Tolerance for clock skew between this host and GCS object timestamps.
@ -44,7 +45,7 @@ class _ServiceAccount(BaseModel):
model_config = ConfigDict(extra="ignore")
client_email: str
private_key: str
private_key: str = Field(repr=False)
class _GcsAuthHeaders(Headers):
@ -146,6 +147,7 @@ class GcsLogReader:
)
return result.body
@step("Read the GCS bucket's log objects for the response {response_id}")
def records_for_response_id(self, response_id: str, *, since: datetime) -> list[GcsLogRecord]:
"""Every payload written for ``response_id``: the direct
``{date}/{response_id}`` object plus any hit inside batch NDJSON
@ -168,6 +170,7 @@ class GcsLogReader:
)
return records
@step("Wait for the log of the response {response_id} to land in the GCS bucket, then watch for duplicates")
def poll_records_for_response_id(self, response_id: str, *, since: datetime) -> list[GcsLogRecord]:
"""Poll until the payload is readable (the gcs_bucket callback flushes
on a ~20s timer), then keep re-reading for GCS_SETTLE_SECONDS - past a

View file

@ -24,6 +24,7 @@ import pytest
from pydantic import BaseModel, ConfigDict, Field, JsonValue, TypeAdapter, ValidationError
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT, settle_propagation
from e2e_metadata import step
from proxy_client import ProxyClient
from e2e_http import (
URL,
@ -326,6 +327,7 @@ def observation_has_guardrail(obs: LangfuseObservation, *, guardrail_name: str)
class LoggingClient:
proxy: ProxyClient
@step("Generate a virtual key named {alias} with models: {models}")
def key_with_alias(
self,
alias: str,
@ -347,9 +349,11 @@ class LoggingClient:
)
)
@step("Delete the virtual key")
def delete_key(self, key: str) -> None:
self.proxy.delete_key(key)
@step("Create the team {alias} with models: {models}")
def create_team(
self,
alias: str,
@ -370,6 +374,7 @@ class LoggingClient:
)
).team_id
@step("Delete the team")
def delete_team(self, team_id: str) -> None:
_ = self.proxy.transport.post(
"/team/delete",
@ -378,6 +383,7 @@ class LoggingClient:
response_type=NoBody,
)
@step("Create the internal user {user_email}")
def create_user(self, *, user_email: str, user_id: str | None = None) -> str:
return unwrap(
self.proxy.transport.post(
@ -392,6 +398,7 @@ class LoggingClient:
)
).user_id
@step("Delete the internal user")
def delete_user(self, user_id: str) -> None:
_ = self.proxy.transport.post(
"/user/delete",
@ -400,6 +407,7 @@ class LoggingClient:
response_type=NoBody,
)
@step("Create the organization {alias} with models: {models}")
def create_org(self, alias: str, *, models: list[str]) -> str:
return unwrap(
self.proxy.transport.post(
@ -410,6 +418,7 @@ class LoggingClient:
)
).organization_id
@step("Delete the organization")
def delete_org(self, organization_id: str) -> None:
_ = self.proxy.transport.delete(
"/organization/delete",
@ -418,6 +427,7 @@ class LoggingClient:
response_type=NoBody,
)
@step("Add a Langfuse OTel logging callback for {callback_type} events to the team")
def add_team_langfuse_callback(
self,
team_id: str,
@ -441,6 +451,7 @@ class LoggingClient:
f"POST /team/{team_id}/callback must return status=success; got {response.status!r}"
)
@step("Create the tool_permission guardrail {name} that allows only the tool {allowed_tool}")
def create_tool_permission_guardrail(self, name: str, *, allowed_tool: str) -> str:
"""Register a tool_permission guardrail that allows one tool and denies the rest."""
response = unwrap(
@ -474,6 +485,7 @@ class LoggingClient:
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}",
@ -482,12 +494,15 @@ class LoggingClient:
response_type=NoBody,
)
@step("Add a 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)
@step("Delete the deployment")
def delete_model(self, model_id: str) -> None:
self.proxy.delete_model(model_id)
@step('Send a /chat/completions request to {model} with the prompt "{text}"')
def chat(self, key: str, model: str, text: str) -> ChatResponse:
return unwrap(
self.proxy.chat(
@ -500,6 +515,7 @@ class LoggingClient:
)
)
@step('Send a /chat/completions request to {model} with stream={stream} and the prompt "{text}"')
def chat_raw(
self,
key: str,
@ -529,6 +545,7 @@ class LoggingClient:
json=body,
)
@step('Send a /v1/messages request to {model} with stream={stream} and the prompt "{text}"')
def messages_raw(
self, key: str, model: str, text: str, *, max_tokens: int = 16, stream: bool = False
) -> StreamingResponse:
@ -545,6 +562,7 @@ class LoggingClient:
return self.proxy.transport.stream("/v1/messages", headers=self.proxy.transport.bearer(key), json=body)
return self.proxy.transport.send("/v1/messages", headers=self.proxy.transport.bearer(key), json=body)
@step('Send a /v1/responses request to {model} with stream={stream} and the prompt "{text}"')
def responses_raw(
self, key: str, model: str, text: str, *, max_output_tokens: int = 64, stream: bool = False
) -> StreamingResponse:
@ -560,9 +578,11 @@ class LoggingClient:
return self.proxy.transport.stream("/v1/responses", headers=self.proxy.transport.bearer(key), json=body)
return self.proxy.transport.send("/v1/responses", headers=self.proxy.transport.bearer(key), json=body)
@step("Scrape the Prometheus metrics from /metrics")
def scrape_metrics(self) -> str:
return self.proxy.probe("/metrics", params=NoBody()).body
@step("Wait for the key's spend log in /spend/logs")
def poll_proxy_spend_for_key(
self,
key: str,
@ -590,6 +610,7 @@ class LoggingClient:
return row
return None
@step("List observations from Langfuse")
def list_langfuse_observations(
self,
creds: LangfuseCreds,
@ -616,6 +637,7 @@ class LoggingClient:
case _:
return []
@step("Look up the Langfuse generation for the key {key_alias}")
def find_langfuse_observation(
self,
creds: LangfuseCreds,
@ -634,6 +656,7 @@ class LoggingClient:
return obs
return None
@step("Wait for the Langfuse generation for the key {key_alias}")
def poll_langfuse_observation(
self,
creds: LangfuseCreds,
@ -653,6 +676,7 @@ class LoggingClient:
time.sleep(POLL_INTERVAL)
return last
@step("Wait for the OTel v2 Langfuse generation for the key {key_alias}")
def poll_langfuse_generation(
self, creds: LangfuseCreds, *, key_alias: str, from_start_time: str
) -> LangfuseObservation | None:
@ -665,6 +689,7 @@ class LoggingClient:
time.sleep(POLL_INTERVAL)
return None
@step("Wait for the Langfuse trace of the key {key_alias} and every observation in it")
def poll_langfuse_trace_observations(
self,
creds: LangfuseCreds,
@ -698,6 +723,7 @@ def build_logging_client(proxy: ProxyClient) -> LoggingClient:
return LoggingClient(proxy=proxy)
@step("Read the proxy's callback list from /health/readiness/details")
def readiness_details_body(client: LoggingClient) -> str:
"""/health/readiness/details, tolerating the 503 it serves while the
ephemeral stack's DB leg blips: the recorded state the logging suites check

View file

@ -24,6 +24,7 @@ import pytest
from pydantic import BaseModel, ConfigDict
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT
from e2e_metadata import step
if TYPE_CHECKING:
from types_boto3_s3.client import S3Client
@ -53,17 +54,21 @@ class S3LogReader:
bucket: str
client: S3Client
@step("List the log objects in the S3 bucket under {prefix}")
def list_keys(self, prefix: str) -> list[str]:
response = self.client.list_objects_v2(Bucket=self.bucket, Prefix=prefix)
return [obj["Key"] for obj in response.get("Contents", []) if "Key" in obj]
@step("Download a log object from the S3 bucket")
def read_record(self, key: str) -> S3LogRecord:
body = self.client.get_object(Bucket=self.bucket, Key=key)["Body"].read()
return S3LogRecord.model_validate_json(body)
@step("Read the log objects in the S3 bucket under {prefix}")
def records_matching(self, *, prefix: str, predicate: Callable[[S3LogRecord], bool]) -> list[S3LogRecord]:
return [record for record in map(self.read_record, self.list_keys(prefix)) if predicate(record)]
@step("Wait for the request's log object to land in the S3 bucket under {prefix}, then watch for duplicates")
def poll_records(self, *, prefix: str, predicate: Callable[[S3LogRecord], bool]) -> list[S3LogRecord]:
"""Poll until at least one matching object is listed (the s3_v2
callback flushes on a ~10s timer), then keep re-reading for

View file

@ -20,10 +20,12 @@ from __future__ import annotations
import math
import time
from typing import Final
import pytest
from datadog_reader import DdLogEvent, DdLogsReader
from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, unique_marker
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok, readiness_details_body
from models import ChatMessage, LiteLLMParamsBody, ReliabilityChatBody, RouterSettingsOverride
@ -33,6 +35,7 @@ pytestmark = pytest.mark.e2e
#: The active DataDog callback's name in /health/readiness/details success_callbacks.
DD_LOGGER_NAME = "DataDogLogger"
FAILING_BACKEND_MODEL: Final = "anthropic/claude-haiku-4-5"
class _DdMessagePayload(BaseModel):
@ -109,6 +112,15 @@ def _assert_exactly_one_event(
class TestDataDogLogDelivery:
@pytest.mark.covers("logging.datadog.success.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_emits_one_log_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -134,6 +146,15 @@ class TestDataDogLogDelivery:
)
@pytest.mark.covers("logging.datadog.success.exports_metric", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_messages_emits_one_log_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -161,6 +182,15 @@ class TestDataDogLogDelivery:
)
@pytest.mark.covers("logging.datadog.success.exports_metric", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_responses_emits_one_log_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -186,6 +216,15 @@ class TestDataDogLogDelivery:
)
@pytest.mark.covers("logging.datadog.stream.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.STREAM,
)
)
def test_chat_completions_stream_emits_one_log_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -232,6 +271,15 @@ class TestDataDogLogDelivery:
)
@pytest.mark.covers("logging.datadog.stream.exports_metric", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.STREAM,
)
)
def test_messages_stream_emits_one_log_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -276,6 +324,15 @@ class TestDataDogLogDelivery:
)
@pytest.mark.covers("logging.datadog.stream.exports_metric", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL,),
mode=Mode.STREAM,
)
)
def test_responses_stream_emits_one_log_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -347,6 +404,15 @@ def _assert_exactly_one_failure_event(events: list[DdLogEvent], *, model_group:
class TestDataDogFailureDelivery:
@pytest.mark.covers("logging.datadog.failure.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(FAILING_BACKEND_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_failed_chat_completions_emits_one_error_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -366,7 +432,7 @@ class TestDataDogFailureDelivery:
model_name = f"dd-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
LiteLLMParamsBody(model=FAILING_BACKEND_MODEL, api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
key = client.key_with_alias(f"dd-err-key-{unique_marker()}", models=[model_name])
@ -399,6 +465,15 @@ class TestDataDogFailureDelivery:
)
@pytest.mark.covers("logging.datadog.stream_failure.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(FAILING_BACKEND_MODEL,),
mode=Mode.STREAM,
)
)
def test_failed_chat_completions_stream_emits_one_error_event(
self, client: LoggingClient, dd_logs: DdLogsReader, resources: ResourceManager
) -> None:
@ -420,7 +495,7 @@ class TestDataDogFailureDelivery:
model_id = client.create_model(
model_name,
LiteLLMParamsBody(
model="anthropic/claude-haiku-4-5",
model=FAILING_BACKEND_MODEL,
api_key=INVALID_UPSTREAM_API_KEY,
api_base="http://localhost:1",
),

View file

@ -22,6 +22,7 @@ import math
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from gcs_reader import GcsLogReader, build_gcs_reader, utc_now
from lifecycle import ResourceManager
from logging_client import LoggingClient, completion_response_id, first_ok, readiness_details_body
@ -52,6 +53,14 @@ def _assert_gcs_configured(client: LoggingClient) -> None:
class TestGcsLogDelivery:
@pytest.mark.covers("logging.gcs_bucket.success.writes_object", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_writes_one_success_record(
self, client: LoggingClient, gcs_logs: GcsLogReader, resources: ResourceManager
) -> None:

View file

@ -21,6 +21,7 @@ from typing import Final
import pytest
from e2e_config import CHEAP_OPENAI_MODEL, POLL_INTERVAL, POLL_TIMEOUT, unique_marker
from e2e_http import Headers, Success, get_external
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from pydantic import BaseModel, ConfigDict, Field, JsonValue
import litellm
@ -90,6 +91,14 @@ def _poll_run(creds: LangsmithCreds, run_id: uuid.UUID) -> LangsmithRun:
class TestLangsmithBatchSerialization:
@pytest.mark.asyncio
@pytest.mark.covers("logging.langsmith.success.serializes_non_native_metadata")
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL,),
mode=Mode.NONSTREAM,
)
)
async def test_non_json_native_metadata_reaches_langsmith(self) -> None:
creds: Final = load_langsmith_creds()
logger: Final = LangsmithLogger(

View file

@ -25,6 +25,7 @@ from typing import Final
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, CHEAP_OPENAI_MODEL, OTEL_EXPORTER_ENDPOINT, unique_marker
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from logging_client import INVALID_UPSTREAM_API_KEY, LoggingClient, first_ok, readiness_details_body
from models import LiteLLMParamsBody
@ -34,6 +35,7 @@ from pydantic import BaseModel, ConfigDict, ValidationError
pytestmark = pytest.mark.e2e
MODEL = CHEAP_ANTHROPIC_MODEL
FAILING_BACKEND_MODEL: Final = "anthropic/claude-haiku-4-5"
DB_SPAN_PREFIX = "postgres."
#: The active OTEL v2 logger's name in /health/readiness/details success_callbacks.
OTEL_V2_LOGGER_NAME = "OpenTelemetryV2"
@ -279,6 +281,15 @@ def _assert_error_span_contract(span: JaegerSpan) -> None:
class TestOtelTraceCompleteness:
@pytest.mark.covers("logging.otel.success.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_exports_complete_trace(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -312,6 +323,14 @@ class TestOtelTraceCompleteness:
_assert_complete_trace(traces, route=route, genai_span=f"chat {MODEL}")
@pytest.mark.covers("logging.otel.success.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
@pytest.mark.otel_tls
def test_otel_export_over_tls_with_internal_ca_reaches_destination(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
@ -337,6 +356,15 @@ class TestOtelTraceCompleteness:
_assert_complete_trace(hits, route=route, genai_span=f"chat {MODEL}")
@pytest.mark.covers("logging.otel.success.exports_metric", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_messages_exports_complete_trace(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -366,6 +394,15 @@ class TestOtelTraceCompleteness:
_assert_complete_trace(traces, route=route, genai_span=f"chat {MODEL}")
@pytest.mark.covers("logging.otel.success.exports_metric", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_responses_exports_complete_trace(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -396,6 +433,15 @@ class TestOtelTraceCompleteness:
_assert_complete_trace(traces, route=route, genai_span=genai_span)
@pytest.mark.covers("logging.otel.stream.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_chat_completions_stream_exports_complete_trace(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -444,6 +490,15 @@ class TestOtelTraceCompleteness:
)
@pytest.mark.covers("logging.otel.stream.exports_metric", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_messages_stream_exports_complete_trace(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -492,6 +547,15 @@ class TestOtelTraceCompleteness:
)
@pytest.mark.covers("logging.otel.stream.exports_metric", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL,),
mode=Mode.STREAM,
)
)
def test_responses_stream_exports_complete_trace(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -544,6 +608,15 @@ class TestOtelTraceCompleteness:
)
@pytest.mark.covers("logging.otel.stream.records_ttft", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_chat_completions_stream_records_real_ttft(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -582,6 +655,15 @@ class TestOtelTraceCompleteness:
_assert_real_ttft(traces.hits, genai_span=genai_span)
@pytest.mark.covers("logging.otel.stream.records_ttft", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(MODEL,),
mode=Mode.STREAM,
)
)
def test_messages_stream_records_real_ttft(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -620,6 +702,15 @@ class TestOtelTraceCompleteness:
_assert_real_ttft(traces.hits, genai_span=genai_span)
@pytest.mark.covers("logging.otel.stream.records_ttft", exercised_on=["responses"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.RESPONSES,
providers=(Provider.OPENAI,),
models=(CHEAP_OPENAI_MODEL,),
mode=Mode.STREAM,
)
)
def test_responses_stream_records_real_ttft(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -658,6 +749,15 @@ class TestOtelTraceCompleteness:
_assert_real_ttft(traces.hits, genai_span=genai_span)
@pytest.mark.covers("logging.otel.failure.exports_metric", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.CHAT_COMPLETIONS,
providers=(Provider.ANTHROPIC,),
models=(FAILING_BACKEND_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_failed_chat_completions_error_span_attributes(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -677,7 +777,7 @@ class TestOtelTraceCompleteness:
model_name = f"otel-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
LiteLLMParamsBody(model=FAILING_BACKEND_MODEL, api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
key = client.key_with_alias(f"otel-err-{unique_marker()}", models=[model_name])
@ -711,6 +811,15 @@ class TestOtelTraceCompleteness:
_assert_error_span_contract(genai)
@pytest.mark.covers("logging.otel.failure.exports_metric", exercised_on=["messages"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MESSAGES,
providers=(Provider.ANTHROPIC,),
models=(FAILING_BACKEND_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_failed_messages_error_span_attributes(
self, client: LoggingClient, otel_reader: OtelReader, resources: ResourceManager
) -> None:
@ -729,7 +838,7 @@ class TestOtelTraceCompleteness:
model_name = f"otel-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
LiteLLMParamsBody(model=FAILING_BACKEND_MODEL, api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
key = client.key_with_alias(f"otel-err-{unique_marker()}", models=[model_name])

View file

@ -24,6 +24,7 @@ from typing import Final
import pytest
from e2e_config import unique_marker
from e2e_http import unwrap
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from logging_client import LangfuseCreds, LangfuseObservation, LoggingClient, load_langfuse_creds
from models import (
@ -69,6 +70,13 @@ RED_SQUARE_PNG: Final = base64.b64decode(
)
BOUNDED_OUTPUT_CHARS: Final = 1024
PLACEHOLDER_INPUT: Final = "default-message-value"
COMPLETION_BACKEND: Final = "openai/gpt-3.5-turbo-instruct"
IMAGE_BACKEND: Final = "openai/gpt-image-1-mini"
SPEECH_BACKEND: Final = "openai/gpt-4o-mini-tts"
TRANSCRIPTION_BACKEND: Final = "openai/gpt-4o-mini-transcribe"
MODERATION_BACKEND: Final = "openai/omni-moderation-latest"
MISTRAL_OCR_BACKEND: Final = "mistral/mistral-ocr-latest"
RERANK_BACKEND: Final = "cohere/rerank-v4.0-fast"
class _OutputMessage(BaseModel):
@ -142,7 +150,7 @@ def _openai(model: str) -> LiteLLMParamsBody:
def _mistral_ocr() -> LiteLLMParamsBody:
return LiteLLMParamsBody(model="mistral/mistral-ocr-latest", api_key="os.environ/MISTRAL_API_KEY")
return LiteLLMParamsBody(model=MISTRAL_OCR_BACKEND, api_key="os.environ/MISTRAL_API_KEY")
def _langfuse_search_tool(
@ -171,10 +179,19 @@ def _langfuse_search_tool(
class TestOtelV2LangfuseGenerationOutput:
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.COMPLETIONS,
providers=(Provider.OPENAI,),
models=(COMPLETION_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_completions_output_is_the_completion_text(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai("openai/gpt-3.5-turbo-instruct"))
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai(COMPLETION_BACKEND))
started: Final = datetime.now(timezone.utc)
response: Final = unwrap(
client.proxy.transport.post(
@ -193,10 +210,19 @@ class TestOtelV2LangfuseGenerationOutput:
)
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["images_generations"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.IMAGES,
providers=(Provider.OPENAI,),
models=(IMAGE_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_images_output_is_a_bounded_summary_without_base64(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai("openai/gpt-image-1-mini"))
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai(IMAGE_BACKEND))
started: Final = datetime.now(timezone.utc)
response: Final = unwrap(
client.proxy.transport.post(
@ -220,10 +246,19 @@ class TestOtelV2LangfuseGenerationOutput:
)
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["audio_speech"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.AUDIO,
providers=(Provider.OPENAI,),
models=(SPEECH_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_speech_output_is_a_bounded_summary_without_audio_bytes(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai("openai/gpt-4o-mini-tts"))
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai(SPEECH_BACKEND))
started: Final = datetime.now(timezone.utc)
audio: Final = client.proxy.transport.stream_binary(
"/v1/audio/speech",
@ -239,10 +274,19 @@ class TestOtelV2LangfuseGenerationOutput:
)
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["audio_transcriptions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.AUDIO,
providers=(Provider.OPENAI,),
models=(TRANSCRIPTION_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_transcription_output_is_the_transcript(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai("openai/gpt-4o-mini-transcribe"))
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai(TRANSCRIPTION_BACKEND))
started: Final = datetime.now(timezone.utc)
response: Final = unwrap(
client.proxy.transport.upload(
@ -262,10 +306,19 @@ class TestOtelV2LangfuseGenerationOutput:
assert transcript in output, f"generation output lacks the transcript {transcript!r}: {output!r}"
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["moderations"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.MODERATIONS,
providers=(Provider.OPENAI,),
models=(MODERATION_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_moderations_output_is_the_verdict(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai("openai/omni-moderation-latest"))
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai(MODERATION_BACKEND))
started: Final = datetime.now(timezone.utc)
response: Final = unwrap(
client.proxy.transport.post(
@ -284,6 +337,15 @@ class TestOtelV2LangfuseGenerationOutput:
)
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["rerank"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.RERANK,
providers=(Provider.COHERE,),
models=(RERANK_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_rerank_output_is_the_ranked_indices_and_scores(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
@ -291,7 +353,7 @@ class TestOtelV2LangfuseGenerationOutput:
client,
langfuse_creds,
resources,
LiteLLMParamsBody(model="cohere/rerank-v4.0-fast", api_key="os.environ/COHERE_API_KEY"),
LiteLLMParamsBody(model=RERANK_BACKEND, api_key="os.environ/COHERE_API_KEY"),
)
query: Final = f"What is the capital of France? {unique_marker()}"
started: Final = datetime.now(timezone.utc)
@ -317,6 +379,15 @@ class TestOtelV2LangfuseGenerationOutput:
assert output == "\n\n".join(ranked), f"generation output is not the ranked indices and scores: {output!r}"
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["ocr"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.OCR,
providers=(Provider.MISTRAL,),
models=(MISTRAL_OCR_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_ocr_input_is_the_document_url(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
@ -334,6 +405,15 @@ class TestOtelV2LangfuseGenerationOutput:
assert response.pages[0].markdown in _output_text(generation)
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["ocr"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.OCR,
providers=(Provider.MISTRAL,),
models=(MISTRAL_OCR_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_ocr_upload_input_is_a_bounded_document_summary_without_base64(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
@ -360,10 +440,19 @@ class TestOtelV2LangfuseGenerationOutput:
)
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["images_edits"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.IMAGES,
providers=(Provider.OPENAI,),
models=(IMAGE_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_image_edit_input_is_the_edit_prompt(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai("openai/gpt-image-1-mini"))
model, key, alias = _langfuse_key(client, langfuse_creds, resources, _openai(IMAGE_BACKEND))
prompt: Final = f"make the square blue {unique_marker()}"
started: Final = datetime.now(timezone.utc)
response: Final = unwrap(
@ -388,6 +477,11 @@ class TestOtelV2LangfuseGenerationOutput:
assert _output_text(generation).startswith("b64_json image (")
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["search"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
)
)
def test_search_input_is_the_query_and_output_the_results(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:

View file

@ -24,6 +24,7 @@ import pytest
from prometheus_client.parser import text_string_to_metric_families
from e2e_config import unique_marker
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from logging_client import LoggingClient
@ -47,6 +48,15 @@ def _aliases_in_metric(exposition: str, metric: str, label: str) -> frozenset[st
class TestPrometheusPerKeyCardinality:
@pytest.mark.covers("logging.prometheus.success.exports_metric", exercised_on=[])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.METRICS,
providers=(Provider.GEMINI,),
models=(DRIVER_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_distinct_key_aliases_produce_distinct_series(
self, client: LoggingClient, resources: ResourceManager
) -> None:

View file

@ -6,6 +6,7 @@ import pytest
from prometheus_client.parser import text_string_to_metric_families
from e2e_config import unique_marker
from e2e_metadata import Domain, Mode, Provider, Route, Subject, meta
from lifecycle import ResourceManager
from logging_client import LoggingClient
@ -30,6 +31,15 @@ def _observation_count(exposition: str, alias: str) -> float | None:
class TestPrometheusRequestQueueTime:
@pytest.mark.covers("logging.prometheus.success.records_queue_time")
@meta(
Subject(
domain=Domain.OBSERVABILITY,
route=Route.METRICS,
providers=(Provider.GEMINI,),
models=(DRIVER_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_queue_time_histogram_records_an_observation(
self, client: LoggingClient, resources: ResourceManager
) -> None:

View file

@ -24,10 +24,12 @@ from __future__ import annotations
import math
import re
import time
from typing import Final
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, S3_PARTITION_GRANULARITY, unique_marker
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from logging_client import (
INVALID_UPSTREAM_API_KEY,
@ -43,6 +45,7 @@ pytestmark = pytest.mark.e2e
#: The active s3_v2 callback's name in /health/readiness/details success_callbacks.
S3_LOGGER_NAME = "S3Logger"
UNREACHABLE_ANTHROPIC_BACKEND: Final = "anthropic/claude-haiku-4-5"
@pytest.fixture(scope="session")
@ -64,6 +67,14 @@ def _assert_s3_configured(client: LoggingClient) -> None:
class TestS3LogDelivery:
@pytest.mark.covers("logging.s3.success.writes_object", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_writes_one_success_object(
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
) -> None:
@ -108,6 +119,14 @@ class TestS3LogDelivery:
), f"payload response_cost {record.response_cost!r} must equal the header cost {outcome.response_cost}"
@pytest.mark.covers("logging.s3.success.partition_layout", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_object_key_follows_the_partition_granularity(
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
) -> None:
@ -148,6 +167,14 @@ class TestS3LogDelivery:
)
@pytest.mark.covers("logging.s3.failure.writes_object", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(UNREACHABLE_ANTHROPIC_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_failure_writes_one_object(
self, client: LoggingClient, s3_logs: S3LogReader, resources: ResourceManager
) -> None:
@ -167,7 +194,7 @@ class TestS3LogDelivery:
model_name = f"s3-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
LiteLLMParamsBody(model=UNREACHABLE_ANTHROPIC_BACKEND, api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
alias = f"s3-err-key-{unique_marker()}"

View file

@ -19,6 +19,7 @@ import time
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from logging_client import (
LangfuseCreds,
@ -46,6 +47,14 @@ def langfuse_creds() -> LangfuseCreds:
class TestTeamLangfuseCallback:
@pytest.mark.covers("logging.langfuse.success.logs_spend", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_team_callback_delivers_and_isolates(
self, client: LoggingClient, langfuse_creds: LangfuseCreds, resources: ResourceManager
) -> None:

View file

@ -21,11 +21,13 @@ query API; nothing is mocked.
from __future__ import annotations
import time
from typing import Final
import pytest
from e2e_config import CHEAP_ANTHROPIC_MODEL, unique_marker
from e2e_http import StreamingResponse
from e2e_metadata import Domain, Mode, Provider, Subject, meta
from lifecycle import ResourceManager
from logging_client import (
INVALID_UPSTREAM_API_KEY,
@ -40,6 +42,8 @@ from weave_reader import WeaveCall, WeaveReader, build_weave_reader
pytestmark = pytest.mark.e2e
UNREACHABLE_ANTHROPIC_BACKEND: Final = "anthropic/claude-haiku-4-5"
@pytest.fixture(scope="session")
def weave_creds() -> WeaveCreds:
@ -79,6 +83,14 @@ WEAVE_STAGE_RED_REASON = (
class TestWeaveLogDelivery:
@pytest.mark.skip(reason=WEAVE_STAGE_RED_REASON)
@pytest.mark.covers("logging.niche_integrations.success.logs_spend", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(CHEAP_ANTHROPIC_MODEL,),
mode=Mode.NONSTREAM,
)
)
def test_chat_completions_delivers_one_call_with_spend(
self,
client: LoggingClient,
@ -121,6 +133,14 @@ class TestWeaveLogDelivery:
@pytest.mark.skip(reason=WEAVE_STAGE_RED_REASON)
@pytest.mark.covers("logging.niche_integrations.failure.logs_spend", exercised_on=["chat_completions"])
@meta(
Subject(
domain=Domain.OBSERVABILITY,
providers=(Provider.ANTHROPIC,),
models=(UNREACHABLE_ANTHROPIC_BACKEND,),
mode=Mode.NONSTREAM,
)
)
def test_failed_chat_completions_delivers_one_error_call(
self,
client: LoggingClient,
@ -135,7 +155,7 @@ class TestWeaveLogDelivery:
model_name = f"weave-err-{unique_marker()}"
model_id = client.create_model(
model_name,
LiteLLMParamsBody(model="anthropic/claude-haiku-4-5", api_key=INVALID_UPSTREAM_API_KEY),
LiteLLMParamsBody(model=UNREACHABLE_ANTHROPIC_BACKEND, api_key=INVALID_UPSTREAM_API_KEY),
)
resources.defer(lambda: client.delete_model(model_id))
key = client.key_with_alias(

View file

@ -28,7 +28,7 @@ import base64
import json
import os
import time
from dataclasses import dataclass
from dataclasses import dataclass, field
from itertools import count, takewhile
from typing import Final
@ -37,6 +37,7 @@ from pydantic import BaseModel, ConfigDict, Field
from e2e_config import POLL_INTERVAL, POLL_TIMEOUT
from e2e_http import URL, AuthHeaders, send
from e2e_metadata import step
_WEAVE_TRACE_API: Final = "https://trace.wandb.ai"
@ -201,7 +202,7 @@ class WeaveCall(BaseModel):
@dataclass(frozen=True, slots=True)
class WeaveReader:
project_id: str
api_key: str
api_key: str = field(repr=False)
@property
def _headers(self) -> AuthHeaders:
@ -232,6 +233,7 @@ class WeaveReader:
)
return tuple(WeaveCall.model_validate_json(line) for line in outcome.body.splitlines() if line.strip())
@step("Read the Weave {op} calls carrying the marker {marker}")
def calls_matching(self, marker: str, *, since: float, op: str = LITELLM_REQUEST_OP) -> tuple[WeaveCall, ...]:
"""Every call under ``op`` started after ``since`` whose inputs carry
``marker``, paging until the window is exhausted.
@ -247,6 +249,7 @@ class WeaveReader:
)
return tuple(call for page in pages for call in page if call.mentions(marker))
@step("Wait for Weave to ingest a {op} call carrying the marker {marker}, then watch for duplicates")
def poll_calls_matching(self, marker: str, *, since: float, op: str = LITELLM_REQUEST_OP) -> tuple[WeaveCall, ...]:
"""Poll until the call is readable, then keep re-reading for
WEAVE_SETTLE_SECONDS so a duplicate exported by a later batch flush