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
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 @@ -292,10 +292,12 @@ def _record_genai_request(self, span: trace.Span) -> None:
output_type=trace_types.GenAIOutputType.TEXT,
)
if self._record_content:
gen_ai_telemetry.record_llm_input_messages(
span, self._chat_ctx, is_delegating=self._genai_operation_name is None
)
gen_ai_telemetry.set_content_attributes(
span,
system_instructions=gen_ai_telemetry.to_system_instructions(self._chat_ctx),
input_messages=gen_ai_telemetry.to_input_messages(self._chat_ctx),
tool_definitions=gen_ai_telemetry.to_tool_definitions(self._tools),
)

Expand Down
219 changes: 209 additions & 10 deletions livekit-agents/livekit/agents/telemetry/gen_ai.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,15 @@
from __future__ import annotations

import contextvars
import hashlib
import json
import os
from collections.abc import Callable, Iterable, Sequence
import weakref
from collections.abc import Callable, Iterable, Mapping, Sequence
from typing import TYPE_CHECKING, Any, TypeAlias

from opentelemetry import trace
from opentelemetry import context as otel_context, trace
from opentelemetry.sdk.trace import Event, ReadableSpan, SpanProcessor
from opentelemetry.util.types import AttributeValue

from . import trace_types
Expand All @@ -33,19 +36,208 @@
os.environ.get("OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT", "").strip().lower()
not in _FALSY
)
_capture_system_instructions = True
_max_input_messages = 0
_capture_input_delta = False
_input_capture_version = 0

_INPUT_MESSAGES_STATE = otel_context.create_key("lk_input_messages_state")
_standalone_input_states: weakref.WeakKeyDictionary[ChatContext, _InputMessagesState] = (
weakref.WeakKeyDictionary()
)


def set_capture_content(enabled: bool) -> None:
"""When off, spans keep every non-content GenAI attribute and omit the message
payloads, tool definitions and tool call arguments/results."""
global _capture_content
global _capture_content, _input_capture_version
if _capture_content != enabled:
_input_capture_version += 1
_capture_content = enabled


def capture_content_enabled() -> bool:
return _capture_content


def set_capture_system_instructions(enabled: bool) -> None:
"""Set process-wide prompt capture (default: enabled).

When disabled, exports omit system instructions and legacy full-chat payloads.
Configure before starting sessions. Model requests are unchanged.
"""
global _capture_system_instructions
_capture_system_instructions = enabled


def set_max_input_messages(count: int) -> None:
"""Keep the newest ``count`` exported GenAI input messages; zero means unlimited.

Full capture is the default. A positive limit omits legacy full-chat payloads
and reports omitted messages in ``lk.gen_ai.input.messages_dropped``. It can
separate tool results from their calls. Model requests are unchanged.
Configure this process-wide setting before starting sessions.
"""
if not isinstance(count, int) or count < 0:
raise ValueError("count must be a non-negative integer")
global _max_input_messages
_max_input_messages = count


def set_capture_input_delta(enabled: bool) -> None:
"""Export only new or changed LLM input messages (default: disabled).

Compare with the previous captured request in the same AgentSession, or on
the same ChatContext for standalone LLM calls. The first request includes
all input messages. An unchanged request exports an empty list. Message IDs
distinguish repeated text from repeated history; edits are exported again.

This omits legacy full-chat payloads. System instructions use their separate
capture setting, and any input-message limit applies after delta selection.
Model requests, speech input and session transcripts are unchanged.
Configure this process-wide setting before starting sessions.
"""
global _capture_input_delta, _input_capture_version
if _capture_input_delta != enabled:
_input_capture_version += 1
_capture_input_delta = enabled
_standalone_input_states.clear()


def legacy_chat_capture_enabled() -> bool:
return (
_capture_content
and _capture_system_instructions
and _max_input_messages == 0
and not _capture_input_delta
)


class _InputMessagesState:
def __init__(self) -> None:
self._previous: dict[tuple[str, ...], bytes] = {}
self._version = _input_capture_version

def delta(self, chat_ctx: ChatContext) -> list[dict[str, Any]]:
if self._version != _input_capture_version:
self._previous.clear()
self._version = _input_capture_version
current: dict[tuple[str, ...], bytes] = {}
changed: list[dict[str, Any]] = []
for ids, message in _input_messages_with_ids(chat_ctx.items):
key = tuple(ids)
fingerprint = hashlib.sha256(_json(message).encode()).digest()
current[key] = fingerprint
if self._previous.get(key) != fingerprint:
changed.append(message)
# Retain fingerprints of the current history, not copies of message content.
self._previous = current
return changed


def _with_input_messages_state(ctx: otel_context.Context) -> otel_context.Context:
return otel_context.set_value(_INPUT_MESSAGES_STATE, _InputMessagesState(), ctx)


def record_llm_input_messages(
span: trace.Span, chat_ctx: ChatContext, *, is_delegating: bool = False
) -> None:
if not _capture_content or not span.is_recording():
return

attrs: dict[str, Any] = {}
if _capture_input_delta:
# A fallback wrapper must not consume the delta before the provider span.
if is_delegating:
return
state = otel_context.get_value(_INPUT_MESSAGES_STATE)
if not isinstance(state, _InputMessagesState):
state = _standalone_input_states.get(chat_ctx)
if state is None:
state = _standalone_input_states[chat_ctx] = _InputMessagesState()
messages = state.delta(chat_ctx)
attrs[trace_types.ATTR_GEN_AI_INPUT_MESSAGES_MODE] = "delta"
else:
messages = to_input_messages(chat_ctx)
if not messages:
return
_set_input_messages(attrs, messages)
span.set_attributes(attrs)


class _ContentFilteringSpanProcessor(SpanProcessor):
"""Apply capture controls before PII filtering can stash content for Cloud export."""

def on_end(self, span: ReadableSpan) -> None:
if legacy_chat_capture_enabled():
return
from . import pii

span._attributes = self._filter_attributes(span.attributes)
# Deprecated content events cannot represent a bounded conversation reliably.
span._events = tuple(
Event(
name=event.name,
attributes=self._filter_attributes(event.attributes),
timestamp=event.timestamp,
)
for event in span.events
if event.name not in pii.PII_EVENT_NAMES
)

@staticmethod
def _filter_attributes(attributes: Mapping[str, Any] | None) -> dict[str, Any]:
from . import pii

attributes = dict(attributes or {})
attributes.pop(trace_types.ATTR_CHAT_CTX, None)
if not _capture_system_instructions:
attributes.pop(trace_types.ATTR_GEN_AI_SYSTEM_INSTRUCTIONS, None)
attributes.pop(trace_types.ATTR_INSTRUCTIONS, None)
if not _capture_content:
legacy_content = {
trace_types.ATTR_INSTRUCTIONS,
trace_types.ATTR_USER_INPUT,
trace_types.ATTR_USER_TRANSCRIPT,
trace_types.ATTR_RESPONSE_TEXT,
trace_types.ATTR_RESPONSE_FUNCTION_CALLS,
trace_types.ATTR_FUNCTION_TOOL_ARGS,
trace_types.ATTR_FUNCTION_TOOL_OUTPUT,
trace_types.ATTR_TTS_INPUT_TEXT,
}
attributes = {
key: value
for key, value in attributes.items()
if key not in legacy_content
and key not in pii.GEN_AI_PII_ATTRIBUTES
and not key.startswith(trace_types.ATTR_GEN_AI_PROMPT_VARIABLE)
}
elif _max_input_messages and (
raw := attributes.get(trace_types.ATTR_GEN_AI_INPUT_MESSAGES)
):
try:
messages = json.loads(raw) if isinstance(raw, str) else None
except (ValueError, TypeError):
messages = None
if isinstance(messages, list):
_set_input_messages(attributes, messages)
else:
# An unparseable payload must not bypass an explicit export limit.
attributes.pop(trace_types.ATTR_GEN_AI_INPUT_MESSAGES, None)
return attributes


def _set_input_messages(attrs: dict[str, Any], messages: list[dict[str, Any]]) -> None:
if _max_input_messages and len(messages) > _max_input_messages:
dropped = len(messages) - _max_input_messages
previous = attrs.get(trace_types.ATTR_GEN_AI_INPUT_MESSAGES_DROPPED, 0)
attrs[trace_types.ATTR_GEN_AI_INPUT_MESSAGES_DROPPED] = dropped + (
previous if isinstance(previous, int) else 0
)
messages = messages[-_max_input_messages:]
attrs[trace_types.ATTR_GEN_AI_INPUT_MESSAGES] = _json(messages)


# A custom `llm_node` may do the inference itself — returning a plain str, streaming its
# own chunks, or calling a third-party engine — and never construct an LLMStream. Those
# paths have no nested `llm_request` span to carry the convention's attributes, so the node
Expand Down Expand Up @@ -173,8 +365,14 @@ def to_input_messages(chat_ctx: ChatContext) -> list[dict[str, Any]]:
"""History in the order it was sent. ``system``/``developer`` messages go to
``gen_ai.system_instructions`` instead, and non-conversational items (agent
handoffs, config updates) are skipped."""
messages: list[dict[str, Any]] = []
for item in chat_ctx.items:
return [message for _, message in _input_messages_with_ids(chat_ctx.items)]


def _input_messages_with_ids(
items: Iterable[ChatItem],
) -> list[tuple[list[str], dict[str, Any]]]:
messages: list[tuple[list[str], dict[str, Any]]] = []
for item in items:
role: str
if item.type == "message":
if item.role in ("system", "developer"):
Expand All @@ -194,13 +392,14 @@ def to_input_messages(chat_ctx: ChatContext) -> list[dict[str, Any]]:
# consecutive tool calls from one assistant turn belong to a single message
if (
messages
and messages[-1]["role"] == role == "assistant"
and messages[-1][1]["role"] == role == "assistant"
and item.type == "function_call"
):
messages[-1]["parts"].extend(parts)
messages[-1][0].append(item.id)
messages[-1][1]["parts"].extend(parts)
continue

messages.append({"role": role, "parts": parts})
messages.append(([item.id], {"role": role, "parts": parts}))
return messages


Expand Down Expand Up @@ -305,10 +504,10 @@ def set_content_attributes(
return

attrs: dict[str, AttributeValue] = {}
if system_instructions:
if system_instructions and _capture_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)
_set_input_messages(attrs, input_messages)
if output_messages:
attrs[trace_types.ATTR_GEN_AI_OUTPUT_MESSAGES] = _json(output_messages)
if tool_definitions:
Expand Down
2 changes: 2 additions & 0 deletions livekit-agents/livekit/agents/telemetry/trace_types.py
Original file line number Diff line number Diff line change
Expand Up @@ -227,6 +227,8 @@

ATTR_GEN_AI_SYSTEM_INSTRUCTIONS = "gen_ai.system_instructions"
ATTR_GEN_AI_INPUT_MESSAGES = "gen_ai.input.messages"
ATTR_GEN_AI_INPUT_MESSAGES_DROPPED = "lk.gen_ai.input.messages_dropped"
ATTR_GEN_AI_INPUT_MESSAGES_MODE = "lk.gen_ai.input.messages_mode"
ATTR_GEN_AI_OUTPUT_MESSAGES = "gen_ai.output.messages"
ATTR_GEN_AI_OUTPUT_TYPE = "gen_ai.output.type"

Expand Down
1 change: 1 addition & 0 deletions livekit-agents/livekit/agents/telemetry/traces.py
Original file line number Diff line number Diff line change
Expand Up @@ -666,6 +666,7 @@ def _install_pii_redaction(
# useful to a backend that can render the conversation
pii._PIIFilteringSpanProcessor(allow_pii=allow_pii if allow_pii is not None else True),
)
_prepend_span_processor(tracer_provider, gen_ai._ContentFilteringSpanProcessor())


def set_tracer_provider(
Expand Down
4 changes: 3 additions & 1 deletion livekit-agents/livekit/agents/voice/agent_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -977,7 +977,9 @@ async def start(
if self._session_ctx_token is not None:
otel_context.detach(self._session_ctx_token)
self._session_ctx_token = None
ctx = trace.set_span_in_context(current_span)
ctx = gen_ai_telemetry._with_input_messages_state(
trace.set_span_in_context(current_span)
)
self._session_ctx_token = otel_context.attach(ctx)

self._recorded_events = []
Expand Down
19 changes: 10 additions & 9 deletions livekit-agents/livekit/agents/voice/generation.py
Original file line number Diff line number Diff line change
Expand Up @@ -194,20 +194,21 @@ async def _llm_inference_task(

if current_span.is_recording():
attrs: dict[str, Any] = {
trace_types.ATTR_CHAT_CTX: json.dumps(
chat_ctx.to_dict(
exclude_audio=True,
exclude_image=True,
exclude_timestamp=True,
exclude_metrics=True,
)
),
trace_types.ATTR_FUNCTION_TOOLS: list(tool_ctx.function_tools.keys()),
trace_types.ATTR_PROVIDER_TOOLS: [
type(tool).__name__ for tool in tool_ctx.provider_tools
],
trace_types.ATTR_TOOL_SETS: [type(tool_set).__name__ for tool_set in tool_ctx.toolsets],
}
if gen_ai_telemetry.legacy_chat_capture_enabled():
attrs[trace_types.ATTR_CHAT_CTX] = json.dumps(
chat_ctx.to_dict(
exclude_audio=True,
exclude_image=True,
exclude_timestamp=True,
exclude_metrics=True,
)
)
current_span.set_attributes(attrs)

# the GenAI inference attributes belong to the nested `llm_request` span, which is the
Expand Down Expand Up @@ -384,10 +385,10 @@ def _record_uninstrumented_inference(
span, finish_reasons=[finish_reason], time_to_first_chunk=data.ttft
)
if span.is_recording() and gen_ai_telemetry.capture_content_enabled():
gen_ai_telemetry.record_llm_input_messages(span, chat_ctx)
gen_ai_telemetry.set_content_attributes(
span,
system_instructions=gen_ai_telemetry.to_system_instructions(chat_ctx),
input_messages=gen_ai_telemetry.to_input_messages(chat_ctx),
tool_definitions=gen_ai_telemetry.to_tool_definitions(tools),
output_messages=gen_ai_telemetry.to_output_messages(
text=data.generated_text,
Expand Down
6 changes: 4 additions & 2 deletions tests/test_llm_telemetry.py
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,9 @@ async def test_llm_stream_capture_requires_enablement_at_start_and_completion(
assert response.text == "hello"
spans = [span for span in span_exporter.get_finished_spans() if span.name == "llm_request"]
assert len(spans) == 1
assert (trace_types.ATTR_GEN_AI_INPUT_MESSAGES in spans[0].attributes) is capture_at_start
assert (trace_types.ATTR_GEN_AI_INPUT_MESSAGES in spans[0].attributes) is (
capture_at_start and capture_during_run
)
assert (trace_types.ATTR_GEN_AI_OUTPUT_MESSAGES in spans[0].attributes) is (
capture_at_start and capture_during_run
)
Expand Down Expand Up @@ -510,7 +512,7 @@ async def test_llm_node_preserves_noncontent_attributes_when_capture_is_disabled

spans = [span for span in span_exporter.get_finished_spans() if span.name == "llm_node"]
assert len(spans) == 1
assert trace_types.ATTR_CHAT_CTX in spans[0].attributes
assert trace_types.ATTR_CHAT_CTX not in spans[0].attributes
assert spans[0].attributes[trace_types.ATTR_GEN_AI_OPERATION_NAME] == "chat"
assert trace_types.ATTR_GEN_AI_INPUT_MESSAGES not in spans[0].attributes
assert trace_types.ATTR_GEN_AI_OUTPUT_MESSAGES not in spans[0].attributes
Expand Down
Loading
Loading