From 7aba77197dc53737f8e882bfceab493397a424b0 Mon Sep 17 00:00:00 2001 From: "devin-ai-integration[bot]" <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sat, 26 Sep 2026 18:15:45 -0700 Subject: [PATCH] feat(otel): add SigNoz preset for OpenTelemetry v2 (#43296) * feat(otel): add SigNoz preset for OpenTelemetry v2 Adds the signoz callback (OTLP/HTTP exporter, GenAI vocabulary, key and team level dynamic ingestion endpoint and key) as an OpenTelemetry v2 preset, with the preset factory accepting the allow_missing_credentials kwarg the V2 registry always passes so construction no longer falls back silently to legacy OpenTelemetry. Ships the deterministic tests/integration/observability/test_signoz_delivery.py audit suite Absorbs the work from https://github.com/BerriAI/litellm/pull/38206 Co-authored-by: Nagesh Bansal Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(otel): drop explanatory comments from the SigNoz preset Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * style(types): keep signoz dynamic param lines within ruff format width Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(signoz): assert the missing-endpoint boot path directly instead of in an except block Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * chore(ui): regenerate schema.d.ts for the signoz health service Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(otel): allowlist SigNoz key/team endpoints and route keyless collectors without the operator key Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(otel): terminate the SigNoz shutdown cell before the flush and drop test docstrings Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(otel): keep the shared tenant routing untouched and require an ingestion key for SigNoz key/team endpoints Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(otel): warn about a keyless SigNoz team endpoint from the header resolver so the shared cache actually reaches it Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: yucheng Co-authored-by: Nagesh Bansal Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/__init__.py | 1 + litellm/integrations/callback_configs.json | 21 + litellm/integrations/otel/model/config.py | 1 + litellm/integrations/otel/presets/__init__.py | 9 + litellm/integrations/otel/presets/signoz.py | 95 ++ .../custom_logger_registry.py | 1 + .../initialize_dynamic_callback_params.py | 6 + litellm/litellm_core_utils/litellm_logging.py | 32 + .../_experimental/out/assets/logos/signoz.svg | 1 + litellm/proxy/_types.py | 6 + .../health_endpoints/_health_endpoints.py | 2 + litellm/proxy/litellm_pre_call_utils.py | 2 + litellm/types/utils.py | 3 + .../observability/test_signoz_delivery.py | 980 ++++++++++++++++++ .../proxy/test_litellm_pre_call_utils.py | 28 + .../integrations/otel/test_otel_v2_dynamic.py | 66 ++ .../integrations/otel/test_otel_v2_presets.py | 66 ++ .../test_litellm_logging.py | 97 ++ .../public/assets/logos/signoz.svg | 1 + .../src/components/callback_info_helpers.tsx | 12 + ui/litellm-dashboard/src/lib/http/schema.d.ts | 2 +- 21 files changed, 1431 insertions(+), 1 deletion(-) create mode 100644 litellm/integrations/otel/presets/signoz.py create mode 100644 litellm/proxy/_experimental/out/assets/logos/signoz.svg create mode 100644 tests/integration/observability/test_signoz_delivery.py create mode 100644 ui/litellm-dashboard/public/assets/logos/signoz.svg diff --git a/litellm/__init__.py b/litellm/__init__.py index 5a7d6e8125d..5d10737e876 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -172,6 +172,7 @@ _custom_logger_compatible_callbacks_literal = Literal[ "levo", "compression_interception", "newrelic", + "signoz", ] cold_storage_custom_logger: Optional[_custom_logger_compatible_callbacks_literal] = None logged_real_time_event_types: Optional[Union[List[str], Literal["*"]]] = None diff --git a/litellm/integrations/callback_configs.json b/litellm/integrations/callback_configs.json index 4e72075dc5c..190c283d087 100644 --- a/litellm/integrations/callback_configs.json +++ b/litellm/integrations/callback_configs.json @@ -502,6 +502,27 @@ }, "description": "S3 Bucket (AWS) Logging Integration" }, + { + "id": "signoz", + "displayName": "SigNoz", + "logo": "signoz.svg", + "supports_key_team_logging": true, + "dynamic_params": { + "signoz_ingestion_endpoint": { + "type": "text", + "ui_name": "SigNoz Ingestion Endpoint", + "description": "Ingestion endpoint for this team, e.g. https://ingest.us.signoz.cloud:443 for SigNoz Cloud or your own collector. Leave blank to use the proxy's configured endpoint. Regions: https://signoz.io/docs/ingestion/signoz-cloud/overview/", + "required": false + }, + "signoz_ingestion_key": { + "type": "password", + "ui_name": "SigNoz Ingestion Key (optional)", + "description": "Ingestion key for this team, so its traces land in its own SigNoz account. Not needed for self-hosted SigNoz. Keys: https://signoz.io/docs/ingestion/signoz-cloud/keys/", + "required": false + } + }, + "description": "SigNoz Logging Integration. Setup: https://signoz.io/docs/litellm-observability/" + }, { "id": "sqs", "displayName": "SQS", diff --git a/litellm/integrations/otel/model/config.py b/litellm/integrations/otel/model/config.py index 5447a8ee80a..5a3965862e0 100644 --- a/litellm/integrations/otel/model/config.py +++ b/litellm/integrations/otel/model/config.py @@ -41,6 +41,7 @@ class ExporterOwner(str, Enum): LEVO = "levo" AGENTOPS = "agentops" NEWRELIC = "newrelic" + SIGNOZ = "signoz" class _OTelV2Flag(BaseSettings): diff --git a/litellm/integrations/otel/presets/__init__.py b/litellm/integrations/otel/presets/__init__.py index a0cd5b3fd98..7c891c29409 100644 --- a/litellm/integrations/otel/presets/__init__.py +++ b/litellm/integrations/otel/presets/__init__.py @@ -30,6 +30,11 @@ from litellm.integrations.otel.presets.phoenix import ( phoenix_preset, phoenix_project_headers, ) +from litellm.integrations.otel.presets.signoz import ( + signoz_dynamic_endpoint, + signoz_dynamic_headers, + signoz_preset, +) from litellm.integrations.otel.presets.weave import weave_dynamic_headers, weave_preset from litellm.types.utils import StandardCallbackDynamicParams @@ -44,6 +49,7 @@ PRESET_BY_CALLBACK: Final[Mapping[str, Preset]] = MappingProxyType( "langtrace": langtrace_preset, "levo": levo_preset, "newrelic": newrelic_preset, + "signoz": signoz_preset, "weave_otel": weave_preset, } ) @@ -58,6 +64,7 @@ DYNAMIC_HEADERS_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynami "arize": arize_dynamic_headers, "langfuse_otel": langfuse_dynamic_headers, "newrelic": newrelic_dynamic_headers, + "signoz": signoz_dynamic_headers, "weave_otel": weave_dynamic_headers, } ) @@ -71,6 +78,7 @@ DYNAMIC_ENDPOINT_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynam MappingProxyType( { "newrelic": newrelic_dynamic_endpoint, + "signoz": signoz_dynamic_endpoint, } ) ) @@ -153,5 +161,6 @@ __all__ = [ "newrelic_preset", "phoenix_preset", "project_routing_headers", + "signoz_preset", "weave_preset", ] diff --git a/litellm/integrations/otel/presets/signoz.py b/litellm/integrations/otel/presets/signoz.py new file mode 100644 index 00000000000..c4d7ed48a38 --- /dev/null +++ b/litellm/integrations/otel/presets/signoz.py @@ -0,0 +1,95 @@ +from functools import lru_cache +from types import MappingProxyType +from typing import Final + +from pydantic import Field +from pydantic_settings import BaseSettings, SettingsConfigDict + +import litellm +from litellm._logging import verbose_logger +from litellm.integrations.otel.model.config import ( + ExporterOwner, + ExporterSpec, + OpenTelemetryV2Config, +) +from litellm.integrations.otel.presets.utils import ensure_mappers +from litellm.litellm_core_utils.url_utils import is_url_destination_allowed_by_host +from litellm.types.utils import StandardCallbackDynamicParams + +SIGNOZ_INGESTION_ENDPOINT_ENV: Final = "SIGNOZ_INGESTION_ENDPOINT" + + +class _SigNozSettings(BaseSettings): + model_config = SettingsConfigDict(case_sensitive=False, extra="ignore") + + endpoint: str | None = Field(default=None, validation_alias=SIGNOZ_INGESTION_ENDPOINT_ENV) + ingestion_key: str | None = Field(default=None, validation_alias="SIGNOZ_INGESTION_KEY") + + +def signoz_preset( + *, + config_overrides: OpenTelemetryV2Config | None = None, + allow_missing_credentials: bool = False, +) -> OpenTelemetryV2Config: + settings: Final = _SigNozSettings() + base: Final = config_overrides or OpenTelemetryV2Config() + key: Final = settings.ingestion_key + spec: Final = ExporterSpec( + kind="otlp_http", + endpoint=settings.endpoint, + headers=(f"signoz-ingestion-key={key}" if key else None), + owner=ExporterOwner.SIGNOZ, + requires_headers=bool(key), + ) + return base.model_copy( + update=MappingProxyType( + { + "exporters": (*base.exporters, spec), + "mapper_names": ensure_mappers(base.mapper_names, "genai"), + } + ) + ) + + +@lru_cache(maxsize=128) +def _warn_host_not_allowlisted(endpoint: str) -> None: + verbose_logger.warning( + "SigNoz: not exporting to key/team endpoint '%s'. Add its host to " + "litellm_settings.provider_url_destination_allowed_hosts to permit it", + endpoint, + ) + + +@lru_cache(maxsize=128) +def _warn_endpoint_without_key(endpoint: str) -> None: + verbose_logger.warning( + "SigNoz: not exporting to key/team endpoint '%s'. Set signoz_ingestion_key alongside it; " + "a keyless collector needs the global callback", + endpoint, + ) + + +def _tenant_endpoint_is_unusable(params: StandardCallbackDynamicParams) -> bool: + return bool(params.get("signoz_ingestion_endpoint")) and signoz_dynamic_endpoint(params) is None + + +def signoz_dynamic_endpoint(params: StandardCallbackDynamicParams) -> str | None: + endpoint: Final = params.get("signoz_ingestion_endpoint") + if not endpoint or not endpoint.startswith(("http://", "https://")): + return None + if not params.get("signoz_ingestion_key"): + _warn_endpoint_without_key(endpoint) + return None + if not is_url_destination_allowed_by_host(endpoint, litellm.provider_url_destination_allowed_hosts): + _warn_host_not_allowlisted(endpoint) + return None + return endpoint + + +def signoz_dynamic_headers( + params: StandardCallbackDynamicParams, +) -> dict[str, str]: # mutable-ok: DYNAMIC_HEADERS_BY_CALLBACK returns a dict + key: Final = params.get("signoz_ingestion_key") + if _tenant_endpoint_is_unusable(params) or not key: + return {} # mutable-ok: same registry contract + return {"signoz-ingestion-key": key} # mutable-ok: same registry contract diff --git a/litellm/litellm_core_utils/custom_logger_registry.py b/litellm/litellm_core_utils/custom_logger_registry.py index 7049fdd1f39..1d277995211 100644 --- a/litellm/litellm_core_utils/custom_logger_registry.py +++ b/litellm/litellm_core_utils/custom_logger_registry.py @@ -89,6 +89,7 @@ class CustomLoggerRegistry: "langtrace": OpenTelemetry, "weave_otel": OpenTelemetry, "levo": OpenTelemetry, + "signoz": OpenTelemetry, "mlflow": MlflowLogger, "langfuse": LangfusePromptManagement, "otel": OpenTelemetry, diff --git a/litellm/litellm_core_utils/initialize_dynamic_callback_params.py b/litellm/litellm_core_utils/initialize_dynamic_callback_params.py index 00ab05aba77..3100ca6fba1 100644 --- a/litellm/litellm_core_utils/initialize_dynamic_callback_params.py +++ b/litellm/litellm_core_utils/initialize_dynamic_callback_params.py @@ -113,6 +113,8 @@ _supported_callback_params: Final[tuple[str, ...]] = ( "dd_agent_port", "newrelic_api_key", "newrelic_region", + "signoz_ingestion_endpoint", + "signoz_ingestion_key", "turn_off_message_logging", ) @@ -126,6 +128,8 @@ _request_blocked_callback_params: Final = frozenset( "dd_agent_port", "newrelic_api_key", "newrelic_region", + "signoz_ingestion_endpoint", + "signoz_ingestion_key", } ) @@ -138,6 +142,8 @@ _trusted_overlay_callback_params: Final = frozenset( { "newrelic_api_key", "newrelic_region", + "signoz_ingestion_endpoint", + "signoz_ingestion_key", } ) diff --git a/litellm/litellm_core_utils/litellm_logging.py b/litellm/litellm_core_utils/litellm_logging.py index 9ee7a7b0a7a..152fd54e55d 100644 --- a/litellm/litellm_core_utils/litellm_logging.py +++ b/litellm/litellm_core_utils/litellm_logging.py @@ -4931,6 +4931,38 @@ def _init_custom_logger_compatible_class( _in_memory_loggers.append(_otel_logger) return _otel_logger + elif logging_integration == "signoz": + from litellm.integrations.otel.presets.signoz import ( + SIGNOZ_INGESTION_ENDPOINT_ENV, + ) + + _signoz_endpoint: Final = os.getenv(SIGNOZ_INGESTION_ENDPOINT_ENV) + if not _signoz_endpoint: + raise ValueError(f"{SIGNOZ_INGESTION_ENDPOINT_ENV} not found in environment variables") + + _signoz_v2: Final = _maybe_construct_otel_v2("signoz", _in_memory_loggers) + if _signoz_v2 is not None: + return _signoz_v2 + + from litellm.integrations.opentelemetry import ( + OpenTelemetry, + OpenTelemetryConfig, + ) + + _signoz_base: Final = _signoz_endpoint.rstrip("/") + _signoz_key: Final = os.getenv("SIGNOZ_INGESTION_KEY") + _signoz_config: Final = OpenTelemetryConfig( + exporter="otlp_http", + endpoint=(_signoz_base if _signoz_base.endswith("/v1/traces") else f"{_signoz_base}/v1/traces"), + headers=(f"signoz-ingestion-key={_signoz_key}" if _signoz_key else None), + ) + for callback in _in_memory_loggers: + if isinstance(callback, OpenTelemetry) and callback.callback_name == "signoz": + return callback + _signoz_logger: Final = OpenTelemetry(config=_signoz_config, callback_name="signoz") + _in_memory_loggers.append(_signoz_logger) + return _signoz_logger + elif logging_integration == "mlflow": for callback in _in_memory_loggers: if isinstance(callback, MlflowLogger): diff --git a/litellm/proxy/_experimental/out/assets/logos/signoz.svg b/litellm/proxy/_experimental/out/assets/logos/signoz.svg new file mode 100644 index 00000000000..9064cb86bd6 --- /dev/null +++ b/litellm/proxy/_experimental/out/assets/logos/signoz.svg @@ -0,0 +1 @@ + \ No newline at end of file diff --git a/litellm/proxy/_types.py b/litellm/proxy/_types.py index 14aa42afefd..34d7fc1e0f0 100644 --- a/litellm/proxy/_types.py +++ b/litellm/proxy/_types.py @@ -4035,6 +4035,12 @@ class AllCallbacks(LiteLLMPydanticObjectBase): ], ) + signoz: CallbackOnUI = CallbackOnUI( + litellm_callback_name="signoz", + ui_callback_name="SigNoz", + litellm_callback_params=("SIGNOZ_INGESTION_ENDPOINT", "SIGNOZ_INGESTION_KEY"), + ) + zerobus: CallbackOnUI = CallbackOnUI( litellm_callback_name="zerobus", ui_callback_name="Databricks Zerobus", diff --git a/litellm/proxy/health_endpoints/_health_endpoints.py b/litellm/proxy/health_endpoints/_health_endpoints.py index fbd4d57bf77..07be73d7573 100644 --- a/litellm/proxy/health_endpoints/_health_endpoints.py +++ b/litellm/proxy/health_endpoints/_health_endpoints.py @@ -221,6 +221,7 @@ services = ( "galileo", "newrelic", "pointfive", + "signoz", "sqs", ] | str @@ -309,6 +310,7 @@ async def health_services_endpoint( "galileo", "newrelic", "pointfive", + "signoz", "sqs", ]: raise HTTPException( diff --git a/litellm/proxy/litellm_pre_call_utils.py b/litellm/proxy/litellm_pre_call_utils.py index 56f647d5acc..866d84ca8f2 100644 --- a/litellm/proxy/litellm_pre_call_utils.py +++ b/litellm/proxy/litellm_pre_call_utils.py @@ -924,6 +924,8 @@ def convert_key_logging_metadata_to_callback( # must not export to it. if var.startswith("newrelic_") and data.callback_name != "newrelic": continue + if var.startswith("signoz_") and data.callback_name != "signoz": + continue if team_callback_settings_obj.callback_vars is None: team_callback_settings_obj.callback_vars = {} team_callback_settings_obj.callback_vars[var] = str(value) diff --git a/litellm/types/utils.py b/litellm/types/utils.py index 4bda1dd53ce..8862df9dc22 100644 --- a/litellm/types/utils.py +++ b/litellm/types/utils.py @@ -3679,6 +3679,9 @@ class StandardCallbackDynamicParams(TypedDict, total=False): newrelic_api_key: str | None # writable-ok: initialize_standard_callback_dynamic_params assigns into the dict newrelic_region: str | None # writable-ok: initialize_standard_callback_dynamic_params assigns into the dict + signoz_ingestion_endpoint: str | None # writable-ok: initialize_standard_callback_dynamic_params assigns it + signoz_ingestion_key: str | None # writable-ok: initialize_standard_callback_dynamic_params assigns it + # Logging settings turn_off_message_logging: bool | None # when true will not log messages litellm_disabled_callbacks: list[str] | None diff --git a/tests/integration/observability/test_signoz_delivery.py b/tests/integration/observability/test_signoz_delivery.py new file mode 100644 index 00000000000..f3d715fe5cf --- /dev/null +++ b/tests/integration/observability/test_signoz_delivery.py @@ -0,0 +1,980 @@ +import asyncio +import base64 +import json +import os +import re +import signal +import threading +import uuid +from collections import deque +from collections.abc import Iterator, Mapping, Sequence +from concurrent.futures import ThreadPoolExecutor +from dataclasses import dataclass +from functools import partial +from pathlib import Path +from typing import Final + +import anthropic +import httpx +import openai +import psutil +import pytest +import yaml +from integration._support.client import Gateway, eventually, gateway_from_environment, object_value, string_value +from integration._support.database import read_rows +from integration._support.process import OwnedProxy, owned_proxy_process +from integration._support.wire import Reply, Request, Wire, wire_server +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest +from opentelemetry.proto.common.v1.common_pb2 import AnyValue +from pydantic import JsonValue, TypeAdapter + +MARKER: Final = re.compile(rb"signoz-[0-9a-f]{32}") +JSON: Final[TypeAdapter[JsonValue]] = TypeAdapter(JsonValue) +RESPONSE_ID: Final = "gen_ai.response.id" +INGESTION_HEADER: Final = "signoz-ingestion-key" +OPERATOR_KEY: Final = "operator-ingestion-" + uuid.uuid4().hex +TENANT_KEY: Final = "tenant-ingestion-" + uuid.uuid4().hex + + +def _marker() -> str: + return "signoz-" + uuid.uuid4().hex + + +def _chat_reply(identity: str, stream: bool) -> Reply: + if not stream: + return Reply( + body=json.dumps( + { + "id": identity, + "object": "chat.completion", + "created": 1, + "model": "gpt-4o-mini", + "choices": [ + {"index": 0, "message": {"role": "assistant", "content": "signoz ok"}, "finish_reason": "stop"} + ], + "usage": {"prompt_tokens": 7, "completion_tokens": 2, "total_tokens": 9}, + } + ).encode() + ) + chunk: Final = {"id": identity, "object": "chat.completion.chunk", "created": 1, "model": "gpt-4o-mini"} + return Reply( + content_type="text/event-stream", + chunks=( + b"data: " + + json.dumps( + {**chunk, "choices": [{"index": 0, "delta": {"role": "assistant", "content": "signoz"}}]} + ).encode() + + b"\n\n", + b"data: " + + json.dumps( + {**chunk, "choices": [{"index": 0, "delta": {"content": " ok"}, "finish_reason": "stop"}]} + ).encode() + + b"\n\n", + b"data: " + + json.dumps( + {**chunk, "choices": [], "usage": {"prompt_tokens": 7, "completion_tokens": 2, "total_tokens": 9}} + ).encode() + + b"\n\n", + b"data: [DONE]\n\n", + ), + ) + + +def _responses_reply(identity: str, stream: bool) -> Reply: + response: Final[dict[str, JsonValue]] = { + "id": identity, + "object": "response", + "created_at": 1, + "status": "completed", + "model": "gpt-4o-mini", + "output": [ + { + "id": "msg_" + identity, + "type": "message", + "role": "assistant", + "status": "completed", + "content": [{"type": "output_text", "text": "signoz ok", "annotations": []}], + } + ], + "usage": {"input_tokens": 7, "output_tokens": 2, "total_tokens": 9}, + } + if not stream: + return Reply(body=json.dumps(response).encode()) + events: Final[tuple[dict[str, JsonValue], ...]] = ( + { + "type": "response.created", + "sequence_number": 0, + "response": {**response, "status": "in_progress", "output": []}, + }, + { + "type": "response.output_text.delta", + "sequence_number": 1, + "item_id": "msg_" + identity, + "output_index": 0, + "content_index": 0, + "delta": "signoz ok", + }, + {"type": "response.completed", "sequence_number": 2, "response": response}, + ) + return Reply( + content_type="text/event-stream", + chunks=tuple(f"event: {event['type']}\ndata: {json.dumps(event)}\n\n".encode() for event in events), + ) + + +def _upstream(request: Request) -> Reply: + found: Final = MARKER.search(request.body) + if found is None: + return Reply(status=404, body=b'{"error":"no marker"}') + if request.headers.get("authorization") == "Bearer revoked-provider-key": + return Reply( + status=401, body=b'{"error":{"message":"Incorrect API key provided","type":"invalid_request_error"}}' + ) + marker: Final = found.group(0).decode() + stream: Final = object_value(JSON.validate_json(request.body)).get("stream") is True + if request.target.endswith("/responses"): + return _responses_reply(f"resp_{marker}", stream) + return _chat_reply(f"chatcmpl-{marker}", stream) + + +def _decoded_responses_id(identity: str) -> str: + try: + return base64.b64decode(identity.removeprefix("resp_").encode()).decode() + except (ValueError, UnicodeDecodeError): + return identity + + +def _canonical_id(identity: str) -> str: + return _decoded_responses_id(identity).rpartition("response_id:")[2] + + +def _sse_events(text: str) -> tuple[dict[str, JsonValue], ...]: + return tuple( + object_value(JSON.validate_json(line[6:])) + for line in text.splitlines() + if line.startswith("data: ") and line != "data: [DONE]" + ) + + +def _text_at(payload: JsonValue, *path: str) -> str: + if not path: + return string_value(payload) + return _text_at(object_value(payload)[path[0]], *path[1:]) + + +def _body_id(response: httpx.Response) -> str: + return _text_at(JSON.validate_json(response.content), "id") + + +@dataclass(frozen=True, slots=True) +class Span: + target: str + ingestion_key: str | None + attributes: Mapping[str, str] + + +def _attribute_text(value: AnyValue) -> str: + match value.WhichOneof("value"): + case "string_value": + return value.string_value + case "int_value": + return str(value.int_value) + case "double_value": + return str(value.double_value) + case "bool_value": + return str(value.bool_value) + case _: + return "" + + +@dataclass(frozen=True, slots=True) +class Collector: + wire: Wire + outage: threading.Event + rejection: threading.Event + missing: threading.Event + slow: threading.Event + release: threading.Event + accepted: Sequence[Request] + refused: Sequence[Request] + guard: threading.Lock + + def refused_batch_carrying(self, response_id: str) -> Request: + def carrying() -> tuple[Request, ...]: + with self.guard: + return tuple(batch for batch in self.refused if response_id.encode() in batch.body) + + return eventually(carrying, lambda found: len(found) >= 1, seconds=30)[0] + + def refused_batches(self) -> tuple[Request, ...]: + def refused() -> tuple[Request, ...]: + with self.guard: + return tuple(self.refused) + + return eventually(refused, lambda found: len(found) >= 1, seconds=30) + + def spans(self) -> tuple[Span, ...]: + with self.guard: + batches: Final = tuple(self.accepted) + return tuple( + Span( + batch.target, + batch.headers.get(INGESTION_HEADER), + {attribute.key: _attribute_text(attribute.value) for attribute in span.attributes}, + ) + for batch in batches + for resource in ExportTraceServiceRequest.FromString(batch.body).resource_spans + for scope in resource.scope_spans + for span in scope.spans + ) + + def spans_for(self, response_id: str) -> tuple[Span, ...]: + return tuple( + span + for span in self.spans() + if RESPONSE_ID in span.attributes + and _canonical_id(span.attributes[RESPONSE_ID]) == _canonical_id(response_id) + ) + + def single_span(self, response_id: str, *, elsewhere: "Collector | None" = None) -> Span: + found: Final = eventually( + lambda: self.spans_for(response_id), lambda spans: len(spans) == 1, seconds=30, return_last_on_timeout=True + ) + assert len(found) == 1, ( + f"{len(found)} spans for {response_id} at this sink; other sink saw " + f"{elsewhere.landed((response_id,)) if elsewhere else 'n/a'}" + ) + return found[0] + + def landed(self, response_ids: Sequence[str]) -> dict[str, int]: + spans: Final = self.spans() + return { + _canonical_id(identity): sum( + 1 + for span in spans + if RESPONSE_ID in span.attributes + and _canonical_id(span.attributes[RESPONSE_ID]) == _canonical_id(identity) + ) + for identity in response_ids + } + + +def _collector() -> Iterator[Collector]: + outage: Final = threading.Event() + rejection: Final = threading.Event() + missing: Final = threading.Event() + slow: Final = threading.Event() + release: Final = threading.Event() + accepted: Final[deque[Request]] = deque() # mutable-ok: the sink thread records each accepted batch as it arrives + refused: Final[deque[Request]] = deque() # mutable-ok: the sink thread records each refused batch as it arrives + guard: Final = threading.Lock() + + def refuse(request: Request, status: int, body: bytes) -> Reply: + with guard: + refused.append(request) + return Reply(status=status, body=body) + + def sink(request: Request) -> Reply: + if slow.is_set(): + release.wait(timeout=30) + if outage.is_set(): + return refuse(request, 503, b'{"error":"sink down"}') + if rejection.is_set(): + return refuse(request, 403, b'{"error":"forbidden"}') + if missing.is_set(): + return refuse(request, 404, b'{"error":"not found"}') + with guard: + accepted.append(request) + return Reply() + + with wire_server(sink) as wire: + yield Collector(wire, outage, rejection, missing, slow, release, accepted, refused, guard) + + +@pytest.fixture(scope="session") +def operator_sink() -> Iterator[Collector]: + yield from _collector() + + +@pytest.fixture(scope="session") +def tenant_sink() -> Iterator[Collector]: + yield from _collector() + + +@pytest.fixture(scope="session") +def provider() -> Iterator[Wire]: + with wire_server(_upstream) as wire: + yield wire + + +@dataclass(frozen=True, slots=True) +class Rig: + proxy: Gateway + process: OwnedProxy + model: str + upstream: Wire + sink: Collector + tenant_sink: Collector + + def openai_client(self) -> openai.OpenAI: + return openai.OpenAI(base_url=str(self.proxy.client.base_url) + "/v1", api_key=self.proxy.key, max_retries=0) + + def async_openai_client(self) -> openai.AsyncOpenAI: + return openai.AsyncOpenAI( + base_url=str(self.proxy.client.base_url) + "/v1", api_key=self.proxy.key, max_retries=0 + ) + + def anthropic_client(self) -> anthropic.Anthropic: + return anthropic.Anthropic(base_url=str(self.proxy.client.base_url), api_key=self.proxy.key, max_retries=0) + + def async_anthropic_client(self) -> anthropic.AsyncAnthropic: + return anthropic.AsyncAnthropic(base_url=str(self.proxy.client.base_url), api_key=self.proxy.key, max_retries=0) + + def chat( + self, marker: str, *, headers: Mapping[str, str] | None = None, key: str | None = None, **extra: JsonValue + ) -> httpx.Response: + return self.proxy.request( + "POST", + "/v1/chat/completions", + { + "model": self.model, + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + **extra, + }, + headers=headers, + key=key, + ) + + def upstream_bodies(self, marker: str) -> tuple[dict[str, JsonValue], ...]: + return tuple( + object_value(JSON.validate_json(request.body)) + for request in self.upstream.drain() + if marker.encode() in request.body + ) + + def spend_rows(self, response_id: str) -> tuple[dict[str, JsonValue], ...]: + return tuple( + eventually( + lambda: read_rows( + 'SELECT request_id, spend FROM "LiteLLM_SpendLogs" WHERE request_id=%s', (response_id,) + ), + lambda values: len(values) == 1, + seconds=70, + ) + ) + + def tenant_logging(self, endpoint: str | None, key: str | None) -> JsonValue: + variables: Final[dict[str, JsonValue]] = { + **({"signoz_ingestion_endpoint": endpoint} if endpoint is not None else {}), + **({"signoz_ingestion_key": key} if key is not None else {}), + } + return [{"callback_name": "signoz", "callback_type": "success", "callback_vars": variables}] + + +@dataclass(frozen=True, slots=True) +class RigFactory: + provider: Wire + sink: Collector + tenant_sink: Collector + directory: Path + otel_v2: bool + workers: int + endpoint: str | None + ingestion_key: str | None = OPERATOR_KEY + + def config_path(self) -> Path: + loaded: Final = object_value( + JSON.validate_python(yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())) + ) + config: Final = { + **loaded, + "litellm_settings": { + **object_value(loaded["litellm_settings"]), + "callbacks": ["signoz"], + "provider_url_destination_allowed_hosts": [self.tenant_sink.wire.url], + }, + "general_settings": {**object_value(loaded["general_settings"]), "disable_model_info_refresh": True}, + } + path: Final = self.directory / f"signoz-{uuid.uuid4().hex}.yaml" + path.write_text(yaml.safe_dump(config)) + return path + + def overrides(self) -> dict[str, str]: + return { + "LITELLM_OTEL_V2": "1" if self.otel_v2 else "0", + "OTEL_BSP_SCHEDULE_DELAY": "300", + **({"SIGNOZ_INGESTION_ENDPOINT": self.endpoint} if self.endpoint is not None else {}), + **({"SIGNOZ_INGESTION_KEY": self.ingestion_key} if self.ingestion_key is not None else {}), + } + + def start(self) -> Iterator[Rig]: + with ( + gateway_from_environment() as gateway, + owned_proxy_process( + gateway, + self.directory, + self.overrides(), + config=self.config_path(), + remove_environment=("SIGNOZ_INGESTION_ENDPOINT", "SIGNOZ_INGESTION_KEY"), + workers=self.workers, + ) as owned, + owned.gateway.scenario() as scenario, + ): + model: Final = scenario.model(api_base=self.provider.url + "/v1") + yield Rig(owned.gateway, owned, model, self.provider, self.sink, self.tenant_sink) + + +@pytest.fixture(scope="session") +def rig( + provider: Wire, operator_sink: Collector, tenant_sink: Collector, tmp_path_factory: pytest.TempPathFactory +) -> Iterator[Rig]: + factory: Final = RigFactory( + provider, operator_sink, tenant_sink, tmp_path_factory.mktemp("signoz"), False, 2, operator_sink.wire.url + ) + yield from factory.start() + + +@pytest.fixture(scope="session") +def v2_rig( + provider: Wire, operator_sink: Collector, tenant_sink: Collector, tmp_path_factory: pytest.TempPathFactory +) -> Iterator[Rig]: + factory: Final = RigFactory( + provider, operator_sink, tenant_sink, tmp_path_factory.mktemp("signoz-v2"), True, 2, operator_sink.wire.url + ) + yield from factory.start() + + +def _assert_operator_span(rig: Rig, response_id: str, marker: str) -> Span: + span: Final = rig.sink.single_span(response_id) + assert span.target == "/v1/traces", span + assert span.ingestion_key == OPERATOR_KEY, span + assert rig.tenant_sink.landed((response_id,)) == {_canonical_id(response_id): 0} + bodies: Final = rig.upstream_bodies(marker) + assert len(bodies) == 1, bodies + assert "signoz" not in json.dumps(bodies[0]).replace(marker, ""), bodies[0] + return span + + +def test_signoz_is_registered_as_an_opentelemetry_callback(rig: Rig) -> None: + listed: Final = rig.proxy.request("GET", "/active/callbacks") + assert listed.status_code == 200, listed.text + assert "OpenTelemetry" in json.dumps(listed.json()), listed.text + log: Final = rig.process.log.read_text() + assert "SIGNOZ_INGESTION_ENDPOINT not found" not in log + + +def test_chat_completion_sdk_span_lands_at_the_operator_sink_with_the_ingestion_key(rig: Rig) -> None: + marker: Final = _marker() + completion: Final = rig.openai_client().chat.completions.create( + model=rig.model, messages=[{"role": "user", "content": marker}] + ) + assert completion.id == f"chatcmpl-{marker}" + span: Final = _assert_operator_span(rig, completion.id, marker) + assert span.attributes.get("gen_ai.request.model") or span.attributes.get("llm.request.model"), span + rows: Final = rig.spend_rows(completion.id) + assert rows[0]["request_id"] == completion.id, rows + + +def test_chat_stream_async_sdk_span_lands_once_after_the_stream_is_consumed(rig: Rig) -> None: + marker: Final = _marker() + + async def consume() -> frozenset[str]: + stream: Final = await rig.async_openai_client().chat.completions.create( + model=rig.model, messages=[{"role": "user", "content": marker}], stream=True + ) + return frozenset([chunk.id async for chunk in stream]) + + identities: Final = asyncio.run(consume()) + assert identities == {f"chatcmpl-{marker}"}, identities + _assert_operator_span(rig, f"chatcmpl-{marker}", marker) + + +def test_messages_sdk_span_lands_at_the_operator_sink(rig: Rig) -> None: + marker: Final = _marker() + message: Final = rig.anthropic_client().messages.create( + model=rig.model, max_tokens=16, messages=[{"role": "user", "content": marker}] + ) + _assert_operator_span(rig, message.id, marker) + + +def test_messages_stream_async_sdk_span_lands_once_after_the_stream_is_consumed(rig: Rig) -> None: + marker: Final = _marker() + + async def consume() -> str: + async with rig.async_anthropic_client().messages.stream( + model=rig.model, max_tokens=16, messages=[{"role": "user", "content": marker}] + ) as stream: + async for _ in stream: + pass + return (await stream.get_final_message()).id + + identity: Final = asyncio.run(consume()) + _assert_operator_span(rig, identity, marker) + + +def test_responses_sdk_span_lands_at_the_operator_sink(rig: Rig) -> None: + marker: Final = _marker() + response: Final = rig.openai_client().responses.create(model=rig.model, input=marker) + _assert_operator_span(rig, response.id, marker) + + +def test_responses_stream_raw_httpx_span_lands_once_after_the_stream_is_consumed(rig: Rig) -> None: + marker: Final = _marker() + response: Final = rig.proxy.request("POST", "/v1/responses", {"model": rig.model, "input": marker, "stream": True}) + assert response.status_code == 200, response.text + _assert_operator_span(rig, _responses_id(response, marker), marker) + + +def test_v2_flag_on_still_delivers_the_operator_span_with_the_ingestion_key(v2_rig: Rig) -> None: + marker: Final = _marker() + response: Final = v2_rig.chat(marker) + assert response.status_code == 200, response.text + _assert_operator_span(v2_rig, _body_id(response), marker) + + +def test_endpoint_already_ending_in_v1_traces_is_not_doubled( + provider: Wire, operator_sink: Collector, tenant_sink: Collector, tmp_path_factory: pytest.TempPathFactory +) -> None: + factory: Final = RigFactory( + provider, + operator_sink, + tenant_sink, + tmp_path_factory.mktemp("signoz-suffixed"), + False, + 2, + operator_sink.wire.url + "/v1/traces", + ) + suffixed: Final = next(started := factory.start()) + marker: Final = _marker() + response: Final = suffixed.chat(marker) + assert response.status_code == 200, response.text + span: Final = suffixed.sink.single_span(_body_id(response)) + assert span.target == "/v1/traces", span + assert tuple(started) == () + + +def test_three_identical_requests_produce_one_span_each(rig: Rig) -> None: + markers: Final = tuple(_marker() for _ in range(3)) + responses: Final = tuple(rig.chat(marker) for marker in markers) + assert all(response.status_code == 200 for response in responses), [response.text for response in responses] + identities: Final = tuple(_body_id(response) for response in responses) + landed: Final = eventually( + lambda: rig.sink.landed(identities), lambda seen: all(count >= 1 for count in seen.values()), seconds=30 + ) + assert landed == {identity: 1 for identity in identities}, landed + assert rig.sink.landed(identities) == landed + + +def test_unauthenticated_request_is_rejected_without_an_upstream_call_and_any_span_records_the_401(rig: Rig) -> None: + marker: Final = _marker() + response: Final = rig.chat(marker, key="sk-not-a-real-key") + assert response.status_code == 401, response.text + later: Final = rig.chat(_marker()) + assert later.status_code == 200, later.text + rig.sink.single_span(_body_id(later)) + assert rig.upstream_bodies(marker) == () + marker_spans: Final = tuple(span for span in rig.sink.spans() if marker in json.dumps(span.attributes)) + assert all(span.attributes.get("error.code") == "401" for span in marker_spans), marker_spans + assert not any(span.attributes.get(RESPONSE_ID, "").startswith("chatcmpl-") for span in marker_spans), marker_spans + + +def test_request_supplied_signoz_variables_are_refused_before_the_upstream_is_called(rig: Rig) -> None: + marker: Final = _marker() + response: Final = rig.chat( + marker, + metadata={"signoz_ingestion_endpoint": rig.tenant_sink.wire.url, "signoz_ingestion_key": TENANT_KEY}, + ) + assert response.status_code == 401, response.text + assert "signoz_ingestion_endpoint is not allowed in request body" in response.text + assert rig.upstream_bodies(marker) == () + assert not any(marker in json.dumps(span.attributes) for span in rig.tenant_sink.spans()) + + +def test_upstream_401_reaches_the_caller_and_unrelated_traffic_keeps_landing(rig: Rig) -> None: + marker: Final = _marker() + with rig.proxy.scenario() as scenario: + broken: Final = scenario.model(api_base=rig.upstream.url + "/v1", api_key="revoked-provider-key") + failed: Final = rig.proxy.request( + "POST", "/v1/chat/completions", {"model": broken, "messages": [{"role": "user", "content": marker}]} + ) + assert failed.status_code == 401, failed.text + assert "Incorrect API key provided" in failed.text + healthy_marker: Final = _marker() + healthy: Final = rig.chat(healthy_marker) + assert healthy.status_code == 200, healthy.text + _assert_operator_span(rig, _body_id(healthy), healthy_marker) + + +def test_health_services_accepts_signoz(rig: Rig) -> None: + response: Final = rig.proxy.request("GET", "/health/services", params={"service": "signoz"}) + assert response.status_code == 200, response.text + + +def test_sink_answering_403_drops_those_spans_and_later_spans_still_land(rig: Rig) -> None: + rig.sink.rejection.set() + try: + rejected: Final = rig.chat(_marker()) + assert rejected.status_code == 200, rejected.text + rig.sink.refused_batch_carrying(_body_id(rejected)) + finally: + rig.sink.rejection.clear() + later: Final = rig.chat(_marker()) + assert later.status_code == 200, later.text + rig.sink.single_span(_body_id(later)) + assert rig.sink.landed((_body_id(rejected),)) == {_body_id(rejected): 0} + + +def test_sink_answering_404_drops_those_spans_and_later_spans_still_land(rig: Rig) -> None: + rig.sink.missing.set() + try: + dropped: Final = rig.chat(_marker()) + assert dropped.status_code == 200, dropped.text + rig.sink.refused_batch_carrying(_body_id(dropped)) + finally: + rig.sink.missing.clear() + later: Final = rig.chat(_marker()) + assert later.status_code == 200, later.text + rig.sink.single_span(_body_id(later)) + assert rig.sink.landed((_body_id(dropped),)) == {_body_id(dropped): 0} + + +def test_key_level_signoz_destination_routes_the_span_to_the_tenant_sink(v2_rig: Rig) -> None: + marker: Final = _marker() + with v2_rig.proxy.scenario() as scenario: + token: Final = scenario.key( + metadata={"logging": v2_rig.tenant_logging(v2_rig.tenant_sink.wire.url, TENANT_KEY)} + ) + response: Final = v2_rig.chat(marker, key=token) + assert response.status_code == 200, response.text + identity: Final = _body_id(response) + span: Final = v2_rig.tenant_sink.single_span(identity, elsewhere=v2_rig.sink) + assert span.ingestion_key == TENANT_KEY, span + assert v2_rig.sink.landed((identity,)) == {identity: 0}, "operator sink also received the tenant span" + + +def test_team_level_signoz_destination_routes_the_span_to_the_tenant_sink(v2_rig: Rig) -> None: + marker: Final = _marker() + with v2_rig.proxy.scenario() as scenario: + team: Final = scenario.team( + metadata={"logging": v2_rig.tenant_logging(v2_rig.tenant_sink.wire.url, TENANT_KEY)} + ) + token: Final = scenario.key(team_id=team) + response: Final = v2_rig.chat(marker, key=token) + assert response.status_code == 200, response.text + identity: Final = _body_id(response) + span: Final = v2_rig.tenant_sink.single_span(identity, elsewhere=v2_rig.sink) + assert span.ingestion_key == TENANT_KEY, span + assert v2_rig.sink.landed((identity,)) == {identity: 0}, "operator sink also received the tenant span" + + +def test_key_level_destination_wins_over_the_team_level_destination(v2_rig: Rig) -> None: + marker: Final = _marker() + team_key: Final = "team-" + TENANT_KEY + with v2_rig.proxy.scenario() as scenario: + team: Final = scenario.team(metadata={"logging": v2_rig.tenant_logging(v2_rig.tenant_sink.wire.url, team_key)}) + token: Final = scenario.key( + team_id=team, metadata={"logging": v2_rig.tenant_logging(v2_rig.tenant_sink.wire.url, TENANT_KEY)} + ) + response: Final = v2_rig.chat(marker, key=token) + assert response.status_code == 200, response.text + span: Final = v2_rig.tenant_sink.single_span(_body_id(response)) + assert span.ingestion_key == TENANT_KEY, span + + +def test_team_endpoint_without_an_ingestion_key_is_ignored_and_the_span_stays_at_the_operator_sink( + v2_rig: Rig, +) -> None: + marker: Final = _marker() + with v2_rig.proxy.scenario() as scenario: + team: Final = scenario.team(metadata={"logging": v2_rig.tenant_logging(v2_rig.tenant_sink.wire.url, None)}) + token: Final = scenario.key(team_id=team) + response: Final = v2_rig.chat(marker, key=token) + assert response.status_code == 200, response.text + _assert_operator_span(v2_rig, _body_id(response), marker) + eventually( + lambda: v2_rig.process.log.read_text(), + lambda text: "Set signoz_ingestion_key alongside it" in text, + seconds=30, + ) + + +def test_team_endpoint_off_the_allowlist_keeps_the_span_at_the_operator_sink(v2_rig: Rig) -> None: + marker: Final = _marker() + with v2_rig.proxy.scenario() as scenario: + team: Final = scenario.team( + metadata={"logging": v2_rig.tenant_logging("http://tenant.invalid:4318/v1/traces", TENANT_KEY)} + ) + token: Final = scenario.key(team_id=team) + response: Final = v2_rig.chat(marker, key=token) + assert response.status_code == 200, response.text + _assert_operator_span(v2_rig, _body_id(response), marker) + eventually( + lambda: v2_rig.process.log.read_text(), + lambda text: "provider_url_destination_allowed_hosts" in text, + seconds=30, + ) + + +def test_legacy_mode_ignores_key_level_signoz_destination_and_keeps_the_operator_sink(rig: Rig) -> None: + marker: Final = _marker() + with rig.proxy.scenario() as scenario: + token: Final = scenario.key(metadata={"logging": rig.tenant_logging(rig.tenant_sink.wire.url, TENANT_KEY)}) + response: Final = rig.chat(marker, key=token) + assert response.status_code == 200, response.text + _assert_operator_span(rig, _body_id(response), marker) + + +@pytest.mark.parametrize( + "endpoint", + ["", "not-a-url", "ftp://tenant.invalid", "x" * 5000, 12345, ["http://tenant.invalid"]], + ids=["empty", "bare", "ftp", "5kb", "int", "list"], +) +def test_hostile_tenant_endpoint_never_breaks_the_request_or_the_operator_sink( + v2_rig: Rig, endpoint: JsonValue +) -> None: + marker: Final = _marker() + created: Final = v2_rig.proxy.request( + "POST", + "/key/generate", + { + "metadata": { + "logging": [ + { + "callback_name": "signoz", + "callback_type": "success", + "callback_vars": {"signoz_ingestion_endpoint": endpoint, "signoz_ingestion_key": TENANT_KEY}, + } + ] + } + }, + ) + assert created.status_code in (200, 400, 422), created.text + if created.status_code != 200: + return + try: + response: Final = v2_rig.chat(marker, key=_text_at(JSON.validate_json(created.content), "key")) + assert response.status_code == 200, response.text + identity: Final = _body_id(response) + eventually( + lambda: v2_rig.sink.landed((identity,))[identity] + v2_rig.tenant_sink.landed((identity,))[identity], + lambda total: total >= 1, + seconds=30, + ) + assert v2_rig.sink.landed((identity,))[identity] + v2_rig.tenant_sink.landed((identity,))[identity] == 1 + assert v2_rig.tenant_sink.landed((identity,)) == {identity: 0}, "unusable endpoint reached the tenant sink" + finally: + v2_rig.proxy.post("/key/delete", {"keys": [created.json()["key"]]}) + + +def test_missing_ingestion_endpoint_fails_loudly_at_boot( + provider: Wire, operator_sink: Collector, tenant_sink: Collector, tmp_path_factory: pytest.TempPathFactory +) -> None: + factory: Final = RigFactory( + provider, operator_sink, tenant_sink, tmp_path_factory.mktemp("signoz-noenv"), False, 2, None + ) + started: Final = factory.start() + broken: Final = next(started) + try: + response: Final = broken.chat(_marker()) + assert response.status_code == 200, response.text + eventually( + lambda: broken.process.log.read_text(), + lambda text: "SIGNOZ_INGESTION_ENDPOINT not found" in text, + seconds=30, + ) + assert operator_sink.landed((_body_id(response),)) == {_body_id(response): 0} + finally: + with pytest.raises(StopIteration): + next(started) + + +def test_empty_ingestion_endpoint_is_treated_as_missing( + provider: Wire, operator_sink: Collector, tenant_sink: Collector, tmp_path_factory: pytest.TempPathFactory +) -> None: + factory: Final = RigFactory( + provider, operator_sink, tenant_sink, tmp_path_factory.mktemp("signoz-empty"), False, 2, "" + ) + started: Final = factory.start() + broken: Final = next(started) + try: + response: Final = broken.chat(_marker()) + assert response.status_code == 200, response.text + eventually( + lambda: broken.process.log.read_text(), + lambda text: "SIGNOZ_INGESTION_ENDPOINT not found" in text, + seconds=30, + ) + finally: + with pytest.raises(StopIteration): + next(started) + + +def _is_event_stream(response: httpx.Response) -> bool: + return "content-type" in response.headers and response.headers["content-type"].startswith("text/event-stream") + + +def _chat_id(response: httpx.Response) -> str: + if not _is_event_stream(response): + return _body_id(response) + identities: Final = frozenset(_text_at(event, "id") for event in _sse_events(response.text)) + assert len(identities) == 1, response.text + return next(iter(identities)) + + +def _responses_id(response: httpx.Response, marker: str) -> str: + if not _is_event_stream(response): + return _body_id(response) + completed: Final = tuple( + _text_at(event, "response", "id") + for event in _sse_events(response.text) + if event.get("type") == "response.completed" + ) + assert len(completed) == 1 and completed[0].startswith("resp_"), response.text + return f"resp_{marker}" + + +def _message_id(response: httpx.Response) -> str: + if not _is_event_stream(response): + return _body_id(response) + starts: Final = tuple( + _text_at(event, "message", "id") for event in _sse_events(response.text) if event.get("type") == "message_start" + ) + assert len(starts) == 1, response.text + return starts[0] + + +def _burst(rig: Rig, count: int) -> tuple[tuple[int, str, str | None], ...]: + markers: Final = tuple(_marker() for _ in range(count)) + + def one(index: int) -> tuple[int, str, str | None]: + marker: Final = markers[index] + headers: Final = {"Authorization": f"Bearer {rig.proxy.key}"} + stream: Final = index % 2 == 0 + path, body, identity_of = ( + ("/v1/chat/completions", {"model": rig.model, "messages": [{"role": "user", "content": marker}]}, _chat_id), + ("/v1/responses", {"model": rig.model, "input": marker}, partial(_responses_id, marker=marker)), + ( + "/v1/messages", + {"model": rig.model, "max_tokens": 16, "messages": [{"role": "user", "content": marker}]}, + _message_id, + ), + )[index % 3] + try: + response: Final = rig.proxy.client.post(path, json={**body, "stream": stream}, headers=headers) + response.read() + except httpx.HTTPError as error: + return index, marker, repr(error) + return (index, marker, response.text) if response.status_code != 200 else (index, identity_of(response), None) + + with ThreadPoolExecutor(max_workers=10) as pool: + return tuple(pool.map(one, range(count))) + + +def _assert_exactly_once(rig: Rig, identities: Sequence[str]) -> None: + landed: Final = eventually( + lambda: rig.sink.landed(identities), lambda seen: all(count >= 1 for count in seen.values()), seconds=80 + ) + assert landed == {_canonical_id(identity): 1 for identity in identities}, landed + charged: Final = frozenset(identity for identity in identities if not identity.startswith("resp_")) + spend: Final = eventually( + lambda: read_rows( + 'SELECT request_id FROM "LiteLLM_SpendLogs" WHERE request_id = ANY(%s::text[])', + ("{" + ",".join(charged) + "}",), + ), + lambda rows: {str(row["request_id"]) for row in rows} >= charged, + seconds=70, + ) + assert {str(row["request_id"]) for row in spend} == charged, spend + + +def test_sink_outage_during_a_mixed_burst_lands_every_response_exactly_once_after_recovery(rig: Rig) -> None: + rig.sink.outage.set() + try: + health_down: Final = rig.proxy.request("GET", "/health/services", params={"service": "signoz"}) + results: Final = _burst(rig, 30) + assert all(error is None for _, _, error in results), [error for _, _, error in results if error] + rig.sink.refused_batches() + finally: + rig.sink.outage.clear() + assert health_down.status_code == 200, health_down.text + _assert_exactly_once(rig, tuple(identity for _, identity, _ in results)) + + +def test_slow_sink_during_a_burst_lands_every_response_exactly_once(rig: Rig) -> None: + rig.sink.release.clear() + rig.sink.slow.set() + try: + results: Final = _burst(rig, 20) + assert all(error is None for _, _, error in results), [error for _, _, error in results if error] + identities: Final = tuple(identity for _, identity, _ in results) + assert all(count == 0 for count in rig.sink.landed(identities).values()), "sink accepted while held" + finally: + rig.sink.slow.clear() + rig.sink.release.set() + _assert_exactly_once(rig, identities) + + +def test_killing_one_of_two_workers_mid_burst_keeps_serving_and_never_duplicates_a_span(rig: Rig) -> None: + root: Final = psutil.Process(rig.process.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, + ) + markers: Final = tuple(_marker() for _ in range(24)) + + def one(index: int) -> tuple[str, str | None]: + if index == 8: + os.kill(workers[0].pid, signal.SIGKILL) + try: + response: Final = rig.chat(markers[index]) + return f"chatcmpl-{markers[index]}", None if response.status_code == 200 else response.text + except httpx.HTTPError as error: + return f"chatcmpl-{markers[index]}", repr(error) + + with ThreadPoolExecutor(max_workers=6) as pool: + results: Final = tuple(pool.map(one, range(24))) + assert rig.process.process.poll() is None, "Proxy root exited after a worker was killed" + after: Final = rig.chat(_marker()) + assert after.status_code == 200, after.text + rig.sink.single_span(_body_id(after)) + failures: Final = tuple(error for _, error in results if error) + assert all(error.startswith(("ReadError(", "RemoteProtocolError(", "ConnectError(")) for error in failures), ( + failures + ) + assert len(failures) <= 6, failures + served: Final = tuple(identity for identity, error in results if error is None) + assert len(served) >= 18, results + settled: Final = tuple(identity for index, (identity, error) in enumerate(results) if index > 14 and not error) + _assert_exactly_once(rig, settled) + assert all(count <= 1 for count in rig.sink.landed(served).values()), rig.sink.landed(served) + + +def test_terminating_the_proxy_right_after_a_burst_flushes_every_span_before_exit( + provider: Wire, operator_sink: Collector, tenant_sink: Collector, tmp_path_factory: pytest.TempPathFactory +) -> None: + pytest.skip("BUG: spans still queued in the OTel batch processor at SIGTERM never reach the sink (4 of 10 lost)") + factory: Final = RigFactory( + provider, + operator_sink, + tenant_sink, + tmp_path_factory.mktemp("signoz-shutdown"), + False, + 2, + operator_sink.wire.url, + ) + started: Final = factory.start() + rig: Final = next(started) + responses: Final = tuple(rig.chat(_marker()) for _ in range(10)) + assert all(response.status_code == 200 for response in responses), [response.text for response in responses] + identities: Final = tuple(_body_id(response) for response in responses) + pending_at_signal: Final = rig.sink.landed(identities) + rig.process.process.terminate() + assert rig.process.process.wait(timeout=40) in (0, -signal.SIGTERM) + with pytest.raises(httpx.ConnectError): + next(started) + landed: Final = rig.sink.landed(identities) + assert landed == {identity: 1 for identity in identities}, ( + f"spans at the sink after exit: {landed}, at the moment of SIGTERM: {pending_at_signal}" + ) diff --git a/tests/test_litellm/proxy/test_litellm_pre_call_utils.py b/tests/test_litellm/proxy/test_litellm_pre_call_utils.py index ee42042bb8e..1a3ffecb0a5 100644 --- a/tests/test_litellm/proxy/test_litellm_pre_call_utils.py +++ b/tests/test_litellm/proxy/test_litellm_pre_call_utils.py @@ -8470,3 +8470,31 @@ async def test_mcp_credentials_only_removed_from_logging_copies(path: str, custo for name, value in secrets.items(): assert updated["secret_fields"]["raw_headers"][name.lower()] == value assert request.headers[name] == value + + +def test_signoz_callback_vars_are_scoped_to_the_signoz_callback(): + from litellm.proxy._types import AddTeamCallback + from litellm.proxy.litellm_pre_call_utils import convert_key_logging_metadata_to_callback + + under_signoz = convert_key_logging_metadata_to_callback( + data=AddTeamCallback( + callback_name="signoz", + callback_type="success", + callback_vars={"signoz_ingestion_key": "team-key", "signoz_ingestion_endpoint": "https://ingest.eu.signoz.cloud:443"}, + ), + team_callback_settings_obj=None, + ) + assert under_signoz.callback_vars == { + "signoz_ingestion_key": "team-key", + "signoz_ingestion_endpoint": "https://ingest.eu.signoz.cloud:443", + } + + under_other = convert_key_logging_metadata_to_callback( + data=AddTeamCallback( + callback_name="langfuse", + callback_type="success", + callback_vars={"signoz_ingestion_key": "team-key", "langfuse_host": "https://cloud.langfuse.com"}, + ), + team_callback_settings_obj=None, + ) + assert under_other.callback_vars == {"langfuse_host": "https://cloud.langfuse.com"} diff --git a/tests/unit/integrations/otel/test_otel_v2_dynamic.py b/tests/unit/integrations/otel/test_otel_v2_dynamic.py index 29772eb92c7..8163ba06317 100644 --- a/tests/unit/integrations/otel/test_otel_v2_dynamic.py +++ b/tests/unit/integrations/otel/test_otel_v2_dynamic.py @@ -1,6 +1,7 @@ """Per-request multi-tenant credential routing (V1 parity).""" import base64 +import logging import pytest from opentelemetry.trace import NoOpTracer @@ -677,3 +678,68 @@ def test_newrelic_key_only_team_routes_to_us_not_operator_region(monkeypatch): ) owned = next(e for e in new_cfg.exporters if e.owner == "newrelic") assert owned.endpoint == "https://otlp.nr-data.net" + + +def test_signoz_dynamic_headers_stamp_ingestion_key(): + from litellm.integrations.otel.presets import dynamic_otlp_headers + + assert dynamic_otlp_headers("signoz", {"signoz_ingestion_key": "team-key"}) == {"signoz-ingestion-key": "team-key"} + # No key means no per-request routing; the caller keeps its default tracer. + assert dynamic_otlp_headers("signoz", {}) is None + + +def test_signoz_dynamic_endpoint_comes_from_team_config_when_its_host_is_allowlisted(monkeypatch): + import litellm + from litellm.integrations.otel.presets import dynamic_otlp_endpoint + + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["ingest.eu.signoz.cloud"]) + assert ( + dynamic_otlp_endpoint( + "signoz", {"signoz_ingestion_endpoint": "https://ingest.eu.signoz.cloud:443", "signoz_ingestion_key": "k"} + ) + == "https://ingest.eu.signoz.cloud:443" + ) + # A team that saved only a key keeps the operator's configured endpoint. + assert dynamic_otlp_endpoint("signoz", {"signoz_ingestion_key": "k"}) is None + assert dynamic_otlp_endpoint("signoz", {}) is None + + +def test_signoz_team_endpoint_off_the_allowlist_is_dropped_along_with_its_key(monkeypatch): + import litellm + from litellm.integrations.otel.presets import dynamic_otlp_endpoint, dynamic_otlp_headers + + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", []) + params = {"signoz_ingestion_endpoint": "http://169.254.169.254/v1/traces", "signoz_ingestion_key": "k"} + assert dynamic_otlp_endpoint("signoz", params) is None + # The tenant key must not ride to the operator's collector either: the request keeps the default tracer. + assert dynamic_otlp_headers("signoz", params) is None + + +def test_signoz_keyless_team_endpoint_is_ignored_so_the_operator_key_never_reaches_it(monkeypatch, caplog): + import litellm + from litellm.integrations.otel.presets import dynamic_otlp_endpoint, dynamic_otlp_headers + from litellm.integrations.otel.presets.signoz import _warn_endpoint_without_key + + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["collector.team.internal"]) + params = {"signoz_ingestion_endpoint": "http://collector.team.internal:4318"} + _warn_endpoint_without_key.cache_clear() + with caplog.at_level(logging.WARNING, logger="LiteLLM"): + assert dynamic_otlp_headers("signoz", params) is None + assert "Set signoz_ingestion_key alongside it" in caplog.text + assert dynamic_otlp_endpoint("signoz", params) is None + cache = _cache( + "signoz", + exporters=[ + ExporterSpec( + kind="otlp_http", + endpoint="https://ingest.us.signoz.cloud:443", + headers="signoz-ingestion-key=OPERATOR", + owner="signoz", + requires_headers=True, + ) + ], + ) + routed = cache._routed_config({}, {}, dynamic_otlp_endpoint("signoz", params), "team-service") + owned = next(e for e in routed.exporters if e.owner == "signoz") + assert owned.endpoint == "https://ingest.us.signoz.cloud:443" + assert owned.headers == "signoz-ingestion-key=OPERATOR" diff --git a/tests/unit/integrations/otel/test_otel_v2_presets.py b/tests/unit/integrations/otel/test_otel_v2_presets.py index 58cfc1ceb3f..a060cdf3648 100644 --- a/tests/unit/integrations/otel/test_otel_v2_presets.py +++ b/tests/unit/integrations/otel/test_otel_v2_presets.py @@ -212,3 +212,69 @@ def test_newrelic_preset_unset_content_knob_keeps_default(monkeypatch): from litellm.integrations.otel.presets.newrelic import newrelic_preset assert newrelic_preset().capture_span_content is False + + +def test_signoz_preset_reads_env_endpoint_and_key(monkeypatch): + monkeypatch.setenv("SIGNOZ_INGESTION_ENDPOINT", "https://ingest.eu.signoz.cloud:443") + monkeypatch.setenv("SIGNOZ_INGESTION_KEY", "env-ingestion-key") + from litellm.integrations.otel.model.config import ExporterOwner + from litellm.integrations.otel.presets.signoz import signoz_preset + + cfg = signoz_preset() + spec = next(e for e in cfg.exporters if e.owner == ExporterOwner.SIGNOZ) + assert spec.kind == "otlp_http" + assert spec.endpoint == "https://ingest.eu.signoz.cloud:443" + assert spec.headers == "signoz-ingestion-key=env-ingestion-key" + assert spec.requires_headers is True + assert "genai" in cfg.mapper_names + + +def test_signoz_preset_without_key_is_self_hosted(monkeypatch): + # A self-hosted collector accepts unauthenticated OTLP, so requiring headers + # would drop exports that would have succeeded. + monkeypatch.setenv("SIGNOZ_INGESTION_ENDPOINT", "http://signoz-collector.internal:4318") + monkeypatch.delenv("SIGNOZ_INGESTION_KEY", raising=False) + from litellm.integrations.otel.model.config import ExporterOwner + from litellm.integrations.otel.presets.signoz import signoz_preset + + cfg = signoz_preset() + spec = next(e for e in cfg.exporters if e.owner == ExporterOwner.SIGNOZ) + assert spec.endpoint == "http://signoz-collector.internal:4318" + assert spec.headers is None + assert spec.requires_headers is False + + +def test_signoz_preset_has_no_default_endpoint(monkeypatch): + # No region table and no default host: the preset never invents a destination. + monkeypatch.delenv("SIGNOZ_INGESTION_ENDPOINT", raising=False) + monkeypatch.delenv("SIGNOZ_INGESTION_KEY", raising=False) + from litellm.integrations.otel.model.config import ExporterOwner + from litellm.integrations.otel.presets.signoz import signoz_preset + + cfg = signoz_preset() + spec = next(e for e in cfg.exporters if e.owner == ExporterOwner.SIGNOZ) + assert spec.endpoint is None + + +def test_signoz_preset_endpoint_passed_through_verbatim(monkeypatch): + # The plumbing appends the signal path, so pre-appending would double it. + monkeypatch.setenv("SIGNOZ_INGESTION_ENDPOINT", "https://ingest.us.signoz.cloud:443/v1/traces") + monkeypatch.delenv("SIGNOZ_INGESTION_KEY", raising=False) + from litellm.integrations.otel.model.config import ExporterOwner + from litellm.integrations.otel.plumbing.providers import _otlp_traces_endpoint + from litellm.integrations.otel.presets.signoz import signoz_preset + + cfg = signoz_preset() + spec = next(e for e in cfg.exporters if e.owner == ExporterOwner.SIGNOZ) + assert spec.endpoint == "https://ingest.us.signoz.cloud:443/v1/traces" + assert _otlp_traces_endpoint(spec.endpoint) == "https://ingest.us.signoz.cloud:443/v1/traces" + + +def test_signoz_preset_accepts_the_factory_call_shape(monkeypatch): + monkeypatch.setenv("SIGNOZ_INGESTION_ENDPOINT", "http://127.0.0.1:1") + monkeypatch.delenv("SIGNOZ_INGESTION_KEY", raising=False) + from litellm.integrations.otel.model.config import ExporterOwner + from litellm.integrations.otel.presets import PRESET_BY_CALLBACK + + cfg = PRESET_BY_CALLBACK["signoz"](allow_missing_credentials=True) + assert any(e.owner == ExporterOwner.SIGNOZ for e in cfg.exporters) diff --git a/tests/unit/litellm_core_utils/test_litellm_logging.py b/tests/unit/litellm_core_utils/test_litellm_logging.py index c8b02ebc790..2fc747e1b48 100644 --- a/tests/unit/litellm_core_utils/test_litellm_logging.py +++ b/tests/unit/litellm_core_utils/test_litellm_logging.py @@ -8932,3 +8932,100 @@ async def test_prompt_management_with_unchanged_variables_replays_a_byte_identic assert json.dumps(messages_n_plus_one[: len(messages_n)], sort_keys=True) == json.dumps(messages_n, sort_keys=True) assert messages_n[0] == {"role": "system", "content": "You are a pirate. Answer in one sentence."} assert len(messages_n_plus_one) == len(messages_n) + 2 + + +def test_signoz_dispatch_prefers_otel_v2_when_flag_on(monkeypatch): + from litellm.integrations.otel.logger import OpenTelemetryV2 + from litellm.integrations.otel.model.config import ExporterOwner, is_otel_v2_enabled + from litellm.litellm_core_utils import litellm_logging as logging_module + + logging_module._in_memory_loggers.clear() + monkeypatch.setenv("LITELLM_OTEL_V2", "true") + monkeypatch.setenv("SIGNOZ_INGESTION_ENDPOINT", "https://ingest.eu.signoz.cloud:443") + monkeypatch.setenv("SIGNOZ_INGESTION_KEY", "test-key") + is_otel_v2_enabled.cache_clear() + try: + v2_logger = logging_module._init_custom_logger_compatible_class( + logging_integration="signoz", + internal_usage_cache=None, + llm_router=None, + custom_logger_init_args={}, + ) + assert isinstance(v2_logger, OpenTelemetryV2) + assert v2_logger.callback_name == "signoz" + spec = next(e for e in v2_logger.config.exporters if e.owner == ExporterOwner.SIGNOZ) + assert spec.endpoint == "https://ingest.eu.signoz.cloud:443" + assert spec.headers == "signoz-ingestion-key=test-key" + again = logging_module._init_custom_logger_compatible_class( + logging_integration="signoz", + internal_usage_cache=None, + llm_router=None, + custom_logger_init_args={}, + ) + assert again is v2_logger + finally: + logging_module._in_memory_loggers.clear() + monkeypatch.delenv("LITELLM_OTEL_V2", raising=False) + is_otel_v2_enabled.cache_clear() + + +def test_signoz_dispatch_keeps_legacy_otel_when_flag_off(monkeypatch): + from litellm.integrations.opentelemetry import OpenTelemetry + from litellm.integrations.otel.model.config import is_otel_v2_enabled + from litellm.litellm_core_utils import litellm_logging as logging_module + + logging_module._in_memory_loggers.clear() + monkeypatch.delenv("LITELLM_OTEL_V2", raising=False) + monkeypatch.setenv("SIGNOZ_INGESTION_ENDPOINT", "http://signoz-collector.internal:4318") + monkeypatch.setenv("SIGNOZ_INGESTION_KEY", "legacy-key") + monkeypatch.delenv("OTEL_EXPORTER_OTLP_TRACES_HEADERS", raising=False) + is_otel_v2_enabled.cache_clear() + try: + legacy = logging_module._init_custom_logger_compatible_class( + logging_integration="signoz", + internal_usage_cache=None, + llm_router=None, + custom_logger_init_args={}, + ) + assert isinstance(legacy, OpenTelemetry) + assert legacy.callback_name == "signoz" + assert legacy.config.endpoint == "http://signoz-collector.internal:4318/v1/traces" + assert legacy.config.headers == "signoz-ingestion-key=legacy-key" + assert "OTEL_EXPORTER_OTLP_TRACES_HEADERS" not in os.environ + # Same name resolves to the same instance, not a second exporter. + again = logging_module._init_custom_logger_compatible_class( + logging_integration="signoz", + internal_usage_cache=None, + llm_router=None, + custom_logger_init_args={}, + ) + assert again is legacy + finally: + logging_module._in_memory_loggers.clear() + is_otel_v2_enabled.cache_clear() + + +def test_signoz_dispatch_requires_an_endpoint(monkeypatch): + from litellm.integrations.otel.model.config import is_otel_v2_enabled + from litellm.litellm_core_utils import litellm_logging as logging_module + + logging_module._in_memory_loggers.clear() + monkeypatch.setenv("LITELLM_OTEL_V2", "true") + monkeypatch.delenv("SIGNOZ_INGESTION_ENDPOINT", raising=False) + monkeypatch.delenv("SIGNOZ_INGESTION_KEY", raising=False) + is_otel_v2_enabled.cache_clear() + try: + created = logging_module._init_custom_logger_compatible_class( + logging_integration="signoz", + internal_usage_cache=None, + llm_router=None, + custom_logger_init_args={}, + ) + assert created is None + assert not [ + cb for cb in logging_module._in_memory_loggers if getattr(cb, "callback_name", None) == "signoz" + ] + finally: + logging_module._in_memory_loggers.clear() + monkeypatch.delenv("LITELLM_OTEL_V2", raising=False) + is_otel_v2_enabled.cache_clear() diff --git a/ui/litellm-dashboard/public/assets/logos/signoz.svg b/ui/litellm-dashboard/public/assets/logos/signoz.svg new file mode 100644 index 00000000000..9064cb86bd6 --- /dev/null +++ b/ui/litellm-dashboard/public/assets/logos/signoz.svg @@ -0,0 +1 @@ + \ No newline at end of file diff --git a/ui/litellm-dashboard/src/components/callback_info_helpers.tsx b/ui/litellm-dashboard/src/components/callback_info_helpers.tsx index bc9889da724..19020d92066 100644 --- a/ui/litellm-dashboard/src/components/callback_info_helpers.tsx +++ b/ui/litellm-dashboard/src/components/callback_info_helpers.tsx @@ -10,6 +10,7 @@ import newrelicLogo from "../../public/assets/logos/newrelic.png"; import openmeterLogo from "../../public/assets/logos/openmeter.png"; import otelLogo from "../../public/assets/logos/otel.png"; import pointfiveLogo from "../../public/assets/logos/pointfive.png"; +import signozLogo from "../../public/assets/logos/signoz.svg"; import databricksLogo from "../../public/assets/logos/databricks.svg"; interface CallbackConfig { @@ -209,6 +210,17 @@ export const CALLBACK_CONFIGS: CallbackConfig[] = [ }, description: "S3 Bucket (AWS) Logging Integration", }, + { + id: "signoz", + displayName: "SigNoz", + logo: signozLogo.src, + supports_key_team_logging: true, + dynamic_params: { + signoz_ingestion_endpoint: "text", + signoz_ingestion_key: "password", + }, + description: "SigNoz Logging Integration. Setup: https://signoz.io/docs/litellm-observability/", + }, { id: "SQS", displayName: "SQS", diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index 6d45f6e691c..2b8ed9aa58d 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -57081,7 +57081,7 @@ export interface operations { parameters: { query: { /** @description Specify the service being hit. */ - service: ("slack_budget_alerts" | "langfuse" | "langfuse_otel" | "slack" | "ms_teams" | "openmeter" | "webhook" | "email" | "braintrust" | "datadog" | "datadog_llm_observability" | "generic_api" | "arize" | "galileo" | "newrelic" | "pointfive" | "sqs") | string; + service: ("slack_budget_alerts" | "langfuse" | "langfuse_otel" | "slack" | "ms_teams" | "openmeter" | "webhook" | "email" | "braintrust" | "datadog" | "datadog_llm_observability" | "generic_api" | "arize" | "galileo" | "newrelic" | "pointfive" | "signoz" | "sqs") | string; }; header?: never; path?: never;