diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index a56188ab3..af06e69ef 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -268,7 +268,8 @@ def mask_otel_spans( additional_headers (Optional[Dict[str, str]]): Additional headers to include in all API requests and in the default OTLPSpanExporter requests. These headers will be merged with default headers. Note: If httpx_client is provided, additional_headers must be set directly on your custom httpx_client as well. If `span_exporter` is provided, these headers are not wired into that exporter and must be configured on the exporter instance directly. tracer_provider(Optional[TracerProvider]): OpenTelemetry TracerProvider to use for Langfuse. This can be useful to set to have disconnected tracing between Langfuse and other OpenTelemetry-span emitting libraries. Note: To track active spans, the context is still shared between TracerProviders. This may lead to broken trace trees. id_generator (Optional[IdGenerator]): OpenTelemetry ID generator to use when Langfuse creates its own TracerProvider. If omitted, the OpenTelemetry SDK default is used. If `tracer_provider` is provided, or an OpenTelemetry TracerProvider is already registered globally, configure the ID generator on that provider instead. - span_exporter (Optional[SpanExporter]): Custom OpenTelemetry span exporter for the Langfuse span processor. If omitted, Langfuse creates an OTLPSpanExporter pointed at the Langfuse OTLP endpoint. If provided, Langfuse does not wire `base_url`, exporter headers, exporter auth, exporter timeout, or the `LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES` request size limit into it. Configure endpoint, headers, and timeout on the exporter instance directly. If you are sending spans to Langfuse v4 or using Langfuse Cloud Fast Preview, include `x-langfuse-ingestion-version=4` on the exporter to enable real time processing of exported spans. + span_exporter (Optional[SpanExporter]): Custom OpenTelemetry span exporter for the Langfuse span processor. If omitted, Langfuse creates an OTLPSpanExporter pointed at the Langfuse OTLP endpoint. If provided, Langfuse does not wire `base_url`, exporter headers, exporter auth, exporter timeout, `otel_compression`, or the `LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES` request size limit into it. Configure endpoint, headers, timeout, and compression on the exporter instance directly. If you are sending spans to Langfuse v4 or using Langfuse Cloud Fast Preview, include `x-langfuse-ingestion-version=4` on the exporter to enable real time processing of exported spans. + otel_compression (Optional[Literal["gzip", "none"]]): Compression for span batches sent by the default OTLP span exporter. Use "gzip" to reduce network bytes or "none" to send uncompressed. Can also be set via LANGFUSE_OTEL_COMPRESSION environment variable. If unset, the standard OTEL_EXPORTER_OTLP_TRACES_COMPRESSION and OTEL_EXPORTER_OTLP_COMPRESSION environment variables apply (no compression by default). "gzip" requires Langfuse server v3.30.0 or later. Example: ```python @@ -335,6 +336,7 @@ def __init__( tracer_provider: Optional[TracerProvider] = None, id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ): self._base_url = ( base_url @@ -435,6 +437,7 @@ def __init__( tracer_provider=tracer_provider, id_generator=id_generator, span_exporter=span_exporter, + otel_compression=otel_compression, ) self._mask = self._resources.mask diff --git a/langfuse/_client/environment_variables.py b/langfuse/_client/environment_variables.py index 424ba53fd..b340c9ac4 100644 --- a/langfuse/_client/environment_variables.py +++ b/langfuse/_client/environment_variables.py @@ -73,6 +73,19 @@ **Default value:** ``67108864`` (64 MiB) """ +LANGFUSE_OTEL_COMPRESSION = "LANGFUSE_OTEL_COMPRESSION" +""" +.. envvar:: LANGFUSE_OTEL_COMPRESSION + +Compression for span batches sent by the default OTLP exporter: ``gzip`` or ``none``. +The ``otel_compression`` client argument takes precedence. If unset, the standard +``OTEL_EXPORTER_OTLP_TRACES_COMPRESSION`` and ``OTEL_EXPORTER_OTLP_COMPRESSION`` +environment variables apply. ``gzip`` requires Langfuse server v3.30.0 or later. +Custom span exporters are not affected. + +**Default value:** unset +""" + LANGFUSE_DEBUG = "LANGFUSE_DEBUG" """ .. envvar:: LANGFUSE_DEBUG diff --git a/langfuse/_client/get_client.py b/langfuse/_client/get_client.py index 0c4ccd321..ff06c7d29 100644 --- a/langfuse/_client/get_client.py +++ b/langfuse/_client/get_client.py @@ -57,6 +57,7 @@ def _create_client_from_instance( tracer_provider=instance.tracer_provider, id_generator=instance.id_generator, span_exporter=instance.span_exporter, + otel_compression=instance.otel_compression, httpx_client=instance.httpx_client, ) diff --git a/langfuse/_client/resource_manager.py b/langfuse/_client/resource_manager.py index a61101fe3..4395872db 100644 --- a/langfuse/_client/resource_manager.py +++ b/langfuse/_client/resource_manager.py @@ -21,7 +21,7 @@ import urllib.request import weakref from queue import Full, Queue -from typing import Any, Callable, Dict, List, Optional, cast +from typing import Any, Callable, Dict, List, Literal, Optional, cast import httpx from opentelemetry import trace as otel_trace_api @@ -133,6 +133,7 @@ def __new__( tracer_provider: Optional[TracerProvider] = None, id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ) -> "LangfuseResourceManager": if public_key in cls._instances: return cls._instances[public_key] @@ -171,6 +172,7 @@ def __new__( tracer_provider=tracer_provider, id_generator=id_generator, span_exporter=span_exporter, + otel_compression=otel_compression, ) cls._instances[public_key] = instance @@ -200,6 +202,7 @@ def _initialize_instance( tracer_provider: Optional[TracerProvider] = None, id_generator: Optional[IdGenerator] = None, span_exporter: Optional[SpanExporter] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ) -> None: self.public_key = public_key self.secret_key = secret_key @@ -222,6 +225,7 @@ def _initialize_instance( self.additional_headers = additional_headers self.id_generator = id_generator self.span_exporter = span_exporter + self.otel_compression = otel_compression self.tracer_provider: Optional[TracerProvider] = None self._custom_httpx_client = httpx_client @@ -261,6 +265,7 @@ def _initialize_instance( span_exporter=span_exporter, media_manager=self._media_manager, mask_otel_spans=mask_otel_spans, + otel_compression=otel_compression, ) tracer_provider.add_span_processor(langfuse_processor) diff --git a/langfuse/_client/span_processor.py b/langfuse/_client/span_processor.py index 524b2996d..b362b8da1 100644 --- a/langfuse/_client/span_processor.py +++ b/langfuse/_client/span_processor.py @@ -15,10 +15,11 @@ import logging import os import threading -from typing import Callable, Dict, List, Optional, cast +from typing import Callable, Dict, List, Literal, Optional, cast from opentelemetry import context as context_api from opentelemetry.context import Context +from opentelemetry.exporter.otlp.proto.http import Compression from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter from opentelemetry.sdk.trace import ReadableSpan, Span from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter @@ -28,6 +29,7 @@ from langfuse._client.environment_variables import ( LANGFUSE_FLUSH_AT, LANGFUSE_FLUSH_INTERVAL, + LANGFUSE_OTEL_COMPRESSION, LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES, LANGFUSE_OTEL_TRACES_EXPORT_PATH, ) @@ -65,6 +67,34 @@ def _resolve_max_batch_size_bytes() -> Optional[int]: return None +# The Langfuse OTLP endpoint only decodes gzip, so deflate is not offered. +_COMPRESSION_BY_NAME = {"gzip": Compression.Gzip, "none": Compression.NoCompression} + + +def _resolve_compression(otel_compression: Optional[str]) -> Optional[Compression]: + """Return the configured compression, or None to defer to OTEL_EXPORTER_OTLP_*COMPRESSION.""" + setting = "otel_compression" + raw_value = otel_compression + if raw_value is None: + setting = LANGFUSE_OTEL_COMPRESSION + raw_value = os.environ.get(LANGFUSE_OTEL_COMPRESSION, "") + + value = raw_value.strip().lower() + if not value: + return None + + compression = _COMPRESSION_BY_NAME.get(value) + if compression is None: + langfuse_logger.warning( + "Invalid %s=%r. Expected 'gzip' or 'none'. Falling back to the " + "OTEL_EXPORTER_OTLP_*COMPRESSION environment variables.", + setting, + raw_value, + ) + + return compression + + class LangfuseSpanProcessor(BatchSpanProcessor): """OpenTelemetry span processor that exports spans to the Langfuse API. @@ -97,6 +127,7 @@ def __init__( span_exporter: Optional[SpanExporter] = None, media_manager: Optional[MediaManager] = None, mask_otel_spans: Optional[MaskOtelSpansFunction] = None, + otel_compression: Optional[Literal["gzip", "none"]] = None, ): self.public_key = public_key self.blocked_instrumentation_scopes = ( @@ -145,8 +176,14 @@ def __init__( endpoint=endpoint, headers=headers, timeout=timeout, + compression=_resolve_compression(otel_compression), max_request_size=_resolve_max_batch_size_bytes(), ) + elif otel_compression is not None: + langfuse_logger.warning( + "otel_compression is ignored because a custom span_exporter was " + "provided. Configure compression on the exporter instead." + ) if media_manager is not None or mask_otel_spans is not None: span_exporter = LangfuseTransformingSpanExporter( diff --git a/tests/unit/test_span_processor.py b/tests/unit/test_span_processor.py index a123fd108..9591929be 100644 --- a/tests/unit/test_span_processor.py +++ b/tests/unit/test_span_processor.py @@ -1,11 +1,19 @@ +import gzip import logging import threading from http.server import BaseHTTPRequestHandler, HTTPServer -from typing import List, Sequence +from typing import List, NamedTuple, Optional, Sequence from unittest.mock import patch import pytest from opentelemetry.exporter.otlp.proto.common.trace_encoder import encode_spans +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ( + ExportTraceServiceRequest, +) +from opentelemetry.sdk.environment_variables import ( + OTEL_EXPORTER_OTLP_COMPRESSION, + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, +) from opentelemetry.sdk.trace import ReadableSpan, TracerProvider from opentelemetry.sdk.trace.export import ( SimpleSpanProcessor, @@ -17,9 +25,11 @@ ) import langfuse._client.span_processor as span_processor_module +from langfuse._client.client import Langfuse from langfuse._client.environment_variables import ( LANGFUSE_FLUSH_AT, LANGFUSE_FLUSH_INTERVAL, + LANGFUSE_OTEL_COMPRESSION, LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES, ) from langfuse._client.span_processor import LangfuseSpanProcessor @@ -71,12 +81,19 @@ def test_span_processor_uses_env_flush_settings_when_constructor_omits_them( processor.shutdown() +class _RecordedRequest(NamedTuple): + content_encoding: Optional[str] + body: bytes + + class _RecordingOTLPHandler(BaseHTTPRequestHandler): - received_body_sizes: List[int] + received_requests: List[_RecordedRequest] def do_POST(self): body = self.rfile.read(int(self.headers.get("Content-Length", 0))) - self.received_body_sizes.append(len(body)) + self.received_requests.append( + _RecordedRequest(self.headers.get("Content-Encoding"), body) + ) self.send_response(200) self.end_headers() @@ -86,17 +103,17 @@ def log_message(self, *args): @pytest.fixture def otlp_http_server(): - received_body_sizes: List[int] = [] + received_requests: List[_RecordedRequest] = [] handler = type( "Handler", (_RecordingOTLPHandler,), - {"received_body_sizes": received_body_sizes}, + {"received_requests": received_requests}, ) server = HTTPServer(("127.0.0.1", 0), handler) thread = threading.Thread(target=server.serve_forever, daemon=True) thread.start() - yield f"http://127.0.0.1:{server.server_port}", received_body_sizes + yield f"http://127.0.0.1:{server.server_port}", received_requests server.shutdown() server.server_close() @@ -136,7 +153,7 @@ def _default_exporter_processor(base_url: str) -> LangfuseSpanProcessor: def test_default_exporter_enforces_max_batch_size_bytes_at_boundary( monkeypatch, otlp_http_server, limit_offset, expected_result ): - base_url, received_body_sizes = otlp_http_server + base_url, received_requests = otlp_http_server spans = _finished_spans("x" * 1_000) request_size = _serialized_request_size(spans) monkeypatch.setenv( @@ -153,13 +170,13 @@ def test_default_exporter_enforces_max_batch_size_bytes_at_boundary( expected_requests = ( [] if expected_result == SpanExportResult.FAILURE else [request_size] ) - assert received_body_sizes == expected_requests + assert [len(request.body) for request in received_requests] == expected_requests def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( monkeypatch, caplog, otlp_http_server ): - base_url, received_body_sizes = otlp_http_server + base_url, received_requests = otlp_http_server oversized_spans = _finished_spans("secret-payload" * 1_000) small_spans = _finished_spans("ok") monkeypatch.setenv( @@ -173,7 +190,7 @@ def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( super(LangfuseSpanProcessor, processor).on_end(oversized_spans[0]) assert processor.force_flush() - assert received_body_sizes == [] + assert received_requests == [] assert "Dropping span batch" in caplog.text assert "secret-payload" not in caplog.text @@ -182,7 +199,9 @@ def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches( finally: processor.shutdown() - assert received_body_sizes == [_serialized_request_size(small_spans)] + assert [len(request.body) for request in received_requests] == [ + _serialized_request_size(small_spans) + ] def test_default_exporter_uses_64_mib_limit_when_env_unset(monkeypatch): @@ -213,6 +232,95 @@ def test_invalid_max_batch_size_bytes_falls_back_to_default_limit( processor.shutdown() +def _decode_request(request: _RecordedRequest) -> ExportTraceServiceRequest: + body = gzip.decompress(request.body) if request.content_encoding else request.body + return ExportTraceServiceRequest.FromString(body) + + +def _span_names(request: ExportTraceServiceRequest) -> List[str]: + return [ + span.name + for resource_spans in request.resource_spans + for scope_spans in resource_spans.scope_spans + for span in scope_spans.spans + ] + + +@pytest.fixture +def compression_env(monkeypatch): + for name in ( + LANGFUSE_OTEL_COMPRESSION, + OTEL_EXPORTER_OTLP_TRACES_COMPRESSION, + OTEL_EXPORTER_OTLP_COMPRESSION, + ): + monkeypatch.delenv(name, raising=False) + + return monkeypatch + + +def _export_one_span(processor: LangfuseSpanProcessor) -> None: + provider = TracerProvider() + provider.add_span_processor(processor) + with provider.get_tracer("test").start_as_current_span( + "llm-call", attributes={"gen_ai.system": "test"} + ): + pass + + try: + assert processor.force_flush() + finally: + provider.shutdown() + + +def test_client_otel_compression_sends_gzip_request(compression_env, otlp_http_server): + base_url, received_requests = otlp_http_server + langfuse = Langfuse( + public_key="pk-test", + secret_key="sk-test", + base_url=base_url, + tracer_provider=TracerProvider(), + otel_compression="gzip", + ) + + with langfuse.start_as_current_observation(name="compressed-span"): + pass + langfuse.flush() + + assert [request.content_encoding for request in received_requests] == ["gzip"] + assert _span_names(_decode_request(received_requests[0])) == ["compressed-span"] + + +@pytest.mark.parametrize( + ("otel_compression", "env", "expected_encoding"), + [ + (None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"), + ("gzip", {LANGFUSE_OTEL_COMPRESSION: "none"}, "gzip"), + ("none", {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, None), + (None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, "gzip"), + ], +) +def test_default_exporter_compression_precedence( + compression_env, otlp_http_server, otel_compression, env, expected_encoding +): + base_url, received_requests = otlp_http_server + for name, value in env.items(): + compression_env.setenv(name, value) + + _export_one_span( + LangfuseSpanProcessor( + public_key="pk-test", + secret_key="sk-test", + base_url=base_url, + otel_compression=otel_compression, + ) + ) + + assert [request.content_encoding for request in received_requests] == [ + expected_encoding + ] + assert _span_names(_decode_request(received_requests[0])) == ["llm-call"] + + @pytest.fixture def tracer_with_processor(): processor = LangfuseSpanProcessor(