From de53a75a98afabd60526e325ffc8aeef305189d8 Mon Sep 17 00:00:00 2001 From: Niklas Semmler Date: Thu, 8 Oct 2026 09:56:32 +0200 Subject: [PATCH 1/5] fix(tracing): mask metadata per key when the mask function fails When a custom mask function raised on dict metadata, the fallback string was written to the bare langfuse.observation.metadata attribute. The server parses that attribute as a JSON object, so the fallback was lost, and nothing showed which keys had been masked. For dict metadata, write the fallback under each key of the failed call instead, so it merges with keys from earlier updates like any other update. Non-dict metadata, input and output keep the plain-string fallback. Co-Authored-By: Claude Opus 5.5 (1M context) --- langfuse/_client/span.py | 20 ++++++++-- tests/unit/test_otel.py | 83 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 100 insertions(+), 3 deletions(-) diff --git a/langfuse/_client/span.py b/langfuse/_client/span.py index 0c1d6db90..756fdb471 100644 --- a/langfuse/_client/span.py +++ b/langfuse/_client/span.py @@ -706,10 +706,17 @@ def _process_media_and_apply_mask( The processed and masked data """ return self._mask_attribute( - data=self._process_media_in_attribute(data=data, field=field) + data=self._process_media_in_attribute(data=data, field=field), field=field ) - def _mask_attribute(self, *, data: Any) -> Any: + def _mask_attribute( + self, + *, + data: Any, + field: Optional[ + Union[Literal["input"], Literal["output"], Literal["metadata"]] + ] = None, + ) -> Any: """Apply the configured mask function to data. Internal method that applies the client's configured masking function to @@ -717,6 +724,7 @@ def _mask_attribute(self, *, data: Any) -> Any: Args: data: The data to mask + field: The attribute the data belongs to Returns: The masked data, or the original data if no mask is configured @@ -732,8 +740,14 @@ def _mask_attribute(self, *, data: Any) -> Any: "data. Using fallback masking. Error: %s", e, ) + fallback = "" + + # Dict metadata is written per key, so mask each key instead of + # writing a plain string to the bare metadata attribute + if field == "metadata" and isinstance(data, dict): + return {key: fallback for key in data} - return "" + return fallback def _process_media_in_attribute( self, diff --git a/tests/unit/test_otel.py b/tests/unit/test_otel.py index 6c4b23d3e..7e7bedb41 100644 --- a/tests/unit/test_otel.py +++ b/tests/unit/test_otel.py @@ -1980,6 +1980,89 @@ def update_random_metadata(thread_id): assert "version" in system_data assert "features" in system_data + MASK_FALLBACK = "" + + def get_metadata_attributes(self, span_data: dict) -> dict: + prefix = LangfuseOtelSpanAttributes.OBSERVATION_METADATA + return { + key: value + for key, value in span_data["attributes"].items() + if key == prefix or key.startswith(f"{prefix}.") + } + + def test_failed_mask_on_update_masks_each_metadata_key( + self, configurable_langfuse_client, memory_exporter + ): + def mask(*, data, **kwargs): + if isinstance(data, dict) and "secret" in data: + raise ValueError("mask failed") + return data + + langfuse_client = configurable_langfuse_client(mask=mask) + span = langfuse_client.start_observation(name="mask-fail", metadata={"a": 1}) + span.update(metadata={"secret": "pw", "b": 2}) + span.end() + + span_data = self.get_spans_by_name(memory_exporter, "mask-fail")[0] + prefix = LangfuseOtelSpanAttributes.OBSERVATION_METADATA + assert self.get_metadata_attributes(span_data) == { + f"{prefix}.a": 1, + f"{prefix}.secret": self.MASK_FALLBACK, + f"{prefix}.b": self.MASK_FALLBACK, + } + + def test_failed_mask_on_start_masks_each_metadata_key( + self, configurable_langfuse_client, memory_exporter + ): + def mask(*, data, **kwargs): + if isinstance(data, dict): + raise ValueError("mask failed") + return data + + langfuse_client = configurable_langfuse_client(mask=mask) + span = langfuse_client.start_observation( + name="mask-fail-start", metadata={"secret": "pw", "b": 2} + ) + span.end() + + span_data = self.get_spans_by_name(memory_exporter, "mask-fail-start")[0] + prefix = LangfuseOtelSpanAttributes.OBSERVATION_METADATA + assert self.get_metadata_attributes(span_data) == { + f"{prefix}.secret": self.MASK_FALLBACK, + f"{prefix}.b": self.MASK_FALLBACK, + } + + def test_failed_mask_keeps_string_fallback_for_non_dict_values( + self, configurable_langfuse_client, memory_exporter + ): + def mask(*, data, **kwargs): + if data is not None: + raise ValueError("mask failed") + return data + + langfuse_client = configurable_langfuse_client(mask=mask) + span = langfuse_client.start_observation( + name="mask-fail-non-dict", + input={"secret": "pw"}, + output={"secret": "pw"}, + metadata="plain-string", + ) + span.end() + + span_data = self.get_spans_by_name(memory_exporter, "mask-fail-non-dict")[0] + attributes = span_data["attributes"] + assert self.get_metadata_attributes(span_data) == { + LangfuseOtelSpanAttributes.OBSERVATION_METADATA: self.MASK_FALLBACK, + } + assert ( + attributes[LangfuseOtelSpanAttributes.OBSERVATION_INPUT] + == self.MASK_FALLBACK + ) + assert ( + attributes[LangfuseOtelSpanAttributes.OBSERVATION_OUTPUT] + == self.MASK_FALLBACK + ) + class TestMultiProjectSetup(TestOTelBase): """Tests for multi-project setup within the same process. From 121c43b43b08f17c0cbc10328ab121ea7f398c07 Mon Sep 17 00:00:00 2001 From: Niklas Semmler Date: Thu, 8 Oct 2026 10:57:29 +0200 Subject: [PATCH 2/5] fix(tracing): apply mask to propagated trace metadata propagate_attributes now applies the client's mask to each propagated trace metadata value (keys are kept), on the raw value before serialization and the 200 character check. If the mask raises, only that value becomes the fallback string. The masked values are stored in the context and baggage, so the current span, child spans, and outgoing baggage all carry them. Matches langfuse-js #1003. Co-Authored-By: Claude Opus 5.5 (1M context) --- langfuse/_client/client.py | 2 +- langfuse/_client/constants.py | 2 + langfuse/_client/propagation.py | 53 +++++++++- langfuse/_client/span.py | 6 +- tests/unit/test_propagate_attributes.py | 130 ++++++++++++++++++++++++ 5 files changed, 188 insertions(+), 5 deletions(-) diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index d8349f802..5ec8cddc7 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -212,7 +212,7 @@ class Langfuse: release (Optional[str]): Release version/hash of your application. Used for grouping analytics by release. media_upload_thread_count (Optional[int]): Number of background threads for handling media uploads. Defaults to 1. Can also be set via LANGFUSE_MEDIA_UPLOAD_THREAD_COUNT environment variable. sample_rate (Optional[float]): Sampling rate for traces (0.0 to 1.0). Defaults to 1.0 (100% of traces are sampled). Can also be set via LANGFUSE_SAMPLE_RATE environment variable. - mask (Optional[MaskFunction]): Function to mask sensitive data synchronously when Langfuse SDK attributes are created. This applies only to data set through Langfuse SDK APIs such as `start_observation()` and `update()`. + mask (Optional[MaskFunction]): Function to mask sensitive data synchronously when Langfuse SDK attributes are created. This applies only to data set through Langfuse SDK APIs such as `start_observation()`, `update()`, and `propagate_attributes()`, where it is applied to each propagated trace metadata value separately (keys are kept). mask_otel_spans (Optional[MaskOtelSpansFunction]): Synchronous export-stage hook for masking raw OpenTelemetry span attributes before this Langfuse client sends them to Langfuse. Use this for spans created by third-party OpenTelemetry instrumentations, or when you need to inspect final span attributes after export filtering and Langfuse media handling. It does not modify spans already exported through other OpenTelemetry exporters. The hook receives one OpenTelemetry export batch. A batch is not guaranteed to contain a complete trace, request, or Langfuse observation tree. The hook usually runs on the OpenTelemetry batch span processor worker thread; during `flush()` and shutdown it may run on the caller thread. Keep it synchronous, deterministic, and fast. diff --git a/langfuse/_client/constants.py b/langfuse/_client/constants.py index c2d0aa7aa..95968507a 100644 --- a/langfuse/_client/constants.py +++ b/langfuse/_client/constants.py @@ -11,6 +11,8 @@ LANGFUSE_SDK_EXPERIMENT_ENVIRONMENT = "sdk-experiment" +MASK_FALLBACK_VALUE = "" + """Note: this type is used with .__args__ / get_args in some cases and therefore must remain flat""" ObservationTypeGenerationLike: TypeAlias = Literal[ "generation", diff --git a/langfuse/_client/propagation.py b/langfuse/_client/propagation.py index a79f57260..3fac90dac 100644 --- a/langfuse/_client/propagation.py +++ b/langfuse/_client/propagation.py @@ -40,11 +40,15 @@ ) from langfuse._client.attributes import LangfuseOtelSpanAttributes -from langfuse._client.constants import LANGFUSE_SDK_EXPERIMENT_ENVIRONMENT +from langfuse._client.constants import ( + LANGFUSE_SDK_EXPERIMENT_ENVIRONMENT, + MASK_FALLBACK_VALUE, +) from langfuse._client.span import _set_span_attributes_within_limit from langfuse._utils.serializer import EventSerializer from langfuse.logger import langfuse_logger from langfuse.model import PromptClient +from langfuse.types import MaskFunction PropagatedKeys = Literal[ "user_id", @@ -157,6 +161,9 @@ def propagate_attributes( metadata: Additional key-value metadata to propagate to all spans. - Keys must be US-ASCII strings - Values are coerced to strings + - If the client has a `mask` function, it is applied to each value + (not the key) before coercion; if it raises, that value becomes + "" - Coerced values must be ≤200 characters - Use for dimensions like internal correlating identifiers - AVOID: large payloads or sensitive data @@ -381,6 +388,10 @@ def _propagate_attributes( as_baggage=as_baggage, ) + # Mask once here: the masked values are stored in the context and baggage, + # so the current span and every child span receive them + mask = _get_current_mask() if metadata else None + for metadata_key, metadata_value in propagated_metadata_attributes.items(): if metadata_value is None: continue @@ -388,6 +399,9 @@ def _propagate_attributes( validated_metadata: Dict[str, str] = {} for key, value in metadata_value.items(): + if mask is not None and metadata_key == "metadata": + value = _mask_propagated_metadata_value(mask=mask, value=value) + serialized_value = _serialize_propagated_metadata_value(value) if serialized_value is None: @@ -436,6 +450,43 @@ def _propagate_attributes( _detach_context_token_safely(token) +def _get_current_mask() -> Optional[MaskFunction]: + """Return the mask of the client that get_client() would resolve, if any. + + Uses the public key in the execution context, otherwise the only initialized + client. With no client, or several clients and no public key, there is no mask. + """ + # Imported here to avoid a circular import via the span processor + from langfuse._client.get_client import _current_public_key + from langfuse._client.resource_manager import LangfuseResourceManager + + with LangfuseResourceManager._lock: + instances = LangfuseResourceManager._instances + public_key = _current_public_key.get(None) + + if public_key: + instance = instances.get(public_key) + elif len(instances) == 1: + instance = next(iter(instances.values())) + else: + instance = None + + return instance.mask if instance is not None else None + + +def _mask_propagated_metadata_value(*, mask: MaskFunction, value: Any) -> Any: + try: + return mask(data=value) + except Exception as e: + langfuse_logger.error( + "Masking error: Custom mask function threw exception when processing " + "propagated metadata. Using fallback masking. Error: %s", + e, + ) + + return MASK_FALLBACK_VALUE + + def _extract_propagated_prompt( prompt: Union[PromptClient, Mapping[str, Any]], ) -> Optional[Tuple[str, int]]: diff --git a/langfuse/_client/span.py b/langfuse/_client/span.py index 756fdb471..52f230dce 100644 --- a/langfuse/_client/span.py +++ b/langfuse/_client/span.py @@ -46,6 +46,7 @@ create_trace_attributes, ) from langfuse._client.constants import ( + MASK_FALLBACK_VALUE, ObservationTypeGenerationLike, ObservationTypeLiteral, ObservationTypeLiteralNoEvent, @@ -740,14 +741,13 @@ def _mask_attribute( "data. Using fallback masking. Error: %s", e, ) - fallback = "" # Dict metadata is written per key, so mask each key instead of # writing a plain string to the bare metadata attribute if field == "metadata" and isinstance(data, dict): - return {key: fallback for key in data} + return {key: MASK_FALLBACK_VALUE for key in data} - return fallback + return MASK_FALLBACK_VALUE def _process_media_in_attribute( self, diff --git a/tests/unit/test_propagate_attributes.py b/tests/unit/test_propagate_attributes.py index 024bac64f..5b4a43701 100644 --- a/tests/unit/test_propagate_attributes.py +++ b/tests/unit/test_propagate_attributes.py @@ -3747,3 +3747,133 @@ def test_prompt_composes_with_outer_propagated_attributes( self.verify_missing_attribute( after_span, LangfuseOtelSpanAttributes.OBSERVATION_PROMPT_NAME ) + + +class TestPropagateAttributesMask(TestPropagateAttributesBase): + """Tests for applying the client's mask to propagated trace metadata.""" + + MASK_FALLBACK = "" + + @pytest.fixture + def masked_langfuse_client(self, monkeypatch, tracer_provider, mock_processor_init): + """Create a mocked Langfuse client with a configurable mask.""" + from langfuse import Langfuse + + def _create_client(mask): + monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "test-public-key") + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "test-secret-key") + + return Langfuse( + public_key="test-public-key", + secret_key="test-secret-key", + host="http://test-host", + tracing_enabled=True, + tracer_provider=tracer_provider, + mask=mask, + ) + + return _create_client + + def get_trace_metadata(self, span_data: dict) -> dict: + prefix = f"{LangfuseOtelSpanAttributes.TRACE_METADATA}." + return { + key[len(prefix) :]: value + for key, value in span_data["attributes"].items() + if key.startswith(prefix) + } + + def test_mask_applies_to_current_and_child_spans( + self, masked_langfuse_client, memory_exporter + ): + """Verify each raw value is masked once and keys are kept.""" + masked_values = [] + + def mask(*, data, **kwargs): + # Observation input/output/metadata are masked too and are None here + if data is None: + return None + masked_values.append(data) + if isinstance(data, dict): + return {k: "***" if k == "token" else v for k, v in data.items()} + return data.replace("secret", "***") + + langfuse_client = masked_langfuse_client(mask) + + with langfuse_client.start_as_current_observation(name="parent-span"): + with propagate_attributes( + metadata={"api_key": "secret-key", "auth": {"token": "secret"}} + ): + child = langfuse_client.start_observation(name="child-span") + child.end() + + expected = {"api_key": "***-key", "auth": '{"token":"***"}'} + for name in ("parent-span", "child-span"): + span_data = self.get_span_by_name(memory_exporter, name) + assert self.get_trace_metadata(span_data) == expected + + assert masked_values == ["secret-key", {"token": "secret"}] + + def test_mask_applies_to_baggage(self, masked_langfuse_client): + """Verify baggage carries masked values.""" + from opentelemetry import baggage + + langfuse_client = masked_langfuse_client(lambda *, data, **kwargs: "***") + + with langfuse_client.start_as_current_observation(name="parent-span"): + with propagate_attributes(metadata={"api_key": "secret"}, as_baggage=True): + assert baggage.get_baggage("langfuse_metadata_api_key") == "***" + + def test_failed_mask_replaces_only_that_value( + self, masked_langfuse_client, memory_exporter + ): + """Verify a throwing mask gives the fallback for that value only.""" + + def mask(*, data, **kwargs): + if data == "secret": + raise ValueError("mask failed") + return data + + langfuse_client = masked_langfuse_client(mask) + + with langfuse_client.start_as_current_observation(name="parent-span"): + with propagate_attributes(metadata={"api_key": "secret", "env": "prod"}): + child = langfuse_client.start_observation(name="child-span") + child.end() + + for name in ("parent-span", "child-span"): + span_data = self.get_span_by_name(memory_exporter, name) + assert self.get_trace_metadata(span_data) == { + "api_key": self.MASK_FALLBACK, + "env": "prod", + } + + def test_length_limit_applies_to_masked_value( + self, masked_langfuse_client, memory_exporter + ): + """Verify the 200 character limit is checked after masking.""" + + def mask(*, data, **kwargs): + if data is None: + return None + return "x" * 201 if data == "expand" else data[:10] + + langfuse_client = masked_langfuse_client(mask) + + with langfuse_client.start_as_current_observation(name="parent-span"): + with propagate_attributes(metadata={"long": "a" * 300, "short": "expand"}): + child = langfuse_client.start_observation(name="child-span") + child.end() + + span_data = self.get_span_by_name(memory_exporter, "child-span") + assert self.get_trace_metadata(span_data) == {"long": "a" * 10} + + def test_no_mask_leaves_metadata_unchanged(self, langfuse_client, memory_exporter): + """Verify metadata is unchanged when no mask is configured.""" + with langfuse_client.start_as_current_observation(name="parent-span"): + with propagate_attributes(metadata={"api_key": "secret", "n": 1}): + child = langfuse_client.start_observation(name="child-span") + child.end() + + for name in ("parent-span", "child-span"): + span_data = self.get_span_by_name(memory_exporter, name) + assert self.get_trace_metadata(span_data) == {"api_key": "secret", "n": "1"} From 37a31a21c0c7582f9ddb3836d717ed1b70eb4c30 Mon Sep 17 00:00:00 2001 From: Niklas Semmler Date: Thu, 8 Oct 2026 11:24:21 +0200 Subject: [PATCH 3/5] docs(tracing): document which client's mask applies to propagated metadata The mask for propagated trace metadata comes from the client for the public key in the execution context, otherwise the only initialized client. With several clients and no public key in context, or with no client yet, propagated metadata is not masked. Document this and the workarounds, and pin the behavior with a test. Co-Authored-By: Claude Opus 5.5 (1M context) --- langfuse/_client/client.py | 2 +- langfuse/_client/propagation.py | 8 +++++- tests/unit/test_propagate_attributes.py | 36 +++++++++++++++++++++++++ 3 files changed, 44 insertions(+), 2 deletions(-) diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index 5ec8cddc7..276818bb9 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -212,7 +212,7 @@ class Langfuse: release (Optional[str]): Release version/hash of your application. Used for grouping analytics by release. media_upload_thread_count (Optional[int]): Number of background threads for handling media uploads. Defaults to 1. Can also be set via LANGFUSE_MEDIA_UPLOAD_THREAD_COUNT environment variable. sample_rate (Optional[float]): Sampling rate for traces (0.0 to 1.0). Defaults to 1.0 (100% of traces are sampled). Can also be set via LANGFUSE_SAMPLE_RATE environment variable. - mask (Optional[MaskFunction]): Function to mask sensitive data synchronously when Langfuse SDK attributes are created. This applies only to data set through Langfuse SDK APIs such as `start_observation()`, `update()`, and `propagate_attributes()`, where it is applied to each propagated trace metadata value separately (keys are kept). + mask (Optional[MaskFunction]): Function to mask sensitive data synchronously when Langfuse SDK attributes are created. This applies only to data set through Langfuse SDK APIs such as `start_observation()`, `update()`, and `propagate_attributes()`, where it is applied to each propagated trace metadata value separately (keys are kept). For propagated metadata, the mask comes from the client for the public key in the execution context, otherwise from the only initialized client. With several clients and no public key in context, or with no client yet, propagated metadata is not masked. To avoid this, pass `langfuse_public_key` via `@observe` or use a single client; `mask_otel_spans` also sees these attributes at export, but not outgoing baggage. mask_otel_spans (Optional[MaskOtelSpansFunction]): Synchronous export-stage hook for masking raw OpenTelemetry span attributes before this Langfuse client sends them to Langfuse. Use this for spans created by third-party OpenTelemetry instrumentations, or when you need to inspect final span attributes after export filtering and Langfuse media handling. It does not modify spans already exported through other OpenTelemetry exporters. The hook receives one OpenTelemetry export batch. A batch is not guaranteed to contain a complete trace, request, or Langfuse observation tree. The hook usually runs on the OpenTelemetry batch span processor worker thread; during `flush()` and shutdown it may run on the caller thread. Keep it synchronous, deterministic, and fast. diff --git a/langfuse/_client/propagation.py b/langfuse/_client/propagation.py index 3fac90dac..81e1b3d68 100644 --- a/langfuse/_client/propagation.py +++ b/langfuse/_client/propagation.py @@ -163,7 +163,11 @@ def propagate_attributes( - Values are coerced to strings - If the client has a `mask` function, it is applied to each value (not the key) before coercion; if it raises, that value becomes - "" + "". The mask comes from + the client for the public key in the execution context, otherwise + the only initialized client. With several clients and no public key + in context, or with no client yet, values are not masked; pass + `langfuse_public_key` via `@observe` or use a single client - Coerced values must be ≤200 characters - Use for dimensions like internal correlating identifiers - AVOID: large payloads or sensitive data @@ -456,6 +460,8 @@ def _get_current_mask() -> Optional[MaskFunction]: Uses the public key in the execution context, otherwise the only initialized client. With no client, or several clients and no public key, there is no mask. """ + # Known limitation: in those cases propagated metadata is exported unmasked + # even if a client has a mask, since there is no single client to pick # Imported here to avoid a circular import via the span processor from langfuse._client.get_client import _current_public_key from langfuse._client.resource_manager import LangfuseResourceManager diff --git a/tests/unit/test_propagate_attributes.py b/tests/unit/test_propagate_attributes.py index 5b4a43701..fe73125b7 100644 --- a/tests/unit/test_propagate_attributes.py +++ b/tests/unit/test_propagate_attributes.py @@ -3867,6 +3867,42 @@ def mask(*, data, **kwargs): span_data = self.get_span_by_name(memory_exporter, "child-span") assert self.get_trace_metadata(span_data) == {"long": "a" * 10} + def test_mask_with_several_clients_needs_public_key_in_context( + self, monkeypatch, tracer_provider, mock_processor_init, memory_exporter + ): + """Verify the known limitation: several clients and no key means no mask.""" + from langfuse import Langfuse + from langfuse._client.get_client import _set_current_public_key + + monkeypatch.setenv("LANGFUSE_SECRET_KEY", "test-secret-key") + masked_client = Langfuse( + public_key="pk-masked", + secret_key="test-secret-key", + host="http://test-host", + tracer_provider=tracer_provider, + mask=lambda *, data, **kwargs: None if data is None else "***", + ) + Langfuse( + public_key="pk-other", + secret_key="test-secret-key", + host="http://test-host", + tracer_provider=tracer_provider, + ) + + with masked_client.start_as_current_observation(name="no-key-span"): + with propagate_attributes(metadata={"api_key": "secret"}): + pass + + with _set_current_public_key("pk-masked"): + with masked_client.start_as_current_observation(name="key-span"): + with propagate_attributes(metadata={"api_key": "secret"}): + pass + + no_key_span = self.get_span_by_name(memory_exporter, "no-key-span") + assert self.get_trace_metadata(no_key_span) == {"api_key": "secret"} + key_span = self.get_span_by_name(memory_exporter, "key-span") + assert self.get_trace_metadata(key_span) == {"api_key": "***"} + def test_no_mask_leaves_metadata_unchanged(self, langfuse_client, memory_exporter): """Verify metadata is unchanged when no mask is configured.""" with langfuse_client.start_as_current_observation(name="parent-span"): From 68f16e439629f0b2e414540be38ecc1492231901 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 9 Oct 2026 09:09:45 +0000 Subject: [PATCH 4/5] revert(tracing): do not mask propagated trace metadata Propagated metadata is passed explicitly by the developer. Masking it adds client-resolution ambiguity and baggage/length-limit ordering complexity, and diverges from the JS SDK. Co-authored-by: Hassieb Pakzad --- langfuse/_client/client.py | 2 +- langfuse/_client/propagation.py | 59 +-------- tests/unit/test_propagate_attributes.py | 166 ------------------------ 3 files changed, 2 insertions(+), 225 deletions(-) diff --git a/langfuse/_client/client.py b/langfuse/_client/client.py index 276818bb9..d8349f802 100644 --- a/langfuse/_client/client.py +++ b/langfuse/_client/client.py @@ -212,7 +212,7 @@ class Langfuse: release (Optional[str]): Release version/hash of your application. Used for grouping analytics by release. media_upload_thread_count (Optional[int]): Number of background threads for handling media uploads. Defaults to 1. Can also be set via LANGFUSE_MEDIA_UPLOAD_THREAD_COUNT environment variable. sample_rate (Optional[float]): Sampling rate for traces (0.0 to 1.0). Defaults to 1.0 (100% of traces are sampled). Can also be set via LANGFUSE_SAMPLE_RATE environment variable. - mask (Optional[MaskFunction]): Function to mask sensitive data synchronously when Langfuse SDK attributes are created. This applies only to data set through Langfuse SDK APIs such as `start_observation()`, `update()`, and `propagate_attributes()`, where it is applied to each propagated trace metadata value separately (keys are kept). For propagated metadata, the mask comes from the client for the public key in the execution context, otherwise from the only initialized client. With several clients and no public key in context, or with no client yet, propagated metadata is not masked. To avoid this, pass `langfuse_public_key` via `@observe` or use a single client; `mask_otel_spans` also sees these attributes at export, but not outgoing baggage. + mask (Optional[MaskFunction]): Function to mask sensitive data synchronously when Langfuse SDK attributes are created. This applies only to data set through Langfuse SDK APIs such as `start_observation()` and `update()`. mask_otel_spans (Optional[MaskOtelSpansFunction]): Synchronous export-stage hook for masking raw OpenTelemetry span attributes before this Langfuse client sends them to Langfuse. Use this for spans created by third-party OpenTelemetry instrumentations, or when you need to inspect final span attributes after export filtering and Langfuse media handling. It does not modify spans already exported through other OpenTelemetry exporters. The hook receives one OpenTelemetry export batch. A batch is not guaranteed to contain a complete trace, request, or Langfuse observation tree. The hook usually runs on the OpenTelemetry batch span processor worker thread; during `flush()` and shutdown it may run on the caller thread. Keep it synchronous, deterministic, and fast. diff --git a/langfuse/_client/propagation.py b/langfuse/_client/propagation.py index 81e1b3d68..a79f57260 100644 --- a/langfuse/_client/propagation.py +++ b/langfuse/_client/propagation.py @@ -40,15 +40,11 @@ ) from langfuse._client.attributes import LangfuseOtelSpanAttributes -from langfuse._client.constants import ( - LANGFUSE_SDK_EXPERIMENT_ENVIRONMENT, - MASK_FALLBACK_VALUE, -) +from langfuse._client.constants import LANGFUSE_SDK_EXPERIMENT_ENVIRONMENT from langfuse._client.span import _set_span_attributes_within_limit from langfuse._utils.serializer import EventSerializer from langfuse.logger import langfuse_logger from langfuse.model import PromptClient -from langfuse.types import MaskFunction PropagatedKeys = Literal[ "user_id", @@ -161,13 +157,6 @@ def propagate_attributes( metadata: Additional key-value metadata to propagate to all spans. - Keys must be US-ASCII strings - Values are coerced to strings - - If the client has a `mask` function, it is applied to each value - (not the key) before coercion; if it raises, that value becomes - "". The mask comes from - the client for the public key in the execution context, otherwise - the only initialized client. With several clients and no public key - in context, or with no client yet, values are not masked; pass - `langfuse_public_key` via `@observe` or use a single client - Coerced values must be ≤200 characters - Use for dimensions like internal correlating identifiers - AVOID: large payloads or sensitive data @@ -392,10 +381,6 @@ def _propagate_attributes( as_baggage=as_baggage, ) - # Mask once here: the masked values are stored in the context and baggage, - # so the current span and every child span receive them - mask = _get_current_mask() if metadata else None - for metadata_key, metadata_value in propagated_metadata_attributes.items(): if metadata_value is None: continue @@ -403,9 +388,6 @@ def _propagate_attributes( validated_metadata: Dict[str, str] = {} for key, value in metadata_value.items(): - if mask is not None and metadata_key == "metadata": - value = _mask_propagated_metadata_value(mask=mask, value=value) - serialized_value = _serialize_propagated_metadata_value(value) if serialized_value is None: @@ -454,45 +436,6 @@ def _propagate_attributes( _detach_context_token_safely(token) -def _get_current_mask() -> Optional[MaskFunction]: - """Return the mask of the client that get_client() would resolve, if any. - - Uses the public key in the execution context, otherwise the only initialized - client. With no client, or several clients and no public key, there is no mask. - """ - # Known limitation: in those cases propagated metadata is exported unmasked - # even if a client has a mask, since there is no single client to pick - # Imported here to avoid a circular import via the span processor - from langfuse._client.get_client import _current_public_key - from langfuse._client.resource_manager import LangfuseResourceManager - - with LangfuseResourceManager._lock: - instances = LangfuseResourceManager._instances - public_key = _current_public_key.get(None) - - if public_key: - instance = instances.get(public_key) - elif len(instances) == 1: - instance = next(iter(instances.values())) - else: - instance = None - - return instance.mask if instance is not None else None - - -def _mask_propagated_metadata_value(*, mask: MaskFunction, value: Any) -> Any: - try: - return mask(data=value) - except Exception as e: - langfuse_logger.error( - "Masking error: Custom mask function threw exception when processing " - "propagated metadata. Using fallback masking. Error: %s", - e, - ) - - return MASK_FALLBACK_VALUE - - def _extract_propagated_prompt( prompt: Union[PromptClient, Mapping[str, Any]], ) -> Optional[Tuple[str, int]]: diff --git a/tests/unit/test_propagate_attributes.py b/tests/unit/test_propagate_attributes.py index fe73125b7..024bac64f 100644 --- a/tests/unit/test_propagate_attributes.py +++ b/tests/unit/test_propagate_attributes.py @@ -3747,169 +3747,3 @@ def test_prompt_composes_with_outer_propagated_attributes( self.verify_missing_attribute( after_span, LangfuseOtelSpanAttributes.OBSERVATION_PROMPT_NAME ) - - -class TestPropagateAttributesMask(TestPropagateAttributesBase): - """Tests for applying the client's mask to propagated trace metadata.""" - - MASK_FALLBACK = "" - - @pytest.fixture - def masked_langfuse_client(self, monkeypatch, tracer_provider, mock_processor_init): - """Create a mocked Langfuse client with a configurable mask.""" - from langfuse import Langfuse - - def _create_client(mask): - monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "test-public-key") - monkeypatch.setenv("LANGFUSE_SECRET_KEY", "test-secret-key") - - return Langfuse( - public_key="test-public-key", - secret_key="test-secret-key", - host="http://test-host", - tracing_enabled=True, - tracer_provider=tracer_provider, - mask=mask, - ) - - return _create_client - - def get_trace_metadata(self, span_data: dict) -> dict: - prefix = f"{LangfuseOtelSpanAttributes.TRACE_METADATA}." - return { - key[len(prefix) :]: value - for key, value in span_data["attributes"].items() - if key.startswith(prefix) - } - - def test_mask_applies_to_current_and_child_spans( - self, masked_langfuse_client, memory_exporter - ): - """Verify each raw value is masked once and keys are kept.""" - masked_values = [] - - def mask(*, data, **kwargs): - # Observation input/output/metadata are masked too and are None here - if data is None: - return None - masked_values.append(data) - if isinstance(data, dict): - return {k: "***" if k == "token" else v for k, v in data.items()} - return data.replace("secret", "***") - - langfuse_client = masked_langfuse_client(mask) - - with langfuse_client.start_as_current_observation(name="parent-span"): - with propagate_attributes( - metadata={"api_key": "secret-key", "auth": {"token": "secret"}} - ): - child = langfuse_client.start_observation(name="child-span") - child.end() - - expected = {"api_key": "***-key", "auth": '{"token":"***"}'} - for name in ("parent-span", "child-span"): - span_data = self.get_span_by_name(memory_exporter, name) - assert self.get_trace_metadata(span_data) == expected - - assert masked_values == ["secret-key", {"token": "secret"}] - - def test_mask_applies_to_baggage(self, masked_langfuse_client): - """Verify baggage carries masked values.""" - from opentelemetry import baggage - - langfuse_client = masked_langfuse_client(lambda *, data, **kwargs: "***") - - with langfuse_client.start_as_current_observation(name="parent-span"): - with propagate_attributes(metadata={"api_key": "secret"}, as_baggage=True): - assert baggage.get_baggage("langfuse_metadata_api_key") == "***" - - def test_failed_mask_replaces_only_that_value( - self, masked_langfuse_client, memory_exporter - ): - """Verify a throwing mask gives the fallback for that value only.""" - - def mask(*, data, **kwargs): - if data == "secret": - raise ValueError("mask failed") - return data - - langfuse_client = masked_langfuse_client(mask) - - with langfuse_client.start_as_current_observation(name="parent-span"): - with propagate_attributes(metadata={"api_key": "secret", "env": "prod"}): - child = langfuse_client.start_observation(name="child-span") - child.end() - - for name in ("parent-span", "child-span"): - span_data = self.get_span_by_name(memory_exporter, name) - assert self.get_trace_metadata(span_data) == { - "api_key": self.MASK_FALLBACK, - "env": "prod", - } - - def test_length_limit_applies_to_masked_value( - self, masked_langfuse_client, memory_exporter - ): - """Verify the 200 character limit is checked after masking.""" - - def mask(*, data, **kwargs): - if data is None: - return None - return "x" * 201 if data == "expand" else data[:10] - - langfuse_client = masked_langfuse_client(mask) - - with langfuse_client.start_as_current_observation(name="parent-span"): - with propagate_attributes(metadata={"long": "a" * 300, "short": "expand"}): - child = langfuse_client.start_observation(name="child-span") - child.end() - - span_data = self.get_span_by_name(memory_exporter, "child-span") - assert self.get_trace_metadata(span_data) == {"long": "a" * 10} - - def test_mask_with_several_clients_needs_public_key_in_context( - self, monkeypatch, tracer_provider, mock_processor_init, memory_exporter - ): - """Verify the known limitation: several clients and no key means no mask.""" - from langfuse import Langfuse - from langfuse._client.get_client import _set_current_public_key - - monkeypatch.setenv("LANGFUSE_SECRET_KEY", "test-secret-key") - masked_client = Langfuse( - public_key="pk-masked", - secret_key="test-secret-key", - host="http://test-host", - tracer_provider=tracer_provider, - mask=lambda *, data, **kwargs: None if data is None else "***", - ) - Langfuse( - public_key="pk-other", - secret_key="test-secret-key", - host="http://test-host", - tracer_provider=tracer_provider, - ) - - with masked_client.start_as_current_observation(name="no-key-span"): - with propagate_attributes(metadata={"api_key": "secret"}): - pass - - with _set_current_public_key("pk-masked"): - with masked_client.start_as_current_observation(name="key-span"): - with propagate_attributes(metadata={"api_key": "secret"}): - pass - - no_key_span = self.get_span_by_name(memory_exporter, "no-key-span") - assert self.get_trace_metadata(no_key_span) == {"api_key": "secret"} - key_span = self.get_span_by_name(memory_exporter, "key-span") - assert self.get_trace_metadata(key_span) == {"api_key": "***"} - - def test_no_mask_leaves_metadata_unchanged(self, langfuse_client, memory_exporter): - """Verify metadata is unchanged when no mask is configured.""" - with langfuse_client.start_as_current_observation(name="parent-span"): - with propagate_attributes(metadata={"api_key": "secret", "n": 1}): - child = langfuse_client.start_observation(name="child-span") - child.end() - - for name in ("parent-span", "child-span"): - span_data = self.get_span_by_name(memory_exporter, name) - assert self.get_trace_metadata(span_data) == {"api_key": "secret", "n": "1"} From 0351764321f7623d8c9a88a7b09633f186473417 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 9 Oct 2026 09:11:12 +0000 Subject: [PATCH 5/5] fix(tracing): keep string fallback when the mask fails on empty dict metadata Co-authored-by: Hassieb Pakzad --- langfuse/_client/span.py | 5 +++-- tests/unit/test_otel.py | 19 +++++++++++++++++++ 2 files changed, 22 insertions(+), 2 deletions(-) diff --git a/langfuse/_client/span.py b/langfuse/_client/span.py index 52f230dce..21bb9fabb 100644 --- a/langfuse/_client/span.py +++ b/langfuse/_client/span.py @@ -743,8 +743,9 @@ def _mask_attribute( ) # Dict metadata is written per key, so mask each key instead of - # writing a plain string to the bare metadata attribute - if field == "metadata" and isinstance(data, dict): + # writing a plain string to the bare metadata attribute. An empty dict + # has no keys to mask, so it keeps the plain string to stay visible + if field == "metadata" and isinstance(data, dict) and data: return {key: MASK_FALLBACK_VALUE for key in data} return MASK_FALLBACK_VALUE diff --git a/tests/unit/test_otel.py b/tests/unit/test_otel.py index 7e7bedb41..286a5e4d0 100644 --- a/tests/unit/test_otel.py +++ b/tests/unit/test_otel.py @@ -2032,6 +2032,25 @@ def mask(*, data, **kwargs): f"{prefix}.b": self.MASK_FALLBACK, } + def test_failed_mask_keeps_string_fallback_for_empty_dict_metadata( + self, configurable_langfuse_client, memory_exporter + ): + def mask(*, data, **kwargs): + if isinstance(data, dict): + raise ValueError("mask failed") + return data + + langfuse_client = configurable_langfuse_client(mask=mask) + span = langfuse_client.start_observation( + name="mask-fail-empty-dict", metadata={} + ) + span.end() + + span_data = self.get_spans_by_name(memory_exporter, "mask-fail-empty-dict")[0] + assert self.get_metadata_attributes(span_data) == { + LangfuseOtelSpanAttributes.OBSERVATION_METADATA: self.MASK_FALLBACK, + } + def test_failed_mask_keeps_string_fallback_for_non_dict_values( self, configurable_langfuse_client, memory_exporter ):