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
18 changes: 15 additions & 3 deletions livekit-agents/livekit/agents/telemetry/traces.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@
recording_enabled,
)
from ..utils import is_given
from . import pii, trace_types, utils as telemetry_utils
from . import gen_ai, pii, trace_types, utils as telemetry_utils

if TYPE_CHECKING:
from ..llm import ChatItem
Expand Down Expand Up @@ -165,9 +165,18 @@ def _serialize_session_options(options: AgentSessionOptions) -> dict[str, Any]:


_USE_SPAN_SIGNATURE = inspect.signature(trace_api.use_span)
_START_SPAN_SIGNATURE = inspect.signature(Tracer.start_span)
_START_AS_CURRENT_SPAN_SIGNATURE = inspect.signature(Tracer.start_as_current_span)


def _with_conversation_id(
attributes: Mapping[str, AttributeValue] | None,
) -> Mapping[str, AttributeValue] | None:
if (conversation_id := gen_ai._conversation_id()) is None:
return attributes
return {trace_types.ATTR_GEN_AI_CONVERSATION_ID: conversation_id, **(attributes or {})}


class _DynamicTracer(Tracer):
def __init__(self, instrumenting_module_name: str) -> None:
self._instrumenting_module_name = instrumenting_module_name
Expand All @@ -182,7 +191,9 @@ def set_provider(self, tracer_provider: trace_api.TracerProvider) -> None:
)

def start_span(self, *args: Any, **kwargs: Any) -> Span:
return self._tracer.start_span(*args, **kwargs)
bound = _START_SPAN_SIGNATURE.bind(self._tracer, *args, **kwargs)
bound.arguments["attributes"] = _with_conversation_id(bound.arguments.get("attributes"))
return self._tracer.start_span(*bound.args[1:], **bound.kwargs)

@_agnosticcontextmanager
def use_span(self, *args: Any, **kwargs: Any) -> Iterator[Span]:
Expand Down Expand Up @@ -214,7 +225,7 @@ def detached_span(
it and becomes the accidental parent of unrelated spans those tasks emit for the rest
of the session. The parent is ``context`` when given, else the ambient context; the
exception, if any, is recorded redaction-aware and the span is ended."""
span = self._tracer.start_span(name, context=context, attributes=attributes)
span = self.start_span(name, context=context, attributes=attributes)
try:
yield span
except Exception as e:
Expand All @@ -226,6 +237,7 @@ def detached_span(
@_agnosticcontextmanager
def start_as_current_span(self, *args: Any, **kwargs: Any) -> Iterator[Span]:
bound = _START_AS_CURRENT_SPAN_SIGNATURE.bind(self._tracer, *args, **kwargs)
bound.arguments["attributes"] = _with_conversation_id(bound.arguments.get("attributes"))
record_exception = bound.arguments.get("record_exception", True)
set_status_on_exception = bound.arguments.get("set_status_on_exception", True)
bound.arguments.update(record_exception=False, set_status_on_exception=False)
Expand Down
75 changes: 75 additions & 0 deletions tests/test_telemetry_conversation_id.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
from __future__ import annotations

from collections.abc import Iterator
from types import MappingProxyType

import pytest
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from opentelemetry.trace import SpanKind

from livekit.agents.telemetry import gen_ai, set_tracer_provider, tracer

pytestmark = [pytest.mark.unit, pytest.mark.no_concurrent]


@pytest.fixture
def span_exporter() -> Iterator[InMemorySpanExporter]:
original_provider = tracer._tracer_provider
provider = TracerProvider()
exporter = InMemorySpanExporter()
provider.add_span_processor(SimpleSpanProcessor(exporter))
set_tracer_provider(provider)
try:
yield exporter
finally:
set_tracer_provider(original_provider)
provider.shutdown()


@pytest.mark.parametrize("conversation_id", [None, "RM_test"])
@pytest.mark.parametrize("explicit_id", [None, "explicit"])
@pytest.mark.parametrize("api", ["start_span", "start_as_current_span", "detached_span"])
def test_conversation_id_preserves_caller_attributes(
span_exporter: InMemorySpanExporter,
monkeypatch: pytest.MonkeyPatch,
conversation_id: str | None,
explicit_id: str | None,
api: str,
) -> None:
monkeypatch.setattr(gen_ai, "_conversation_id", lambda: conversation_id)
attributes = {"test.attribute": "value"}
if explicit_id:
attributes["gen_ai.conversation.id"] = explicit_id
original = attributes.copy()
# The standard tracer APIs also allow positional, read-only attribute mappings.
if api == "start_span":
tracer.start_span("test", None, SpanKind.INTERNAL, MappingProxyType(attributes)).end()
elif api == "start_as_current_span":
with tracer.start_as_current_span(
"test", None, SpanKind.INTERNAL, MappingProxyType(attributes)
):
pass
else:
with tracer.detached_span("test", attributes=attributes):
pass
assert attributes == original
expected = dict(original)
if explicit_id or conversation_id:
expected["gen_ai.conversation.id"] = explicit_id or conversation_id
assert dict(span_exporter.get_finished_spans()[0].attributes) == expected


def test_spans_without_attributes_receive_conversation_id(
span_exporter: InMemorySpanExporter, monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setattr(gen_ai, "_conversation_id", lambda: "RM_test")
tracer.start_span("test").end()
with tracer.start_as_current_span("test"):
pass
with tracer.detached_span("test"):
pass
assert len(span_exporter.get_finished_spans()) == 3
for span in span_exporter.get_finished_spans():
assert span.attributes["gen_ai.conversation.id"] == "RM_test"
2 changes: 2 additions & 0 deletions tests/test_telemetry_metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ def test_exported_spans_keep_provider_metadata(
"langfuse.session.id": "customer-session",
"job_id": "provider-job" if state == "outside" else "job-a",
"room_id": "provider-room" if state == "outside" else "room-a",
**({"gen_ai.conversation.id": "room-a"} if state != "outside" else {}),
}


Expand All @@ -107,6 +108,7 @@ def test_job_fallback_metadata_does_not_cross_jobs(
"langfuse.session.id": "customer-session",
"job_id": "job-a",
"room_id": "room-a",
"gen_ai.conversation.id": "room-a",
}

with tracer.start_as_current_span("worker"):
Expand Down
Loading