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
4 changes: 3 additions & 1 deletion livekit-agents/livekit/agents/llm/llm.py
Original file line number Diff line number Diff line change
Expand Up @@ -433,7 +433,9 @@ async def _metrics_monitor_task(self, event_aiter: AsyncIterable[ChatChunk]) ->
)

# the GenAI response side; the request side was recorded at span creation
gen_ai_telemetry.set_usage_attributes(self._llm_request_span, metrics)
# a delegating span leaves the token counts to the provider call beneath it
if self._genai_operation_name is not None:
gen_ai_telemetry.set_usage_attributes(self._llm_request_span, metrics)
finish_reason = gen_ai_telemetry.finish_reason_for(
function_calls=tool_calls, interrupted=metrics.cancelled
)
Expand Down
81 changes: 73 additions & 8 deletions livekit-agents/livekit/agents/telemetry/gen_ai.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
from opentelemetry import trace
from opentelemetry.util.types import AttributeValue

from ..log import logger
from . import trace_types

if TYPE_CHECKING:
Expand Down Expand Up @@ -86,6 +87,44 @@ def mark_inference_span_recorded() -> None:
on_created()


def _env_max_input_messages() -> int:
raw = os.environ.get("OTEL_INSTRUMENTATION_GENAI_MAX_INPUT_MESSAGES", "").strip()
try:
return max(0, int(raw or 0))
except ValueError:
# a malformed telemetry setting must not stop the process from starting
logger.warning(
"ignoring non-integer OTEL_INSTRUMENTATION_GENAI_MAX_INPUT_MESSAGES, keeping all"
)
return 0


# opt-in truncation, both off by default: an inference span repeats the prompt and history
_capture_system_instructions: bool = (
os.environ.get("OTEL_INSTRUMENTATION_GENAI_CAPTURE_SYSTEM_INSTRUCTIONS", "").strip().lower()
not in _FALSY
)
_max_input_messages: int = _env_max_input_messages()


def set_capture_system_instructions(enabled: bool) -> None:
"""Set whether ``gen_ai.system_instructions`` is recorded; False omits it.

The instructions are the largest payload in a session and identical on every
inference span."""
global _capture_system_instructions
_capture_system_instructions = enabled


def set_max_input_messages(count: int) -> None:
"""Keep only the ``count`` most recent ``gen_ai.input.messages``, or all when 0.

Negative is treated as 0, and a span that drops any reports how many in
``lk.gen_ai.input.messages_dropped``."""
global _max_input_messages
_max_input_messages = max(0, count)


def _text_part(content: str) -> dict[str, Any]:
return {"type": "text", "content": content}

Expand Down Expand Up @@ -130,13 +169,15 @@ def _message_parts(item: ChatItem) -> list[dict[str, Any]]:
}
)
elif item.type == "function_call_output":
parts.append(
{
"type": "tool_call_response",
"id": item.call_id,
"response": _maybe_json(item.output),
}
)
response_part: dict[str, Any] = {
"type": "tool_call_response",
"id": item.call_id,
"response": _maybe_json(item.output),
}
# the part's schema is open, and a backend labelling from it alone shows "unknown"
if item.name:
response_part["name"] = item.name
parts.append(response_part)
return parts


Expand Down Expand Up @@ -224,6 +265,16 @@ def to_output_messages(
return [message]


def to_speech_messages(text: str, *, role: str) -> list[dict[str, Any]]:
"""An STT transcript or the words handed to a TTS, in the convention's message shape.

``role`` is ``"user"`` for a transcript and ``"assistant"`` for synthesized speech.
Empty text yields no message."""
if not text:
return []
return [{"role": role, "parts": [_text_part(text)]}]


def to_tool_definitions(tools: Iterable[Tool]) -> list[dict[str, Any]]:
"""``parameters`` is deliberately omitted: the convention marks it NOT RECOMMENDED by
default because a schema is large, and building one per request would be pure
Expand Down Expand Up @@ -291,17 +342,31 @@ def set_content_attributes(
input_messages: list[dict[str, Any]] | None = None,
output_messages: list[dict[str, Any]] | None = None,
tool_definitions: list[dict[str, Any]] | None = None,
truncate: bool = True,
) -> None:
"""Values are JSON strings: OpenTelemetry attributes cannot hold structured values
yet, which the convention explicitly allows for spans."""
yet, which the convention explicitly allows for spans.

``truncate`` applies the session-wide content limits; the span carrying a whole
session's conversation passes False, since trimming it defeats the purpose."""
if not _capture_content or not span.is_recording():
return

dropped = 0
if truncate:
if not _capture_system_instructions:
system_instructions = None
if input_messages and _max_input_messages:
dropped = max(0, len(input_messages) - _max_input_messages)
input_messages = input_messages[len(input_messages) - _max_input_messages :]

attrs: dict[str, AttributeValue] = {}
if system_instructions:
attrs[trace_types.ATTR_GEN_AI_SYSTEM_INSTRUCTIONS] = _json(system_instructions)
if input_messages:
attrs[trace_types.ATTR_GEN_AI_INPUT_MESSAGES] = _json(input_messages)
if dropped:
attrs[trace_types.ATTR_INPUT_MESSAGES_DROPPED] = dropped
if output_messages:
attrs[trace_types.ATTR_GEN_AI_OUTPUT_MESSAGES] = _json(output_messages)
if tool_definitions:
Expand Down
7 changes: 7 additions & 0 deletions livekit-agents/livekit/agents/telemetry/trace_types.py
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,9 @@

ATTR_ERROR_TYPE = "error.type"

# how many leading messages `gen_ai.input.messages` left out of a truncated history
ATTR_INPUT_MESSAGES_DROPPED = "lk.gen_ai.input.messages_dropped"


class GenAIOperationName:
"""Well-known ``gen_ai.operation.name`` values."""
Expand All @@ -259,6 +262,10 @@ class GenAIOperationName:
CREATE_MEMORY_STORE = "create_memory_store"
DELETE_MEMORY_STORE = "delete_memory_store"

# not in the registry, but the enum is open and a span without one is not GenAI
TRANSCRIBE = "transcribe"
SYNTHESIZE = "synthesize"


class GenAIOutputType:
"""Well-known ``gen_ai.output.type`` values."""
Expand Down
24 changes: 21 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 @@ -181,8 +181,21 @@ def set_provider(self, tracer_provider: trace_api.TracerProvider) -> None:
tracer_provider=self._tracer_provider,
)

def _with_conversation_id(self, kwargs: dict[str, Any]) -> dict[str, Any]:
# a backend grouping a session by the id needs every span to carry it, not just GenAI
conv = gen_ai._conversation_id()
if conv is None:
return kwargs
attributes = kwargs.get("attributes") or {}
if trace_types.ATTR_GEN_AI_CONVERSATION_ID in attributes:
return kwargs
return {
**kwargs,
"attributes": {**attributes, trace_types.ATTR_GEN_AI_CONVERSATION_ID: conv},
}

def start_span(self, *args: Any, **kwargs: Any) -> Span:
return self._tracer.start_span(*args, **kwargs)
return self._tracer.start_span(*args, **self._with_conversation_id(kwargs))

@_agnosticcontextmanager
def use_span(self, *args: Any, **kwargs: Any) -> Iterator[Span]:
Expand Down Expand Up @@ -214,7 +227,9 @@ 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._tracer.start_span(
name, **self._with_conversation_id({"context": context, "attributes": attributes})
)
try:
yield span
except Exception as e:
Expand All @@ -229,6 +244,9 @@ def start_as_current_span(self, *args: Any, **kwargs: Any) -> Iterator[Span]:
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)
bound.arguments.update(
self._with_conversation_id({"attributes": bound.arguments.get("attributes")})
)
with self._tracer.start_as_current_span(*bound.args[1:], **bound.kwargs) as span:
try:
yield span
Expand Down
35 changes: 35 additions & 0 deletions livekit-agents/livekit/agents/voice/agent_activity.py

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.

🟡 Realtime fallback spans lose served model

With RealtimeModelFallbackAdapter, realtime spans report the adapter's model and provider instead of the active provider. Its metrics_metadata exposes the serving instance, while these properties stay RealtimeModelFallbackAdapter and livekit. Backends therefore attribute every fallback response to the adapter.

(Refers to this code)

Learn more

A realtime fallback adapter delegates each live session to one underlying realtime model. The adapter keeps that serving model in _active_instance, and metrics_metadata exposes its model and provider. By contrast, the adapter's model and provider are constant adapter labels. Reading those constants records neither the primary model nor a model selected after failover.

Example: An adapter backed by OpenAI and Google starts on OpenAI. Its realtime inference span records model RealtimeModelFallbackAdapter and provider livekit, rather than OpenAI's model and provider. After a swap to Google, the span records the same adapter labels again.

Recommended fix: Read the model and provider from self.llm.metrics_metadata when creating the realtime inference span, or provide an equivalent API that exposes the adapter's active serving instance.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Original file line number Diff line number Diff line change
Expand Up @@ -3312,6 +3312,10 @@ def _on_first_frame(fut: asyncio.Future[float] | asyncio.Future[None]) -> None:
else:
forwarded_text = ""
current_span.set_attribute(trace_types.ATTR_RESPONSE_TEXT, forwarded_text)
gen_ai_telemetry.set_content_attributes(
current_span,
output_messages=gen_ai_telemetry.to_output_messages(text=forwarded_text),
)

assistant_metrics: llm.MetricsReport = {}

Expand Down Expand Up @@ -3502,6 +3506,13 @@ async def _pipeline_reply_task_impl(
[task.ctx.function_call for task in _RunningTasks.get(self._session, {}).values()],
)

# the turn's own input, not the whole context the `llm_request` span beneath carries
if new_message is not None:
gen_ai_telemetry.set_content_attributes(
current_span,
input_messages=gen_ai_telemetry.to_input_messages(llm.ChatContext([new_message])),
)

tasks: list[asyncio.Task[Any]] = []
llm_task, llm_gen_data = perform_llm_inference(
node=self._agent.llm_node,
Expand Down Expand Up @@ -3863,6 +3874,19 @@ async def _next_segment() -> _SpeechSegment | None:
speech_handle._item_added([msg])
current_span.set_attribute(trace_types.ATTR_RESPONSE_TEXT, forwarded_text)

# outside the block above: a turn that only called tools has no forwarded text
gen_ai_telemetry.set_content_attributes(
current_span,
output_messages=gen_ai_telemetry.to_output_messages(
text=forwarded_text,
function_calls=llm_gen_data.generated_functions,
finish_reason=gen_ai_telemetry.finish_reason_for(
function_calls=llm_gen_data.generated_functions,
interrupted=speech_handle.interrupted,
),
),
)

if not speech_handle.interrupted and len(tool_output.output) > 0:
self._session._update_agent_state("thinking")
self._on_end_of_agent_speech(ended_at=time.time())
Expand Down Expand Up @@ -4558,6 +4582,17 @@ def _create_assistant_message(
if trace_text_parts:
current_span.set_attribute(trace_types.ATTR_RESPONSE_TEXT, "\n".join(trace_text_parts))

gen_ai_telemetry.set_content_attributes(
current_span,
output_messages=gen_ai_telemetry.to_output_messages(
text="\n".join(trace_text_parts),
function_calls=function_calls,
finish_reason=gen_ai_telemetry.finish_reason_for(
function_calls=function_calls, interrupted=speech_handle.interrupted
),
),
)

# sync local chat ctx to the realtime server to remove any items the
# model added but the user never heard (interrupted before we pulled
# them, or message_outputs entries left in "skipped")
Expand Down
15 changes: 15 additions & 0 deletions livekit-agents/livekit/agents/voice/agent_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -1341,12 +1341,27 @@ async def _aclose_locked(
close_span.end()
otel_context.detach(close_token)
if self._session_span:
self._record_session_content(self._session_span)
self._session_span.end()
self._session_span = None
self._root_span_context = None

logger.debug("session closed", extra={"reason": reason.value, "error": error})

def _record_session_content(self, span: trace.Span) -> None:
"""The whole conversation on the ``agent_session`` span, recorded once at close.

A turn span carries only its own messages, so this is the one span a chat renderer
can read end to end; instructions stay on the inference spans that were given them."""
messages = gen_ai_telemetry.to_input_messages(self._chat_ctx)
output_messages: list[dict[str, Any]] = []
if messages and messages[-1]["role"] == "assistant":
output_messages = [messages.pop()]

gen_ai_telemetry.set_content_attributes(
span, input_messages=messages, output_messages=output_messages, truncate=False
)

async def _teardown_activity(self, *, reason: CloseReason, drain: bool) -> None:
"""Stop the activity and the models; the first step of closing."""
self._closing = True
Expand Down
23 changes: 14 additions & 9 deletions livekit-agents/livekit/agents/voice/audio_recognition.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
from ..language import LanguageCode
from ..log import logger
from ..stt import SpeechEvent
from ..telemetry import trace_types, tracer
from ..telemetry import gen_ai as gen_ai_telemetry, trace_types, tracer
from ..types import NOT_GIVEN, NotGivenOr
from ..utils import aio, is_given
from ..vad import VADStream
Expand Down Expand Up @@ -1769,6 +1769,12 @@ async def _bounce_eou_task(
trace_types.ATTR_END_OF_TURN_DELAY: metrics.end_of_turn_delay or 0,
}
)
gen_ai_telemetry.set_content_attributes(
user_turn_span,
output_messages=gen_ai_telemetry.to_speech_messages(
self._audio_transcript, role="user"
),
)
if self._stt_request_ids:
user_turn_span.set_attribute(
trace_types.ATTR_PROVIDER_REQUEST_IDS, self._stt_request_ids
Expand Down Expand Up @@ -1991,14 +1997,13 @@ def _ensure_user_turn_span(self, start_time: float | None = None) -> trace.Span:
_set_participant_attributes(self._user_turn_span, room_io.linked_participant)

# add STT model/provider attributes
if self._stt_model:
self._user_turn_span.set_attribute(
trace_types.ATTR_GEN_AI_REQUEST_MODEL, self._stt_model
)
if self._stt_provider:
self._user_turn_span.set_attribute(
trace_types.ATTR_GEN_AI_PROVIDER_NAME, self._stt_provider
)
gen_ai_telemetry.set_request_attributes(
self._user_turn_span,
operation=trace_types.GenAIOperationName.TRANSCRIBE,
provider=self._stt_provider,
model=self._stt_model,
output_type=trace_types.GenAIOutputType.TEXT,
)

return self._user_turn_span

Expand Down
Loading
Loading