parlot.instrumentation.livekit
parlot-instrumentation-livekit: OTel instrumentation for LiveKit Agents.
1"""parlot-instrumentation-livekit: OTel instrumentation for LiveKit Agents.""" 2 3from parlot.core import ( 4 add_platform_ref, 5 human_escalation, 6 record_human_rep, 7 set_session_attribute, 8 set_session_metadata, 9 stamp_platform_refs, 10) 11 12from ._auto import parlotize 13from ._events import install_session_hooks 14from ._processor import LiveKitGenAIProcessor 15 16__all__ = [ 17 "parlotize", 18 "install_session_hooks", 19 "LiveKitGenAIProcessor", 20 "add_platform_ref", 21 "human_escalation", 22 "record_human_rep", 23 "set_session_attribute", 24 "set_session_metadata", 25 "stamp_platform_refs", 26]
66def parlotize( 67 agent_id: str, 68 *, 69 endpoint: Optional[str] = None, 70 api_key: Optional[str] = None, 71 capture_genai_content: Optional[bool] = None, 72 service_name: Optional[str] = None, 73 tracer_provider: TracerProvider | None = None, 74 auto_escalate_sip: bool = False, 75 escalation_metadata_match: dict[str, str] | None = None, 76 version: Optional[str] = None, 77 record: bool | list[str] | None = None, 78 capture_logs: bool | list[str] | None = None, 79 log_level: Optional[str] = None, 80) -> ParlotContext: 81 """Parlotize LiveKit instrumentation and OTLP export. 82 83 Call before constructing ``AgentSession``. Builds a ``TracerProvider`` with 84 an OTLP exporter, registers it with ``livekit.agents.telemetry``, fetches 85 telemetry bootstrap when ``PARLOT_API_KEY`` is set, patches 86 ``JobContext.connect`` / ``AgentSession.__init__``, and installs session 87 event hooks. 88 89 Shared parameters (``agent_id``, plus keyword-only ``endpoint``, ``api_key``, 90 ``version``, ``capture_genai_content``, ``capture_logs``, ``log_level``, 91 ``service_name``, ``tracer_provider``) match every adapter — see 92 ``parlot.core.ParlotizeProtocol``. 93 94 Under LiveKit ``dev`` / job workers, ``__main__`` is often LiveKit's IPC 95 entrypoint, not your agent file — pass ``agent_id`` / ``version=`` 96 explicitly when you care about stable deployment identity. 97 98 Args: 99 agent_id: Shared — required canonical ``session.agent_id``. Must be 100 non-empty after stripping whitespace. 101 endpoint: Shared — Parlot OTLP base URL (or ``PARLOT_ENDPOINT``). 102 api_key: Shared — org API key (or ``PARLOT_API_KEY``). 103 capture_genai_content: Shared — GenAI payload capture override. 104 Precedence: job metadata → this kwarg → Settings → Generative AI → on. 105 service_name: Shared — OTel ``service.name`` (defaults to ``agent_id`` 106 or ``"unknown"``). 107 tracer_provider: Shared — existing ``TracerProvider``, or build one with 108 Parlot's OTLP exporter and ``LiveKitGenAIProcessor``. 109 auto_escalate_sip: When ``True``, mark the session escalated when a SIP 110 participant joins the room. 111 escalation_metadata_match: Participant metadata key/value pairs that 112 classify joining participants as human representatives. 113 version: Shared — ``gen_ai.agent.version``. Pass ``version=`` to set it; 114 otherwise stamped as ``\"unknown\"``. 115 record: Audio recording policy. Boolean or agent-id glob patterns 116 (e.g. ``["support-*", "billing"]``). Precedence: LiveKit job 117 metadata ``record`` → this kwarg → Settings → Recording. 118 capture_logs: Shared — session log capture (bool or globs). Precedence: 119 job metadata → this kwarg → Settings → Logs → on. 120 log_level: Shared — minimum level for session log capture. 121 122 Returns: 123 The ``ParlotContext`` created for this process (or the prior one if 124 already configured). 125 """ 126 global _configured, _auto_escalate_sip, _escalation_metadata_match 127 global _configured_agent_id, _configured_agent_version 128 global _configured_record 129 global _configured_capture_genai_content, _configured_capture_logs, _configured_log_level 130 global _parlot_context 131 if _configured: 132 logger.debug("parlot-instrumentation.livekit already configured — skipping") 133 assert _parlot_context is not None 134 return _parlot_context 135 136 if _is_livekit_dev_watch_parent(): 137 logger.debug( 138 "Skipping parlot-instrumentation.livekit parlotize in LiveKit dev " 139 "watcher parent (worker child will parlotize)" 140 ) 141 return ParlotContext() 142 143 _auto_escalate_sip = auto_escalate_sip 144 _escalation_metadata_match = ( 145 dict(escalation_metadata_match) if escalation_metadata_match else None 146 ) 147 if isinstance(record, list): 148 _configured_record = [str(item).strip() for item in record if str(item).strip()] 149 else: 150 _configured_record = record 151 from parlot.core.provider import resolve_capture_genai_content 152 153 _configured_capture_genai_content = resolve_capture_genai_content(capture_genai_content) 154 if isinstance(capture_logs, list): 155 _configured_capture_logs = [ 156 str(item).strip() for item in capture_logs if str(item).strip() 157 ] 158 else: 159 _configured_capture_logs = capture_logs 160 _configured_log_level = log_level.strip().upper() if log_level else None 161 162 from ._agent_version import resolve_agent_version 163 _configured_agent_version = resolve_agent_version(version) 164 165 from parlot.core import base_parlotize 166 167 res = base_parlotize( 168 agent_id, 169 endpoint=endpoint, 170 api_key=api_key, 171 capture_genai_content=capture_genai_content, 172 tracer_provider=tracer_provider, 173 version=_configured_agent_version, 174 capture_logs=capture_logs, 175 log_level=log_level, 176 ) 177 _parlot_context = res.context 178 _configured_agent_id = res.agent_id 179 180 from ._recording_guard import ( 181 current_job_capture_logs_metadata, 182 livekit_session_log_fields, 183 recording_agent_id_from_ctx, 184 ) 185 186 def _agent_id_resolver() -> str: 187 try: 188 from livekit.agents.job import get_job_context 189 190 ctx = get_job_context() 191 if ctx is not None: 192 return recording_agent_id_from_ctx(ctx) 193 except Exception: 194 pass 195 return _configured_agent_id or "" 196 197 res.context.session_logs.set_resolvers( 198 session_resolver=livekit_session_log_fields, 199 agent_id_resolver=_agent_id_resolver, 200 metadata_resolver=current_job_capture_logs_metadata, 201 ) 202 203 tracer_provider = res.tracer_provider 204 if tracer_provider is None: 205 tracer_provider = _build_provider( 206 endpoint=res.endpoint, 207 api_key=res.api_key, 208 service_name=service_name, 209 service_version=_configured_agent_version or None, 210 context=res.context, 211 ) 212 213 from parlot.core.processor import assert_sync_span_processors 214 215 from ._session import set_span_context_attach_enabled 216 217 attach_ok = assert_sync_span_processors(tracer_provider) 218 set_span_context_attach_enabled(attach_ok) 219 220 _register_with_livekit(tracer_provider) 221 _patch_agent_session_with_tracer(tracer_provider) 222 _patch_job_context_connect() 223 _install_telemetry_compare() 224 225 _configured = True 226 logger.debug("parlot-instrumentation.livekit configured (endpoint=%s)", res.endpoint) 227 return res.context
Parlotize LiveKit instrumentation and OTLP export.
Call before constructing AgentSession. Builds a TracerProvider with
an OTLP exporter, registers it with livekit.agents.telemetry, fetches
telemetry bootstrap when PARLOT_API_KEY is set, patches
JobContext.connect / AgentSession.__init__, and installs session
event hooks.
Shared parameters (agent_id, plus keyword-only endpoint, api_key,
version, capture_genai_content, capture_logs, log_level,
service_name, tracer_provider) match every adapter — see
parlot.core.ParlotizeProtocol.
Under LiveKit dev / job workers, __main__ is often LiveKit's IPC
entrypoint, not your agent file — pass agent_id / version=
explicitly when you care about stable deployment identity.
Arguments:
- agent_id: Shared — required canonical
session.agent_id. Must be non-empty after stripping whitespace. - endpoint: Shared — Parlot OTLP base URL (or
PARLOT_ENDPOINT). - api_key: Shared — org API key (or
PARLOT_API_KEY). - capture_genai_content: Shared — GenAI payload capture override. Precedence: job metadata → this kwarg → Settings → Generative AI → on.
- service_name: Shared — OTel
service.name(defaults toagent_idor"unknown"). - tracer_provider: Shared — existing
TracerProvider, or build one with Parlot's OTLP exporter andLiveKitGenAIProcessor. - auto_escalate_sip: When
True, mark the session escalated when a SIP participant joins the room. - escalation_metadata_match: Participant metadata key/value pairs that classify joining participants as human representatives.
- version: Shared —
gen_ai.agent.version. Passversion=to set it; otherwise stamped as"unknown". - record: Audio recording policy. Boolean or agent-id glob patterns
(e.g.
["support-*", "billing"]). Precedence: LiveKit job metadatarecord→ this kwarg → Settings → Recording. - capture_logs: Shared — session log capture (bool or globs). Precedence: job metadata → this kwarg → Settings → Logs → on.
- log_level: Shared — minimum level for session log capture.
Returns:
The
ParlotContextcreated for this process (or the prior one if already configured).
483def install_session_hooks( 484 session: Any, 485 processor: "LiveKitGenAIProcessor", 486 tracer: "Tracer", 487) -> None: 488 """Install all AgentSession event listeners for Parlot instrumentation. 489 490 ``parlotize()`` patches ``AgentSession.__init__`` to invoke this 491 automatically. Call directly only when manually instantiating unpatched 492 sessions. 493 494 Args: 495 session: LiveKit ``AgentSession`` instance to instrument. 496 processor: Span processor that owns session/turn enrichment state. 497 tracer: OpenTelemetry tracer used for Parlot contract spans. 498 """ 499 LiveKitEventBridge(processor, tracer).install(session)
Install all AgentSession event listeners for Parlot instrumentation.
parlotize() patches AgentSession.__init__ to invoke this
automatically. Call directly only when manually instantiating unpatched
sessions.
Arguments:
- session: LiveKit
AgentSessioninstance to instrument. - processor: Span processor that owns session/turn enrichment state.
- tracer: OpenTelemetry tracer used for Parlot contract spans.
113class LiveKitGenAIProcessor(ParlotBaseProcessor): 114 """Enriches LiveKit Agent spans in-place with Parlot conventions. 115 116 Intercepts spans from the LiveKit Agents SDK and normalizes them into 117 Parlot's three-layer semantic vocabulary: 118 119 1. **Conversation Contract** — session start/close, turns, agent handoffs. 120 2. **OTel GenAI (v1.41.0)** — LLM inference, tool executions, workflows. 121 3. **Voice spans** — TTS, STT, and end-of-utterance operational timings. 122 """ 123 124 def __init__( 125 self, 126 capture_genai_content: Optional[bool] = None, 127 handoff_tool_names: Optional[set[str]] = None, 128 *, 129 context: Optional["ParlotContext"] = None, 130 ) -> None: 131 from parlot.core.context import ParlotContext 132 133 # Explicit override for tests / rare call sites; None → policy at emit time. 134 self._capture_genai_content_override = capture_genai_content 135 self._handoff_tools = handoff_tool_names or set() 136 self._context = context if context is not None else ParlotContext() 137 self._sessions: dict[str, _LiveKitSessionState] = {} 138 self._turn_trace_registry: dict[str, dict[int, tuple[str, str]]] = {} 139 self._tracer: Tracer | None = None 140 self._metrics = None 141 self._turn_source: str = "spans" 142 143 self._handoff = HandoffTracker( 144 set_attr=self._set, 145 handoff_tools=self._handoff_tools, 146 ) 147 self._tokens = TokenAggregator( 148 set_attr=self._set, 149 plugin_host=self, 150 ) 151 self._recording = RecordingCoordinator( 152 get_capture_override=lambda: self._capture_genai_content_override, 153 active_agent_id_fn=self._active_agent_id, 154 get_context=lambda: self._context, 155 ) 156 self._turns = TurnEnricher( 157 set_attr=self._set, 158 add_event=self._add_event, 159 maybe_update=self._maybe_update, 160 get_turn_source=lambda: self._turn_source, 161 get_tracer=lambda: self._tracer, 162 get_metrics=lambda: self._metrics, 163 genai_content_enabled=self._recording.genai_content_enabled, 164 active_agent_id=self._active_agent_id, 165 stamp_agent_identity=self._stamp_agent_identity, 166 stamp_response_model_if_distinct=self._stamp_response_model_if_distinct, 167 handoff_tracker=self._handoff, 168 token_aggregator=self._tokens, 169 recording=self._recording, 170 turn_trace_registry=self._turn_trace_registry, 171 ) 172 173 def on_start(self, span, parent_context=None) -> None: 174 super().on_start(span, parent_context) 175 # Capture turn.index while open_agent_turn_index is still set — child 176 # llm/tts spans often end after conversation_item_added clears it. 177 bootstrap = get_job_bootstrap() 178 if bootstrap is None or not bootstrap.state.parlot_session_id: 179 return 180 self._turns.stamp_turn_index_at_start(span, bootstrap.state) 181 182 def on_end(self, span: ReadableSpan) -> None: 183 super().on_end(span) 184 try: 185 name = span.name 186 if name == SPAN_CONVERSATION_SESSION: 187 handle_conversation_session_on_end() 188 return 189 self._enrich(span) 190 self._log_compare_span(span) 191 except Exception as exc: 192 logger.error( 193 "parlot: LiveKitGenAIProcessor failed on span %r — %s", 194 span.name, 195 exc, 196 ) 197 logger.debug( 198 "parlot: LiveKitGenAIProcessor failed on span %r", 199 span.name, 200 exc_info=True, 201 ) 202 203 def set_tracer(self, tracer: Tracer) -> None: 204 self._tracer = tracer 205 206 def set_metrics(self, metrics) -> None: 207 self._metrics = metrics 208 209 def set_turn_source(self, source: str) -> None: 210 if source in ("spans", "events"): 211 self._turn_source = source 212 213 @property 214 def turn_source(self) -> str: 215 return self._turn_source 216 217 def enrich_spans_for_export(self, spans: list[ReadableSpan]) -> None: 218 """Deferred enrichment, turn-attr copy, then GenAI/voice rename for export.""" 219 self._turns.correct_function_tool_timing(spans) 220 221 llm_spans = sorted( 222 (s for s in spans if s.name == "llm_node"), 223 key=lambda s: s.end_time or 0, 224 ) 225 for span in llm_spans: 226 # FIFO token attach before stamping active_speech_id (which would 227 # pin every span to the latest speech and break multi-node batches). 228 self._tokens.apply_plugin_llm_usage_to_span( 229 span, dict(span.attributes or {}), prefer_fifo=True 230 ) 231 attrs = span.attributes or {} 232 if not attrs.get(ATTR_LK_SPEECH_ID): 233 bootstrap = get_job_bootstrap() 234 speech_id = "" 235 if bootstrap is not None: 236 speech_id = bootstrap.state.active_speech_id.strip() 237 if speech_id: 238 self._set(span, ATTR_LK_SPEECH_ID, speech_id) 239 240 self._turns.merge_native_turn_attrs_onto_parlot_turns(spans) 241 242 from ._telemetry_compare import compare_enabled, get_compare_logger 243 244 if compare_enabled(): 245 bootstrap = get_job_bootstrap() 246 if bootstrap is not None and bootstrap.state.parlot_session_id: 247 session_id = bootstrap.state.parlot_session_id 248 for span in spans: 249 if span.name == "llm_node": 250 get_compare_logger().accumulate_export_tokens( 251 session_id, 252 span_name=span.name or "", 253 attrs=dict(span.attributes or {}), 254 ) 255 256 # Rename LiveKit-native ops → GenAI/voice names before export filter. 257 for span in spans: 258 native = span.name or "" 259 if native in NATIVE_TURN_SPANS: 260 # Dropped by export filter (remap returns None); attrs already merged. 261 continue 262 if remap_livekit_span_name(native, span.attributes or {}) is not None: 263 apply_livekit_span_rename(span) 264 role = livekit_agent_role_for_span(span.name or "") 265 if role and not (span.attributes or {}).get(ATTR_AGENT_ROLE): 266 self._set(span, ATTR_AGENT_ROLE, role) 267 268 def _log_compare_span(self, span: ReadableSpan) -> None: 269 from ._session import get_job_bootstrap 270 from ._telemetry_compare import get_compare_logger 271 272 bootstrap = get_job_bootstrap() 273 if bootstrap is None or not bootstrap.state.parlot_session_id: 274 return 275 attrs = span.attributes or {} 276 get_compare_logger().log_span( 277 bootstrap.state.parlot_session_id, 278 span_name=span.name or "", 279 attrs=dict(attrs), 280 ) 281 282 # ------------------------------------------------------------------ 283 # Public facades (called by _events / _egress) 284 # ------------------------------------------------------------------ 285 286 def mark_conversation_item_committed(self, item_id: str) -> bool: 287 return self._turns.mark_conversation_item_committed(item_id) 288 289 def mark_handoff_item_committed(self, item_id: str) -> None: 290 self._handoff.mark_handoff_item_committed(item_id) 291 292 def committed_handoff_item_ids(self) -> set[str]: 293 return self._handoff.committed_handoff_item_ids() 294 295 def record_handoff_from_event( 296 self, *, from_agent: str = "", to_agent: str = "" 297 ) -> None: 298 self._handoff.record_handoff_from_event( 299 from_agent=from_agent, to_agent=to_agent 300 ) 301 302 def note_user_transcription_meta( 303 self, *, speaker_id: str = "", language: str = "" 304 ) -> None: 305 self._turns.note_user_transcription_meta( 306 speaker_id=speaker_id, language=language 307 ) 308 309 def note_function_tools_executed(self, count: int) -> None: 310 self._turns.note_function_tools_executed(count) 311 312 def apply_session_usage(self, total_in: int, total_out: int) -> None: 313 self._tokens.apply_session_usage(total_in, total_out) 314 315 def note_session_error(self, message: str, *, recoverable: bool) -> None: 316 if not recoverable and message: 317 bootstrap = get_job_bootstrap() 318 if bootstrap is not None: 319 bootstrap.state.pending_close_error = message 320 321 def pop_pending_close_error(self) -> str | None: 322 bootstrap = get_job_bootstrap() 323 if bootstrap is None: 324 return None 325 err = bootstrap.state.pending_close_error.strip() 326 bootstrap.state.pending_close_error = "" 327 return err or None 328 329 def commit_user_message( 330 self, 331 text: str, 332 *, 333 interrupted: bool = False, 334 metrics: dict[str, float] | None = None, 335 ) -> None: 336 self._turns.commit_user_message( 337 text, interrupted=interrupted, metrics=metrics 338 ) 339 340 def commit_agent_message( 341 self, 342 text: str, 343 *, 344 interrupted: bool = False, 345 metrics: dict[str, float] | None = None, 346 ) -> None: 347 self._turns.commit_agent_message( 348 text, interrupted=interrupted, metrics=metrics 349 ) 350 351 def lookup_turn_trace( 352 self, session_id: str, turn_index: int 353 ) -> tuple[str, str] | None: 354 return self._turns.lookup_turn_trace(session_id, turn_index) 355 356 def set_recording_anchor_wall_ms( 357 self, state: _LiveKitSessionState, anchor_wall_ms: int 358 ) -> None: 359 self._recording.set_recording_anchor_wall_ms(state, anchor_wall_ms) 360 361 def _apply_plugin_llm_usage_to_span( 362 self, 363 span: ReadableSpan, 364 attrs: Mapping[str, AttributeValue], 365 *, 366 prefer_fifo: bool = False, 367 ) -> None: 368 self._tokens.apply_plugin_llm_usage_to_span( 369 span, attrs, prefer_fifo=prefer_fifo 370 ) 371 372 # ------------------------------------------------------------------ 373 # Session / agent identity helpers (kept on processor) 374 # ------------------------------------------------------------------ 375 376 def _resolve_session_state( 377 self, span: ReadableSpan, attrs: Mapping[str, AttributeValue] 378 ) -> _LiveKitSessionState | None: 379 bootstrap = get_job_bootstrap() 380 if bootstrap is not None: 381 return bootstrap.state 382 sid = attrs.get(ATTR_SESSION_ID) 383 if sid: 384 state = self._sessions.get(str(sid)) 385 if state is not None: 386 return state 387 job_id = attrs.get(ATTR_LK_JOB_ID) or attrs.get(METADATA_JOB_ID) 388 if job_id: 389 state = self._sessions.get(str(job_id)) 390 if state is not None: 391 return state 392 name = span.name or "" 393 if ( 394 name in _AGENT_PIPELINE_SPANS 395 or name in _AGENT_LABEL_SPANS 396 or job_id is not None 397 ): 398 logger.debug( 399 "parlot: no session bootstrap for span %s (trace=%s); skipping enrich", 400 name, 401 self._trace_id_hex(span), 402 ) 403 return None 404 logger.error( 405 "parlot: no session bootstrap for span %s (trace=%s)", 406 name, 407 self._trace_id_hex(span), 408 ) 409 return None 410 411 def _track_agent_label( 412 self, state: _LiveKitSessionState, attrs: Mapping[str, AttributeValue] 413 ) -> None: 414 raw_label = attrs.get(ATTR_LK_AGENT_LABEL) 415 raw_name = attrs.get(ATTR_LK_AGENT_NAME) 416 label_str = raw_label if isinstance(raw_label, str) else None 417 name_str = raw_name if isinstance(raw_name, str) else None 418 label = topology_agent_name(label_str) or topology_agent_name(name_str) 419 if label: 420 self._maybe_update(state, "agent_label", label) 421 canonical = explicit_agent_id() 422 if label != canonical: 423 append_agent_chain_step(state, label) 424 425 def _active_agent_id( 426 self, 427 state: _LiveKitSessionState, 428 attrs: Mapping[str, AttributeValue], 429 *, 430 label_override: str | None = None, 431 ) -> str: 432 if label_override: 433 override = topology_agent_name(label_override) 434 if override: 435 return override 436 # Prefer sub-agent labels over worker name; never use LiveKit AD_* dispatch ids. 437 for candidate in ( 438 attrs.get(ATTR_LK_AGENT_LABEL), 439 state.agent_label, 440 attrs.get(ATTR_LK_AGENT_NAME), 441 state.worker_agent_name, 442 explicit_agent_id(), 443 state.agent_chain[-1] if state.agent_chain else "", 444 ): 445 cand_str = candidate if isinstance(candidate, str) else None 446 name = topology_agent_name(cand_str) 447 if name: 448 return name 449 return "unknown" 450 451 def _stamp_agent_identity( 452 self, 453 span: ReadableSpan, 454 state: _LiveKitSessionState, 455 attrs: Mapping[str, AttributeValue] | None = None, 456 *, 457 label_override: str | None = None, 458 ) -> None: 459 resolved: Mapping[str, AttributeValue] = ( 460 attrs if attrs is not None else (span.attributes or {}) 461 ) 462 if not resolved.get(ATTR_GEN_AI_AGENT_NAME): 463 agent_id = self._active_agent_id( 464 state, resolved, label_override=label_override 465 ) 466 if agent_id == "unknown": 467 agent_id = configured_agent_id() or "unknown" 468 if agent_id != "unknown": 469 self._set(span, ATTR_GEN_AI_AGENT_NAME, agent_id) 470 version = configured_agent_version() 471 if version and not resolved.get(ATTR_GEN_AI_AGENT_VERSION): 472 self._set(span, ATTR_GEN_AI_AGENT_VERSION, version) 473 474 def _stamp_response_model_if_distinct( 475 self, 476 span: ReadableSpan, 477 attrs: Mapping[str, AttributeValue], 478 ) -> None: 479 response_model = str(attrs.get(ATTR_GEN_AI_RESPONSE_MODEL, "") or "").strip() 480 if not response_model: 481 return 482 request_model = str(attrs.get(ATTR_GEN_AI_MODEL, "") or "").strip() 483 if request_model and response_model == request_model: 484 return 485 self._set(span, ATTR_GEN_AI_RESPONSE_MODEL, response_model) 486 487 def _stamp_session_turn_attrs( 488 self, span: ReadableSpan, state: _LiveKitSessionState 489 ) -> None: 490 if state.parlot_session_id: 491 self._set(span, ATTR_SESSION_ID, state.parlot_session_id) 492 self._set(span, ATTR_SESSION_CONVERSATION_ID, state.conversation_id) 493 self._set(span, ATTR_GEN_AI_CONVERSATION_ID, state.conversation_id) 494 if (span.attributes or {}).get(ATTR_TURN_INDEX) is not None: 495 return 496 active = self._turns.active_turn_index(state, span.name) 497 if active is not None: 498 self._set(span, ATTR_TURN_INDEX, active) 499 500 def _stamp_error_status_if_needed(self, span: ReadableSpan) -> None: 501 """Normalize failed pipeline spans to OTel ERROR status for OTLP export.""" 502 attrs = span.attributes or {} 503 if attrs.get(ATTR_EXCEPTION_TYPE) or attrs.get(ATTR_LK_FNC_TOOL_ERROR): 504 span._status = Status(StatusCode.ERROR) 505 506 # ------------------------------------------------------------------ 507 # Enrichment dispatcher 508 # ------------------------------------------------------------------ 509 510 def _enrich(self, span: ReadableSpan) -> None: 511 name = span.name 512 if name in ("parlot.turn", SPAN_PARLOT_SESSION_CLOSE, SPAN_CONVERSATION_SESSION): 513 return 514 attrs = span.attributes or {} 515 state = self._resolve_session_state(span, attrs) 516 if state is None: 517 # Post-close JudgeGroup LLM: sticky closed session only — no aggregates. 518 if name in ("llm_request", "llm_request_run", "llm_node"): 519 self._enrich_late_evaluation_llm(span, attrs) 520 return 521 522 explicit_job = attrs.get(ATTR_LK_JOB_ID) or attrs.get(METADATA_JOB_ID) 523 if explicit_job: 524 self._maybe_update(state, "session_id", explicit_job) 525 self._maybe_update( 526 state, 527 "room_name", 528 attr_get(attrs, ATTR_LK_ROOM_NAME, ATTR_ROOM_NAME_LEGACY), 529 ) 530 self._maybe_update(state, "room_sid", attrs.get(ATTR_LK_ROOM_SID) or attrs.get(METADATA_ROOM_ID)) 531 532 if name in _AGENT_LABEL_SPANS: 533 self._track_agent_label(state, attrs) 534 535 if state.session_id and not state.room_sid: 536 rn, rs = lookup_room_context(state.session_id) 537 self._maybe_update(state, "room_name", rn or None) 538 self._maybe_update(state, "room_sid", rs or None) 539 540 if state.session_id: 541 self._set(span, ATTR_LK_JOB_ID, state.session_id) 542 if state.room_name: 543 self._set(span, ATTR_LK_ROOM_NAME, state.room_name) 544 if state.room_sid: 545 self._set(span, ATTR_LK_ROOM_SID, state.room_sid) 546 547 stamp_livekit_platform_refs( 548 span, 549 job_id=state.session_id, 550 room_name=state.room_name, 551 room_sid=state.room_sid, 552 ) 553 554 self._stamp_session_turn_attrs(span, state) 555 self._set(span, ATTR_AGENT_FRAMEWORK, "livekit") 556 stage = livekit_agent_stage_for_span(name) 557 if stage and not (span.attributes or {}).get(ATTR_AGENT_STAGE): 558 self._set(span, ATTR_AGENT_STAGE, stage) 559 560 if name in ("llm_request", "llm_request_run"): 561 self._enrich_llm_request(span, state) 562 elif name == "llm_node": 563 self._turns.enrich_llm_node(span, state) 564 elif name == "tts_node": 565 self._turns.enrich_tts_node(span, state) 566 elif name == "tts_request_run": 567 self._turns.enrich_tts_request(span, state) 568 elif name == "function_tool": 569 self._turns.enrich_function_tool(span, state) 570 elif name == "user_turn": 571 self._turns.enrich_user_turn(span, state) 572 elif name == "agent_turn": 573 self._turns.enrich_agent_turn(span, state) 574 elif name == "drain_agent_activity": 575 self._turns.enrich_drain(span, state) 576 elif name == "eou_detection": 577 self._turns.enrich_eou(span, state) 578 elif name == "amd": 579 self._turns.enrich_amd(span, state) 580 elif name == SPAN_AGENT_HANDOFF: 581 self._handoff.enrich_handoff( 582 span, 583 state, 584 stamp_agent_identity=self._stamp_agent_identity, 585 turn_source=self._turn_source, 586 ) 587 588 if name in _AGENT_PIPELINE_SPANS: 589 self._stamp_error_status_if_needed(span) 590 591 def _enrich_late_evaluation_llm( 592 self, 593 span: ReadableSpan, 594 attrs: Mapping[str, AttributeValue], 595 ) -> None: 596 """Stamp closed-session ids + evaluation attrs; do not mutate turn/token state.""" 597 job_id = attrs.get(ATTR_LK_JOB_ID) or attrs.get(METADATA_JOB_ID) 598 if not job_id: 599 try: 600 from livekit.agents.job import get_job_context 601 602 ctx = get_job_context() 603 if ctx is not None: 604 job_id = str(ctx.job.id) 605 except Exception: 606 job_id = None 607 if not job_id: 608 return 609 sticky = get_sticky_closed_session(str(job_id)) 610 if sticky is None: 611 return 612 self._set(span, ATTR_SESSION_ID, sticky.session_id) 613 self._set(span, ATTR_SESSION_CONVERSATION_ID, sticky.conversation_id) 614 self._set(span, ATTR_GEN_AI_CONVERSATION_ID, sticky.conversation_id) 615 self._set(span, ATTR_LK_JOB_ID, str(job_id)) 616 self._set(span, ATTR_AGENT_FRAMEWORK, "livekit") 617 # Evaluation is a Parlot Layer-3 signal (parlot.span.kind), not a GenAI op. 618 self._set(span, ATTR_PARLOT_SPAN_KIND, PARLOT_SPAN_KIND_EVALUATION) 619 stage = livekit_agent_stage_for_span(span.name or "") 620 if stage and not attrs.get(ATTR_AGENT_STAGE): 621 self._set(span, ATTR_AGENT_STAGE, stage) 622 self._stamp_error_status_if_needed(span) 623 624 def _enrich_llm_request( 625 self, span: ReadableSpan, state: _LiveKitSessionState 626 ) -> None: 627 attrs = span.attributes or {} 628 self._stamp_agent_identity(span, state, attrs) 629 630 if not attrs.get(ATTR_GEN_AI_SYSTEM): 631 system = provider_to_system(str(attrs.get(ATTR_GEN_AI_PROVIDER, ""))) 632 if system: 633 self._set(span, ATTR_GEN_AI_SYSTEM, system) 634 635 self._tokens.accumulate_llm_request_tokens(span, state, attrs) 636 self._handoff.stamp_transfer_latency_if_pending(span, state) 637 self._stamp_response_model_if_distinct(span, span.attributes or attrs) 638 639 def _apply_root_to_live_span( 640 self, session_span: Any, state: _LiveKitSessionState 641 ) -> None: 642 """Stamp session aggregates on the live ``parlot.session`` span.""" 643 if session_span is None or not hasattr(session_span, "set_attribute"): 644 return 645 session_span.set_attribute(ATTR_SESSION_TURN_COUNT, state.turn_count) 646 session_span.set_attribute(ATTR_SESSION_TOOL_CALL_COUNT, state.tool_call_count) 647 session_span.set_attribute(ATTR_SESSION_HANDOFF_COUNT, state.handoff_count) 648 session_span.set_attribute( 649 ATTR_SESSION_TOTAL_INPUT_TOKENS, state.total_input_tokens 650 ) 651 session_span.set_attribute( 652 ATTR_SESSION_TOTAL_OUTPUT_TOKENS, state.total_output_tokens 653 ) 654 if state.user_id: 655 session_span.set_attribute(ATTR_SESSION_USER_ID, state.user_id) 656 ensure_agent_chain_seeded(state) 657 for agent in state.agent_chain: 658 state.topology.push_agent_chain(agent) 659 state.topology.stamp_session_span(session_span, final_turn=state.turn_count) 660 stamp_session_agent_identity(session_span, state) 661 stamp_session_sdk_version(session_span) 662 if state.amd: 663 session_span.set_attribute(ATTR_SESSION_AMD, state.amd) 664 if state.recording_anchor_wall_ms is not None: 665 session_span.set_attribute( 666 ATTR_SESSION_RECORDING_ANCHOR_WALL_MS, 667 state.recording_anchor_wall_ms, 668 ) 669 if state.languages_seen: 670 session_span.set_attribute( 671 ATTR_SESSION_LANGUAGES, 672 json.dumps(sorted(state.languages_seen)), 673 ) 674 if state.session_id: 675 stamp_livekit_platform_refs( 676 session_span, 677 job_id=state.session_id, 678 room_name=state.room_name, 679 room_sid=state.room_sid, 680 ) 681 if self._metrics and state.parlot_session_id: 682 self._metrics.record_session_close(state)
Enriches LiveKit Agent spans in-place with Parlot conventions.
Intercepts spans from the LiveKit Agents SDK and normalizes them into Parlot's three-layer semantic vocabulary:
- Conversation Contract — session start/close, turns, agent handoffs.
- OTel GenAI (v1.41.0) — LLM inference, tool executions, workflows.
- Voice spans — TTS, STT, and end-of-utterance operational timings.
124 def __init__( 125 self, 126 capture_genai_content: Optional[bool] = None, 127 handoff_tool_names: Optional[set[str]] = None, 128 *, 129 context: Optional["ParlotContext"] = None, 130 ) -> None: 131 from parlot.core.context import ParlotContext 132 133 # Explicit override for tests / rare call sites; None → policy at emit time. 134 self._capture_genai_content_override = capture_genai_content 135 self._handoff_tools = handoff_tool_names or set() 136 self._context = context if context is not None else ParlotContext() 137 self._sessions: dict[str, _LiveKitSessionState] = {} 138 self._turn_trace_registry: dict[str, dict[int, tuple[str, str]]] = {} 139 self._tracer: Tracer | None = None 140 self._metrics = None 141 self._turn_source: str = "spans" 142 143 self._handoff = HandoffTracker( 144 set_attr=self._set, 145 handoff_tools=self._handoff_tools, 146 ) 147 self._tokens = TokenAggregator( 148 set_attr=self._set, 149 plugin_host=self, 150 ) 151 self._recording = RecordingCoordinator( 152 get_capture_override=lambda: self._capture_genai_content_override, 153 active_agent_id_fn=self._active_agent_id, 154 get_context=lambda: self._context, 155 ) 156 self._turns = TurnEnricher( 157 set_attr=self._set, 158 add_event=self._add_event, 159 maybe_update=self._maybe_update, 160 get_turn_source=lambda: self._turn_source, 161 get_tracer=lambda: self._tracer, 162 get_metrics=lambda: self._metrics, 163 genai_content_enabled=self._recording.genai_content_enabled, 164 active_agent_id=self._active_agent_id, 165 stamp_agent_identity=self._stamp_agent_identity, 166 stamp_response_model_if_distinct=self._stamp_response_model_if_distinct, 167 handoff_tracker=self._handoff, 168 token_aggregator=self._tokens, 169 recording=self._recording, 170 turn_trace_registry=self._turn_trace_registry, 171 )
173 def on_start(self, span, parent_context=None) -> None: 174 super().on_start(span, parent_context) 175 # Capture turn.index while open_agent_turn_index is still set — child 176 # llm/tts spans often end after conversation_item_added clears it. 177 bootstrap = get_job_bootstrap() 178 if bootstrap is None or not bootstrap.state.parlot_session_id: 179 return 180 self._turns.stamp_turn_index_at_start(span, bootstrap.state)
Called when a opentelemetry.trace.Span is started.
This method is called synchronously on the thread that starts the span, therefore it should not block or throw an exception.
Arguments:
- span: The
opentelemetry.trace.Spanthat just started. - parent_context: The parent context of the span that just started.
182 def on_end(self, span: ReadableSpan) -> None: 183 super().on_end(span) 184 try: 185 name = span.name 186 if name == SPAN_CONVERSATION_SESSION: 187 handle_conversation_session_on_end() 188 return 189 self._enrich(span) 190 self._log_compare_span(span) 191 except Exception as exc: 192 logger.error( 193 "parlot: LiveKitGenAIProcessor failed on span %r — %s", 194 span.name, 195 exc, 196 ) 197 logger.debug( 198 "parlot: LiveKitGenAIProcessor failed on span %r", 199 span.name, 200 exc_info=True, 201 )
Called when a opentelemetry.trace.Span is ended.
This method is called synchronously on the thread that ends the span, therefore it should not block or throw an exception.
Arguments:
- span: The
opentelemetry.trace.Spanthat just ended.
217 def enrich_spans_for_export(self, spans: list[ReadableSpan]) -> None: 218 """Deferred enrichment, turn-attr copy, then GenAI/voice rename for export.""" 219 self._turns.correct_function_tool_timing(spans) 220 221 llm_spans = sorted( 222 (s for s in spans if s.name == "llm_node"), 223 key=lambda s: s.end_time or 0, 224 ) 225 for span in llm_spans: 226 # FIFO token attach before stamping active_speech_id (which would 227 # pin every span to the latest speech and break multi-node batches). 228 self._tokens.apply_plugin_llm_usage_to_span( 229 span, dict(span.attributes or {}), prefer_fifo=True 230 ) 231 attrs = span.attributes or {} 232 if not attrs.get(ATTR_LK_SPEECH_ID): 233 bootstrap = get_job_bootstrap() 234 speech_id = "" 235 if bootstrap is not None: 236 speech_id = bootstrap.state.active_speech_id.strip() 237 if speech_id: 238 self._set(span, ATTR_LK_SPEECH_ID, speech_id) 239 240 self._turns.merge_native_turn_attrs_onto_parlot_turns(spans) 241 242 from ._telemetry_compare import compare_enabled, get_compare_logger 243 244 if compare_enabled(): 245 bootstrap = get_job_bootstrap() 246 if bootstrap is not None and bootstrap.state.parlot_session_id: 247 session_id = bootstrap.state.parlot_session_id 248 for span in spans: 249 if span.name == "llm_node": 250 get_compare_logger().accumulate_export_tokens( 251 session_id, 252 span_name=span.name or "", 253 attrs=dict(span.attributes or {}), 254 ) 255 256 # Rename LiveKit-native ops → GenAI/voice names before export filter. 257 for span in spans: 258 native = span.name or "" 259 if native in NATIVE_TURN_SPANS: 260 # Dropped by export filter (remap returns None); attrs already merged. 261 continue 262 if remap_livekit_span_name(native, span.attributes or {}) is not None: 263 apply_livekit_span_rename(span) 264 role = livekit_agent_role_for_span(span.name or "") 265 if role and not (span.attributes or {}).get(ATTR_AGENT_ROLE): 266 self._set(span, ATTR_AGENT_ROLE, role)
Deferred enrichment, turn-attr copy, then GenAI/voice rename for export.
71def add_platform_ref( 72 kind: str, 73 value: str, 74 *, 75 framework: str = "custom", 76) -> None: 77 """ 78 Attach a searchable external ID to the active Parlot session. 79 80 Stamps ``platform.ref.{kind}`` (and the primary triple when this is the 81 first ref) on the live session span so the session can be found by that 82 value in Parlot search / resolve. 83 84 Args: 85 kind: Identifier type (e.g. ``"crm_ticket"``, ``"order_number"``, 86 ``"call_sid"``). 87 value: Unique identifier value (e.g. ``"TKT-9921"``). 88 framework: Originating framework name. Defaults to ``"custom"``. 89 90 Example:: 91 92 from parlot.instrumentation.livekit import add_platform_ref 93 94 add_platform_ref("crm_ticket", "TKT-9") 95 """ 96 kind = str(kind or "").strip() 97 value = str(value or "").strip() 98 framework = str(framework or "custom").strip() or "custom" 99 if not kind or not value: 100 return 101 span = get_active_session_span() 102 if span is None: 103 return 104 stamp_platform_refs(span, [(framework, kind, value)])
Attach a searchable external ID to the active Parlot session.
Stamps platform.ref.{kind} (and the primary triple when this is the
first ref) on the live session span so the session can be found by that
value in Parlot search / resolve.
Arguments:
- kind: Identifier type (e.g.
"crm_ticket","order_number","call_sid"). - value: Unique identifier value (e.g.
"TKT-9921"). - framework: Originating framework name. Defaults to
"custom".
Example::
from parlot.instrumentation.livekit import add_platform_ref
add_platform_ref("crm_ticket", "TKT-9")
50@contextmanager 51def human_escalation(label: str | None = None) -> Iterator[None]: 52 """Mark the next participant who joins the active session as a human rep. 53 54 Args: 55 label: Optional role or team label for the incoming human representative. 56 """ 57 token = _pending_escalation_label.set(label) 58 try: 59 yield 60 finally: 61 _pending_escalation_label.reset(token)
Mark the next participant who joins the active session as a human rep.
Arguments:
- label: Optional role or team label for the incoming human representative.
19def record_human_rep(participant_id: str, *, label: str | None = None) -> None: 20 """ 21 Mark a participant as a human representative. 22 23 Stamps ``session.topology.agents`` on the active session span and registers 24 the participant for ``turn.participant_role=human_rep`` on future turns. 25 26 Args: 27 participant_id: Participant identifier within the room/session. 28 label: Optional human-readable name (e.g. ``"Tier 2 Escalation Desk"``). 29 """ 30 participant_id = str(participant_id or "").strip() 31 if not participant_id: 32 return 33 34 state = get_active_session() 35 if state is not None: 36 state.human_rep_participant_ids.add(participant_id) 37 entry: dict[str, str] = {"id": participant_id, "role": "human_rep"} 38 if label: 39 entry["label"] = label 40 if not any(a.get("id") == participant_id for a in state.topology_agents): 41 state.topology_agents.append(entry) 42 43 span = get_active_session_span() 44 if span is not None and state is not None and state.topology_agents: 45 span.set_attribute( 46 ATTR_SESSION_TOPOLOGY_AGENTS, json.dumps(state.topology_agents) 47 )
Mark a participant as a human representative.
Stamps session.topology.agents on the active session span and registers
the participant for turn.participant_role=human_rep on future turns.
Arguments:
- participant_id: Participant identifier within the room/session.
- label: Optional human-readable name (e.g.
"Tier 2 Escalation Desk").
20def set_session_attribute(key: str, value: str | int | float | bool) -> None: 21 """ 22 Stamp one custom attribute on the active session under ``session.metadata.*``. 23 24 Keys are normalized to ``session.metadata.<key>``. Values are stored as 25 strings on the live session span and remembered on session state so they 26 are also present on ``parlot.session.close``. 27 28 Args: 29 key: Attribute name. If not prefixed with ``session.metadata.``, the 30 prefix is added automatically. 31 value: Value to record (``str``, ``int``, ``float``, or ``bool``). 32 """ 33 full_key = session_metadata_key(key) 34 if not full_key or full_key == ATTR_SESSION_METADATA_PREFIX: 35 return 36 str_value = value if isinstance(value, str) else str(value) 37 38 state = get_active_session() 39 if state is not None: 40 state.custom_metadata[full_key] = str_value 41 42 span = get_active_session_span() 43 if span is not None and hasattr(span, "set_attribute"): 44 span.set_attribute(full_key, str_value)
Stamp one custom attribute on the active session under session.metadata.*.
Keys are normalized to session.metadata.<key>. Values are stored as
strings on the live session span and remembered on session state so they
are also present on parlot.session.close.
Arguments:
- key: Attribute name. If not prefixed with
session.metadata., the prefix is added automatically. - value: Value to record (
str,int,float, orbool).
47def set_session_metadata(**pairs: str | int | float | bool) -> None: 48 """ 49 Attach custom key/value metadata to the active Parlot session. 50 51 Each keyword becomes ``session.metadata.<name>`` on the session span and 52 appears in session detail in the Parlot UI. 53 54 Args: 55 **pairs: Keyword metadata pairs (values ``str``, ``int``, ``float``, 56 or ``bool``). 57 58 Example:: 59 60 from parlot.instrumentation.livekit import set_session_metadata 61 62 set_session_metadata(order_id="12345", crm_ticket="TKT-9") 63 """ 64 for key, value in pairs.items(): 65 set_session_attribute(key, value)
Attach custom key/value metadata to the active Parlot session.
Each keyword becomes session.metadata.<name> on the session span and
appears in session detail in the Parlot UI.
Arguments:
- **pairs: Keyword metadata pairs (values
str,int,float, orbool).
Example::
from parlot.instrumentation.livekit import set_session_metadata
set_session_metadata(order_id="12345", crm_ticket="TKT-9")
28def stamp_platform_refs( 29 span: Any, 30 refs: list[tuple[str, str, str]], 31) -> None: 32 """Stamp ``platform.ref.*`` triples onto a span. 33 34 ``refs`` is a list of ``(framework, kind, value)`` tuples, e.g. 35 ``("livekit", "room_sid", "RM_abc")``. The first tuple is also written 36 to the canonical triple attributes (``platform.ref.framework/kind/value``) 37 so backends can pivot on a single primary ref. 38 39 Every tuple is also written as a flat ``platform.ref.{kind} = value`` key 40 for ingestion fallbacks. 41 42 Args: 43 span: A live OTel span (``set_attribute``) or a ``ReadableSpan`` with a 44 mutable ``_attributes`` dict. 45 refs: Non-empty list of reference triples. 46 """ 47 if not refs: 48 return 49 50 def _write(key: str, value: str) -> None: 51 if hasattr(span, "set_attribute") and callable(span.set_attribute): 52 try: 53 span.set_attribute(key, value) 54 return 55 except Exception: 56 pass 57 try: 58 ParlotBaseProcessor._set(span, key, value) 59 except Exception: 60 return 61 62 fw, kind, val = refs[0] 63 _write(ATTR_PLATFORM_FRAMEWORK, fw) 64 _write(ATTR_PLATFORM_KIND, kind) 65 _write(ATTR_PLATFORM_VALUE, val) 66 67 for _fw_i, kind_i, val_i in refs: 68 _write(platform_ref_flat_key(kind_i), val_i)
Stamp platform.ref.* triples onto a span.
refs is a list of (framework, kind, value) tuples, e.g.
("livekit", "room_sid", "RM_abc"). The first tuple is also written
to the canonical triple attributes (platform.ref.framework/kind/value)
so backends can pivot on a single primary ref.
Every tuple is also written as a flat platform.ref.{kind} = value key
for ingestion fallbacks.
Arguments:
- span: A live OTel span (
set_attribute) or aReadableSpanwith a mutable_attributesdict. - refs: Non-empty list of reference triples.