Skip to content
Draft
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
31 changes: 28 additions & 3 deletions langfuse/_client/propagation.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
propagate to all child spans within the context.
"""

import ast
import json
import re
from typing import (
Any,
Expand Down Expand Up @@ -158,7 +160,8 @@ def propagate_attributes(
- Use for dimensions like internal correlating identifiers
- AVOID: large payloads or sensitive data
version: Version identfier for parts of your application that are independently versioned, e.g. agents
tags: List of tags to categorize the group of observations
tags: List of tags to categorize the group of observations. Appended to tags
inherited from the current context, including cross-service baggage.
trace_name: Name to assign to the trace. Must be US-ASCII string, ≤200 characters.
Use this to set a consistent trace name for all spans created within this context.
prompt: Langfuse prompt to link to generations created within this context.
Expand Down Expand Up @@ -492,6 +495,12 @@ def _get_propagated_attributes_from_context(
propagated_attributes[span_key] = int(baggage_value)
continue

if span_key == LangfuseOtelSpanAttributes.TRACE_TAGS and isinstance(
baggage_value, str
):
propagated_attributes[span_key] = _parse_baggage_tags(baggage_value)
continue

propagated_attributes[span_key] = (
baggage_value
if isinstance(baggage_value, (str, list))
Expand Down Expand Up @@ -542,6 +551,22 @@ def _get_propagated_attributes_from_context(
return propagated_attributes


def _parse_baggage_tags(value: str) -> List[str]:
# The Python SDK writes str(list), e.g. "['a', 'b']", and the JS SDK writes "a,b".
# JSON goes first because literal_eval does not join JSON's escaped surrogate pairs.
if value.startswith("["):
for parse in (json.loads, ast.literal_eval):
try:
tags = parse(value)
except Exception:
continue

if isinstance(tags, list) and all(isinstance(tag, str) for tag in tags):
return tags

return value.split(",")


def _set_propagated_attribute(
*,
key: str,
Expand All @@ -562,10 +587,10 @@ def _set_propagated_attribute(
)
value = existing_metadata_in_context | value

# Merge tags with previously set tags
# Merge with inherited tags, including baggage after a process boundary.
if isinstance(value, list):
existing_tags_in_context = cast(
list, otel_context_api.get_value(context_key) or []
list, _get_propagated_attributes_from_context(context).get(span_key) or []
)
merged_tags = list(existing_tags_in_context)
merged_tags.extend(tag for tag in value if tag not in existing_tags_in_context)
Expand Down
148 changes: 148 additions & 0 deletions tests/unit/test_propagate_attributes.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
"""

import concurrent.futures
import json
from datetime import datetime

import pytest
Expand Down Expand Up @@ -1849,6 +1850,153 @@ def test_baggage_survives_context_isolation(self, langfuse_client, memory_export
"cross_process_session",
)

def test_tags_survive_w3c_baggage_header(self, langfuse_client, memory_exporter):
"""Verify tags keep their list shape after crossing a W3C baggage header."""
from opentelemetry import context as otel_context
from opentelemetry.baggage.propagation import W3CBaggagePropagator

propagator = W3CBaggagePropagator()
carrier: dict = {}

with langfuse_client.start_as_current_observation(name="upstream"):
with propagate_attributes(tags=["tag-a", "comma,tag"], as_baggage=True):
propagator.inject(carrier)

# Only the header reaches the downstream service
downstream_context = propagator.extract(carrier, context=otel_context.Context())

token = otel_context.attach(downstream_context)
try:
child = langfuse_client.start_observation(name="downstream")
child.end()
finally:
otel_context.detach(token)

downstream_span = self.get_span_by_name(memory_exporter, "downstream")
self.verify_span_attribute(
downstream_span,
LangfuseOtelSpanAttributes.TRACE_TAGS,
tuple(["tag-a", "comma,tag"]),
)

@pytest.mark.parametrize("as_baggage", [False, True])
def test_tags_append_after_baggage_header_extraction(
self, langfuse_client, memory_exporter, as_baggage
):
"""Append inherited tags locally and forward them only when requested."""
from opentelemetry import context as otel_context
from opentelemetry.baggage.propagation import W3CBaggagePropagator

propagator = W3CBaggagePropagator()
incoming: dict = {}
outgoing: dict = {}
inherited_tags = ["upstream", "shared", "comma,tag"]
merged_tags = (*inherited_tags, "worker")

with propagate_attributes(tags=inherited_tags, as_baggage=True):
propagator.inject(incoming)

received_context = propagator.extract(incoming, context=otel_context.Context())
token = otel_context.attach(received_context)
try:
with langfuse_client.start_as_current_observation(name="receiver"):
with propagate_attributes(
tags=["shared", "worker"], as_baggage=as_baggage
):
with langfuse_client.start_as_current_observation(name="child"):
propagator.inject(outgoing)
with langfuse_client.start_as_current_observation(name="restored"):
pass
finally:
otel_context.detach(token)

forwarded_context = propagator.extract(outgoing, context=otel_context.Context())
token = otel_context.attach(forwarded_context)
try:
with langfuse_client.start_as_current_observation(name="next-service"):
pass
finally:
otel_context.detach(token)

for name in ("receiver", "child"):
self.verify_span_attribute(
self.get_span_by_name(memory_exporter, name),
LangfuseOtelSpanAttributes.TRACE_TAGS,
merged_tags,
)
self.verify_span_attribute(
self.get_span_by_name(memory_exporter, "restored"),
LangfuseOtelSpanAttributes.TRACE_TAGS,
tuple(inherited_tags),
)
self.verify_span_attribute(
self.get_span_by_name(memory_exporter, "next-service"),
LangfuseOtelSpanAttributes.TRACE_TAGS,
merged_tags if as_baggage else tuple(inherited_tags),
)

def test_tags_append_prefers_local_context_over_baggage(
self, langfuse_client, memory_exporter
):
"""Keep local tags authoritative when baggage has different values."""
from opentelemetry import baggage
from opentelemetry import context as otel_context

with propagate_attributes(tags=["local"]):
token = otel_context.attach(
baggage.set_baggage("langfuse_tags", "['baggage-only']")
)
try:
with propagate_attributes(tags=["worker"], as_baggage=True):
with langfuse_client.start_as_current_observation(name="child"):
pass
finally:
otel_context.detach(token)

self.verify_span_attribute(
self.get_span_by_name(memory_exporter, "child"),
LangfuseOtelSpanAttributes.TRACE_TAGS,
("local", "worker"),
)

@pytest.mark.parametrize(
("tags_value", "expected_tags"),
[
# JS SDK
("tag-a,tag-b", ["tag-a", "tag-b"]),
# JSON array; json.dumps escapes the emoji as a surrogate pair
(json.dumps(["tag-a", "\U0001f600"]), ["tag-a", "\U0001f600"]),
],
ids=["js-sdk", "json-array"],
)
def test_tags_read_from_js_and_json_baggage(
self, langfuse_client, memory_exporter, tags_value, expected_tags
):
"""Verify tags from JS SDK or JSON-array baggage are read as a list."""
from urllib.parse import quote_plus

from opentelemetry import context as otel_context
from opentelemetry.baggage.propagation import W3CBaggagePropagator

downstream_context = W3CBaggagePropagator().extract(
{"baggage": f"langfuse_tags={quote_plus(tags_value)}"},
context=otel_context.Context(),
)

token = otel_context.attach(downstream_context)
try:
child = langfuse_client.start_observation(name="downstream")
child.end()
finally:
otel_context.detach(token)

downstream_span = self.get_span_by_name(memory_exporter, "downstream")
self.verify_span_attribute(
downstream_span,
LangfuseOtelSpanAttributes.TRACE_TAGS,
tuple(expected_tags),
)


class TestPropagateAttributesEnvironment(TestPropagateAttributesBase):
"""Tests for first-class Langfuse environment propagation."""
Expand Down