Skip to content
Open
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
5 changes: 4 additions & 1 deletion langfuse/_client/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
13 changes: 13 additions & 0 deletions langfuse/_client/environment_variables.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions langfuse/_client/get_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)

Expand Down
7 changes: 6 additions & 1 deletion langfuse/_client/resource_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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]
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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)

Expand Down
39 changes: 38 additions & 1 deletion langfuse/_client/span_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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,
)
Expand Down Expand Up @@ -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
Comment on lines +78 to +84

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Blank argument skips environment setting

If otel_compression is empty or contains only spaces while LANGFUSE_OTEL_COMPRESSION=gzip, this return skips the valid environment setting and silently falls back to the standard OTEL settings. Treat a blank argument as unset, or reject it consistently with other invalid values.

Suggested change
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
if raw_value is None or not raw_value.strip():
setting = LANGFUSE_OTEL_COMPRESSION
raw_value = os.environ.get(LANGFUSE_OTEL_COMPRESSION, "")
value = raw_value.strip().lower()
if not value:
return None
Prompt To Fix With AI
This is a comment left during a code review.
Path: langfuse/_client/span_processor.py
Line: 78-84

Comment:
**Blank argument skips environment setting**

If `otel_compression` is empty or contains only spaces while `LANGFUSE_OTEL_COMPRESSION=gzip`, this return skips the valid environment setting and silently falls back to the standard OTEL settings. Treat a blank argument as unset, or reject it consistently with other invalid values.

```suggestion
    if raw_value is None or not raw_value.strip():
        setting = LANGFUSE_OTEL_COMPRESSION
        raw_value = os.environ.get(LANGFUSE_OTEL_COMPRESSION, "")

    value = raw_value.strip().lower()
    if not value:
        return None
```

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Leaving this as is. otel_compression is typed Literal["gzip", "none"], so a blank string is not a supported value for the argument. The blank check exists for LANGFUSE_OTEL_COMPRESSION set to an empty string in the environment. The JS SDK treats its compression option the same way.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That distinction is reasonable. Since otel_compression is typed as Literal["gzip", "none"], blank values are only expected from LANGFUSE_OTEL_COMPRESSION, where the existing empty-string handling is sufficient. I’ll withdraw this concern; no change is needed.

Tip: You can customize Greptile's behavior for this repo with .greptile/rules.md and .greptile/config.json.


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.

Expand Down Expand Up @@ -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 = (
Expand Down Expand Up @@ -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(
Expand Down
130 changes: 119 additions & 11 deletions tests/unit/test_span_processor.py
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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()

Expand All @@ -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()
Expand Down Expand Up @@ -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(
Expand All @@ -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(
Expand All @@ -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

Expand All @@ -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):
Expand Down Expand Up @@ -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(
Expand Down