test(integration): callback credential canary slots C1-C3 and D5 (#43630)

* test(integration): credential canary suite harness

Adds tests/integration/security with canary generation and search, sweeps over the database, GET routes, client responses, sink doubles and Redis, an owned proxy rig, a sweep sensitivity self-test and the config deployment api_key slot. Registers the security group in run.py, the manifest and the CircleCI integration matrix.

* test(integration): widen canary route sweep and harden the rig

Enumerate lazily registered feature routers, call parameterized routes with placeholder ids, fail on routes that return no response, skip provider pass-through routes, add an explicit admin-only route allowance, let the sink double use a configurable token, inflate gzip members anywhere in a blob, sweep Redis before the route walk, and trap outbound connections from the owned proxy.

* test(integration): descend into any decoded value that can still hold an encoded canary

* test(integration): bound canary decoding by depth and decoded bytes

* test(integration): scope log-table and spend-log reads to the scenario window

* test(integration): sweep spend-log rows in the scenario date window

* test(integration): keep spend-log date window summarized

* test(integration): resolve deployment ids, scope paginated log lists, key allowances by slot

* test(integration): expect 404 from the caller-scoped team membership route

* test(integration): use the rig's own master key and expect 404 from submission lookups

* test(integration): check the overridden rig key without assuming the default key is unknown

* test(integration): callback credential canary slots C1-C3 and D5

Team callback, team callback_settings, config default_team_settings and key metadata.logging Langfuse secrets, a team Datadog dd_api_key, and request-body Langfuse keys (allow_client_side_credentials) must reach only their sink. Each scenario checks its sink received the canary as auth and that the marker is visible at the stored body, the Logs drawer route and the sink. Adds a unit test that the stored request body snapshot carries no callback parameter.

* test(integration): give the callback sink waits a wider bound

* test(integration): sweep provider requests for callback credentials
This commit is contained in:
yucheng-berri 2026-09-29 10:49:30 -07:00 • committed by GitHub
parent a3552c451b
commit b3dcf8208d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 616 additions and 0 deletions

View file

@ -0,0 +1,173 @@
"""Traffic matrix and sink doubles for the callback credential slots.
- ``upstream(request)``: provider double for every endpoint in ``ENDPOINTS``: OpenAI chat (plain
and SSE) and OpenAI Responses (``/v1/messages`` reaches it as chat). A body carrying
``PROVIDER_4XX`` gets HTTP 400 and one carrying ``PROVIDER_5XX`` gets HTTP 500. The sensitivity
marker found in the body is echoed back.
- ``langfuse_sink`` / ``datadog_sink``: Langfuse OTLP ingest and Datadog intake doubles.
- ``send(gateway, key, endpoint, model, text, extra)``: one client call per endpoint.
- ``spend_request_id(marker)``: the spend row written for the request carrying ``marker``.
- ``wait_for_sink(recorder, marker)``: bounded wait until a sink received the marker (gzip aware).
"""
from __future__ import annotations
import json
import re
import uuid
from collections.abc import Mapping
from typing import Final
import httpx
from integration._support.client import Gateway, eventually, string_value
from integration._support.database import read_rows
from integration._support.wire import Reply, Request
from integration.security._canary import Canary, find_canary
from integration.security._sinks import PROVIDER_4XX, Recorder
from pydantic import JsonValue
PROVIDER_5XX: Final = "canary-provider-5xx"
ENDPOINTS: Final = ("chat", "chat_stream", "messages", "responses")
OUTCOMES: Final = ("success", "provider_4xx", "provider_5xx")
EXPECTED_STATUS: Final = {"success": 200, "provider_4xx": 400, "provider_5xx": 500}
LANGFUSE_PUBLIC_KEY: Final = "pk-lf-canary-public"
_MARKER: Final = re.compile(rb"lkc-M0-[0-9a-f]{32}")
def _echo(body: bytes) -> str:
found: Final = _MARKER.search(body)
return "echo " + (found.group().decode() if found else "none")
def _failure(body: bytes) -> Reply | None:
if PROVIDER_4XX.encode() in body:
return Reply(
status=400,
body=b'{"error":{"type":"invalid_request_error","code":"canary_rejected","message":"rejected"}}',
)
if PROVIDER_5XX.encode() in body:
return Reply(status=500, body=b'{"error":{"type":"server_error","message":"upstream exploded"}}')
return None
def _chat(body: Mapping[str, JsonValue], text: str) -> Reply:
identity: Final = f"chatcmpl-{uuid.uuid4().hex}"
usage: Final = {"prompt_tokens": 7, "completion_tokens": 3, "total_tokens": 10}
if body.get("stream") is True:
chunks: Final = (
{"choices": [{"index": 0, "delta": {"role": "assistant", "content": text}, "finish_reason": None}]},
{"choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}]},
{"choices": [], "usage": usage},
)
events: Final = b"".join(
b"data: "
+ json.dumps(
{"id": identity, "object": "chat.completion.chunk", "created": 1, "model": "gpt-4o-mini", **chunk}
).encode()
+ b"\n\n"
for chunk in chunks
)
return Reply(body=events + b"data: [DONE]\n\n", content_type="text/event-stream")
return Reply(
body=json.dumps(
{
"id": identity,
"object": "chat.completion",
"created": 1,
"model": "gpt-4o-mini",
"choices": [{"index": 0, "message": {"role": "assistant", "content": text}, "finish_reason": "stop"}],
"usage": usage,
}
).encode()
)
def _responses(text: str) -> Reply:
return Reply(
body=json.dumps(
{
"id": f"resp_{uuid.uuid4().hex}",
"object": "response",
"created_at": 1,
"status": "completed",
"model": "gpt-4o-mini",
"output": [
{
"type": "message",
"id": f"msg_{uuid.uuid4().hex}",
"status": "completed",
"role": "assistant",
"content": [{"type": "output_text", "text": text, "annotations": []}],
}
],
"parallel_tool_calls": True,
"tool_choice": "auto",
"tools": [],
"usage": {"input_tokens": 7, "output_tokens": 3, "total_tokens": 10},
}
).encode()
)
def upstream(request: Request) -> Reply:
failure: Final = _failure(request.body)
if failure is not None:
return failure
text: Final = _echo(request.body)
if request.target.split("?", 1)[0].endswith("/responses"):
return _responses(text)
return _chat(json.loads(request.body or b"{}"), text)
def langfuse_sink(request: Request) -> Reply:
if request.method == "GET" and request.target.startswith("/api/public/projects"):
return Reply(body=b'{"data":[{"id":"canary-project","name":"canary"}]}')
return Reply(body=b"", content_type="application/x-protobuf")
def datadog_sink(request: Request) -> Reply:
return Reply(status=202, body=b"{}")
def body_for(endpoint: str, model: str, text: str) -> dict[str, JsonValue]:
if endpoint == "responses":
return {"model": model, "input": text}
if endpoint == "messages":
return {"model": model, "max_tokens": 16, "messages": [{"role": "user", "content": text}]}
return {
"model": model,
"messages": [{"role": "user", "content": text}],
**({"stream": True, "stream_options": {"include_usage": True}} if endpoint == "chat_stream" else {}),
}
def send(
gateway: Gateway, key: str, endpoint: str, model: str, text: str, extra: Mapping[str, JsonValue] | None = None
) -> httpx.Response:
path: Final = {"responses": "/v1/responses", "messages": "/v1/messages"}.get(endpoint, "/v1/chat/completions")
return gateway.request("POST", path, {**body_for(endpoint, model, text), **(extra or {})}, key=key)
def outcome_text(slot: str, marker: Canary, outcome: str) -> str:
trigger: Final = {"success": "", "provider_4xx": f" {PROVIDER_4XX}", "provider_5xx": f" {PROVIDER_5XX}"}[outcome]
return f"slot {slot} {marker.value}{trigger}"
def spend_request_id(marker: Canary) -> str:
rows: Final = eventually(
lambda: read_rows(
'SELECT request_id FROM "LiteLLM_SpendLogs" WHERE proxy_server_request::text LIKE %s',
(f"%{marker.core}%",),
),
lambda found: len(found) >= 1,
seconds=70,
)
return string_value(rows[0]["request_id"])
def wait_for_sink(recorder: Recorder, marker: Canary, seconds: float = 90) -> tuple[Request, ...]:
return eventually(
lambda: tuple(request for request in recorder.requests() if find_canary(request.body, (marker,))),
bool,
seconds=seconds,
)

View file

@ -83,6 +83,12 @@ SLOTS: Final = MappingProxyType(
"A1": Slot("A1", "Virtual key raw value, set as a custom key through /key/generate", prefix="sk-"),
"A2": Slot("A2", "Proxy master key from the LITELLM_MASTER_KEY environment variable", prefix="sk-"),
"B1": Slot("B1", "Deployment api_key declared in the proxy config.yaml model_list"),
"C1": Slot(
"C1", "Team callback langfuse_secret_key (team callback API, config team settings, callback_settings)"
),
"C2": Slot("C2", "Key-level callback langfuse_secret_key in key metadata.logging"),
"C3": Slot("C3", "Team callback dd_api_key for the Datadog sink"),
"D5": Slot("D5", "Request-supplied langfuse_secret_key in the request body"),
"B2": Slot("B2", "Deployment api_key added through /model/new and stored encrypted"),
"B3": Slot("B3", "Credentials table api_key referenced by a deployment's litellm_credential_name"),
"B4": Slot("B4", "Deployment aws_secret_access_key added through /model/new"),

View file

@ -0,0 +1,401 @@
"""Slots C1, C2, C3 and D5: callback credentials must reach only their sink.
C1 is the team callback ``langfuse_secret_key`` (team callback API, the deprecated team
``metadata.callback_settings`` and the config ``default_team_settings``), C2 the key-level
``metadata.logging`` Langfuse key, C3 a team callback ``dd_api_key`` for Datadog, and D5 a
``langfuse_secret_key`` the caller sends in the request body (``langfuse_host`` in a body is
rejected without an admin opt-in, so D5 runs on its own proxy with
``general_settings.allow_client_side_credentials`` on).
Positive control: the owning sink double must receive the request's marker under an auth
header built from the canary (Langfuse ``Basic pk:sk``, Datadog ``DD-API-KEY``), or the test
fails before sweeping. Sensitivity control: the marker must be seen in the stored request body,
the Logs drawer route and the owning sink. Then no sweep may find the canary anywhere else,
including every request the provider double received (swept as the ``provider`` sink, with no
header allowance; the provider's own key is slot B1, which these tests do not search for).
"""
from __future__ import annotations
import base64
from collections.abc import Callable, Iterator, Mapping
from contextlib import contextmanager
from dataclasses import dataclass
from datetime import UTC, datetime
from pathlib import Path
from typing import Final
from urllib.parse import quote
import pytest
from integration._support.client import Scenario
from integration._support.wire import Request, wire_server
from integration.security._callback_traffic import (
ENDPOINTS,
EXPECTED_STATUS,
LANGFUSE_PUBLIC_KEY,
OUTCOMES,
datadog_sink,
langfuse_sink,
outcome_text,
send,
spend_request_id,
upstream,
wait_for_sink,
)
from integration.security._canary import MARKER, Canary, canary, find_canary
from integration.security._sinks import CONFIG_MODEL, GENERIC_SINK, Caller, Recorder, Rig, canary_rig
from integration.security._sweeps import assert_marker_seen, assert_no_hits, record_route_sweep, sweep_all
from pydantic import JsonValue
LANGFUSE: Final = "langfuse"
DATADOG: Final = "datadog"
PROVIDER: Final = "provider"
BOTH: Final = "success_and_failure"
@dataclass(frozen=True, slots=True)
class CallbackRig:
rig: Rig
langfuse: Recorder
datadog: Recorder
def sinks(self) -> dict[str, tuple[Request, ...]]:
return {
**{name: sink.requests() for name, sink in self.rig.sinks.items()},
LANGFUSE: self.langfuse.requests(),
DATADOG: self.datadog.requests(),
PROVIDER: self.rig.provider.requests(),
}
def datadog_port(self) -> str:
return self.datadog.url.rsplit(":", 1)[1]
@contextmanager
def callback_rig(
root: Path, configure: Callable[[dict[str, object], str, str], None] | None = None
) -> Iterator[CallbackRig]:
with (
wire_server(langfuse_sink) as langfuse,
wire_server(datadog_sink) as datadog,
canary_rig(
root,
configure=(lambda config, provider: configure(config, provider, langfuse.url)) if configure else None,
environment={"LANGFUSE_FLUSH_INTERVAL": "1"},
upstream=upstream,
) as rig,
):
yield CallbackRig(rig, Recorder(langfuse), Recorder(datadog))
def _allow_client_side_credentials(config: dict[str, object], _provider: str, _langfuse: str) -> None:
settings: Final = config["general_settings"]
assert isinstance(settings, dict)
settings["allow_client_side_credentials"] = True
@pytest.fixture(scope="module")
def client_side(tmp_path_factory: pytest.TempPathFactory) -> Iterator[CallbackRig]:
with callback_rig(tmp_path_factory.mktemp("canary-client-side"), _allow_client_side_credentials) as value:
yield value
@pytest.fixture(scope="module")
def shared(tmp_path_factory: pytest.TempPathFactory) -> Iterator[CallbackRig]:
with callback_rig(tmp_path_factory.mktemp("canary-callbacks")) as value:
yield value
def langfuse_vars(secret: Canary, host: str) -> dict[str, JsonValue]:
return {"langfuse_public_key": LANGFUSE_PUBLIC_KEY, "langfuse_secret_key": secret.value, "langfuse_host": host}
def caller(
scenario: Scenario,
*,
team_id: str | None = None,
team_metadata: Mapping[str, JsonValue] | None = None,
key_metadata: Mapping[str, JsonValue] | None = None,
) -> Caller:
team: Final = scenario.team(
**({"team_id": team_id} if team_id else {}), **({"metadata": dict(team_metadata)} if team_metadata else {})
)
user: Final = scenario.member(team)
key: Final = scenario.key(
team_id=team, user_id=user, models=[CONFIG_MODEL], **({"metadata": dict(key_metadata)} if key_metadata else {})
)
return Caller(team, user, key)
def langfuse_control(secret: Canary) -> Callable[[CallbackRig, Canary], None]:
expected: Final = "Basic " + base64.b64encode(f"{LANGFUSE_PUBLIC_KEY}:{secret.value}".encode()).decode()
def check(rig: CallbackRig, marker: Canary) -> None:
delivered: Final = wait_for_sink(rig.langfuse, marker)
assert {request.headers.get("authorization") for request in delivered} == {expected}, (
f"Positive control: the Langfuse double never received the {secret.slot} canary as its Basic auth"
)
return check
def datadog_control(secret: Canary) -> Callable[[CallbackRig, Canary], None]:
def check(rig: CallbackRig, marker: Canary) -> None:
delivered: Final = wait_for_sink(rig.datadog, marker)
assert {request.headers.get("dd-api-key") for request in delivered} == {secret.value}, (
"Positive control: the Datadog double never received the C3 canary as DD-API-KEY"
)
return check
def run_scenario(
cb: CallbackRig,
scenario: Scenario,
who: Caller,
secret: Canary,
endpoint: str,
outcome: str,
*,
control: Callable[[CallbackRig, Canary], None],
sink: str,
own_header: tuple[str, str],
node: str,
extra: Mapping[str, JsonValue] | None = None,
) -> None:
marker: Final = canary(MARKER)
started: Final = datetime.now(UTC)
response: Final = send(
cb.rig.proxy, who.key, endpoint, CONFIG_MODEL, outcome_text(secret.slot, marker, outcome), extra
)
assert response.status_code == EXPECTED_STATUS[outcome], response.text
control(cb, marker)
request_id: Final = spend_request_id(marker)
wait_for_sink(cb.rig.sinks[GENERIC_SINK], marker)
report: Final = sweep_all(
cb.rig.proxy,
(marker, secret),
responses=(response,),
sinks=cb.sinks(),
ids={
"request_id": request_id,
"team_id": who.team_id,
"user_id": who.user_id,
"model_id": cb.rig.model_id,
"model": CONFIG_MODEL,
},
callers=who.callers(cb.rig),
own_headers={**cb.rig.own_headers, sink: own_header},
since=started,
)
record_route_sweep(report.routes, node)
assert_marker_seen(
report,
{
"S1": "LiteLLM_SpendLogs.proxy_server_request",
"S2": f"GET /spend/logs/ui/{quote(request_id, safe='')} as admin -> 200",
"S4": f"{sink}[",
},
)
assert_marker_seen(report, {"S2": f"GET /spend/logs?request_id={quote(request_id, safe='')} as admin -> 200"})
assert_marker_seen(report, {"S4": f"{PROVIDER}["})
assert_no_hits(report.credential_hits(), f"slot {secret.slot}, {endpoint}, {outcome}")
MATRIX: Final = [
pytest.param(endpoint, outcome, id=f"{endpoint}-{outcome}") for endpoint in ENDPOINTS for outcome in OUTCOMES
]
@pytest.mark.timeout(240) # full S1/S2 walk: every table and ~430 GET routes as two callers
@pytest.mark.parametrize(("endpoint", "outcome"), MATRIX)
def test_c1_team_callback_api_langfuse_secret_reaches_only_langfuse(
shared: CallbackRig, endpoint: str, outcome: str, request: pytest.FixtureRequest
) -> None:
secret: Final = canary("C1")
with shared.rig.proxy.scenario() as scenario:
who: Final = caller(scenario)
shared.rig.proxy.post(
f"/team/{who.team_id}/callback",
{
"callback_name": "langfuse",
"callback_type": BOTH,
"callback_vars": langfuse_vars(secret, shared.langfuse.url),
},
)
run_scenario(
shared,
scenario,
who,
secret,
endpoint,
outcome,
control=langfuse_control(secret),
sink=LANGFUSE,
own_header=("authorization", "C1"),
node=request.node.nodeid,
)
@pytest.mark.timeout(240) # full S1/S2 walk: every table and ~430 GET routes as two callers
@pytest.mark.parametrize("endpoint", ENDPOINTS)
def test_c1_deprecated_team_callback_settings_langfuse_secret_reaches_only_langfuse(
shared: CallbackRig, endpoint: str, request: pytest.FixtureRequest
) -> None:
secret: Final = canary("C1")
settings: Final = {
"success_callback": ["langfuse"],
"failure_callback": ["langfuse"],
"callback_vars": langfuse_vars(secret, shared.langfuse.url),
}
with shared.rig.proxy.scenario() as scenario:
who: Final = caller(scenario, team_metadata={"callback_settings": settings})
run_scenario(
shared,
scenario,
who,
secret,
endpoint,
"success",
control=langfuse_control(secret),
sink=LANGFUSE,
own_header=("authorization", "C1"),
node=request.node.nodeid,
)
@pytest.mark.timeout(240) # full S1/S2 walk: every table and ~430 GET routes as two callers
@pytest.mark.parametrize("endpoint", ENDPOINTS)
def test_c1_config_default_team_settings_langfuse_secret_reaches_only_langfuse(
tmp_path: Path, endpoint: str, request: pytest.FixtureRequest
) -> None:
"""The team callback comes from ``litellm_settings.default_team_settings`` in config.yaml."""
secret: Final = canary("C1")
team_id: Final = f"canary-config-team-{secret.core[:12]}"
def configure(config: dict[str, object], _provider: str, langfuse_url: str) -> None:
settings: Final = config["litellm_settings"]
assert isinstance(settings, dict)
settings["default_team_settings"] = [
{
"team_id": team_id,
"success_callback": ["langfuse"],
"failure_callback": ["langfuse"],
"langfuse_public_key": LANGFUSE_PUBLIC_KEY,
"langfuse_secret": secret.value,
"langfuse_host": langfuse_url,
}
]
with callback_rig(tmp_path, configure) as cb, cb.rig.proxy.scenario() as scenario:
who: Final = caller(scenario, team_id=team_id)
run_scenario(
cb,
scenario,
who,
secret,
endpoint,
"success",
control=langfuse_control(secret),
sink=LANGFUSE,
own_header=("authorization", "C1"),
node=request.node.nodeid,
)
@pytest.mark.timeout(240) # full S1/S2 walk: every table and ~430 GET routes as two callers
@pytest.mark.parametrize(("endpoint", "outcome"), MATRIX)
def test_c2_key_logging_langfuse_secret_reaches_only_langfuse(
shared: CallbackRig, endpoint: str, outcome: str, request: pytest.FixtureRequest
) -> None:
secret: Final = canary("C2")
logging: Final = [
{
"callback_name": "langfuse",
"callback_type": BOTH,
"callback_vars": langfuse_vars(secret, shared.langfuse.url),
}
]
with shared.rig.proxy.scenario() as scenario:
who: Final = caller(scenario, key_metadata={"logging": logging})
run_scenario(
shared,
scenario,
who,
secret,
endpoint,
outcome,
control=langfuse_control(secret),
sink=LANGFUSE,
own_header=("authorization", "C2"),
node=request.node.nodeid,
)
@pytest.mark.timeout(240) # full S1/S2 walk: every table and ~430 GET routes as two callers
@pytest.mark.parametrize(("endpoint", "outcome"), MATRIX)
def test_c3_team_callback_datadog_api_key_reaches_only_datadog(
shared: CallbackRig, endpoint: str, outcome: str, request: pytest.FixtureRequest
) -> None:
secret: Final = canary("C3")
with shared.rig.proxy.scenario() as scenario:
who: Final = caller(scenario)
shared.rig.proxy.post(
f"/team/{who.team_id}/callback",
{
"callback_name": "datadog",
"callback_type": BOTH,
"callback_vars": {
"dd_api_key": secret.value,
"dd_agent_host": "127.0.0.1",
"dd_agent_port": shared.datadog_port(),
},
},
)
run_scenario(
shared,
scenario,
who,
secret,
endpoint,
outcome,
control=datadog_control(secret),
sink=DATADOG,
own_header=("dd-api-key", "C3"),
node=request.node.nodeid,
)
@pytest.mark.timeout(240) # full S1/S2 walk: every table and ~430 GET routes as two callers
@pytest.mark.parametrize(("endpoint", "outcome"), MATRIX)
def test_d5_request_body_langfuse_secret_reaches_only_langfuse(
client_side: CallbackRig, endpoint: str, outcome: str, request: pytest.FixtureRequest
) -> None:
secret: Final = canary("D5")
with client_side.rig.proxy.scenario() as scenario:
who: Final = caller(scenario)
run_scenario(
client_side,
scenario,
who,
secret,
endpoint,
outcome,
control=langfuse_control(secret),
sink=LANGFUSE,
own_header=("authorization", "D5"),
node=request.node.nodeid,
extra={
**langfuse_vars(secret, client_side.langfuse.url),
"success_callback": ["langfuse"],
"failure_callback": ["langfuse"],
},
)
def test_find_canary_sees_the_langfuse_basic_auth_header() -> None:
"""The Langfuse positive control and own-header rule depend on decoding ``Basic pk:sk``."""
secret: Final = canary("C1")
header: Final = "Basic " + base64.b64encode(f"{LANGFUSE_PUBLIC_KEY}:{secret.value}".encode()).decode()
assert [match.slot for match in find_canary(header, (secret,))] == ["C1"]

View file

@ -0,0 +1,36 @@
"""The stored request body never carries callback parameters.
Every ``StandardCallbackDynamicParams`` key and ``litellm_trusted_callback_vars`` is set on the
request dict with a unique value, the body snapshot is refreshed, and none of the keys or values
may be in ``proxy_server_request["body"]``. A control key proves the snapshot was rebuilt.
"""
from __future__ import annotations
import json
import uuid
from typing import Final
from litellm.proxy.litellm_pre_call_utils import refresh_proxy_server_request_body_snapshot
from litellm.types.utils import TRUSTED_CALLBACK_VARS_FIELD, StandardCallbackDynamicParams
def test_body_snapshot_excludes_every_callback_dynamic_param_and_the_trusted_vars() -> None:
core: Final = uuid.uuid4().hex
params: Final = {name: f"lkc-{name}-{core}" for name in StandardCallbackDynamicParams.__annotations__}
control: Final = f"control-{uuid.uuid4().hex}"
data: Final = {
"model": "gpt-4o-mini",
"messages": [{"role": "user", "content": control}],
**params,
TRUSTED_CALLBACK_VARS_FIELD: dict(params),
"proxy_server_request": {"url": "http://proxy/v1/chat/completions", "body": {}},
}
refresh_proxy_server_request_body_snapshot(data)
body: Final = data["proxy_server_request"]["body"]
assert control in json.dumps(body), "Sensitivity control: the snapshot was not rebuilt from the request"
present: Final = sorted({*params, TRUSTED_CALLBACK_VARS_FIELD} & set(body))
assert present == [], f"Callback parameters copied into the stored request body: {present}"
assert core not in json.dumps(body, default=str)