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.
110class LiveKitGenAIProcessor(ParlotBaseProcessor): 111 """Enriches LiveKit Agent spans in-place with Parlot conventions. 112 113 Intercepts spans from the LiveKit Agents SDK and normalizes them into 114 Parlot's three-layer semantic vocabulary: 115 116 1. **Conversation Contract** — session start/close, turns, agent handoffs. 117 2. **OTel GenAI (v1.41.0)** — LLM inference, tool executions, workflows. 118 3. **Voice spans** — TTS, STT, and end-of-utterance operational timings. 119 """ 120 121 def __init__( 122 self, 123 capture_genai_content: Optional[bool] = None, 124 handoff_tool_names: Optional[set[str]] = None, 125 *, 126 context: Optional["ParlotContext"] = None, 127 ) -> None: 128 from parlot.core.context import ParlotContext 129 130 # Explicit override for tests / rare call sites; None → policy at emit time. 131 self._capture_genai_content_override = capture_genai_content 132 self._handoff_tools = handoff_tool_names or set() 133 self._context = context if context is not None else ParlotContext() 134 self._sessions: dict[str, _LiveKitSessionState] = {} 135 self._turn_trace_registry: dict[str, dict[int, tuple[str, str]]] = {} 136 self._tracer: Tracer | None = None 137 self._metrics = None 138 self._turn_source: str = "spans" 139 140 self._handoff = HandoffTracker( 141 set_attr=self._set, 142 handoff_tools=self._handoff_tools, 143 ) 144 self._tokens = TokenAggregator( 145 set_attr=self._set, 146 plugin_host=self, 147 ) 148 self._recording = RecordingCoordinator( 149 get_capture_override=lambda: self._capture_genai_content_override, 150 active_agent_id_fn=self._active_agent_id, 151 get_context=lambda: self._context, 152 ) 153 self._turns = TurnEnricher( 154 set_attr=self._set, 155 add_event=self._add_event, 156 maybe_update=self._maybe_update, 157 get_turn_source=lambda: self._turn_source, 158 get_tracer=lambda: self._tracer, 159 get_metrics=lambda: self._metrics, 160 genai_content_enabled=self._recording.genai_content_enabled, 161 active_agent_id=self._active_agent_id, 162 stamp_agent_identity=self._stamp_agent_identity, 163 stamp_response_model_if_distinct=self._stamp_response_model_if_distinct, 164 handoff_tracker=self._handoff, 165 token_aggregator=self._tokens, 166 recording=self._recording, 167 turn_trace_registry=self._turn_trace_registry, 168 ) 169 170 def on_start(self, span, parent_context=None) -> None: 171 super().on_start(span, parent_context) 172 # Capture turn.index while open_agent_turn_index is still set — child 173 # llm/tts spans often end after conversation_item_added clears it. 174 bootstrap = get_job_bootstrap() 175 if bootstrap is None or not bootstrap.state.parlot_session_id: 176 return 177 self._turns.stamp_turn_index_at_start(span, bootstrap.state) 178 179 def on_end(self, span: ReadableSpan) -> None: 180 super().on_end(span) 181 try: 182 name = span.name 183 if name == SPAN_CONVERSATION_SESSION: 184 handle_conversation_session_on_end() 185 return 186 self._enrich(span) 187 self._log_compare_span(span) 188 except Exception as exc: 189 logger.error( 190 "parlot: LiveKitGenAIProcessor failed on span %r — %s", 191 span.name, 192 exc, 193 ) 194 logger.debug( 195 "parlot: LiveKitGenAIProcessor failed on span %r", 196 span.name, 197 exc_info=True, 198 ) 199 200 def set_tracer(self, tracer: Tracer) -> None: 201 self._tracer = tracer 202 203 def set_metrics(self, metrics) -> None: 204 self._metrics = metrics 205 206 def set_turn_source(self, source: str) -> None: 207 if source in ("spans", "events"): 208 self._turn_source = source 209 210 @property 211 def turn_source(self) -> str: 212 return self._turn_source 213 214 def enrich_spans_for_export(self, spans: list[ReadableSpan]) -> None: 215 """Deferred enrichment, turn-attr copy, then GenAI/voice rename for export.""" 216 self._turns.correct_function_tool_timing(spans) 217 218 llm_spans = sorted( 219 (s for s in spans if s.name == "llm_node"), 220 key=lambda s: s.end_time or 0, 221 ) 222 for span in llm_spans: 223 # FIFO token attach before stamping active_speech_id (which would 224 # pin every span to the latest speech and break multi-node batches). 225 self._tokens.apply_plugin_llm_usage_to_span( 226 span, dict(span.attributes or {}), prefer_fifo=True 227 ) 228 attrs = span.attributes or {} 229 if not attrs.get(ATTR_LK_SPEECH_ID): 230 bootstrap = get_job_bootstrap() 231 speech_id = "" 232 if bootstrap is not None: 233 speech_id = bootstrap.state.active_speech_id.strip() 234 if speech_id: 235 self._set(span, ATTR_LK_SPEECH_ID, speech_id) 236 237 self._turns.merge_native_turn_attrs_onto_parlot_turns(spans) 238 239 from ._telemetry_compare import compare_enabled, get_compare_logger 240 241 if compare_enabled(): 242 bootstrap = get_job_bootstrap() 243 if bootstrap is not None and bootstrap.state.parlot_session_id: 244 session_id = bootstrap.state.parlot_session_id 245 for span in spans: 246 if span.name == "llm_node": 247 get_compare_logger().accumulate_export_tokens( 248 session_id, 249 span_name=span.name or "", 250 attrs=dict(span.attributes or {}), 251 ) 252 253 # Rename LiveKit-native ops → GenAI/voice names before export filter. 254 for span in spans: 255 native = span.name or "" 256 if native in NATIVE_TURN_SPANS: 257 # Dropped by export filter (remap returns None); attrs already merged. 258 continue 259 if remap_livekit_span_name(native, span.attributes or {}) is not None: 260 apply_livekit_span_rename(span) 261 role = livekit_agent_role_for_span(span.name or "") 262 if role and not (span.attributes or {}).get(ATTR_AGENT_ROLE): 263 self._set(span, ATTR_AGENT_ROLE, role) 264 265 def _log_compare_span(self, span: ReadableSpan) -> None: 266 from ._session import get_job_bootstrap 267 from ._telemetry_compare import get_compare_logger 268 269 bootstrap = get_job_bootstrap() 270 if bootstrap is None or not bootstrap.state.parlot_session_id: 271 return 272 attrs = span.attributes or {} 273 get_compare_logger().log_span( 274 bootstrap.state.parlot_session_id, 275 span_name=span.name or "", 276 attrs=dict(attrs), 277 ) 278 279 # ------------------------------------------------------------------ 280 # Public facades (called by _events / _egress) 281 # ------------------------------------------------------------------ 282 283 def mark_conversation_item_committed(self, item_id: str) -> bool: 284 return self._turns.mark_conversation_item_committed(item_id) 285 286 def mark_handoff_item_committed(self, item_id: str) -> None: 287 self._handoff.mark_handoff_item_committed(item_id) 288 289 def committed_handoff_item_ids(self) -> set[str]: 290 return self._handoff.committed_handoff_item_ids() 291 292 def record_handoff_from_event( 293 self, *, from_agent: str = "", to_agent: str = "" 294 ) -> None: 295 self._handoff.record_handoff_from_event( 296 from_agent=from_agent, to_agent=to_agent 297 ) 298 299 def note_user_transcription_meta( 300 self, *, speaker_id: str = "", language: str = "" 301 ) -> None: 302 self._turns.note_user_transcription_meta( 303 speaker_id=speaker_id, language=language 304 ) 305 306 def note_function_tools_executed(self, count: int) -> None: 307 self._turns.note_function_tools_executed(count) 308 309 def apply_session_usage(self, total_in: int, total_out: int) -> None: 310 self._tokens.apply_session_usage(total_in, total_out) 311 312 def note_session_error(self, message: str, *, recoverable: bool) -> None: 313 if not recoverable and message: 314 bootstrap = get_job_bootstrap() 315 if bootstrap is not None: 316 bootstrap.state.pending_close_error = message 317 318 def pop_pending_close_error(self) -> str | None: 319 bootstrap = get_job_bootstrap() 320 if bootstrap is None: 321 return None 322 err = bootstrap.state.pending_close_error.strip() 323 bootstrap.state.pending_close_error = "" 324 return err or None 325 326 def commit_user_message( 327 self, 328 text: str, 329 *, 330 interrupted: bool = False, 331 metrics: dict[str, float] | None = None, 332 ) -> None: 333 self._turns.commit_user_message( 334 text, interrupted=interrupted, metrics=metrics 335 ) 336 337 def commit_agent_message( 338 self, 339 text: str, 340 *, 341 interrupted: bool = False, 342 metrics: dict[str, float] | None = None, 343 ) -> None: 344 self._turns.commit_agent_message( 345 text, interrupted=interrupted, metrics=metrics 346 ) 347 348 def lookup_turn_trace( 349 self, session_id: str, turn_index: int 350 ) -> tuple[str, str] | None: 351 return self._turns.lookup_turn_trace(session_id, turn_index) 352 353 def set_recording_anchor_wall_ms( 354 self, state: _LiveKitSessionState, anchor_wall_ms: int 355 ) -> None: 356 self._recording.set_recording_anchor_wall_ms(state, anchor_wall_ms) 357 358 def _apply_plugin_llm_usage_to_span( 359 self, 360 span: ReadableSpan, 361 attrs: Mapping[str, AttributeValue], 362 *, 363 prefer_fifo: bool = False, 364 ) -> None: 365 self._tokens.apply_plugin_llm_usage_to_span( 366 span, attrs, prefer_fifo=prefer_fifo 367 ) 368 369 # ------------------------------------------------------------------ 370 # Session / agent identity helpers (kept on processor) 371 # ------------------------------------------------------------------ 372 373 def _resolve_session_state( 374 self, span: ReadableSpan, attrs: Mapping[str, AttributeValue] 375 ) -> _LiveKitSessionState | None: 376 bootstrap = get_job_bootstrap() 377 if bootstrap is not None: 378 return bootstrap.state 379 sid = attrs.get(ATTR_SESSION_ID) 380 if sid: 381 state = self._sessions.get(str(sid)) 382 if state is not None: 383 return state 384 job_id = attrs.get(ATTR_LK_JOB_ID) or attrs.get(METADATA_JOB_ID) 385 if job_id: 386 state = self._sessions.get(str(job_id)) 387 if state is not None: 388 return state 389 name = span.name or "" 390 if ( 391 name in _AGENT_PIPELINE_SPANS 392 or name in _AGENT_LABEL_SPANS 393 or job_id is not None 394 ): 395 logger.debug( 396 "parlot: no session bootstrap for span %s (trace=%s); skipping enrich", 397 name, 398 self._trace_id_hex(span), 399 ) 400 return None 401 logger.error( 402 "parlot: no session bootstrap for span %s (trace=%s)", 403 name, 404 self._trace_id_hex(span), 405 ) 406 return None 407 408 def _track_agent_label( 409 self, state: _LiveKitSessionState, attrs: Mapping[str, AttributeValue] 410 ) -> None: 411 raw_label = attrs.get(ATTR_LK_AGENT_LABEL) 412 raw_name = attrs.get(ATTR_LK_AGENT_NAME) 413 label_str = raw_label if isinstance(raw_label, str) else None 414 name_str = raw_name if isinstance(raw_name, str) else None 415 label = topology_agent_name(label_str) or topology_agent_name(name_str) 416 if label: 417 self._maybe_update(state, "agent_label", label) 418 canonical = explicit_agent_id() 419 if label != canonical: 420 append_agent_chain_step(state, label) 421 422 def _active_agent_id( 423 self, 424 state: _LiveKitSessionState, 425 attrs: Mapping[str, AttributeValue], 426 *, 427 label_override: str | None = None, 428 ) -> str: 429 if label_override: 430 override = topology_agent_name(label_override) 431 if override: 432 return override 433 # Prefer sub-agent labels over worker name; never use LiveKit AD_* dispatch ids. 434 for candidate in ( 435 attrs.get(ATTR_LK_AGENT_LABEL), 436 state.agent_label, 437 attrs.get(ATTR_LK_AGENT_NAME), 438 state.worker_agent_name, 439 explicit_agent_id(), 440 state.agent_chain[-1] if state.agent_chain else "", 441 ): 442 cand_str = candidate if isinstance(candidate, str) else None 443 name = topology_agent_name(cand_str) 444 if name: 445 return name 446 return "unknown" 447 448 def _stamp_agent_identity( 449 self, 450 span: ReadableSpan, 451 state: _LiveKitSessionState, 452 attrs: Mapping[str, AttributeValue] | None = None, 453 *, 454 label_override: str | None = None, 455 ) -> None: 456 resolved: Mapping[str, AttributeValue] = ( 457 attrs if attrs is not None else (span.attributes or {}) 458 ) 459 if not resolved.get(ATTR_GEN_AI_AGENT_NAME): 460 agent_id = self._active_agent_id( 461 state, resolved, label_override=label_override 462 ) 463 if agent_id == "unknown": 464 agent_id = configured_agent_id() or "unknown" 465 if agent_id != "unknown": 466 self._set(span, ATTR_GEN_AI_AGENT_NAME, agent_id) 467 version = configured_agent_version() 468 if version and not resolved.get(ATTR_GEN_AI_AGENT_VERSION): 469 self._set(span, ATTR_GEN_AI_AGENT_VERSION, version) 470 471 def _stamp_response_model_if_distinct( 472 self, 473 span: ReadableSpan, 474 attrs: Mapping[str, AttributeValue], 475 ) -> None: 476 response_model = str(attrs.get(ATTR_GEN_AI_RESPONSE_MODEL, "") or "").strip() 477 if not response_model: 478 return 479 request_model = str(attrs.get(ATTR_GEN_AI_MODEL, "") or "").strip() 480 if request_model and response_model == request_model: 481 return 482 self._set(span, ATTR_GEN_AI_RESPONSE_MODEL, response_model) 483 484 def _stamp_session_turn_attrs( 485 self, span: ReadableSpan, state: _LiveKitSessionState 486 ) -> None: 487 if state.parlot_session_id: 488 self._set(span, ATTR_SESSION_ID, state.parlot_session_id) 489 self._set(span, ATTR_SESSION_CONVERSATION_ID, state.conversation_id) 490 self._set(span, ATTR_GEN_AI_CONVERSATION_ID, state.conversation_id) 491 if (span.attributes or {}).get(ATTR_TURN_INDEX) is not None: 492 return 493 active = self._turns.active_turn_index(state, span.name) 494 if active is not None: 495 self._set(span, ATTR_TURN_INDEX, active) 496 497 def _stamp_error_status_if_needed(self, span: ReadableSpan) -> None: 498 """Normalize failed pipeline spans to OTel ERROR status for OTLP export.""" 499 attrs = span.attributes or {} 500 if attrs.get(ATTR_EXCEPTION_TYPE) or attrs.get(ATTR_LK_FNC_TOOL_ERROR): 501 span._status = Status(StatusCode.ERROR) 502 503 # ------------------------------------------------------------------ 504 # Enrichment dispatcher 505 # ------------------------------------------------------------------ 506 507 def _enrich(self, span: ReadableSpan) -> None: 508 name = span.name 509 if name in ("parlot.turn", SPAN_PARLOT_SESSION_CLOSE, SPAN_CONVERSATION_SESSION): 510 return 511 attrs = span.attributes or {} 512 state = self._resolve_session_state(span, attrs) 513 if state is None: 514 return 515 516 explicit_job = attrs.get(ATTR_LK_JOB_ID) or attrs.get(METADATA_JOB_ID) 517 if explicit_job: 518 self._maybe_update(state, "session_id", explicit_job) 519 self._maybe_update( 520 state, 521 "room_name", 522 attr_get(attrs, ATTR_LK_ROOM_NAME, ATTR_ROOM_NAME_LEGACY), 523 ) 524 self._maybe_update(state, "room_sid", attrs.get(ATTR_LK_ROOM_SID) or attrs.get(METADATA_ROOM_ID)) 525 526 if name in _AGENT_LABEL_SPANS: 527 self._track_agent_label(state, attrs) 528 529 if state.session_id and not state.room_sid: 530 rn, rs = lookup_room_context(state.session_id) 531 self._maybe_update(state, "room_name", rn or None) 532 self._maybe_update(state, "room_sid", rs or None) 533 534 if state.session_id: 535 self._set(span, ATTR_LK_JOB_ID, state.session_id) 536 if state.room_name: 537 self._set(span, ATTR_LK_ROOM_NAME, state.room_name) 538 if state.room_sid: 539 self._set(span, ATTR_LK_ROOM_SID, state.room_sid) 540 541 stamp_livekit_platform_refs( 542 span, 543 job_id=state.session_id, 544 room_name=state.room_name, 545 room_sid=state.room_sid, 546 ) 547 548 self._stamp_session_turn_attrs(span, state) 549 self._set(span, ATTR_AGENT_FRAMEWORK, "livekit") 550 stage = livekit_agent_stage_for_span(name) 551 if stage and not (span.attributes or {}).get(ATTR_AGENT_STAGE): 552 self._set(span, ATTR_AGENT_STAGE, stage) 553 554 if name in ("llm_request", "llm_request_run"): 555 self._enrich_llm_request(span, state) 556 elif name == "llm_node": 557 self._turns.enrich_llm_node(span, state) 558 elif name == "tts_node": 559 self._turns.enrich_tts_node(span, state) 560 elif name == "tts_request_run": 561 self._turns.enrich_tts_request(span, state) 562 elif name == "function_tool": 563 self._turns.enrich_function_tool(span, state) 564 elif name == "user_turn": 565 self._turns.enrich_user_turn(span, state) 566 elif name == "agent_turn": 567 self._turns.enrich_agent_turn(span, state) 568 elif name == "drain_agent_activity": 569 self._turns.enrich_drain(span, state) 570 elif name == "eou_detection": 571 self._turns.enrich_eou(span, state) 572 elif name == "amd": 573 self._turns.enrich_amd(span, state) 574 elif name == SPAN_AGENT_HANDOFF: 575 self._handoff.enrich_handoff( 576 span, 577 state, 578 stamp_agent_identity=self._stamp_agent_identity, 579 turn_source=self._turn_source, 580 ) 581 582 if name in _AGENT_PIPELINE_SPANS: 583 self._stamp_error_status_if_needed(span) 584 585 def _enrich_llm_request( 586 self, span: ReadableSpan, state: _LiveKitSessionState 587 ) -> None: 588 attrs = span.attributes or {} 589 self._stamp_agent_identity(span, state, attrs) 590 591 if not attrs.get(ATTR_GEN_AI_SYSTEM): 592 system = provider_to_system(str(attrs.get(ATTR_GEN_AI_PROVIDER, ""))) 593 if system: 594 self._set(span, ATTR_GEN_AI_SYSTEM, system) 595 596 self._tokens.accumulate_llm_request_tokens(span, state, attrs) 597 self._handoff.stamp_transfer_latency_if_pending(span, state) 598 self._stamp_response_model_if_distinct(span, span.attributes or attrs) 599 600 def _apply_root_to_live_span( 601 self, session_span: Any, state: _LiveKitSessionState 602 ) -> None: 603 """Stamp session aggregates on the live ``parlot.session`` span.""" 604 if session_span is None or not hasattr(session_span, "set_attribute"): 605 return 606 session_span.set_attribute(ATTR_SESSION_TURN_COUNT, state.turn_count) 607 session_span.set_attribute(ATTR_SESSION_TOOL_CALL_COUNT, state.tool_call_count) 608 session_span.set_attribute(ATTR_SESSION_HANDOFF_COUNT, state.handoff_count) 609 session_span.set_attribute( 610 ATTR_SESSION_TOTAL_INPUT_TOKENS, state.total_input_tokens 611 ) 612 session_span.set_attribute( 613 ATTR_SESSION_TOTAL_OUTPUT_TOKENS, state.total_output_tokens 614 ) 615 if state.user_id: 616 session_span.set_attribute(ATTR_SESSION_USER_ID, state.user_id) 617 ensure_agent_chain_seeded(state) 618 for agent in state.agent_chain: 619 state.topology.push_agent_chain(agent) 620 state.topology.stamp_session_span(session_span, final_turn=state.turn_count) 621 stamp_session_agent_identity(session_span, state) 622 stamp_session_sdk_version(session_span) 623 if state.amd: 624 session_span.set_attribute(ATTR_SESSION_AMD, state.amd) 625 if state.recording_anchor_wall_ms is not None: 626 session_span.set_attribute( 627 ATTR_SESSION_RECORDING_ANCHOR_WALL_MS, 628 state.recording_anchor_wall_ms, 629 ) 630 if state.languages_seen: 631 session_span.set_attribute( 632 ATTR_SESSION_LANGUAGES, 633 json.dumps(sorted(state.languages_seen)), 634 ) 635 if state.session_id: 636 stamp_livekit_platform_refs( 637 session_span, 638 job_id=state.session_id, 639 room_name=state.room_name, 640 room_sid=state.room_sid, 641 ) 642 if self._metrics and state.parlot_session_id: 643 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.
121 def __init__( 122 self, 123 capture_genai_content: Optional[bool] = None, 124 handoff_tool_names: Optional[set[str]] = None, 125 *, 126 context: Optional["ParlotContext"] = None, 127 ) -> None: 128 from parlot.core.context import ParlotContext 129 130 # Explicit override for tests / rare call sites; None → policy at emit time. 131 self._capture_genai_content_override = capture_genai_content 132 self._handoff_tools = handoff_tool_names or set() 133 self._context = context if context is not None else ParlotContext() 134 self._sessions: dict[str, _LiveKitSessionState] = {} 135 self._turn_trace_registry: dict[str, dict[int, tuple[str, str]]] = {} 136 self._tracer: Tracer | None = None 137 self._metrics = None 138 self._turn_source: str = "spans" 139 140 self._handoff = HandoffTracker( 141 set_attr=self._set, 142 handoff_tools=self._handoff_tools, 143 ) 144 self._tokens = TokenAggregator( 145 set_attr=self._set, 146 plugin_host=self, 147 ) 148 self._recording = RecordingCoordinator( 149 get_capture_override=lambda: self._capture_genai_content_override, 150 active_agent_id_fn=self._active_agent_id, 151 get_context=lambda: self._context, 152 ) 153 self._turns = TurnEnricher( 154 set_attr=self._set, 155 add_event=self._add_event, 156 maybe_update=self._maybe_update, 157 get_turn_source=lambda: self._turn_source, 158 get_tracer=lambda: self._tracer, 159 get_metrics=lambda: self._metrics, 160 genai_content_enabled=self._recording.genai_content_enabled, 161 active_agent_id=self._active_agent_id, 162 stamp_agent_identity=self._stamp_agent_identity, 163 stamp_response_model_if_distinct=self._stamp_response_model_if_distinct, 164 handoff_tracker=self._handoff, 165 token_aggregator=self._tokens, 166 recording=self._recording, 167 turn_trace_registry=self._turn_trace_registry, 168 )
170 def on_start(self, span, parent_context=None) -> None: 171 super().on_start(span, parent_context) 172 # Capture turn.index while open_agent_turn_index is still set — child 173 # llm/tts spans often end after conversation_item_added clears it. 174 bootstrap = get_job_bootstrap() 175 if bootstrap is None or not bootstrap.state.parlot_session_id: 176 return 177 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.
179 def on_end(self, span: ReadableSpan) -> None: 180 super().on_end(span) 181 try: 182 name = span.name 183 if name == SPAN_CONVERSATION_SESSION: 184 handle_conversation_session_on_end() 185 return 186 self._enrich(span) 187 self._log_compare_span(span) 188 except Exception as exc: 189 logger.error( 190 "parlot: LiveKitGenAIProcessor failed on span %r — %s", 191 span.name, 192 exc, 193 ) 194 logger.debug( 195 "parlot: LiveKitGenAIProcessor failed on span %r", 196 span.name, 197 exc_info=True, 198 )
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.
214 def enrich_spans_for_export(self, spans: list[ReadableSpan]) -> None: 215 """Deferred enrichment, turn-attr copy, then GenAI/voice rename for export.""" 216 self._turns.correct_function_tool_timing(spans) 217 218 llm_spans = sorted( 219 (s for s in spans if s.name == "llm_node"), 220 key=lambda s: s.end_time or 0, 221 ) 222 for span in llm_spans: 223 # FIFO token attach before stamping active_speech_id (which would 224 # pin every span to the latest speech and break multi-node batches). 225 self._tokens.apply_plugin_llm_usage_to_span( 226 span, dict(span.attributes or {}), prefer_fifo=True 227 ) 228 attrs = span.attributes or {} 229 if not attrs.get(ATTR_LK_SPEECH_ID): 230 bootstrap = get_job_bootstrap() 231 speech_id = "" 232 if bootstrap is not None: 233 speech_id = bootstrap.state.active_speech_id.strip() 234 if speech_id: 235 self._set(span, ATTR_LK_SPEECH_ID, speech_id) 236 237 self._turns.merge_native_turn_attrs_onto_parlot_turns(spans) 238 239 from ._telemetry_compare import compare_enabled, get_compare_logger 240 241 if compare_enabled(): 242 bootstrap = get_job_bootstrap() 243 if bootstrap is not None and bootstrap.state.parlot_session_id: 244 session_id = bootstrap.state.parlot_session_id 245 for span in spans: 246 if span.name == "llm_node": 247 get_compare_logger().accumulate_export_tokens( 248 session_id, 249 span_name=span.name or "", 250 attrs=dict(span.attributes or {}), 251 ) 252 253 # Rename LiveKit-native ops → GenAI/voice names before export filter. 254 for span in spans: 255 native = span.name or "" 256 if native in NATIVE_TURN_SPANS: 257 # Dropped by export filter (remap returns None); attrs already merged. 258 continue 259 if remap_livekit_span_name(native, span.attributes or {}) is not None: 260 apply_livekit_span_rename(span) 261 role = livekit_agent_role_for_span(span.name or "") 262 if role and not (span.attributes or {}).get(ATTR_AGENT_ROLE): 263 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.