Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion langfuse/_client/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -251,7 +251,7 @@ def mask_otel_spans(
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, `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.
otel_compression (Optional[Literal["gzip", "none"]]): Compression for span batches sent by the default OTLP span exporter: "gzip" (default) or "none" to send uncompressed. If unset, the LANGFUSE_OTEL_COMPRESSION, OTEL_EXPORTER_OTLP_TRACES_COMPRESSION and OTEL_EXPORTER_OTLP_COMPRESSION environment variables apply, in that order, then gzip. Values are case-insensitive; an invalid value falls through to the next setting.

Example:
```python
Expand Down
10 changes: 5 additions & 5 deletions langfuse/_client/environment_variables.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,13 +77,13 @@
"""
.. 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.
Compression for span batches sent by the default OTLP exporter: ``gzip`` or ``none``
(case-insensitive). The ``otel_compression`` client argument takes precedence. If this
variable is unset or invalid, ``OTEL_EXPORTER_OTLP_TRACES_COMPRESSION`` and then
``OTEL_EXPORTER_OTLP_COMPRESSION`` apply, and otherwise span exports are gzip-compressed.
Custom span exporters are not affected.

**Default value:** unset
**Default value:** unset (span exports use gzip)
"""

LANGFUSE_DEBUG = "LANGFUSE_DEBUG"
Expand Down
53 changes: 41 additions & 12 deletions langfuse/_client/span_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,10 @@
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.environment_variables import (
OTEL_EXPORTER_OTLP_COMPRESSION,
OTEL_EXPORTER_OTLP_TRACES_COMPRESSION,
)
from opentelemetry.sdk.trace import ReadableSpan, Span
from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter
from opentelemetry.trace import format_span_id, format_trace_id
Expand Down Expand Up @@ -71,30 +75,55 @@ def _resolve_max_batch_size_bytes() -> Optional[int]:
_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()
def _parse_compression(
raw_value: Optional[str], *, setting: str, warn: bool
) -> Optional[Compression]:
value = (raw_value or "").strip().lower()
if not value:
return None

compression = _COMPRESSION_BY_NAME.get(value)
if compression is None:
if compression is None and warn:
langfuse_logger.warning(
"Invalid %s=%r. Expected 'gzip' or 'none'. Falling back to the "
"OTEL_EXPORTER_OTLP_*COMPRESSION environment variables.",
"Invalid %s=%r. Expected 'gzip' or 'none'. Falling back to the next "
"compression setting, then gzip.",
setting,
raw_value,
)

return compression


def _resolve_compression(otel_compression: Optional[str]) -> Compression:
"""Resolve compression for the default exporter; gzip unless configured otherwise.

Order: the ``otel_compression`` argument, ``LANGFUSE_OTEL_COMPRESSION``,
``OTEL_EXPORTER_OTLP_TRACES_COMPRESSION``, ``OTEL_EXPORTER_OTLP_COMPRESSION``,
then gzip. Values are case-insensitive, and an invalid value falls through to
the next setting, matching the JS SDK.
"""
candidates = (
("otel_compression", otel_compression, True),
(LANGFUSE_OTEL_COMPRESSION, os.environ.get(LANGFUSE_OTEL_COMPRESSION), True),
(
OTEL_EXPORTER_OTLP_TRACES_COMPRESSION,
os.environ.get(OTEL_EXPORTER_OTLP_TRACES_COMPRESSION),
False,
),
(
OTEL_EXPORTER_OTLP_COMPRESSION,
os.environ.get(OTEL_EXPORTER_OTLP_COMPRESSION),
False,
),
)
for setting, raw_value, warn in candidates:
compression = _parse_compression(raw_value, setting=setting, warn=warn)
if compression is not None:
return compression

return Compression.Gzip


class LangfuseSpanProcessor(BatchSpanProcessor):
"""OpenTelemetry span processor that exports spans to the Langfuse API.

Expand Down
51 changes: 47 additions & 4 deletions tests/unit/test_span_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,11 @@ def _serialized_request_size(spans: List[ReadableSpan]) -> int:
return len(encode_spans(spans).SerializePartialToString())


def _uncompressed_size(request: _RecordedRequest) -> int:
body = gzip.decompress(request.body) if request.content_encoding else request.body
return len(body)


def _default_exporter_processor(base_url: str) -> LangfuseSpanProcessor:
return LangfuseSpanProcessor(
public_key="pk-test",
Expand Down Expand Up @@ -170,7 +175,9 @@ def test_default_exporter_enforces_max_batch_size_bytes_at_boundary(
expected_requests = (
[] if expected_result == SpanExportResult.FAILURE else [request_size]
)
assert [len(request.body) for request in received_requests] == expected_requests
assert [_uncompressed_size(request) for request in received_requests] == (
expected_requests
)


def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches(
Expand Down Expand Up @@ -199,7 +206,7 @@ def test_oversized_batch_is_dropped_on_flush_without_blocking_later_batches(
finally:
processor.shutdown()

assert [len(request.body) for request in received_requests] == [
assert [_uncompressed_size(request) for request in received_requests] == [
_serialized_request_size(small_spans)
]

Expand Down Expand Up @@ -293,10 +300,46 @@ def test_client_otel_compression_sends_gzip_request(compression_env, otlp_http_s
@pytest.mark.parametrize(
("otel_compression", "env", "expected_encoding"),
[
(None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"),
# Nothing configured: gzip by default.
(None, {}, "gzip"),
# The argument wins over every environment variable.
("gzip", {LANGFUSE_OTEL_COMPRESSION: "none"}, "gzip"),
("none", {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, None),
(None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip"}, "gzip"),
# LANGFUSE_OTEL_COMPRESSION wins over the OTEL variables (case-insensitive).
(None, {LANGFUSE_OTEL_COMPRESSION: " GZIP "}, "gzip"),
(
None,
{LANGFUSE_OTEL_COMPRESSION: "None", OTEL_EXPORTER_OTLP_COMPRESSION: "gzip"},
None,
),
# The traces-specific OTEL variable wins over the generic one.
(
None,
{
OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "gzip",
OTEL_EXPORTER_OTLP_COMPRESSION: "none",
},
"gzip",
),
(None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "NONE"}, None),
(None, {OTEL_EXPORTER_OTLP_COMPRESSION: "none"}, None),
# Invalid values fall through to the next setting, then gzip.
("brotli", {LANGFUSE_OTEL_COMPRESSION: "none"}, None),
("brotli", {}, "gzip"),
(
None,
{LANGFUSE_OTEL_COMPRESSION: "zstd", OTEL_EXPORTER_OTLP_COMPRESSION: "none"},
None,
),
(
None,
{
OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "deflate",
OTEL_EXPORTER_OTLP_COMPRESSION: "none",
},
None,
),
(None, {OTEL_EXPORTER_OTLP_TRACES_COMPRESSION: "bogus"}, "gzip"),
],
)
def test_default_exporter_compression_precedence(
Expand Down
Loading