mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
fix(otel): tolerate non-dict callback_settings.otel and ignore bare EXCLUDED_SERVICES env (#44086)
* test(otel): cover non-dict callback_settings.otel and bare EXCLUDED_SERVICES env Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(otel): tolerate non-dict callback_settings.otel and ignore bare EXCLUDED_SERVICES env Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(otel): simplify settings_customise_sources signature Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel): poll for present spans instead of waiting the full window Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel): audit null otel block and bare EXCLUDED_SERVICES across the excluded-services matrix Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel): assert per-trace datastore spans in the unconfigured burst cells Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(otel): keep pydantic-settings runtime options on the OTel v2 env source Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel): require post-auth datastore spans at the tenant on the cache-hit twin Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel): cover pydantic-settings runtime options in callback_settings.otel through the proxy Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: yucheng <yucheng@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
b1e0e9e84b
commit
f95446ea39
6 changed files with 602 additions and 5 deletions
|
|
@ -915,7 +915,8 @@ def _excluded_db_systems(logger: "OpenTelemetryV2") -> frozenset[str]:
|
|||
logger got published: with ``callbacks: [langfuse_otel, otel]`` the ``otel``
|
||||
callback folds into the preset, whose config is env-only.
|
||||
"""
|
||||
configured: Final = litellm.callback_settings.get("otel", {}).get("excluded_services")
|
||||
otel_settings: Final = (litellm.callback_settings or {}).get("otel")
|
||||
configured: Final = otel_settings.get("excluded_services") if isinstance(otel_settings, dict) else None
|
||||
if configured is None:
|
||||
return logger.config.excluded_services
|
||||
return excluded_db_systems_from(configured)
|
||||
|
|
|
|||
|
|
@ -5,7 +5,8 @@ from functools import lru_cache
|
|||
from typing import Annotated, Any, Final
|
||||
|
||||
from pydantic import AliasChoices, BaseModel, Field, TypeAdapter, ValidationError, field_validator, model_validator
|
||||
from pydantic_settings import BaseSettings, NoDecode, SettingsConfigDict
|
||||
from pydantic.fields import FieldInfo
|
||||
from pydantic_settings import BaseSettings, NoDecode, PydanticBaseSettingsSource, SettingsConfigDict
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.model.baggage import (
|
||||
|
|
@ -121,9 +122,37 @@ class ExporterSpec(BaseModel):
|
|||
)
|
||||
|
||||
|
||||
class _EnvWithoutBareExcludedServices(PydanticBaseSettingsSource):
|
||||
def __init__(self, settings_cls: type[BaseSettings], env_settings: PydanticBaseSettingsSource) -> None:
|
||||
super().__init__(settings_cls)
|
||||
self._env_settings: Final = env_settings
|
||||
|
||||
def get_field_value(self, field: FieldInfo, field_name: str) -> tuple[object, str, bool]:
|
||||
return self._env_settings.get_field_value(field, field_name)
|
||||
|
||||
def __call__(self) -> dict[str, object]:
|
||||
return {key: value for key, value in self._env_settings().items() if key != "excluded_services"}
|
||||
|
||||
|
||||
class OpenTelemetryV2Config(BaseSettings):
|
||||
model_config = SettingsConfigDict(populate_by_name=True, extra="ignore")
|
||||
|
||||
@classmethod
|
||||
def settings_customise_sources(
|
||||
cls,
|
||||
settings_cls: type[BaseSettings],
|
||||
init_settings: PydanticBaseSettingsSource,
|
||||
env_settings: PydanticBaseSettingsSource,
|
||||
dotenv_settings: PydanticBaseSettingsSource,
|
||||
file_secret_settings: PydanticBaseSettingsSource,
|
||||
) -> tuple[PydanticBaseSettingsSource, ...]:
|
||||
return (
|
||||
init_settings,
|
||||
_EnvWithoutBareExcludedServices(settings_cls, env_settings),
|
||||
dotenv_settings,
|
||||
file_secret_settings,
|
||||
)
|
||||
|
||||
# ----- single-destination shorthand, read from standard OTEL_* envs ----- #
|
||||
exporter: str = Field(
|
||||
default="console",
|
||||
|
|
@ -178,7 +207,7 @@ class OpenTelemetryV2Config(BaseSettings):
|
|||
)
|
||||
excluded_services: Annotated[frozenset[str], NoDecode] = Field(
|
||||
default_factory=frozenset,
|
||||
validation_alias=AliasChoices("excluded_services", "LITELLM_OTEL_EXCLUDED_SERVICES"),
|
||||
validation_alias=AliasChoices("LITELLM_OTEL_EXCLUDED_SERVICES"),
|
||||
description=(
|
||||
"Datastore services whose spans are withheld from key/team ``callback_vars`` "
|
||||
"OTel destinations (the operator's own exporters still receive them). Accepted "
|
||||
|
|
|
|||
|
|
@ -124,6 +124,20 @@ def _trace_spans(sink_url: str, trace_id: str, seconds: float = 30) -> tuple[Spa
|
|||
return group
|
||||
|
||||
|
||||
def _trace_spans_when(
|
||||
sink_url: str,
|
||||
trace_id: str,
|
||||
ready: Callable[[tuple[Span, ...]], bool],
|
||||
seconds: float = 30,
|
||||
) -> tuple[Span, ...]:
|
||||
spans: Final = eventually(
|
||||
lambda: spans_for_trace(recorded_spans(sink_url)[1], trace_id),
|
||||
ready,
|
||||
seconds=seconds,
|
||||
)
|
||||
return spans
|
||||
|
||||
|
||||
def _await_db_span(sink_url: str, trace_id: str | None, needle: str, seconds: float = 40, since: int = 0) -> None:
|
||||
def seen() -> bool:
|
||||
_, spans = recorded_spans(sink_url, since)
|
||||
|
|
@ -242,6 +256,196 @@ def test_without_excluded_services_the_tenant_still_gets_redis_and_postgres_span
|
|||
assert {"redis", "postgresql"} <= systems, f"datastore spans missing at tenant: {systems}"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("otel", [None, True, "on", "", []], ids=["null", "true", "on", "empty_string", "empty_list"])
|
||||
@pytest.mark.timeout(180)
|
||||
def test_a_non_mapping_otel_block_still_publishes_the_tenant_fan_out(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: dict[str, JsonValue],
|
||||
tmp_path: Path,
|
||||
otel: JsonValue,
|
||||
) -> None:
|
||||
def with_callback_settings(config: dict) -> None:
|
||||
config["litellm_settings"]["callbacks"] = ["langfuse_otel"]
|
||||
config["callback_settings"]["otel"] = otel
|
||||
|
||||
config: Final = _config_with(tmp_path, otel_audit_config, extra=with_callback_settings)
|
||||
overrides: Final = {"LITELLM_OTEL_V2": "1", **_operator_langfuse(audit_sinks)}
|
||||
with owned_proxy(gateway, tmp_path, overrides, config=config, workers=2) as candidate:
|
||||
traffic: Final = _drive(candidate, langfuse_vars)
|
||||
tenant_trace: Final = _trace_id(audit_sinks.tenant, traffic)
|
||||
_await_db_span(audit_sinks.tenant, tenant_trace, "redis")
|
||||
tenant_spans: Final = _trace_spans_when(
|
||||
audit_sinks.tenant,
|
||||
tenant_trace,
|
||||
lambda spans: any(span["kind"] == 2 for span in spans) and "redis" in _db_systems(spans),
|
||||
seconds=15,
|
||||
)
|
||||
assert any(span["kind"] == 2 for span in tenant_spans), "tenant SERVER root span missing"
|
||||
assert "redis" in _db_systems(tenant_spans), f"tenant redis span missing: {_db_systems(tenant_spans)}"
|
||||
operator_trace: Final = _trace_id(audit_sinks.operator, traffic)
|
||||
operator_spans: Final = _trace_spans_when(
|
||||
audit_sinks.operator,
|
||||
operator_trace,
|
||||
lambda spans: any(span["kind"] == 2 for span in spans),
|
||||
seconds=15,
|
||||
)
|
||||
assert any(span["kind"] == 2 for span in operator_spans), "operator SERVER root span missing"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", ["EXCLUDED_SERVICES", "excluded_services"])
|
||||
@pytest.mark.timeout(180)
|
||||
def test_a_bare_excluded_services_env_var_is_ignored(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: dict[str, JsonValue],
|
||||
tmp_path: Path,
|
||||
name: str,
|
||||
) -> None:
|
||||
config: Final = _config_with(tmp_path, otel_audit_config)
|
||||
overrides: Final = {"LITELLM_OTEL_V2": "1", name: "redis,postgres"}
|
||||
with owned_proxy(gateway, tmp_path, overrides, config=config, workers=2) as candidate:
|
||||
tenant_start, _ = recorded_spans(audit_sinks.tenant)
|
||||
traffic: Final = _drive(candidate, langfuse_vars)
|
||||
tenant_trace: Final = _trace_id(audit_sinks.tenant, traffic)
|
||||
_await_db_span(audit_sinks.tenant, tenant_trace, "redis")
|
||||
tenant_spans: Final = _trace_spans_when(
|
||||
audit_sinks.tenant,
|
||||
tenant_trace,
|
||||
lambda spans: "redis" in _db_systems(spans),
|
||||
seconds=15,
|
||||
)
|
||||
assert "redis" in _db_systems(tenant_spans), f"redis span missing at tenant: {_db_systems(tenant_spans)}"
|
||||
_await_db_span(audit_sinks.tenant, None, "postgresql", since=tenant_start)
|
||||
_, all_tenant = recorded_spans(audit_sinks.tenant, tenant_start)
|
||||
systems: Final = _db_systems(all_tenant)
|
||||
assert {"redis", "postgresql"} <= systems, f"datastore spans missing at tenant: {systems}"
|
||||
|
||||
|
||||
@pytest.mark.timeout(180)
|
||||
def test_the_documented_env_var_wins_over_a_bare_excluded_services(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: dict[str, JsonValue],
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
config: Final = _config_with(tmp_path, otel_audit_config)
|
||||
overrides: Final = {
|
||||
"LITELLM_OTEL_V2": "1",
|
||||
"LITELLM_OTEL_EXCLUDED_SERVICES": "redis",
|
||||
"EXCLUDED_SERVICES": "postgres",
|
||||
}
|
||||
with owned_proxy(gateway, tmp_path, overrides, config=config, workers=2) as candidate:
|
||||
tenant_start, _ = recorded_spans(audit_sinks.tenant)
|
||||
traffic: Final = _drive(candidate, langfuse_vars)
|
||||
tenant_trace: Final = _trace_id(audit_sinks.tenant, traffic)
|
||||
_trace_spans(audit_sinks.tenant, tenant_trace, seconds=15)
|
||||
_await_db_span(audit_sinks.tenant, None, "postgresql", since=tenant_start)
|
||||
_, tenant_spans = recorded_spans(audit_sinks.tenant, tenant_start)
|
||||
systems: Final = _db_systems(tenant_spans)
|
||||
assert "postgresql" in systems, f"postgresql spans missing at tenant: {systems}"
|
||||
assert "redis" not in systems, f"redis spans reached tenant: {systems}"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("env_name", "redis_reaches_tenant"),
|
||||
[
|
||||
pytest.param("LITELLM_OTEL_EXCLUDED_SERVICES", False, id="exact-case"),
|
||||
pytest.param("litellm_otel_excluded_services", True, id="wrong-case"),
|
||||
],
|
||||
)
|
||||
@pytest.mark.timeout(180)
|
||||
def test_case_sensitive_otel_settings_read_only_the_exact_env_name(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: Mapping[str, JsonValue],
|
||||
tmp_path: Path,
|
||||
env_name: str,
|
||||
redis_reaches_tenant: bool,
|
||||
) -> None:
|
||||
config: Final = _config_with(tmp_path, otel_audit_config, otel={"_case_sensitive": True})
|
||||
overrides: Final = {"LITELLM_OTEL_V2": "1", env_name: "redis"}
|
||||
with owned_proxy(
|
||||
gateway,
|
||||
tmp_path,
|
||||
overrides,
|
||||
config=config,
|
||||
remove_environment=("LITELLM_OTEL_EXCLUDED_SERVICES", "litellm_otel_excluded_services"),
|
||||
workers=2,
|
||||
) as candidate:
|
||||
tenant_start, _ = recorded_spans(audit_sinks.tenant)
|
||||
traffic: Final = _drive(candidate, langfuse_vars)
|
||||
tenant_trace: Final = _trace_id(audit_sinks.tenant, traffic)
|
||||
if redis_reaches_tenant:
|
||||
_await_db_span(audit_sinks.tenant, tenant_trace, "redis")
|
||||
_trace_spans_when(
|
||||
audit_sinks.tenant,
|
||||
tenant_trace,
|
||||
lambda spans: "redis" in _db_systems(spans),
|
||||
seconds=15,
|
||||
)
|
||||
else:
|
||||
_trace_spans(audit_sinks.tenant, tenant_trace, seconds=15)
|
||||
_await_db_span(audit_sinks.tenant, None, "postgresql", seconds=60, since=tenant_start)
|
||||
_, tenant_spans = recorded_spans(audit_sinks.tenant, tenant_start)
|
||||
systems: Final = _db_systems(tenant_spans)
|
||||
assert "postgresql" in systems, f"postgresql spans missing at tenant: {systems}"
|
||||
assert ("redis" in systems) is redis_reaches_tenant, (
|
||||
f"tenant redis presence={('redis' in systems)}; expected={redis_reaches_tenant}; systems={systems}"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.timeout(180)
|
||||
def test_env_ignore_empty_keeps_the_default_service_name(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: Mapping[str, JsonValue],
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
config: Final = _config_with(tmp_path, otel_audit_config, otel={"_env_ignore_empty": True})
|
||||
overrides: Final = {"LITELLM_OTEL_V2": "1", "OTEL_SERVICE_NAME": ""}
|
||||
with owned_proxy(gateway, tmp_path, overrides, config=config, workers=2) as candidate:
|
||||
traffic: Final = _drive(candidate, langfuse_vars)
|
||||
operator_trace: Final = _trace_id(audit_sinks.operator, traffic)
|
||||
operator_spans: Final = _trace_spans_when(
|
||||
audit_sinks.operator,
|
||||
operator_trace,
|
||||
lambda spans: any(span["kind"] == 2 for span in spans),
|
||||
seconds=15,
|
||||
)
|
||||
service_names: Final = tuple(span["resource"].get("service.name") for span in operator_spans)
|
||||
assert service_names and all(name == "litellm" for name in service_names), (
|
||||
f"operator service.name values={service_names}"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.timeout(180)
|
||||
def test_env_parse_none_str_reads_a_null_traces_endpoint_as_unset(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: Mapping[str, JsonValue],
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
config: Final = _config_with(tmp_path, otel_audit_config, otel={"_env_parse_none_str": "null"})
|
||||
overrides: Final = {"LITELLM_OTEL_V2": "1", "OTEL_TRACES_ENDPOINT": "null"}
|
||||
with owned_proxy(gateway, tmp_path, overrides, config=config, workers=2) as candidate:
|
||||
traffic: Final = _drive(candidate, langfuse_vars)
|
||||
operator_trace: Final = _trace_id(audit_sinks.operator, traffic)
|
||||
operator_spans: Final = _trace_spans_when(
|
||||
audit_sinks.operator,
|
||||
operator_trace,
|
||||
lambda spans: any(span["kind"] == 2 for span in spans),
|
||||
seconds=15,
|
||||
)
|
||||
assert any(span["kind"] == 2 for span in operator_spans), "operator SERVER root span missing"
|
||||
|
||||
|
||||
def test_env_excluded_services_drops_only_redis(
|
||||
gateway: Gateway,
|
||||
audit_sinks: SpanSinks,
|
||||
|
|
|
|||
|
|
@ -9,7 +9,8 @@ from concurrent.futures import ThreadPoolExecutor
|
|||
from contextlib import contextmanager
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Final, Literal
|
||||
from types import MappingProxyType
|
||||
from typing import Final, Literal, Protocol, cast
|
||||
|
||||
import anthropic
|
||||
import httpx
|
||||
|
|
@ -26,6 +27,9 @@ from pydantic import JsonValue, TypeAdapter
|
|||
MARKER: Final = re.compile(rb"excl-[0-9a-f]{32}")
|
||||
FAILING: Final = re.compile(rb"excl-fail-[0-9a-f]{32}")
|
||||
JSON: Final[TypeAdapter[JsonValue]] = TypeAdapter(JsonValue)
|
||||
UNCONFIGURED_VARIANT: Final[TypeAdapter[Literal["null_block", "bare_env"]]] = TypeAdapter(
|
||||
Literal["null_block", "bare_env"]
|
||||
)
|
||||
REPLY_TEXT: Final = "excluded ok"
|
||||
SERVER: Final = 2
|
||||
INVALID_NAME_LOG: Final = "is not a datastore service"
|
||||
|
|
@ -37,6 +41,11 @@ CLIENTS: Final[tuple[Client, ...]] = ("raw", "sdk", "async_sdk")
|
|||
AuditConfigWriter = Callable[[Path, Mapping[str, JsonValue]], Path]
|
||||
|
||||
|
||||
class _FixtureRequestParam(Protocol):
|
||||
@property
|
||||
def param(self) -> object: ...
|
||||
|
||||
|
||||
def _marker() -> str:
|
||||
return "excl-" + uuid.uuid4().hex
|
||||
|
||||
|
|
@ -339,6 +348,11 @@ def _db_systems(spans: tuple[Span, ...]) -> set[str]:
|
|||
}
|
||||
|
||||
|
||||
def _post_auth_datastore_spans(spans: tuple[Span, ...]) -> tuple[Span, ...]:
|
||||
auth_ids: Final = frozenset(span["span_id"] for span in spans if span["name"].startswith("auth "))
|
||||
return tuple(span for span in spans if _db_systems((span,)) and span["parent_span_id"] not in auth_ids)
|
||||
|
||||
|
||||
def _names(spans: tuple[Span, ...]) -> list[str]:
|
||||
return sorted(span["name"] for span in spans)
|
||||
|
||||
|
|
@ -408,6 +422,29 @@ def _config(directory: Path, otel_audit_config: AuditConfigWriter, otel: Mapping
|
|||
return path
|
||||
|
||||
|
||||
def _null_otel_config(directory: Path, otel_audit_config: AuditConfigWriter, name: str) -> Path:
|
||||
written: Final = otel_audit_config(directory, {})
|
||||
loaded: Final = object_value(JSON.validate_python(yaml.safe_load(written.read_text())))
|
||||
config: Final = {
|
||||
**loaded,
|
||||
"litellm_settings": {**object_value(loaded["litellm_settings"]), "callbacks": ["langfuse_otel"]},
|
||||
"callback_settings": {**object_value(loaded["callback_settings"]), "otel": None},
|
||||
}
|
||||
path: Final = directory / f"{name}.yaml"
|
||||
path.write_text(yaml.safe_dump(config))
|
||||
return path
|
||||
|
||||
|
||||
def _operator_langfuse(sinks: SpanSinks) -> dict[str, str]:
|
||||
return {
|
||||
"LANGFUSE_HOST": sinks.operator,
|
||||
"LANGFUSE_PUBLIC_KEY": "pk-lf-operator",
|
||||
"LANGFUSE_SECRET_KEY": "sk-lf-operator",
|
||||
"OTEL_EXPORTER": "http/json",
|
||||
"OTEL_ENDPOINT": sinks.operator,
|
||||
}
|
||||
|
||||
|
||||
@contextmanager
|
||||
def _started(
|
||||
provider: Wire,
|
||||
|
|
@ -416,13 +453,14 @@ def _started(
|
|||
directory: Path,
|
||||
langfuse_vars: Mapping[str, JsonValue],
|
||||
workers: int,
|
||||
environment: Mapping[str, str] = MappingProxyType({}),
|
||||
) -> Generator[Rig]:
|
||||
with (
|
||||
gateway_from_environment() as gateway,
|
||||
owned_proxy_process(
|
||||
gateway,
|
||||
directory,
|
||||
{"LITELLM_OTEL_V2": "1", "OTEL_BSP_SCHEDULE_DELAY": "300"},
|
||||
{"LITELLM_OTEL_V2": "1", "OTEL_BSP_SCHEDULE_DELAY": "300", **environment},
|
||||
config=config,
|
||||
remove_environment=("LITELLM_OTEL_EXCLUDED_SERVICES",),
|
||||
workers=workers,
|
||||
|
|
@ -458,6 +496,32 @@ def rig(
|
|||
yield started
|
||||
|
||||
|
||||
@pytest.fixture(scope="module", params=["null_block", "bare_env"], ids=["null_block", "bare_env"])
|
||||
def unconfigured_rig(
|
||||
request: pytest.FixtureRequest,
|
||||
provider: Wire,
|
||||
audit_sinks: SpanSinks,
|
||||
otel_audit_config: AuditConfigWriter,
|
||||
langfuse_vars: dict[str, JsonValue],
|
||||
tmp_path_factory: pytest.TempPathFactory,
|
||||
) -> Iterator[Rig]:
|
||||
parameter: Final = cast(_FixtureRequestParam, request).param
|
||||
variant: Final = UNCONFIGURED_VARIANT.validate_python(parameter)
|
||||
directory: Final = tmp_path_factory.mktemp(f"excluded-{variant}")
|
||||
config: Final = (
|
||||
_null_otel_config(directory, otel_audit_config, variant)
|
||||
if variant == "null_block"
|
||||
else _config(directory, otel_audit_config, {}, variant)
|
||||
)
|
||||
environment: Final = (
|
||||
_operator_langfuse(audit_sinks) if variant == "null_block" else {"EXCLUDED_SERVICES": "redis,postgres"}
|
||||
)
|
||||
with _started(
|
||||
provider, audit_sinks, config, directory, langfuse_vars, workers=2, environment=environment
|
||||
) as started:
|
||||
yield started
|
||||
|
||||
|
||||
@pytest.mark.timeout(120)
|
||||
@pytest.mark.parametrize("stream", [False, True], ids=["unary", "stream"])
|
||||
@pytest.mark.parametrize("client", CLIENTS)
|
||||
|
|
@ -654,6 +718,237 @@ def test_killing_one_of_two_workers_mid_burst_keeps_the_filter_on_the_survivor(r
|
|||
_assert_withheld(rig, rig.raw("chat", _marker(), stream=False), after)
|
||||
|
||||
|
||||
def _assert_tenant_kept(
|
||||
rig: Rig, trace_id: str, cursors: Cursors, *, needs_model_span: bool = False
|
||||
) -> tuple[Span, ...]:
|
||||
def ready(spans: tuple[Span, ...]) -> bool:
|
||||
return (
|
||||
sum(1 for span in spans if span["kind"] == SERVER) == 1
|
||||
and "redis" in _db_systems(spans)
|
||||
and (not needs_model_span or any("gen_ai.operation.name" in span["attributes"] for span in spans))
|
||||
)
|
||||
|
||||
tenant: Final = eventually(
|
||||
lambda: spans_for_trace(recorded_spans(rig.sinks.tenant, cursors.tenant)[1], trace_id),
|
||||
ready,
|
||||
seconds=40,
|
||||
return_last_on_timeout=True,
|
||||
)
|
||||
assert sum(1 for span in tenant if span["kind"] == SERVER) == 1, _names(tenant)
|
||||
assert "redis" in _db_systems(tenant), f"redis spans missing at the tenant: {_names(tenant)}"
|
||||
assert not needs_model_span or any("gen_ai.operation.name" in span["attributes"] for span in tenant), _names(tenant)
|
||||
return tenant
|
||||
|
||||
|
||||
def _assert_kept(rig: Rig, sent: Sent, cursors: Cursors) -> tuple[Span, ...]:
|
||||
operator: Final = _operator_trace(rig, sent, cursors)
|
||||
return _assert_tenant_kept(rig, operator[0]["trace_id"], cursors, needs_model_span=True)
|
||||
|
||||
|
||||
@pytest.mark.timeout(120)
|
||||
@pytest.mark.parametrize("stream", [False, True], ids=["unary", "stream"])
|
||||
@pytest.mark.parametrize("client", CLIENTS)
|
||||
@pytest.mark.parametrize("endpoint", ENDPOINTS)
|
||||
def test_unconfigured_tenant_trace_keeps_datastore_spans(
|
||||
unconfigured_rig: Rig, endpoint: Endpoint, client: Client, stream: bool
|
||||
) -> None:
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
marker: Final = _marker()
|
||||
sent: Final = unconfigured_rig.send(endpoint, client, marker, stream)
|
||||
assert sent.text == REPLY_TEXT, sent
|
||||
assert unconfigured_rig.upstream_hits(marker) == 1
|
||||
_assert_kept(unconfigured_rig, sent, cursors)
|
||||
|
||||
|
||||
@pytest.mark.timeout(120)
|
||||
@pytest.mark.parametrize("endpoint", ["chat", "messages"])
|
||||
def test_unconfigured_cache_hit_twin_keeps_datastore_spans(unconfigured_rig: Rig, endpoint: Endpoint) -> None:
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
marker: Final = _marker()
|
||||
first_result: Final = _traced_raw(unconfigured_rig, endpoint, marker)
|
||||
first: Final = first_result[1]
|
||||
assert first.text == REPLY_TEXT, first
|
||||
assert unconfigured_rig.upstream_hits(marker) == 1
|
||||
_assert_kept(unconfigured_rig, first, cursors)
|
||||
hit_cursors: Final = unconfigured_rig.cursors()
|
||||
|
||||
def read_hit() -> tuple[str, Sent, tuple[Span, ...]]:
|
||||
trace_id, sent = _traced_raw(unconfigured_rig, endpoint, marker)
|
||||
return trace_id, sent, _operator_trace_by_id(unconfigured_rig, trace_id, hit_cursors)
|
||||
|
||||
trace_id, hit, operator = eventually(
|
||||
read_hit,
|
||||
lambda result: (
|
||||
unconfigured_rig.upstream_hits(marker) == 0
|
||||
and "redis" in _db_systems(_post_auth_datastore_spans(result[2]))
|
||||
),
|
||||
seconds=60,
|
||||
)
|
||||
assert hit.text == REPLY_TEXT, hit
|
||||
post_auth_datastore: Final = _post_auth_datastore_spans(operator)
|
||||
post_auth_span_ids: Final = frozenset(span["span_id"] for span in post_auth_datastore)
|
||||
non_datastore_names: Final = frozenset(span["name"] for span in operator if not _db_systems((span,)))
|
||||
tenant: Final = eventually(
|
||||
lambda: spans_for_trace(recorded_spans(unconfigured_rig.sinks.tenant, hit_cursors.tenant)[1], trace_id),
|
||||
lambda spans: (
|
||||
non_datastore_names <= frozenset(span["name"] for span in spans)
|
||||
and post_auth_span_ids <= frozenset(span["span_id"] for span in spans)
|
||||
),
|
||||
seconds=40,
|
||||
return_last_on_timeout=True,
|
||||
)
|
||||
tenant_span_ids: Final = frozenset(span["span_id"] for span in tenant)
|
||||
missing_post_auth_names: Final = tuple(
|
||||
span["name"] for span in post_auth_datastore if span["span_id"] not in tenant_span_ids
|
||||
)
|
||||
assert sum(1 for span in tenant if span["kind"] == SERVER) == 1, _names(tenant)
|
||||
assert "redis" in _db_systems(tenant), (
|
||||
f"operator datastore systems={sorted(_db_systems(post_auth_datastore))}; "
|
||||
f"tenant datastore systems={sorted(_db_systems(tenant))}; tenant spans={_names(tenant)}"
|
||||
)
|
||||
assert not missing_post_auth_names, (
|
||||
f"missing post-auth datastore span names={missing_post_auth_names}; "
|
||||
f"operator={_names(post_auth_datastore)}; tenant={_names(tenant)}"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.timeout(120)
|
||||
@pytest.mark.parametrize("endpoint", ENDPOINTS)
|
||||
def test_unconfigured_failed_upstream_keeps_datastore_spans(unconfigured_rig: Rig, endpoint: Endpoint) -> None:
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
marker: Final = "excl-fail-" + uuid.uuid4().hex
|
||||
trace_id: Final = uuid.uuid4().hex
|
||||
path, body = _body(unconfigured_rig.model, endpoint, marker, stream=False)
|
||||
failed: Final = unconfigured_rig.proxy.client.post(
|
||||
path,
|
||||
json=body,
|
||||
headers={
|
||||
"Authorization": f"Bearer {unconfigured_rig.key}",
|
||||
"traceparent": f"00-{trace_id}-{uuid.uuid4().hex[:16]}-01",
|
||||
},
|
||||
)
|
||||
assert failed.status_code == 500, failed.text
|
||||
assert unconfigured_rig.upstream_hits(marker) >= 1
|
||||
operator: Final = eventually(
|
||||
lambda: spans_for_trace(recorded_spans(unconfigured_rig.sinks.operator, cursors.operator)[1], trace_id),
|
||||
lambda spans: _has_root(spans) and "redis" in _db_systems(spans),
|
||||
seconds=40,
|
||||
)
|
||||
assert "redis" in _db_systems(operator), _names(operator)
|
||||
_assert_tenant_kept(unconfigured_rig, trace_id, cursors)
|
||||
|
||||
|
||||
@pytest.mark.timeout(120)
|
||||
@pytest.mark.parametrize("status", [403, 404])
|
||||
def test_unconfigured_rejecting_tenant_destination_recovers(unconfigured_rig: Rig, status: int) -> None:
|
||||
configure_sink(unconfigured_rig.sinks.tenant, status=status)
|
||||
try:
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
marker: Final = _marker()
|
||||
sent: Final = unconfigured_rig.raw("chat", marker, stream=True)
|
||||
assert sent.text == REPLY_TEXT, sent
|
||||
assert unconfigured_rig.upstream_hits(marker) == 1
|
||||
_operator_trace(unconfigured_rig, sent, cursors)
|
||||
finally:
|
||||
configure_sink(unconfigured_rig.sinks.tenant, status=200)
|
||||
after: Final = unconfigured_rig.cursors()
|
||||
recovered: Final = unconfigured_rig.raw("responses", _marker(), stream=False)
|
||||
assert recovered.text == REPLY_TEXT, recovered
|
||||
_assert_kept(unconfigured_rig, recovered, after)
|
||||
|
||||
|
||||
@pytest.mark.timeout(120)
|
||||
def test_unconfigured_key_level_destination_keeps_datastore_spans(
|
||||
unconfigured_rig: Rig, langfuse_vars: dict[str, JsonValue]
|
||||
) -> None:
|
||||
key: Final = unconfigured_rig.scenario.key(
|
||||
metadata={
|
||||
"logging": [
|
||||
{"callback_name": "langfuse_otel", "callback_type": "success", "callback_vars": dict(langfuse_vars)}
|
||||
]
|
||||
}
|
||||
)
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
marker: Final = _marker()
|
||||
sent: Final = unconfigured_rig.raw("chat", marker, stream=False, key=key)
|
||||
assert sent.text == REPLY_TEXT, sent
|
||||
assert unconfigured_rig.upstream_hits(marker) == 1
|
||||
_assert_kept(unconfigured_rig, sent, cursors)
|
||||
|
||||
|
||||
def _assert_tenant_kept_the_burst(rig: Rig, cursors: Cursors, traces: set[str]) -> None:
|
||||
def ready(spans: tuple[Span, ...]) -> bool:
|
||||
def trace_kept(trace: str) -> bool:
|
||||
trace_spans: Final = spans_for_trace(spans, trace)
|
||||
return any(span["kind"] == SERVER for span in trace_spans) and "redis" in _db_systems(trace_spans)
|
||||
|
||||
return all(trace_kept(trace) for trace in traces)
|
||||
|
||||
tenant: Final = eventually(
|
||||
lambda: recorded_spans(rig.sinks.tenant, cursors.tenant)[1],
|
||||
ready,
|
||||
seconds=90,
|
||||
return_last_on_timeout=True,
|
||||
)
|
||||
burst: Final = tuple(span for span in tenant if span["trace_id"] in traces)
|
||||
missing_roots: Final = tuple(
|
||||
trace for trace in traces if not any(span["kind"] == SERVER for span in spans_for_trace(burst, trace))
|
||||
)
|
||||
missing_redis: Final = tuple(trace for trace in traces if "redis" not in _db_systems(spans_for_trace(burst, trace)))
|
||||
assert not missing_roots, f"SERVER root missing from tenant burst traces: {missing_roots}, {_names(burst)}"
|
||||
assert not missing_redis, f"redis spans missing from tenant burst traces: {missing_redis}, {_names(burst)}"
|
||||
|
||||
|
||||
@pytest.mark.timeout(300)
|
||||
def test_unconfigured_tenant_outage_during_a_mixed_burst(unconfigured_rig: Rig) -> None:
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
configure_sink(unconfigured_rig.sinks.tenant, status=503)
|
||||
try:
|
||||
results: Final = _burst(unconfigured_rig, 30)
|
||||
finally:
|
||||
configure_sink(unconfigured_rig.sinks.tenant, status=200)
|
||||
served: Final = _served(results)
|
||||
assert len(served) == 30, [result for result in results if isinstance(result, str)]
|
||||
assert all(sent.text == REPLY_TEXT for sent in served), served
|
||||
traces: Final = _assert_operator_exactly_once(unconfigured_rig, served, cursors)
|
||||
_assert_tenant_kept_the_burst(unconfigured_rig, cursors, traces)
|
||||
after: Final = unconfigured_rig.cursors()
|
||||
_assert_kept(unconfigured_rig, unconfigured_rig.raw("messages", _marker(), stream=True), after)
|
||||
|
||||
|
||||
@pytest.mark.timeout(300)
|
||||
def test_unconfigured_killing_one_of_two_workers_keeps_the_fan_out(unconfigured_rig: Rig) -> None:
|
||||
root: Final = psutil.Process(unconfigured_rig.owned.process.pid)
|
||||
workers: Final = eventually(
|
||||
lambda: tuple(child for child in root.children() if "resource_tracker" not in " ".join(child.cmdline())),
|
||||
lambda found: len(found) == 2,
|
||||
seconds=30,
|
||||
)
|
||||
cursors: Final = unconfigured_rig.cursors()
|
||||
|
||||
def one(index: int) -> Sent | str:
|
||||
if index == 6:
|
||||
os.kill(workers[0].pid, signal.SIGKILL)
|
||||
try:
|
||||
return unconfigured_rig.raw("chat", _marker(), stream=index % 2 == 0)
|
||||
except (httpx.HTTPError, AssertionError) as error:
|
||||
return repr(error)
|
||||
|
||||
with ThreadPoolExecutor(max_workers=6) as pool:
|
||||
results: Final = tuple(pool.map(one, range(18)))
|
||||
assert unconfigured_rig.owned.process.poll() is None, "Proxy root exited after a worker was killed"
|
||||
failures: Final = tuple(result for result in results if isinstance(result, str))
|
||||
assert all(failure.startswith(("ReadError(", "RemoteProtocolError(", "ConnectError(")) for failure in failures), (
|
||||
failures
|
||||
)
|
||||
assert len(failures) <= 6, failures
|
||||
settled: Final = tuple(result for index, result in enumerate(results) if index > 12 and isinstance(result, Sent))
|
||||
traces: Final = _assert_operator_exactly_once(unconfigured_rig, settled, cursors)
|
||||
_assert_tenant_kept_the_burst(unconfigured_rig, cursors, traces)
|
||||
after: Final = unconfigured_rig.cursors()
|
||||
_assert_kept(unconfigured_rig, unconfigured_rig.raw("chat", _marker(), stream=False), after)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Setting:
|
||||
otel: Mapping[str, JsonValue]
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@
|
|||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Final
|
||||
|
||||
import pytest
|
||||
|
||||
|
|
@ -110,6 +111,61 @@ def test_excluded_services_from_env_csv(monkeypatch):
|
|||
assert OpenTelemetryV2Config().excluded_services == frozenset({"redis", "postgresql"})
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", ["EXCLUDED_SERVICES", "excluded_services", "Excluded_Services"])
|
||||
def test_a_bare_excluded_services_env_var_is_ignored(monkeypatch, name):
|
||||
for env_name in ("LITELLM_OTEL_EXCLUDED_SERVICES", "EXCLUDED_SERVICES", "excluded_services", "Excluded_Services"):
|
||||
monkeypatch.delenv(env_name, raising=False)
|
||||
monkeypatch.setenv(name, "redis,postgres")
|
||||
assert OpenTelemetryV2Config().excluded_services == frozenset()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("set_env_name", "env_value", "case_sensitive", "env_ignore_empty", "env_parse_none_str"),
|
||||
[
|
||||
pytest.param("otel_service_name", "lower", True, False, None, id="case-sensitive"),
|
||||
pytest.param("OTEL_SERVICE_NAME", "", False, True, None, id="ignore-empty"),
|
||||
pytest.param("OTEL_ENDPOINT", "null", False, False, "null", id="parse-none"),
|
||||
pytest.param("excluded_services", "redis", True, False, None, id="bare-exclusion"),
|
||||
],
|
||||
)
|
||||
def test_env_source_preserves_runtime_options(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
set_env_name: str,
|
||||
env_value: str,
|
||||
case_sensitive: bool,
|
||||
env_ignore_empty: bool,
|
||||
env_parse_none_str: str | None,
|
||||
) -> None:
|
||||
for env_name in (
|
||||
"OTEL_SERVICE_NAME",
|
||||
"otel_service_name",
|
||||
"OTEL_ENDPOINT",
|
||||
"OTEL_EXPORTER_OTLP_ENDPOINT",
|
||||
"LITELLM_OTEL_EXCLUDED_SERVICES",
|
||||
"EXCLUDED_SERVICES",
|
||||
"excluded_services",
|
||||
"Excluded_Services",
|
||||
):
|
||||
monkeypatch.delenv(env_name, raising=False)
|
||||
monkeypatch.setenv(set_env_name, env_value)
|
||||
config: Final = OpenTelemetryV2Config(
|
||||
_case_sensitive=case_sensitive,
|
||||
_env_ignore_empty=env_ignore_empty,
|
||||
_env_parse_none_str=env_parse_none_str,
|
||||
)
|
||||
assert config.service_name == "litellm"
|
||||
assert config.endpoint is None
|
||||
assert config.excluded_services == frozenset()
|
||||
|
||||
|
||||
def test_the_documented_env_var_wins_over_a_bare_excluded_services(monkeypatch):
|
||||
for env_name in ("LITELLM_OTEL_EXCLUDED_SERVICES", "EXCLUDED_SERVICES", "excluded_services", "Excluded_Services"):
|
||||
monkeypatch.delenv(env_name, raising=False)
|
||||
monkeypatch.setenv("EXCLUDED_SERVICES", "postgres")
|
||||
monkeypatch.setenv("LITELLM_OTEL_EXCLUDED_SERVICES", "redis")
|
||||
assert OpenTelemetryV2Config().excluded_services == frozenset({"redis"})
|
||||
|
||||
|
||||
def test_excluded_services_config_wins_over_env(monkeypatch):
|
||||
monkeypatch.setenv("LITELLM_OTEL_EXCLUDED_SERVICES", "redis")
|
||||
assert OpenTelemetryV2Config(excluded_services=["postgres"]).excluded_services == frozenset({"postgresql"})
|
||||
|
|
|
|||
|
|
@ -1116,6 +1116,18 @@ class TestProviderWiring:
|
|||
|
||||
assert self._fan_out_of(preset)._excluded_db_systems == frozenset({"redis"})
|
||||
|
||||
@pytest.mark.parametrize("otel", [None, True, "on", "", []], ids=["null", "true", "on", "empty_string", "empty_list"])
|
||||
def test_a_non_mapping_otel_block_falls_back_to_the_published_logger_config(self, monkeypatch, otel):
|
||||
monkeypatch.setattr(litellm, "callback_settings", {"otel": otel}, raising=False)
|
||||
preset = OpenTelemetryV2(
|
||||
config=OpenTelemetryV2Config(exporters=[ExporterSpec(kind="in_memory")], excluded_services=["redis"]),
|
||||
callback_name="langfuse_otel",
|
||||
)
|
||||
|
||||
publish_global_otel_v2_provider([], lambda _p: None, registered=preset)
|
||||
|
||||
assert self._fan_out_of(preset)._excluded_db_systems == frozenset({"redis"})
|
||||
|
||||
def test_otel_after_a_preset_reuses_it_and_still_takes_callback_settings_exclusions(self, monkeypatch):
|
||||
"""``callbacks: [langfuse_otel, otel]`` keeps one v2 logger, exactly as
|
||||
before ``excluded_services`` existed, and the exclusion still comes from
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue