Skip to content
Merged
26 changes: 0 additions & 26 deletions langfuse/_client/attributes.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,32 +160,6 @@ def _serialize(obj: Any) -> Optional[str]:
return json.dumps(obj, cls=EventSerializer)


def _flatten_and_serialize_metadata_values(
metadata: Optional[Dict[str, Any]],
) -> Optional[Dict[str, str]]:
if metadata is None:
return None

flattened_metadata: Dict[str, str] = {}

def flatten_value(path: str, value: Any) -> None:
if isinstance(value, dict):
for nested_key, nested_value in value.items():
flatten_value(f"{path}.{nested_key}", nested_value)

return

serialized_value = _serialize(value)

if serialized_value is not None:
flattened_metadata[path] = serialized_value

for key, value in metadata.items():
flatten_value(str(key), value)

return flattened_metadata


def _flatten_and_serialize_metadata(
metadata: Any, type: Literal["observation", "trace"]
) -> dict:
Expand Down
63 changes: 41 additions & 22 deletions langfuse/_client/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@

from langfuse._client.attributes import (
LangfuseOtelSpanAttributes,
_flatten_and_serialize_metadata_values,
_serialize,
)
from langfuse._client.constants import (
Expand Down Expand Up @@ -83,6 +82,7 @@
LangfuseRetriever,
LangfuseSpan,
LangfuseTool,
_set_span_attributes_within_limit,
)
from langfuse._client.utils import (
get_sha256_hash_hex,
Expand Down Expand Up @@ -694,7 +694,9 @@ def start_observation(
cast(otel_trace_api.Span, remote_parent_span)
):
otel_span = self._otel_tracer.start_span(name=name)
otel_span.set_attribute(LangfuseOtelSpanAttributes.AS_ROOT, True)
_set_span_attributes_within_limit(
otel_span, {LangfuseOtelSpanAttributes.AS_ROOT: True}
)

return self._create_observation_from_otel_span(
otel_span=otel_span,
Expand Down Expand Up @@ -1284,8 +1286,9 @@ def _create_span_with_parent_context(
prompt=prompt,
) as langfuse_span:
if remote_parent_span is not None:
langfuse_span._otel_span.set_attribute(
LangfuseOtelSpanAttributes.AS_ROOT, True
_set_span_attributes_within_limit(
langfuse_span._otel_span,
{LangfuseOtelSpanAttributes.AS_ROOT: True},
)

yield langfuse_span
Expand Down Expand Up @@ -1615,7 +1618,9 @@ def create_event(
otel_span = self._otel_tracer.start_span(
name=name, start_time=timestamp
)
otel_span.set_attribute(LangfuseOtelSpanAttributes.AS_ROOT, True)
_set_span_attributes_within_limit(
otel_span, {LangfuseOtelSpanAttributes.AS_ROOT: True}
)

return cast(
LangfuseEvent,
Expand Down Expand Up @@ -2886,17 +2891,15 @@ async def _process_experiment_item(
else getattr(item, "metadata", None)
)

final_observation_metadata = {
**(item_metadata if isinstance(item_metadata, dict) else {}),
**(experiment_metadata or {}),
"experiment_name": experiment_name,
"experiment_run_name": experiment_run_name,
}

trace_id = span.trace_id
dataset_id = None
dataset_item_id = None

experiment_run_metadata: Dict[str, Any] = {
"experiment_name": experiment_name,
"experiment_run_name": experiment_run_name,
}

if (
not isinstance(item, dict)
and hasattr(item, "dataset_id")
Expand All @@ -2905,10 +2908,20 @@ async def _process_experiment_item(
dataset_id = item.dataset_id
dataset_item_id = item.id

final_observation_metadata.update(
experiment_run_metadata.update(
{"dataset_id": dataset_id, "dataset_item_id": dataset_item_id}
)

# Experiment run keys go first so the span attribute limit drops
# user metadata before them, and last so they still win over
# user keys.
final_observation_metadata = {
**experiment_run_metadata,
**(item_metadata if isinstance(item_metadata, dict) else {}),
**(experiment_metadata or {}),
**experiment_run_metadata,
}

experiment_item_id = (
dataset_item_id or get_sha256_hash_hex(_serialize(input_data))[:16]
)
Expand All @@ -2926,25 +2939,26 @@ async def _process_experiment_item(
}.items()
if v is not None
}
span._otel_span.set_attributes(experiment_span_attributes)
_set_span_attributes_within_limit(
span._otel_span, experiment_span_attributes
)

with span.start_as_current_observation(
name="experiment-item-task",
as_type="span",
input=input_data,
metadata=final_observation_metadata,
) as task_span:
task_span._otel_span.set_attributes(experiment_span_attributes)
_set_span_attributes_within_limit(
task_span._otel_span, experiment_span_attributes
)

propagated_experiment_attributes = PropagatedExperimentAttributes(
experiment_id=experiment_id,
experiment_name=experiment_run_name,
experiment_metadata=_flatten_and_serialize_metadata_values(
experiment_metadata
),
experiment_metadata=_serialize(experiment_metadata),
experiment_dataset_id=dataset_id,
experiment_item_id=experiment_item_id,
experiment_item_metadata=_flatten_and_serialize_metadata_values(
experiment_item_metadata=_serialize(
item_metadata if isinstance(item_metadata, dict) else None
),
experiment_item_root_observation_id=task_span.id,
Expand All @@ -2955,11 +2969,16 @@ async def _process_experiment_item(
):
# _propagate_attributes updates the current task span and future children.
# Explicitly backfill the parent item-run span to preserve experiment association.
span._otel_span.set_attributes(
_set_span_attributes_within_limit(
span._otel_span,
_get_propagated_attributes_from_context(
otel_context_api.get_current()
)
),
)
# Write the observation metadata after the experiment
# attributes, so the span attribute limit trims the
# metadata instead of leaving no room for the output.
task_span.update(metadata=final_observation_metadata)
try:
output = await _run_task(task, item)
except Exception as e:
Expand Down
42 changes: 22 additions & 20 deletions langfuse/_client/propagation.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@

from langfuse._client.attributes import LangfuseOtelSpanAttributes
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
Expand Down Expand Up @@ -90,10 +91,10 @@
class PropagatedExperimentAttributes(TypedDict):
experiment_id: str
experiment_name: str
experiment_metadata: Optional[Dict[str, str]]
experiment_metadata: Optional[str] # serialized JSON
experiment_dataset_id: Optional[str]
experiment_item_id: str
experiment_item_metadata: Optional[Dict[str, str]]
experiment_item_metadata: Optional[str] # serialized JSON
experiment_item_root_observation_id: str


Expand Down Expand Up @@ -363,17 +364,6 @@ def _propagate_attributes(
"metadata": metadata,
}

if experiment:
for key, value in experiment.items():
if key in ("experiment_metadata", "experiment_item_metadata"):
propagated_metadata_attributes[key] = cast(
Optional[Dict[str, str]], value
)
else:
propagated_string_attributes[key] = cast(
Optional[Union[str, List[str]]], value
)

# Filter out None values
propagated_string_attributes = {
k: v for k, v in propagated_string_attributes.items() if v is not None
Expand Down Expand Up @@ -423,6 +413,19 @@ def _propagate_attributes(
as_baggage=as_baggage,
)

# Experiment attributes are set by the SDK and already serialized, so they
# skip validation. Mirrors langfuse-js.
if experiment:
for experiment_key, experiment_value in experiment.items():
if experiment_value is not None:
context = _set_propagated_attribute(
key=experiment_key,
value=cast(str, experiment_value),
context=context,
span=current_span,
as_baggage=as_baggage,
)

# Activate context, execute, and detach context
token = otel_context_api.attach(context=context)

Expand Down Expand Up @@ -599,14 +602,13 @@ def _set_propagated_attribute(
if span is not None and span.is_recording():
if isinstance(value, dict):
# Handle metadata
for k, v in value.items():
span.set_attribute(
key=f"{span_key}.{k}",
value=v,
)

span_attributes: Dict[str, Any] = {
f"{span_key}.{k}": v for k, v in value.items()
}
else:
span.set_attribute(key=span_key, value=value)
span_attributes = {span_key: value}

_set_span_attributes_within_limit(span, span_attributes)

# Set on baggage
if as_baggage:
Expand Down
Loading
Loading