test(integration): audit cells for responses guardrail block contract

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-23 09:50:02 +00:00
parent ee71241d53
commit 1a4c94a268
2 changed files with 659 additions and 16 deletions

View file

@ -2033,6 +2033,66 @@
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_returns_json_with_typed_message_item": [
"other.observability.guardrails.responses_pre_call_denial_returns_typed_message"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_openai_sdk_streams_typed_message": [
"other.observability.guardrails.responses_pre_call_denial_openai_sdk_streams_typed_message"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_openai_async_sdk_streams_typed_message": [
"other.observability.guardrails.responses_pre_call_denial_openai_async_sdk_streams_typed_message"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_openai_sdk_returns_typed_message": [
"other.observability.guardrails.responses_pre_call_denial_openai_sdk_returns_typed_message"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_reports_zero_usage": [
"other.observability.guardrails.responses_pre_call_denial_reports_zero_usage"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_string_true_returns_json": [
"other.observability.guardrails.responses_pre_call_denial_stream_string_true_returns_json"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_event_vocabulary": [
"other.observability.guardrails.responses_pre_call_denial_stream_event_vocabulary"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_large_denial_text": [
"other.observability.guardrails.responses_pre_call_denial_stream_large_denial_text"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_requests_have_distinct_ids": [
"other.observability.guardrails.responses_pre_call_denial_stream_requests_have_distinct_ids"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_post_call_pipeline_denial_streams_real_usage": [
"other.observability.guardrails.responses_post_call_pipeline_denial_streams_real_usage"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_post_call_pipeline_denial_returns_real_usage": [
"other.observability.guardrails.responses_post_call_pipeline_denial_returns_real_usage"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_denial_requires_authentication": [
"other.observability.guardrails.responses_denial_requires_authentication"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_does_not_reach_upstream": [
"other.observability.guardrails.responses_pre_call_denial_stream_does_not_reach_upstream"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_unguarded_stream_reaches_upstream": [
"other.observability.guardrails.responses_unguarded_stream_reaches_upstream"
],
"tests/integration/observability/test_guardrail_effects.py::test_chat_pre_call_denial_streams_content_filter": [
"other.observability.guardrails.chat_pre_call_denial_streams_content_filter"
],
"tests/integration/observability/test_guardrail_effects.py::test_chat_pre_call_denial_returns_content_filter": [
"other.observability.guardrails.chat_pre_call_denial_returns_content_filter"
],
"tests/integration/observability/test_guardrail_effects.py::test_messages_pre_call_denial_returns_message": [
"other.observability.guardrails.messages_pre_call_denial_returns_message"
],
"tests/integration/observability/test_guardrail_effects.py::test_messages_pre_call_denial_streams_message": [
"other.observability.guardrails.messages_pre_call_denial_streams_message"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_writes_zero_spend_row": [
"other.observability.guardrails.responses_pre_call_denial_writes_zero_spend_row"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_survives_worker_burst": [
"other.observability.guardrails.responses_pre_call_denial_stream_survives_worker_burst"
],
"tests/integration/observability/test_guardrail_effects.py::test_responses_pre_call_denial_stream_survives_worker_kill": [
"other.observability.guardrails.responses_pre_call_denial_stream_survives_worker_kill"
]
},
"browser": {

View file

@ -1,15 +1,21 @@
import json
import os
import signal
import socket
import uuid
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from typing import Final
import httpx
import pytest
import yaml
from integration._support.client import Gateway
from integration._support.client import Gateway, eventually
from integration._support.database import read_rows
from integration._support.mcp import mcp_peer, register_mcp, tool_names
from integration._support.process import owned_proxy
from integration._support.process import group_members, owned_proxy, owned_proxy_process
from integration._support.wire import Reply, Request, wire_server
from openai import AsyncOpenAI, OpenAI
@pytest.mark.covers("other.observability.guardrails.rewrite_reaches_correct_anthropic_positions")
@ -314,21 +320,21 @@ def test_request_selected_mcp_guardrail_blocks_direct_and_virtual_calls(gateway:
_RESPONSES_DENIAL: Final = "This model is not currently available."
def _responses_denial_config(tmp_path: Path, identity: str) -> Path:
def _deny_guardrail(name: str, denial: str = _RESPONSES_DENIAL) -> dict[str, object]:
return {
"guardrail_name": name,
"litellm_params": {
"guardrail": "custom_code",
"mode": "pre_call",
"default_on": False,
"custom_code": (f"def apply_guardrail(inputs, request_data, input_type):\n return block({denial!r})\n"),
},
}
def _responses_denial_config(tmp_path: Path, identity: str, denial: str = _RESPONSES_DENIAL) -> Path:
config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())
config["guardrails"] = [
{
"guardrail_name": identity,
"litellm_params": {
"guardrail": "custom_code",
"mode": "pre_call",
"default_on": False,
"custom_code": (
f"def apply_guardrail(inputs, request_data, input_type):\n return block({_RESPONSES_DENIAL!r})\n"
),
},
}
]
config["guardrails"] = [_deny_guardrail(identity, denial)]
path: Final = tmp_path / "responses-deny.yaml"
path.write_text(yaml.safe_dump(config))
return path
@ -406,3 +412,580 @@ def test_responses_pre_call_denial_returns_json_with_typed_message_item(gateway:
assert len(body["output"]) == 1, body
_assert_blocked_message_item(body["output"][0], body)
assert observed.get("/__observations").json()["requests"] == []
_RESPONSES_OUTPUT_DENIAL: Final = "Output withheld by policy."
_UPSTREAM_INPUT_TOKENS: Final = 20
_UPSTREAM_OUTPUT_TOKENS: Final = 20
_UPSTREAM_TOTAL_TOKENS: Final = 40
def _responses_output_denial_config(tmp_path: Path, identity: str, model: str) -> Path:
config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())
config["guardrails"] = [
{
"guardrail_name": identity,
"litellm_params": {
"guardrail": "custom_code",
"mode": "post_call",
"default_on": False,
"custom_code": (
"def apply_guardrail(inputs, request_data, input_type):\n"
f" return block({_RESPONSES_OUTPUT_DENIAL!r})\n"
),
},
}
]
config["policies"] = {
f"{identity}-pipeline": {
"guardrails": {"add": [identity]},
"pipeline": {
"mode": "post_call",
"steps": [
{
"guardrail": identity,
"on_pass": "allow",
"on_fail": "modify_response",
"modify_response_message": _RESPONSES_OUTPUT_DENIAL,
}
],
},
}
}
config["policy_attachments"] = [{"policy": f"{identity}-pipeline", "models": [model]}]
path: Final = tmp_path / "responses-output-deny.yaml"
path.write_text(yaml.safe_dump(config))
return path
def _blocked_stream_events(text: str) -> tuple[dict[str, object], ...]:
lines: Final = tuple(line for line in text.split("\n") if line.startswith("data: "))
assert lines[-1] == "data: [DONE]", text
return tuple(json.loads(line.removeprefix("data: ")) for line in lines[:-1])
def _dead_api_base() -> str:
with socket.socket() as reserve:
reserve.bind(("127.0.0.1", 0))
port: Final = reserve.getsockname()[1]
return f"http://127.0.0.1:{port}/v1"
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_openai_sdk_streams_typed_message")
def test_responses_pre_call_denial_openai_sdk_streams_typed_message(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
client: Final = OpenAI(
base_url=f"{candidate.client.base_url}/v1", api_key=candidate.key, max_retries=0, timeout=15
)
events: Final = tuple(
client.responses.create(model=model, input="say hi", stream=True, extra_body={"guardrails": [identity]})
)
assert events[-1].type == "response.completed", [event.type for event in events]
completed: Final = events[-1].response
assert completed is not None and len(completed.output) == 1, completed
item: Final = completed.output[0]
assert item.type == "message", item
assert item.role == "assistant" and item.status == "completed", item
assert item.content[0].type == "output_text" and item.content[0].text == _RESPONSES_DENIAL, item.content
assert completed.usage is not None and completed.usage.total_tokens == 0, completed.usage
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_openai_async_sdk_streams_typed_message")
async def test_responses_pre_call_denial_openai_async_sdk_streams_typed_message(
gateway: Gateway, tmp_path: Path
) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
client: Final = AsyncOpenAI(
base_url=f"{candidate.client.base_url}/v1", api_key=candidate.key, max_retries=0, timeout=15
)
stream: Final = await client.responses.create(
model=model, input="say hi", stream=True, extra_body={"guardrails": [identity]}
)
kinds: Final = [event.type async for event in stream]
assert kinds[-1] == "response.completed", kinds
assert "response.output_text.delta" in kinds, kinds
assert "response.in_progress" in kinds, kinds
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_openai_sdk_returns_typed_message")
def test_responses_pre_call_denial_openai_sdk_returns_typed_message(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
client: Final = OpenAI(
base_url=f"{candidate.client.base_url}/v1", api_key=candidate.key, max_retries=0, timeout=15
)
body: Final = client.responses.create(model=model, input="say hi", extra_body={"guardrails": [identity]})
assert body.object == "response" and body.status == "completed", body
assert len(body.output) == 1, body.output
item: Final = body.output[0]
assert item.type == "message" and item.role == "assistant", item
assert item.content[0].type == "output_text" and item.content[0].text == _RESPONSES_DENIAL, item.content
assert body.output_text == _RESPONSES_DENIAL, body
assert body.usage is not None and body.usage.total_tokens == 0, body.usage
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_reports_zero_usage")
def test_responses_pre_call_denial_reports_zero_usage(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST", "/v1/responses", {"model": model, "input": "say hi", "stream": False, "guardrails": [identity]}
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("application/json"), response.text
body: Final = response.json()
_assert_blocked_message_item(body["output"][0], body)
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_string_true_returns_json")
def test_responses_pre_call_denial_stream_string_true_returns_json(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST", "/v1/responses", {"model": model, "input": "say hi", "stream": "true", "guardrails": [identity]}
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("application/json"), (
response.headers["content-type"],
response.text,
)
body: Final = response.json()
_assert_blocked_message_item(body["output"][0], body)
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_event_vocabulary")
def test_responses_pre_call_denial_stream_event_vocabulary(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
identity: Final = "guardrail" + uuid.uuid4().hex
second: Final = "guardrail-2-" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
loaded: Final = yaml.safe_load(config.read_text())
loaded["guardrails"].append(_deny_guardrail(second))
config.write_text(yaml.safe_dump(loaded))
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
with httpx.Client(base_url=gateway.upstream_url, timeout=5, trust_env=False) as observed:
observed.get("/__observations")
response: Final = candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": "say hi", "stream": True, "guardrails": [identity, second]},
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
events: Final = _blocked_stream_events(response.text)
kinds: Final = {event["type"] for event in events}
assert kinds == {
"response.created",
"response.in_progress",
"response.output_item.added",
"response.content_part.added",
"response.output_text.delta",
"response.output_text.done",
"response.content_part.done",
"response.output_item.done",
"response.completed",
}, kinds
item_done: Final = tuple(event for event in events if event["type"] == "response.output_item.done")
assert len(item_done) == 1, events
assert len(events[-1]["response"]["output"]) == 1, events[-1]
assert observed.get("/__observations").json()["requests"] == []
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_large_denial_text")
def test_responses_pre_call_denial_stream_large_denial_text(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
denial: Final = ("Denied: " + "mixed ascii and unicode text " * 200 + "fin")[:5000]
config: Final = _responses_denial_config(tmp_path, identity, denial)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": "say hi", "stream": True, "guardrails": [identity]},
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
events: Final = _blocked_stream_events(response.text)
assert "".join(event["delta"] for event in events if event["type"] == "response.output_text.delta") == denial
done: Final = next(event for event in events if event["type"] == "response.output_text.done")
assert done["text"] == denial, done
completed: Final = events[-1]["response"]
assert completed["output"][0]["content"][0]["text"] == denial, completed
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_requests_have_distinct_ids")
def test_responses_pre_call_denial_stream_requests_have_distinct_ids(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
with httpx.Client(base_url=gateway.upstream_url, timeout=5, trust_env=False) as observed:
observed.get("/__observations")
responses: Final = tuple(
candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": "say hi", "stream": True, "guardrails": [identity]},
)
for _ in range(2)
)
completed: Final = tuple(_blocked_stream_events(response.text)[-1]["response"] for response in responses)
for response in responses:
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
assert completed[0]["id"] != completed[1]["id"], completed
assert completed[0]["output"][0]["id"] != completed[1]["output"][0]["id"], completed
assert observed.get("/__observations").json()["requests"] == []
def _register_named_model(candidate: Gateway, name: str, api_base: str | None = None, **parameters: object) -> str:
created: Final = candidate.post(
"/model/new",
{
"model_name": name,
"litellm_params": {
"model": "openai/gpt-4o-mini",
"api_key": "integration-provider-key",
"api_base": api_base or f"{candidate.upstream_url}/v1",
**parameters,
},
},
)
return str(created["model_info"]["id"])
@pytest.mark.covers("other.observability.guardrails.responses_post_call_pipeline_denial_streams_real_usage")
def test_responses_post_call_pipeline_denial_streams_real_usage(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
model: Final = f"integration-{uuid.uuid4().hex}"
config: Final = _responses_output_denial_config(tmp_path, identity, model)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate:
model_id: Final = _register_named_model(candidate, model, use_chat_completions_api=True)
try:
response: Final = candidate.request(
"POST", "/v1/responses", {"model": model, "input": f"say hi {uuid.uuid4().hex}", "stream": True}
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
events: Final = _blocked_stream_events(response.text)
assert events[-1]["type"] == "response.completed", events
completed: Final = events[-1]["response"]
item: Final = completed["output"][0]
assert item["type"] == "message" and item["role"] == "assistant", item
assert item["content"][0]["type"] == "output_text", item
assert item["content"][0]["text"] == _RESPONSES_OUTPUT_DENIAL, item
usage: Final = completed["usage"]
assert (
usage["input_tokens"],
usage["output_tokens"],
usage["total_tokens"],
) == (_UPSTREAM_INPUT_TOKENS, _UPSTREAM_OUTPUT_TOKENS, _UPSTREAM_TOTAL_TOKENS), usage
finally:
candidate.post("/model/delete", {"id": model_id})
@pytest.mark.covers("other.observability.guardrails.responses_post_call_pipeline_denial_returns_real_usage")
def test_responses_post_call_pipeline_denial_returns_real_usage(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
model: Final = f"integration-{uuid.uuid4().hex}"
config: Final = _responses_output_denial_config(tmp_path, identity, model)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate:
model_id: Final = _register_named_model(candidate, model, use_chat_completions_api=True)
try:
response: Final = candidate.request(
"POST", "/v1/responses", {"model": model, "input": f"say hi {uuid.uuid4().hex}"}
)
assert response.status_code == 200, response.text
body: Final = response.json()
item: Final = body["output"][0]
assert item["type"] == "message" and item["role"] == "assistant", item
assert item["content"][0]["type"] == "output_text", item
assert item["content"][0]["text"] == _RESPONSES_OUTPUT_DENIAL, item
usage: Final = body["usage"]
assert (
usage["input_tokens"],
usage["output_tokens"],
usage["total_tokens"],
) == (_UPSTREAM_INPUT_TOKENS, _UPSTREAM_OUTPUT_TOKENS, _UPSTREAM_TOTAL_TOKENS), usage
finally:
candidate.post("/model/delete", {"id": model_id})
@pytest.mark.covers("other.observability.guardrails.responses_denial_requires_authentication")
def test_responses_denial_requires_authentication(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST", "/v1/responses", {"model": model, "input": "say hi", "guardrails": [identity]}, key="sk-invalid"
)
assert response.status_code == 401, (response.status_code, response.text)
assert response.json()["error"]["type"] == "token_not_found_in_db", response.text
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_does_not_reach_upstream")
def test_responses_pre_call_denial_stream_does_not_reach_upstream(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model(api_base=_dead_api_base())
response: Final = candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": "say hi", "stream": True, "guardrails": [identity]},
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
events: Final = _blocked_stream_events(response.text)
completed: Final = events[-1]["response"]
_assert_blocked_message_item(completed["output"][0], completed)
@pytest.mark.covers("other.observability.guardrails.responses_unguarded_stream_reaches_upstream")
def test_responses_unguarded_stream_reaches_upstream(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
dead: Final = scenario.model(api_base=_dead_api_base())
denied: Final = candidate.request(
"POST",
"/v1/responses",
{"model": dead, "input": "say hi", "stream": True, "guardrails": [identity]},
)
assert denied.status_code == 200, denied.text
model: Final = scenario.model(use_chat_completions_api=True)
with httpx.Client(base_url=gateway.upstream_url, timeout=5, trust_env=False) as observed:
observed.get("/__observations")
response: Final = candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": f"say hi {uuid.uuid4().hex}", "stream": True, "guardrails": []},
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
assert "response.completed" in response.text, response.text
batches: Final[list[dict[str, object]]] = []
def _drain() -> list[dict[str, object]]:
batches.extend(observed.get("/__observations").json()["requests"])
return batches
eventually(_drain, lambda values: len(values) >= 1, seconds=30, return_last_on_timeout=True)
assert len(batches) == 1, batches
assert batches[0]["path"] == "/v1/chat/completions", batches
@pytest.mark.covers("other.observability.guardrails.chat_pre_call_denial_streams_content_filter")
def test_chat_pre_call_denial_streams_content_filter(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": "say hi"}],
"stream": True,
"guardrails": [identity],
},
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
lines: Final = tuple(line for line in response.text.split("\n") if line.startswith("data: "))
assert lines[-1] == "data: [DONE]", response.text
chunks: Final = tuple(json.loads(line.removeprefix("data: ")) for line in lines[:-1])
assert chunks[0]["choices"][0]["delta"]["content"] == _RESPONSES_DENIAL, chunks
assert chunks[-1]["choices"][0]["finish_reason"] == "stop", chunks
@pytest.mark.covers("other.observability.guardrails.chat_pre_call_denial_returns_content_filter")
def test_chat_pre_call_denial_returns_content_filter(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST",
"/v1/chat/completions",
{
"model": model,
"messages": [{"role": "user", "content": "say hi"}],
"guardrails": [identity],
},
)
assert response.status_code == 200, response.text
body: Final = response.json()
choice: Final = body["choices"][0]
assert choice["finish_reason"] == "content_filter", body
assert choice["message"]["content"] == _RESPONSES_DENIAL, body
assert (
body["usage"]["prompt_tokens"],
body["usage"]["completion_tokens"],
body["usage"]["total_tokens"],
) == (0, 0, 0), body["usage"]
@pytest.mark.covers("other.observability.guardrails.messages_pre_call_denial_returns_message")
def test_messages_pre_call_denial_returns_message(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST",
"/v1/messages",
{
"model": model,
"messages": [{"role": "user", "content": "say hi"}],
"max_tokens": 16,
"guardrails": [identity],
},
)
assert response.status_code == 200, response.text
body: Final = response.json()
assert body["type"] == "message" and body["role"] == "assistant", body
assert body["content"] == [{"type": "text", "text": _RESPONSES_DENIAL}], body
assert body["stop_reason"] == "end_turn", body
assert (body["usage"]["input_tokens"], body["usage"]["output_tokens"]) == (0, 0), body
@pytest.mark.covers("other.observability.guardrails.messages_pre_call_denial_streams_message")
def test_messages_pre_call_denial_streams_message(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST",
"/v1/messages",
{
"model": model,
"messages": [{"role": "user", "content": "say hi"}],
"max_tokens": 16,
"stream": True,
"guardrails": [identity],
},
)
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text
lines: Final = tuple(line for line in response.text.split("\n") if line.startswith("data: "))
assert len(lines) == 1, response.text
body: Final = json.loads(lines[0].removeprefix("data: "))
assert body["type"] == "message" and body["role"] == "assistant", body
assert body["content"] == [{"type": "text", "text": _RESPONSES_DENIAL}], body
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_writes_zero_spend_row")
def test_responses_pre_call_denial_writes_zero_spend_row(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model()
response: Final = candidate.request(
"POST", "/v1/responses", {"model": model, "input": "say hi", "guardrails": [identity]}
)
assert response.status_code == 200, response.text
rows: Final = eventually(
lambda: read_rows(
'SELECT spend, total_tokens FROM "LiteLLM_SpendLogs" WHERE model=%s AND call_type=%s',
(model, "aresponses"),
),
lambda values: len(values) == 1,
seconds=70,
)
assert float(rows[0]["spend"]) == 0, rows
assert rows[0]["total_tokens"] == 0, rows
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_survives_worker_burst")
def test_responses_pre_call_denial_stream_survives_worker_burst(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy(gateway, tmp_path, {}, config=config, workers=2) as candidate, candidate.scenario() as scenario:
model: Final = scenario.model(api_base=_dead_api_base())
healthy: Final = scenario.model(use_chat_completions_api=True)
with httpx.Client(base_url=gateway.upstream_url, timeout=5, trust_env=False) as observed:
observed.get("/__observations")
def burst(index: int) -> httpx.Response:
if index % 3 == 0:
return candidate.request(
"POST",
"/v1/responses",
{"model": healthy, "input": f"say hi {uuid.uuid4().hex} {index}", "stream": True},
)
stream: Final = index % 3 == 1
return candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": f"say hi {index}", "stream": stream, "guardrails": [identity]},
)
with ThreadPoolExecutor(max_workers=8) as pool:
responses: Final = tuple(pool.map(burst, range(30)))
assert {response.status_code for response in responses} == {200}
response_ids: Final = set()
for index, response in enumerate(responses):
assert response.status_code == 200, (index, response.text)
if index % 3 == 0:
assert response.headers["content-type"].startswith("text/event-stream"), response.text
completed: Final = _blocked_stream_events(response.text)[-1]["response"]
response_ids.add(completed["id"])
elif index % 3 == 1:
assert response.headers["content-type"].startswith("text/event-stream"), response.text
blocked: Final = _blocked_stream_events(response.text)[-1]["response"]
_assert_blocked_message_item(blocked["output"][0], blocked)
response_ids.add(blocked["id"])
else:
assert response.headers["content-type"].startswith("application/json"), response.text
body: Final = response.json()
_assert_blocked_message_item(body["output"][0], body)
response_ids.add(body["id"])
assert len(response_ids) == 30, response_ids
assert len(observed.get("/__observations").json()["requests"]) == 10
@pytest.mark.covers("other.observability.guardrails.responses_pre_call_denial_stream_survives_worker_kill")
def test_responses_pre_call_denial_stream_survives_worker_kill(gateway: Gateway, tmp_path: Path) -> None:
identity: Final = "guardrail" + uuid.uuid4().hex
config: Final = _responses_denial_config(tmp_path, identity)
with owned_proxy_process(gateway, tmp_path, {}, config=config, workers=2) as owned:
candidate: Final = owned.gateway
with candidate.scenario() as scenario:
model: Final = scenario.model(api_base=_dead_api_base())
children: Final = tuple(
member.pid for member in group_members(owned.process.pid) if member.pid != owned.process.pid
)
assert len(children) >= 2, children
os.kill(children[0], signal.SIGKILL)
eventually(lambda: group_members(owned.process.pid), lambda members: all(m.is_running() for m in members))
def burst(index: int) -> httpx.Response:
return candidate.request(
"POST",
"/v1/responses",
{"model": model, "input": f"say hi {index}", "stream": True, "guardrails": [identity]},
)
with ThreadPoolExecutor(max_workers=5) as pool:
responses: Final = tuple(pool.map(burst, range(10)))
for response in responses:
assert response.status_code == 200, response.text
assert response.headers["content-type"].startswith("text/event-stream"), response.text