diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index 55eb8e8fb71..e1983a44451 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -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) diff --git a/litellm/integrations/otel/model/config.py b/litellm/integrations/otel/model/config.py index 9eb29157d6f..ddb8e127408 100644 --- a/litellm/integrations/otel/model/config.py +++ b/litellm/integrations/otel/model/config.py @@ -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 " diff --git a/tests/integration/observability/test_otel_excluded_services.py b/tests/integration/observability/test_otel_excluded_services.py index 48b20651b97..8768052b488 100644 --- a/tests/integration/observability/test_otel_excluded_services.py +++ b/tests/integration/observability/test_otel_excluded_services.py @@ -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, diff --git a/tests/integration/observability/test_otel_excluded_services_matrix.py b/tests/integration/observability/test_otel_excluded_services_matrix.py index 0d4b5d087c6..7bca7ffc0dc 100644 --- a/tests/integration/observability/test_otel_excluded_services_matrix.py +++ b/tests/integration/observability/test_otel_excluded_services_matrix.py @@ -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] diff --git a/tests/unit/integrations/otel/test_otel_v2_config_baggage_parenting_guardrails.py b/tests/unit/integrations/otel/test_otel_v2_config_baggage_parenting_guardrails.py index 86837f7f46c..20f8c89b5ff 100644 --- a/tests/unit/integrations/otel/test_otel_v2_config_baggage_parenting_guardrails.py +++ b/tests/unit/integrations/otel/test_otel_v2_config_baggage_parenting_guardrails.py @@ -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"}) diff --git a/tests/unit/integrations/otel/test_otel_v2_destinations.py b/tests/unit/integrations/otel/test_otel_v2_destinations.py index 5a7057203e4..cb986e61229 100644 --- a/tests/unit/integrations/otel/test_otel_v2_destinations.py +++ b/tests/unit/integrations/otel/test_otel_v2_destinations.py @@ -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