fix(langfuse_otel): build per-request OTLP exporter from key and team dynamic Langfuse credentials (#32437)

* fix(langfuse_otel): build per-request OTLP exporter from key/team dynamic Langfuse credentials

Key-scoped langfuse_otel callbacks only injected Authorization headers into the
init-time exporter, so a proxy without global LANGFUSE_* env vars kept its
fallback exporter and never exported traces to Langfuse. Dynamic params now
build a full per-request OTLP config (endpoint from the key's langfuse_host,
otlp_http, basic auth from the key's credentials).

Resolves LIT-3976

* fix(otel): log dynamic config endpoint in span processor debug output

* fix(otel): redact authorization headers in exporter debug logs
This commit is contained in:
Yassin Kortam 2026-07-16 13:39:10 -07:00 • committed by GitHub
parent 903219a8b1
commit 2162da5015
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 293 additions and 62 deletions

View file

@ -267,29 +267,10 @@ class LangfuseOtelLogger(OpenTelemetry):
# If no keys, return default from env (likely logging to console or something else)
return OpenTelemetryConfig.from_env()
# Determine endpoint - default to US cloud
langfuse_host = LangfuseOtelLogger._get_langfuse_otel_host()
if langfuse_host:
# If LANGFUSE_HOST is provided, construct OTEL endpoint from it
if not langfuse_host.startswith("http"):
langfuse_host = "https://" + langfuse_host
endpoint = f"{langfuse_host.rstrip('/')}/api/public/otel"
verbose_logger.debug(f"Using Langfuse OTEL endpoint from host: {endpoint}")
else:
# Default to US cloud endpoint
endpoint = LANGFUSE_CLOUD_US_ENDPOINT
verbose_logger.debug(f"Using Langfuse US cloud endpoint: {endpoint}")
auth_header = LangfuseOtelLogger._get_langfuse_authorization_header(
public_key=public_key, secret_key=secret_key
)
otlp_auth_headers = f"Authorization={auth_header}"
return OpenTelemetryConfig(
exporter="otlp_http",
endpoint=endpoint,
headers=otlp_auth_headers,
return LangfuseOtelLogger._build_langfuse_otel_config(
public_key=public_key,
secret_key=secret_key,
langfuse_host=LangfuseOtelLogger._get_langfuse_otel_host(),
)
@staticmethod
@ -316,33 +297,36 @@ class LangfuseOtelLogger(OpenTelemetry):
"LANGFUSE_PUBLIC_KEY and LANGFUSE_SECRET_KEY must be set for Langfuse OpenTelemetry integration."
)
# Determine endpoint - default to US cloud
langfuse_host = LangfuseOtelLogger._get_langfuse_otel_host()
return LangfuseOtelLogger._build_langfuse_otel_config(
public_key=public_key,
secret_key=secret_key,
langfuse_host=LangfuseOtelLogger._get_langfuse_otel_host(),
)
@staticmethod
def _build_langfuse_otel_config(
public_key: str, secret_key: str, langfuse_host: Optional[str]
) -> "OpenTelemetryConfig":
"""
Builds an OTLP HTTP config pointing at the Langfuse OTEL endpoint for the
given host (US cloud when no host is provided), authorized with the given keys.
"""
if langfuse_host:
# If LANGFUSE_HOST is provided, construct OTEL endpoint from it
if not langfuse_host.startswith("http"):
langfuse_host = "https://" + langfuse_host
endpoint = f"{langfuse_host.rstrip('/')}/api/public/otel"
normalized_host = langfuse_host if langfuse_host.startswith("http") else f"https://{langfuse_host}"
endpoint = f"{normalized_host.rstrip('/')}/api/public/otel"
verbose_logger.debug(f"Using Langfuse OTEL endpoint from host: {endpoint}")
else:
# Default to US cloud endpoint
endpoint = LANGFUSE_CLOUD_US_ENDPOINT
verbose_logger.debug(f"Using Langfuse US cloud endpoint: {endpoint}")
auth_header = LangfuseOtelLogger._get_langfuse_authorization_header(
public_key=public_key, secret_key=secret_key
)
otlp_auth_headers = f"Authorization={auth_header}"
# Prevent modification of global env vars which causes leakage
# os.environ["OTEL_EXPORTER_OTLP_ENDPOINT"] = endpoint
# os.environ["OTEL_EXPORTER_OTLP_HEADERS"] = otlp_auth_headers
return OpenTelemetryConfig(
exporter="otlp_http",
endpoint=endpoint,
headers=otlp_auth_headers,
headers=f"Authorization={auth_header}",
)
@staticmethod
@ -378,6 +362,29 @@ class LangfuseOtelLogger(OpenTelemetry):
return dynamic_headers
def construct_dynamic_otel_config(
self, standard_callback_dynamic_params: StandardCallbackDynamicParams
) -> Optional["OpenTelemetryConfig"]:
"""
Build a full per-request OTLP config from team/key dynamic Langfuse credentials.
Key-scoped credentials must define the export target, not just the auth
headers: without this, a proxy with no global LANGFUSE_* env vars keeps its
init-time fallback exporter (console), so key-level langfuse_otel silently
never reaches Langfuse.
"""
public_key = standard_callback_dynamic_params.get("langfuse_public_key")
secret_key = standard_callback_dynamic_params.get("langfuse_secret_key")
if not public_key or not secret_key:
return None
langfuse_host = standard_callback_dynamic_params.get("langfuse_host") or self._get_langfuse_otel_host()
return LangfuseOtelLogger._build_langfuse_otel_config(
public_key=public_key,
secret_key=secret_key,
langfuse_host=langfuse_host,
)
def create_litellm_proxy_request_started_span(
self,
start_time: datetime,

View file

@ -28,6 +28,7 @@ from litellm.integrations.opentelemetry_utils.gen_ai_semconv import (
parse_semconv_opt_in,
)
from litellm.litellm_core_utils.safe_json_dumps import safe_dumps
from litellm.litellm_core_utils.secret_redaction import redact_string
from litellm.secret_managers.main import get_secret_bool, str_to_bool
from litellm.types.services import ServiceLoggerPayload
from litellm.types.utils import (
@ -948,12 +949,22 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
Returns:
Tracer: The tracer to use for this request
"""
dynamic_config = self._get_dynamic_otel_config_from_kwargs(kwargs)
if dynamic_config is not None:
verbose_logger.debug(
"[OTEL DEBUG] Using DYNAMIC config tracer with endpoint: %s",
dynamic_config.endpoint,
)
return self._get_tracer_with_dynamic_config(dynamic_config)
dynamic_headers = self._get_dynamic_otel_headers_from_kwargs(kwargs)
if dynamic_headers is not None:
# Create spans using a temporary tracer with dynamic headers
tracer_to_use = self._get_tracer_with_dynamic_headers(dynamic_headers)
verbose_logger.debug("[OTEL DEBUG] Using DYNAMIC tracer with headers: %s", dynamic_headers)
verbose_logger.debug(
"[OTEL DEBUG] Using DYNAMIC tracer with headers: %s", redact_string(str(dynamic_headers))
)
else:
# For langfuse_otel without dynamic headers, create a provider with env var credentials
if hasattr(self, "callback_name") and self.callback_name == "langfuse_otel":
@ -989,6 +1000,32 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
return dynamic_headers if dynamic_headers else None
def _get_dynamic_otel_config_from_kwargs(self, kwargs: dict) -> Optional[OpenTelemetryConfig]:
"""Extract a full dynamic exporter config from kwargs if available."""
standard_callback_dynamic_params: Optional[StandardCallbackDynamicParams] = kwargs.get(
"standard_callback_dynamic_params"
)
if not standard_callback_dynamic_params:
return None
return self.construct_dynamic_otel_config(standard_callback_dynamic_params=standard_callback_dynamic_params)
def _get_tracer_with_dynamic_config(self, dynamic_config: OpenTelemetryConfig):
"""Create (or reuse) a tracer whose exporter target comes from a per-request config."""
from opentelemetry.sdk.trace import TracerProvider
cache_key = f"dynamic_config:{dynamic_config.exporter}:{dynamic_config.endpoint}:{dynamic_config.headers}"
if cache_key in self._tracer_provider_cache:
return self._tracer_provider_cache[cache_key].get_tracer(LITELLM_TRACER_NAME)
temp_provider = TracerProvider(resource=self._get_litellm_resource(self.config))
temp_provider.add_span_processor(self._get_span_processor(config_override=dynamic_config))
self._tracer_provider_cache[cache_key] = temp_provider
return temp_provider.get_tracer(LITELLM_TRACER_NAME)
def _get_tracer_with_dynamic_headers(self, dynamic_headers: dict):
"""Create a temporary tracer with dynamic headers for this request only."""
from opentelemetry.sdk.trace import TracerProvider
@ -1020,6 +1057,19 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
"""
return None
def construct_dynamic_otel_config(
self, standard_callback_dynamic_params: StandardCallbackDynamicParams
) -> Optional[OpenTelemetryConfig]:
"""
Construct a full exporter config from standard callback dynamic params.
Override this when team/key dynamic params must control the export
target (exporter kind + endpoint), not just the request headers. When
this returns a config, it takes precedence over
construct_dynamic_otel_headers for the request.
"""
return None
#########################################################
# End of Team/Key Based Logging Control Flow
#########################################################
@ -2747,7 +2797,11 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
verbose_logger.debug("OpenTelemetry: No parent context found, creating root span")
return None, None
def _get_span_processor(self, dynamic_headers: Optional[dict] = None):
def _get_span_processor(
self,
dynamic_headers: Optional[dict] = None,
config_override: Optional[OpenTelemetryConfig] = None,
):
from opentelemetry.sdk.trace.export import (
BatchSpanProcessor,
ConsoleSpanExporter,
@ -2755,40 +2809,45 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
SpanExporter,
)
otel_exporter = config_override.exporter if config_override else self.OTEL_EXPORTER
otel_endpoint = config_override.endpoint if config_override else self.OTEL_ENDPOINT
otel_headers = config_override.headers if config_override else self.OTEL_HEADERS
verbose_logger.debug(
"OpenTelemetry Logger, initializing span processor \nself.OTEL_EXPORTER: %s\nself.OTEL_ENDPOINT: %s\nself.OTEL_HEADERS: %s",
self.OTEL_EXPORTER,
self.OTEL_ENDPOINT,
self.OTEL_HEADERS,
"OpenTelemetry Logger, initializing span processor \nexporter: %s\nendpoint: %s\nheaders: %s",
otel_exporter,
otel_endpoint,
redact_string(str(otel_headers)),
)
_split_otel_headers = OpenTelemetry._get_headers_dictionary(headers=dynamic_headers or self.OTEL_HEADERS)
_split_otel_headers = OpenTelemetry._get_headers_dictionary(headers=dynamic_headers or otel_headers)
if dynamic_headers:
verbose_logger.debug(
"[OTEL DEBUG] Creating span processor with DYNAMIC headers: %s",
{k: v[:20] + "..." if len(str(v)) > 20 else v for k, v in _split_otel_headers.items()},
redact_string(str(_split_otel_headers)),
)
elif config_override:
verbose_logger.debug(
"[OTEL DEBUG] Creating span processor with DYNAMIC config, endpoint: %s",
otel_endpoint,
)
else:
verbose_logger.debug("[OTEL DEBUG] Creating span processor with GLOBAL headers")
if hasattr(self.OTEL_EXPORTER, "export"): # Check if it has the export method that SpanExporter requires
if hasattr(otel_exporter, "export"): # Check if it has the export method that SpanExporter requires
verbose_logger.debug(
"OpenTelemetry: intiializing SpanExporter. Value of OTEL_EXPORTER: %s",
self.OTEL_EXPORTER,
otel_exporter,
)
return SimpleSpanProcessor(cast(SpanExporter, self.OTEL_EXPORTER))
return SimpleSpanProcessor(cast(SpanExporter, otel_exporter))
if self.OTEL_EXPORTER == "console":
if otel_exporter == "console":
verbose_logger.debug(
"OpenTelemetry: intiializing console exporter. Value of OTEL_EXPORTER: %s",
self.OTEL_EXPORTER,
otel_exporter,
)
return BatchSpanProcessor(ConsoleSpanExporter())
elif (
self.OTEL_EXPORTER == "otlp_http"
or self.OTEL_EXPORTER == "http/protobuf"
or self.OTEL_EXPORTER == "http/json"
):
elif otel_exporter == "otlp_http" or otel_exporter == "http/protobuf" or otel_exporter == "http/json":
try:
from opentelemetry.exporter.otlp.proto.http.trace_exporter import (
OTLPSpanExporter as OTLPSpanExporterHTTP,
@ -2801,13 +2860,13 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
verbose_logger.debug(
"OpenTelemetry: intiializing http exporter. Value of OTEL_EXPORTER: %s",
self.OTEL_EXPORTER,
otel_exporter,
)
normalized_endpoint = self._normalize_otel_endpoint(self.OTEL_ENDPOINT, "traces")
normalized_endpoint = self._normalize_otel_endpoint(otel_endpoint, "traces")
return BatchSpanProcessor(
OTLPSpanExporterHTTP(endpoint=normalized_endpoint, headers=_split_otel_headers),
)
elif self.OTEL_EXPORTER == "otlp_grpc" or self.OTEL_EXPORTER == "grpc":
elif otel_exporter == "otlp_grpc" or otel_exporter == "grpc":
try:
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import (
OTLPSpanExporter as OTLPSpanExporterGRPC,
@ -2820,16 +2879,16 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
verbose_logger.debug(
"OpenTelemetry: intiializing grpc exporter. Value of OTEL_EXPORTER: %s",
self.OTEL_EXPORTER,
otel_exporter,
)
normalized_endpoint = self._normalize_otel_endpoint(self.OTEL_ENDPOINT, "traces")
normalized_endpoint = self._normalize_otel_endpoint(otel_endpoint, "traces")
return BatchSpanProcessor(
OTLPSpanExporterGRPC(endpoint=normalized_endpoint, headers=_split_otel_headers),
)
else:
verbose_logger.debug(
"OpenTelemetry: intiializing console exporter. Value of OTEL_EXPORTER: %s",
self.OTEL_EXPORTER,
otel_exporter,
)
return BatchSpanProcessor(ConsoleSpanExporter())
@ -2841,7 +2900,7 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
"OpenTelemetry Logger, initializing log exporter \nself.OTEL_EXPORTER: %s\nself.OTEL_ENDPOINT: %s\nself.OTEL_HEADERS: %s",
self.OTEL_EXPORTER,
self.OTEL_ENDPOINT,
self.OTEL_HEADERS,
redact_string(str(self.OTEL_HEADERS)),
)
_split_otel_headers = OpenTelemetry._get_headers_dictionary(self.OTEL_HEADERS)
@ -2928,7 +2987,7 @@ class OpenTelemetry(OTELGenAISemconvMixin, CustomLogger):
"OpenTelemetry Logger, initializing metric reader\nself.OTEL_EXPORTER: %s\nself.OTEL_ENDPOINT: %s\nself.OTEL_HEADERS: %s",
self.OTEL_EXPORTER,
self.OTEL_ENDPOINT,
self.OTEL_HEADERS,
redact_string(str(self.OTEL_HEADERS)),
)
_split_otel_headers = OpenTelemetry._get_headers_dictionary(self.OTEL_HEADERS)

View file

@ -412,6 +412,171 @@ class TestLangfuseOtelIntegration:
# Endpoint assertion removed as side effect is gone
class TestLangfuseOtelKeyDynamicConfig:
"""Key/team-scoped Langfuse credentials must define the full export target
(OTLP endpoint + auth), not just auth headers on the init-time exporter."""
CLEAN_ENV_VARS = [
"LANGFUSE_PUBLIC_KEY",
"LANGFUSE_SECRET_KEY",
"LANGFUSE_HOST",
"LANGFUSE_OTEL_HOST",
"OTEL_EXPORTER",
"OTEL_EXPORTER_OTLP_PROTOCOL",
"OTEL_ENDPOINT",
"OTEL_EXPORTER_OTLP_ENDPOINT",
"OTEL_HEADERS",
"OTEL_EXPORTER_OTLP_HEADERS",
]
def _clean_env(self):
cleaned = {k: v for k, v in os.environ.items() if k not in self.CLEAN_ENV_VARS}
return patch.dict(os.environ, cleaned, clear=True)
def _dynamic_params(self, **overrides):
from litellm.types.utils import StandardCallbackDynamicParams
params = {
"langfuse_public_key": "key_public",
"langfuse_secret_key": "key_secret",
"langfuse_host": "https://langfuse.example.com",
}
params.update(overrides)
return StandardCallbackDynamicParams(**{k: v for k, v in params.items() if v is not None})
def test_construct_dynamic_otel_config_with_key_credentials(self):
with self._clean_env():
logger = LangfuseOtelLogger()
config = logger.construct_dynamic_otel_config(self._dynamic_params())
assert config is not None
assert config.exporter == "otlp_http"
assert config.endpoint == "https://langfuse.example.com/api/public/otel"
import base64
expected_auth = base64.b64encode(b"key_public:key_secret").decode()
assert config.headers == f"Authorization=Basic {expected_auth}"
def test_construct_dynamic_otel_config_host_without_protocol(self):
with self._clean_env():
logger = LangfuseOtelLogger()
config = logger.construct_dynamic_otel_config(self._dynamic_params(langfuse_host="langfuse.example.com"))
assert config is not None
assert config.endpoint == "https://langfuse.example.com/api/public/otel"
def test_construct_dynamic_otel_config_defaults_to_us_cloud(self):
with self._clean_env():
logger = LangfuseOtelLogger()
config = logger.construct_dynamic_otel_config(self._dynamic_params(langfuse_host=None))
assert config is not None
assert config.endpoint == "https://us.cloud.langfuse.com/api/public/otel"
def test_construct_dynamic_otel_config_falls_back_to_env_host(self):
with self._clean_env():
with patch.dict(os.environ, {"LANGFUSE_HOST": "https://env-host.example.com"}):
logger = LangfuseOtelLogger()
config = logger.construct_dynamic_otel_config(self._dynamic_params(langfuse_host=None))
assert config is not None
assert config.endpoint == "https://env-host.example.com/api/public/otel"
def test_construct_dynamic_otel_config_requires_both_keys(self):
with self._clean_env():
logger = LangfuseOtelLogger()
assert logger.construct_dynamic_otel_config(self._dynamic_params(langfuse_secret_key=None)) is None
assert logger.construct_dynamic_otel_config(self._dynamic_params(langfuse_public_key=None)) is None
def test_key_dynamic_params_create_otlp_exporter_without_global_env(self):
"""Without global LANGFUSE_* env vars, a request carrying key-scoped Langfuse
credentials must get a tracer exporting via OTLP HTTP to that key's host."""
from opentelemetry.exporter.otlp.proto.http.trace_exporter import (
OTLPSpanExporter,
)
from opentelemetry.sdk.trace.export import BatchSpanProcessor
with self._clean_env():
logger = LangfuseOtelLogger()
assert logger.OTEL_EXPORTER == "console"
tracer = logger.get_tracer_to_use_for_request(
{"standard_callback_dynamic_params": self._dynamic_params()}
)
assert tracer is not logger.tracer
assert len(logger._tracer_provider_cache) == 1
provider = next(iter(logger._tracer_provider_cache.values()))
span_processors = provider._active_span_processor._span_processors
assert len(span_processors) == 1
assert isinstance(span_processors[0], BatchSpanProcessor)
exporter = span_processors[0].span_exporter
assert isinstance(exporter, OTLPSpanExporter)
assert exporter._endpoint == "https://langfuse.example.com/api/public/otel/v1/traces"
import base64
expected_auth = base64.b64encode(b"key_public:key_secret").decode()
assert exporter._headers == {"Authorization": f"Basic {expected_auth}"}
def test_key_dynamic_params_reuse_cached_provider(self):
with self._clean_env():
logger = LangfuseOtelLogger()
kwargs = {"standard_callback_dynamic_params": self._dynamic_params()}
logger.get_tracer_to_use_for_request(kwargs)
logger.get_tracer_to_use_for_request(kwargs)
assert len(logger._tracer_provider_cache) == 1
def test_no_dynamic_params_keeps_default_tracer(self):
with self._clean_env():
logger = LangfuseOtelLogger()
tracer = logger.get_tracer_to_use_for_request({})
assert tracer is logger.tracer
assert logger._tracer_provider_cache == {}
def test_key_credentials_never_passed_to_debug_logger(self):
"""The span-processor debug logs must receive a redacted header value, so the
key-scoped Langfuse secret never enters a log record regardless of downstream
handler configuration, while the exporter still gets the real header."""
import base64
from opentelemetry.exporter.otlp.proto.http.trace_exporter import (
OTLPSpanExporter,
)
from litellm.integrations import opentelemetry as otel_module
secret = base64.b64encode(b"key_public:key_secret").decode()
recorded_arguments = []
def _spy(message, *args, **kwargs):
recorded_arguments.append(" ".join(str(part) for part in (message, *args)))
with self._clean_env():
logger = LangfuseOtelLogger()
with patch.object(otel_module.verbose_logger, "debug", side_effect=_spy):
logger.get_tracer_to_use_for_request(
{"standard_callback_dynamic_params": self._dynamic_params()}
)
logged = "\n".join(recorded_arguments)
assert "initializing span processor" in logged
assert secret not in logged
assert f"Basic {secret}" not in logged
provider = next(iter(logger._tracer_provider_cache.values()))
exporter = provider._active_span_processor._span_processors[0].span_exporter
assert isinstance(exporter, OTLPSpanExporter)
assert exporter._headers == {"Authorization": f"Basic {secret}"}
class TestLangfuseOtelResponsesAPI:
"""Test suite for Langfuse OTEL integration with ResponsesAPI"""