parlot.core
parlot-core: shared semantic conventions, base processor, and utilities.
1"""parlot-core: shared semantic conventions, base processor, and utilities.""" 2 3from .attrs import * # noqa: F401, F403 — re-export all attribute constants 4from .context import ParlotContext 5from .escalation import human_escalation, record_human_rep 6from .metadata import set_session_attribute, set_session_metadata 7from .platform_refs import add_platform_ref, stamp_platform_refs 8from .processor import ParlotBaseProcessor, assert_sync_span_processors 9from .provider import ( 10 adopt_existing_tracer_provider, 11 build_otlp_http_exporter, 12 build_parlot_client_headers, 13 build_resource, 14 build_tracer_provider, 15 resolve_api_key, 16 resolve_capture_genai_content, 17 resolve_endpoint, 18 HEADER_SDK_NAME, 19 HEADER_SDK_VERSION, 20 HEADER_INGESTION_VERSION, 21 INGESTION_PROTOCOL_VERSION, 22) 23from .bootstrap import fetch_telemetry_bootstrap 24from .parlotize import BaseParlotizeResult, ParlotizeProtocol, base_parlotize, configure_parlot_logging 25from .genai_content_capture import should_capture_genai_content 26from .export import ExportFilterSpanExporter 27from .intent import derive_intent 28from .logs_capture import should_capture_logs 29from .runtime import ParlotRuntimeContext, runtime_from_bootstrap 30from .sdk_version import resolve_parlot_sdk_version, stamp_session_sdk_version 31from .session import ( 32 SessionState, 33 clear_active_session, 34 get_active_session, 35 get_active_session_span, 36 session_owned, 37 set_active_session, 38) 39from .session_logs import SessionLogsCollector 40from .topology import SessionTopology 41from .turn_emit import stamp_turn_utterance_text 42 43__all__ = [ 44 "BaseParlotizeResult", 45 "ParlotizeProtocol", 46 "ParlotContext", 47 "ParlotRuntimeContext", 48 "SessionLogsCollector", 49 "base_parlotize", 50 "configure_parlot_logging", 51 "ExportFilterSpanExporter", 52 "ParlotBaseProcessor", 53 "assert_sync_span_processors", 54 "SessionState", 55 "SessionTopology", 56 "add_platform_ref", 57 "adopt_existing_tracer_provider", 58 "build_otlp_http_exporter", 59 "build_parlot_client_headers", 60 "build_resource", 61 "build_tracer_provider", 62 "HEADER_SDK_NAME", 63 "HEADER_SDK_VERSION", 64 "HEADER_INGESTION_VERSION", 65 "INGESTION_PROTOCOL_VERSION", 66 "clear_active_session", 67 "derive_intent", 68 "fetch_telemetry_bootstrap", 69 "get_active_session", 70 "get_active_session_span", 71 "human_escalation", 72 "record_human_rep", 73 "resolve_api_key", 74 "resolve_capture_genai_content", 75 "resolve_endpoint", 76 "resolve_parlot_sdk_version", 77 "runtime_from_bootstrap", 78 "session_owned", 79 "set_active_session", 80 "set_session_attribute", 81 "set_session_metadata", 82 "should_capture_genai_content", 83 "should_capture_logs", 84 "stamp_platform_refs", 85 "stamp_session_sdk_version", 86 "stamp_turn_utterance_text", 87]
77@dataclass 78class BaseParlotizeResult: 79 """Resolved context from ``base_parlotize`` for framework adapters.""" 80 81 context: ParlotContext 82 endpoint: str 83 api_key: str 84 tracer_provider: Any 85 agent_id: str 86 agent_version: str 87 capture_genai_content: Optional[bool] 88 capture_logs: bool | list[str] | None 89 log_level: Optional[str]
Resolved context from base_parlotize for framework adapters.
23@runtime_checkable 24class ParlotizeProtocol(Protocol): 25 """Common parameters every adapter ``parlotize()`` must accept. 26 27 Framework packages may add extra keyword-only parameters (e.g. LiveKit 28 ``record=`` / ``auto_escalate_sip=``). 29 """ 30 31 def __call__( 32 self, 33 agent_id: str, 34 *, 35 endpoint: Optional[str] = None, 36 api_key: Optional[str] = None, 37 capture_genai_content: Optional[bool] = None, 38 service_name: Optional[str] = None, 39 tracer_provider: TracerProvider | None = None, 40 version: Optional[str] = None, 41 capture_logs: bool | list[str] | None = None, 42 log_level: Optional[str] = None, 43 **kwargs: Any, 44 ) -> ParlotContext: 45 """Shared parameter surface for every adapter ``parlotize()``. 46 47 Args: 48 agent_id: Required canonical deployment identity stamped on 49 ``session.agent_id``. Must be non-empty after stripping whitespace. 50 endpoint: Parlot OTLP collector base URL (e.g. 51 ``https://ingest.parlot.ai``). Spans export to 52 ``{endpoint}/v1/traces``. If omitted, reads ``PARLOT_ENDPOINT``. 53 api_key: Org-scoped API key minted in Parlot **Settings → API Keys**. 54 If omitted, reads ``PARLOT_API_KEY``. Required for remote telemetry 55 bootstrap and recording grants. 56 capture_genai_content: Process-wide override for LLM message bodies and 57 tool input/output payloads. If ``False``, payloads are omitted while 58 preserving span durations, tokens, and turn text. Precedence: job 59 metadata → this kwarg → Settings → Generative AI (default: on). 60 service_name: OpenTelemetry resource ``service.name``. Defaults to 61 ``agent_id`` or a framework-specific fallback. 62 tracer_provider: Existing OpenTelemetry ``TracerProvider`` to adopt. If 63 omitted, adapters build one with Parlot's OTLP exporter (or adopt an 64 already-registered provider when another adapter configured first). 65 version: Deployment version stamped on ``gen_ai.agent.version``. 66 Pass ``version=`` to set it; otherwise ``\"unknown\"``. 67 capture_logs: Intercept Python ``logging`` during active sessions and 68 stream to the session Logs tab. Boolean or agent-id glob patterns. 69 Precedence: job metadata → this kwarg → Settings → Logs (default: on). 70 log_level: Minimum level for session log capture (e.g. ``"INFO"``, 71 ``"WARNING"``). Defaults to ``"INFO"``. 72 **kwargs: Framework-specific options (ignored by the shared surface). 73 """ 74 ...
Common parameters every adapter parlotize() must accept.
Framework packages may add extra keyword-only parameters (e.g. LiveKit
record= / auto_escalate_sip=).
1739def _no_init_or_replace_init(self, *args, **kwargs): 1740 cls = type(self) 1741 1742 if cls._is_protocol: 1743 raise TypeError('Protocols cannot be instantiated') 1744 1745 # Already using a custom `__init__`. No need to calculate correct 1746 # `__init__` to call. This can lead to RecursionError. See bpo-45121. 1747 if cls.__init__ is not _no_init_or_replace_init: 1748 return 1749 1750 # Initially, `__init__` of a protocol subclass is set to `_no_init_or_replace_init`. 1751 # The first instantiation of the subclass will call `_no_init_or_replace_init` which 1752 # searches for a proper new `__init__` in the MRO. The new `__init__` 1753 # replaces the subclass' old `__init__` (ie `_no_init_or_replace_init`). Subsequent 1754 # instantiation of the protocol subclass will thus use the new 1755 # `__init__` and no longer call `_no_init_or_replace_init`. 1756 for base in cls.__mro__: 1757 init = base.__dict__.get('__init__', _no_init_or_replace_init) 1758 if init is not _no_init_or_replace_init: 1759 cls.__init__ = init 1760 break 1761 else: 1762 # should not happen 1763 cls.__init__ = object.__init__ 1764 1765 cls.__init__(self, *args, **kwargs)
13@dataclass 14class ParlotContext: 15 """Owns bootstrap runtime and process-local collectors for one SDK instance.""" 16 17 runtime: Optional[ParlotRuntimeContext] = None 18 session_logs: SessionLogsCollector = field(default_factory=SessionLogsCollector) 19 20 def __post_init__(self) -> None: 21 self.session_logs.bind_context(self) 22 23 def shutdown(self) -> None: 24 self.session_logs.shutdown() 25 self.runtime = None
Owns bootstrap runtime and process-local collectors for one SDK instance.
12@dataclass(frozen=True) 13class ParlotRuntimeContext: 14 endpoint: str 15 api_key: str 16 tenant_id: str 17 content_bucket: str 18 r2_endpoint: str 19 recording_globs: tuple[str, ...] = () 20 recording_agents: tuple[tuple[str, bool], ...] = () 21 logs_globs: tuple[str, ...] = () 22 logs_agents: tuple[tuple[str, bool], ...] = () 23 logs_agent_min_levels: tuple[tuple[str, str], ...] = () 24 logs_min_level: str = DEFAULT_LOGS_MIN_LEVEL 25 # True when bootstrap JSON included a ``logs`` object (even empty). 26 logs_policy_present: bool = False 27 capture_genai_content_globs: tuple[str, ...] = () 28 capture_genai_content_agents: tuple[tuple[str, bool], ...] = () 29 # True when bootstrap JSON included a ``capture_genai_content`` object (even empty). 30 capture_genai_content_policy_present: bool = False 31 32 def recording_agents_map(self) -> dict[str, bool]: 33 return dict(self.recording_agents) 34 35 def logs_agents_map(self) -> dict[str, bool]: 36 return dict(self.logs_agents) 37 38 def logs_agent_min_levels_map(self) -> dict[str, str]: 39 return dict(self.logs_agent_min_levels) 40 41 def capture_genai_content_agents_map(self) -> dict[str, bool]: 42 return dict(self.capture_genai_content_agents)
241class SessionLogsCollector: 242 """Per-``ParlotContext`` session log capture pipeline.""" 243 244 def __init__(self) -> None: 245 self._lock = threading.Lock() 246 self._buffer: deque[SessionLogRecord] = deque(maxlen=_MAX_BUFFER) 247 self._consecutive_failures = 0 248 self._circuit_open_until = 0.0 249 self._last_flush_at = 0.0 250 self._endpoint = "" 251 self._api_key = "" 252 self._flush_thread: threading.Thread | None = None 253 self._stop = threading.Event() 254 self._handler: Optional[SessionLogHandler] = None 255 self._installed = False 256 self._atexit_registered = False 257 self._capture_logs_config: CaptureLogsConfig = None 258 self._capture_logs_level: Optional[str] = None 259 self._session_resolver: Optional[SessionResolver] = None 260 self._agent_id_resolver: Optional[Callable[[], str]] = None 261 self._metadata_resolver: Optional[Callable[[], Optional[bool]]] = None 262 self._context_ref: Optional[weakref.ref[ParlotContext]] = None 263 264 def bind_context(self, context: ParlotContext) -> None: 265 self._context_ref = weakref.ref(context) 266 267 def _runtime(self): 268 ctx = self._context_ref() if self._context_ref is not None else None 269 return ctx.runtime if ctx is not None else None 270 271 def set_capture_logs_config( 272 self, 273 capture_logs: CaptureLogsConfig = None, 274 *, 275 log_level: Optional[str] = None, 276 ) -> None: 277 self._capture_logs_config = capture_logs 278 if log_level is not None: 279 self._capture_logs_level = normalize_log_level(log_level) 280 281 def set_resolvers( 282 self, 283 *, 284 session_resolver: Optional[SessionResolver] = None, 285 agent_id_resolver: Optional[Callable[[], str]] = None, 286 metadata_resolver: Optional[Callable[[], Optional[bool]]] = None, 287 ) -> None: 288 """Optional framework hooks (e.g. LiveKit job bootstrap fallback).""" 289 if session_resolver is not None: 290 self._session_resolver = session_resolver 291 if agent_id_resolver is not None: 292 self._agent_id_resolver = agent_id_resolver 293 if metadata_resolver is not None: 294 self._metadata_resolver = metadata_resolver 295 296 def _resolve_min_level(self) -> str: 297 if self._capture_logs_level: 298 return self._capture_logs_level 299 runtime = self._runtime() 300 if runtime is not None: 301 agent_id = "" 302 if self._agent_id_resolver is not None: 303 try: 304 agent_id = self._agent_id_resolver() or "" 305 except Exception: 306 agent_id = "" 307 if agent_id: 308 agent_levels = runtime.logs_agent_min_levels_map() 309 if agent_id in agent_levels: 310 return normalize_log_level(agent_levels[agent_id]) 311 if runtime.logs_min_level: 312 return normalize_log_level(runtime.logs_min_level) 313 return DEFAULT_LOGS_MIN_LEVEL 314 315 def capture_logs_enabled_for_agent(self, agent_name: str = "") -> bool: 316 runtime = self._runtime() 317 metadata = None 318 if self._metadata_resolver is not None: 319 try: 320 metadata = self._metadata_resolver() 321 except Exception: 322 metadata = None 323 name = agent_name 324 if not name and self._agent_id_resolver is not None: 325 try: 326 name = self._agent_id_resolver() or "" 327 except Exception: 328 name = "" 329 return should_capture_logs( 330 name, 331 metadata_capture_logs=metadata, 332 capture_logs_config=self._capture_logs_config, 333 bootstrap_globs=list(runtime.logs_globs) if runtime else None, 334 bootstrap_agents=runtime.logs_agents_map() if runtime else None, 335 bootstrap_present=runtime is not None and runtime.logs_policy_present, 336 ) 337 338 def handle_record( 339 self, handler: logging.Handler, record: logging.LogRecord 340 ) -> None: 341 name = record.name or "" 342 if any(name == p or name.startswith(p + ".") for p in _SKIP_LOGGER_PREFIXES): 343 return 344 level_name = record.levelname or "INFO" 345 if not level_at_least(level_name, self._resolve_min_level()): 346 return 347 if not self.capture_logs_enabled_for_agent(): 348 return 349 session = _resolve_session_fields(self._session_resolver) 350 if not session or not session.get("session_id"): 351 return 352 message = handler.format(record) if handler.formatter else record.getMessage() 353 if not isinstance(message, str): 354 message = str(message) 355 trace_id, span_id = _current_trace_span() 356 ts = datetime.fromtimestamp(record.created, tz=timezone.utc).strftime( 357 "%Y-%m-%d %H:%M:%S.%f" 358 )[:-3] 359 attrs: dict[str, str] = {} 360 if record.pathname: 361 attrs["pathname"] = str(record.pathname)[:512] 362 if record.lineno: 363 attrs["lineno"] = str(record.lineno) 364 if record.funcName: 365 attrs["funcName"] = str(record.funcName)[:128] 366 event = SessionLogRecord( 367 session_id=str(session["session_id"]), 368 conversation_id=str(session.get("conversation_id") or session["session_id"]), 369 ts=ts, 370 level=normalize_log_level(level_name), 371 logger_name=name, 372 message=message, 373 turn_index=int(session.get("turn_index") or 0), 374 trace_id=trace_id, 375 span_id=span_id, 376 attributes=attrs, 377 ) 378 with self._lock: 379 self._buffer.append(event) 380 381 def init(self, *, endpoint: str = "", api_key: str = "") -> None: 382 """Install root handler + start background flusher.""" 383 with self._lock: 384 self._endpoint = (endpoint or "").rstrip("/") 385 self._api_key = api_key or "" 386 if self._handler is None: 387 self._handler = SessionLogHandler(self) 388 self._handler.setLevel(logging.DEBUG) 389 self._handler.setFormatter(logging.Formatter("%(message)s")) 390 if not self._installed: 391 logging.root.addHandler(self._handler) 392 self._installed = True 393 if self._flush_thread is None or not self._flush_thread.is_alive(): 394 self._stop.clear() 395 self._flush_thread = threading.Thread( 396 target=self._flush_loop, 397 name="parlot-session-logs-flush", 398 daemon=True, 399 ) 400 self._flush_thread.start() 401 if not self._atexit_registered: 402 atexit.register(self.shutdown) 403 self._atexit_registered = True 404 405 def drain(self) -> list[SessionLogRecord]: 406 with self._lock: 407 events = list(self._buffer) 408 self._buffer.clear() 409 return events 410 411 def snapshot(self) -> list[SessionLogRecord]: 412 with self._lock: 413 return list(self._buffer) 414 415 def _circuit_open(self) -> bool: 416 return time.time() < self._circuit_open_until 417 418 def _flush_loop(self) -> None: 419 while not self._stop.wait(_FLUSH_INTERVAL_S): 420 try: 421 self._try_flush() 422 except Exception: 423 continue 424 425 def _try_flush(self, *, force: bool = False) -> None: 426 if self._circuit_open(): 427 with self._lock: 428 self._buffer.clear() 429 return 430 if not self._endpoint or not self._api_key: 431 return 432 433 if not force: 434 now = time.time() 435 if now - self._last_flush_at < _FLUSH_INTERVAL_S * 0.5: 436 return 437 438 events = self.drain() 439 if not events: 440 return 441 442 ok = self._post_events(events) 443 self._last_flush_at = time.time() 444 if ok: 445 self._consecutive_failures = 0 446 return 447 448 self._consecutive_failures += 1 449 if self._consecutive_failures >= _CIRCUIT_FAILURES: 450 self._circuit_open_until = time.time() + _CIRCUIT_COOLDOWN_S 451 self._consecutive_failures = 0 452 453 def _post_events(self, events: list[SessionLogRecord]) -> bool: 454 try: 455 import httpx 456 except Exception: 457 return False 458 459 url = f"{self._endpoint}/v1/logs" 460 try: 461 payload = _encode_otlp_logs(events[:_MAX_BATCH]) 462 except Exception: 463 return False 464 if not payload: 465 return True 466 from parlot.core.provider import build_parlot_client_headers 467 468 headers = build_parlot_client_headers(self._api_key) 469 headers["Content-Type"] = "application/x-protobuf" 470 try: 471 with httpx.Client(timeout=3.0) as client: 472 res = client.post( 473 url, 474 content=payload, 475 headers=headers, 476 ) 477 return 200 <= res.status_code < 300 478 except Exception: 479 return False 480 481 def shutdown(self) -> None: 482 """Stop capture and push any remaining buffered logs. Never raises.""" 483 try: 484 self._stop.set() 485 thread = self._flush_thread 486 if ( 487 thread is not None 488 and thread.is_alive() 489 and thread is not threading.current_thread() 490 ): 491 thread.join(timeout=1.0) 492 self._flush_thread = None 493 with self._lock: 494 if self._installed and self._handler is not None: 495 try: 496 logging.root.removeHandler(self._handler) 497 except Exception: 498 pass 499 self._installed = False 500 self._try_flush(force=True) 501 except Exception: 502 return
Per-ParlotContext session log capture pipeline.
281 def set_resolvers( 282 self, 283 *, 284 session_resolver: Optional[SessionResolver] = None, 285 agent_id_resolver: Optional[Callable[[], str]] = None, 286 metadata_resolver: Optional[Callable[[], Optional[bool]]] = None, 287 ) -> None: 288 """Optional framework hooks (e.g. LiveKit job bootstrap fallback).""" 289 if session_resolver is not None: 290 self._session_resolver = session_resolver 291 if agent_id_resolver is not None: 292 self._agent_id_resolver = agent_id_resolver 293 if metadata_resolver is not None: 294 self._metadata_resolver = metadata_resolver
Optional framework hooks (e.g. LiveKit job bootstrap fallback).
315 def capture_logs_enabled_for_agent(self, agent_name: str = "") -> bool: 316 runtime = self._runtime() 317 metadata = None 318 if self._metadata_resolver is not None: 319 try: 320 metadata = self._metadata_resolver() 321 except Exception: 322 metadata = None 323 name = agent_name 324 if not name and self._agent_id_resolver is not None: 325 try: 326 name = self._agent_id_resolver() or "" 327 except Exception: 328 name = "" 329 return should_capture_logs( 330 name, 331 metadata_capture_logs=metadata, 332 capture_logs_config=self._capture_logs_config, 333 bootstrap_globs=list(runtime.logs_globs) if runtime else None, 334 bootstrap_agents=runtime.logs_agents_map() if runtime else None, 335 bootstrap_present=runtime is not None and runtime.logs_policy_present, 336 )
338 def handle_record( 339 self, handler: logging.Handler, record: logging.LogRecord 340 ) -> None: 341 name = record.name or "" 342 if any(name == p or name.startswith(p + ".") for p in _SKIP_LOGGER_PREFIXES): 343 return 344 level_name = record.levelname or "INFO" 345 if not level_at_least(level_name, self._resolve_min_level()): 346 return 347 if not self.capture_logs_enabled_for_agent(): 348 return 349 session = _resolve_session_fields(self._session_resolver) 350 if not session or not session.get("session_id"): 351 return 352 message = handler.format(record) if handler.formatter else record.getMessage() 353 if not isinstance(message, str): 354 message = str(message) 355 trace_id, span_id = _current_trace_span() 356 ts = datetime.fromtimestamp(record.created, tz=timezone.utc).strftime( 357 "%Y-%m-%d %H:%M:%S.%f" 358 )[:-3] 359 attrs: dict[str, str] = {} 360 if record.pathname: 361 attrs["pathname"] = str(record.pathname)[:512] 362 if record.lineno: 363 attrs["lineno"] = str(record.lineno) 364 if record.funcName: 365 attrs["funcName"] = str(record.funcName)[:128] 366 event = SessionLogRecord( 367 session_id=str(session["session_id"]), 368 conversation_id=str(session.get("conversation_id") or session["session_id"]), 369 ts=ts, 370 level=normalize_log_level(level_name), 371 logger_name=name, 372 message=message, 373 turn_index=int(session.get("turn_index") or 0), 374 trace_id=trace_id, 375 span_id=span_id, 376 attributes=attrs, 377 ) 378 with self._lock: 379 self._buffer.append(event)
381 def init(self, *, endpoint: str = "", api_key: str = "") -> None: 382 """Install root handler + start background flusher.""" 383 with self._lock: 384 self._endpoint = (endpoint or "").rstrip("/") 385 self._api_key = api_key or "" 386 if self._handler is None: 387 self._handler = SessionLogHandler(self) 388 self._handler.setLevel(logging.DEBUG) 389 self._handler.setFormatter(logging.Formatter("%(message)s")) 390 if not self._installed: 391 logging.root.addHandler(self._handler) 392 self._installed = True 393 if self._flush_thread is None or not self._flush_thread.is_alive(): 394 self._stop.clear() 395 self._flush_thread = threading.Thread( 396 target=self._flush_loop, 397 name="parlot-session-logs-flush", 398 daemon=True, 399 ) 400 self._flush_thread.start() 401 if not self._atexit_registered: 402 atexit.register(self.shutdown) 403 self._atexit_registered = True
Install root handler + start background flusher.
481 def shutdown(self) -> None: 482 """Stop capture and push any remaining buffered logs. Never raises.""" 483 try: 484 self._stop.set() 485 thread = self._flush_thread 486 if ( 487 thread is not None 488 and thread.is_alive() 489 and thread is not threading.current_thread() 490 ): 491 thread.join(timeout=1.0) 492 self._flush_thread = None 493 with self._lock: 494 if self._installed and self._handler is not None: 495 try: 496 logging.root.removeHandler(self._handler) 497 except Exception: 498 pass 499 self._installed = False 500 self._try_flush(force=True) 501 except Exception: 502 return
Stop capture and push any remaining buffered logs. Never raises.
100def base_parlotize( 101 agent_id: str, 102 *, 103 endpoint: Optional[str] = None, 104 api_key: Optional[str] = None, 105 capture_genai_content: Optional[bool] = None, 106 tracer_provider: TracerProvider | None = None, 107 version: Optional[str] = None, 108 capture_logs: bool | list[str] | None = None, 109 log_level: Optional[str] = None, 110) -> BaseParlotizeResult: 111 """Execute shared telemetry configuration common across all adapters. 112 113 Creates a ``ParlotContext``, initializes the session log 114 collector, and fetches remote bootstrap into ``context.runtime``. 115 116 See ``ParlotizeProtocol`` for the shared parameter surface. 117 118 Args: 119 agent_id: Required canonical deployment identity for ``session.agent_id``. 120 endpoint: Parlot OTLP collector base URL. If omitted, reads 121 ``PARLOT_ENDPOINT``. 122 api_key: Org-scoped API key. If omitted, reads ``PARLOT_API_KEY``. 123 capture_genai_content: Process-wide GenAI content capture override. 124 tracer_provider: Existing ``TracerProvider`` to adopt, if any. 125 version: Deployment version for ``gen_ai.agent.version``. Defaults to 126 ``\"unknown\"`` when omitted or blank. 127 capture_logs: Session log capture policy (bool or agent-id globs). 128 log_level: Minimum level for session log capture. 129 130 Returns: 131 Resolved endpoint, credentials, provider, and capture settings for the 132 calling adapter. 133 """ 134 from parlot.core.bootstrap import fetch_telemetry_bootstrap 135 from parlot.core.context import ParlotContext 136 from parlot.core.provider import ( 137 adopt_existing_tracer_provider, 138 resolve_api_key, 139 resolve_capture_genai_content, 140 resolve_endpoint, 141 ) 142 143 configure_parlot_logging() 144 145 resolved_agent_id = agent_id.strip() 146 if not resolved_agent_id: 147 raise ValueError("parlotize(agent_id) is required") 148 resolved_version = (version or "").strip() or "unknown" 149 resolved_endpoint = resolve_endpoint(endpoint) 150 resolved_api_key = resolve_api_key(api_key) 151 resolved_capture_content = resolve_capture_genai_content(capture_genai_content) 152 153 if isinstance(capture_logs, list): 154 resolved_capture_logs: bool | list[str] | None = [ 155 str(item).strip() for item in capture_logs if str(item).strip() 156 ] 157 else: 158 resolved_capture_logs = capture_logs 159 160 resolved_log_level = log_level.strip().upper() if log_level else None 161 162 context = ParlotContext() 163 context.session_logs.set_capture_logs_config( 164 resolved_capture_logs, 165 log_level=resolved_log_level, 166 ) 167 context.session_logs.init(endpoint=resolved_endpoint, api_key=resolved_api_key) 168 169 if resolved_api_key: 170 fetch_telemetry_bootstrap(resolved_endpoint, resolved_api_key, context) 171 172 if tracer_provider is None: 173 tracer_provider = adopt_existing_tracer_provider() 174 175 return BaseParlotizeResult( 176 context=context, 177 endpoint=resolved_endpoint, 178 api_key=resolved_api_key, 179 tracer_provider=tracer_provider, 180 agent_id=resolved_agent_id, 181 agent_version=resolved_version, 182 capture_genai_content=resolved_capture_content, 183 capture_logs=resolved_capture_logs, 184 log_level=resolved_log_level, 185 )
Execute shared telemetry configuration common across all adapters.
Creates a ParlotContext, initializes the session log
collector, and fetches remote bootstrap into context.runtime.
See ParlotizeProtocol for the shared parameter surface.
Arguments:
- agent_id: Required canonical deployment identity for
session.agent_id. - endpoint: Parlot OTLP collector base URL. If omitted, reads
PARLOT_ENDPOINT. - api_key: Org-scoped API key. If omitted, reads
PARLOT_API_KEY. - capture_genai_content: Process-wide GenAI content capture override.
- tracer_provider: Existing
TracerProviderto adopt, if any. - version: Deployment version for
gen_ai.agent.version. Defaults to"unknown"when omitted or blank. - capture_logs: Session log capture policy (bool or agent-id globs).
- log_level: Minimum level for session log capture.
Returns:
Resolved endpoint, credentials, provider, and capture settings for the calling adapter.
92def configure_parlot_logging() -> None: 93 """Configure root logging from ``PARLOT_DEBUG_LEVEL`` when no handlers exist.""" 94 level_name = os.getenv("PARLOT_DEBUG_LEVEL", "INFO").upper() 95 level = getattr(logging, level_name, logging.INFO) 96 if not logging.root.handlers: 97 logging.basicConfig(level=level)
Configure root logging from PARLOT_DEBUG_LEVEL when no handlers exist.
14class ExportFilterSpanExporter(SpanExporter): 15 """Pass through only Conversation Contract + GenAI + voice spans.""" 16 17 def __init__(self, exporter: SpanExporter) -> None: 18 self._exporter = exporter 19 20 def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: 21 filtered = [span for span in spans if is_exportable_span_name(span.name)] 22 if not filtered: 23 return SpanExportResult.SUCCESS 24 return self._exporter.export(filtered) 25 26 def shutdown(self) -> None: 27 self._exporter.shutdown() 28 29 def force_flush(self, timeout_millis: int = 30000) -> bool: 30 return self._exporter.force_flush(timeout_millis)
Pass through only Conversation Contract + GenAI + voice spans.
20 def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: 21 filtered = [span for span in spans if is_exportable_span_name(span.name)] 22 if not filtered: 23 return SpanExportResult.SUCCESS 24 return self._exporter.export(filtered)
Exports a batch of telemetry data.
Arguments:
- spans: The list of
opentelemetry.trace.Spanobjects to be exported
Returns:
The result of the export
29 def force_flush(self, timeout_millis: int = 30000) -> bool: 30 return self._exporter.force_flush(timeout_millis)
Hint to ensure that the export of any spans the exporter has received prior to the call to ForceFlush SHOULD be completed as soon as possible, preferably before returning from this method.
24class ParlotBaseProcessor(SpanProcessor): 25 """ 26 Subclass and override ``on_end`` (or the ``_enrich`` dispatch method) to 27 add framework-specific span enrichment. Call ``super().on_start()`` / 28 ``super().on_end()`` so shared debug logging runs. 29 30 Do NOT wrap a downstream processor — add alongside BatchSpanProcessor: 31 32 provider.add_span_processor(MyFrameworkProcessor()) 33 provider.add_span_processor(BatchSpanProcessor(your_exporter)) 34 """ 35 36 def on_start(self, span, parent_context=None) -> None: 37 log_span_event("on_start", span) 38 39 def on_end(self, span: ReadableSpan) -> None: 40 log_span_event("on_end", span) 41 42 def shutdown(self) -> None: 43 pass 44 45 def force_flush(self, timeout_millis: int = 30_000) -> bool: 46 return True 47 48 # ------------------------------------------------------------------ 49 # Shared helpers — used by all subclasses 50 # ------------------------------------------------------------------ 51 52 @staticmethod 53 def _set(span: ReadableSpan, key: str, value) -> None: 54 """Write an attribute into a ReadableSpan after it has ended. 55 56 OTel SDK ≥1.44 stores attributes in an immutable ``BoundedAttributes`` 57 by default; mutating it raises ``TypeError``. Copy to a plain dict 58 when needed so enrichment can still stamp session/turn fields. 59 """ 60 attrs = span._attributes 61 if attrs is None: 62 span._attributes = {key: value} 63 return 64 if getattr(attrs, "_immutable", False) or not isinstance(attrs, dict): 65 span._attributes = dict(attrs) 66 span._attributes[key] = value 67 return 68 try: 69 attrs[key] = value 70 except TypeError: 71 span._attributes = dict(attrs) 72 span._attributes[key] = value 73 74 @staticmethod 75 def _add_event(span: ReadableSpan, name: str, attributes: dict) -> None: 76 """Append a span event to a ReadableSpan after it has ended.""" 77 from opentelemetry.sdk.trace import Event 78 evt = Event(name=name, attributes=attributes, timestamp=time.time_ns()) 79 if hasattr(span, "_events") and isinstance(span._events, list): 80 span._events.append(evt) 81 82 @staticmethod 83 def _trace_id_hex(span: ReadableSpan) -> str: 84 return format(span.context.trace_id, "032x") 85 86 @staticmethod 87 def _maybe_update(state: SessionState, attr: str, value) -> None: 88 """Set a state attribute only if it is currently falsy.""" 89 if value and not getattr(state, attr, ""): 90 setattr(state, attr, str(value))
Subclass and override on_end (or the _enrich dispatch method) to
add framework-specific span enrichment. Call super().on_start() /
super().on_end() so shared debug logging runs.
Do NOT wrap a downstream processor — add alongside BatchSpanProcessor:
provider.add_span_processor(MyFrameworkProcessor())
provider.add_span_processor(BatchSpanProcessor(your_exporter))
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.
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.
93def assert_sync_span_processors(provider) -> bool: 94 """Require ``SynchronousMultiSpanProcessor`` on the provider. 95 96 Parlot processors rely on synchronous ``on_start`` / ``on_end`` ordering and 97 often on ``contextvars`` copied at task creation. ``ConcurrentMultiSpanProcessor`` 98 runs handlers on worker threads and breaks that model. 99 100 Returns ``False`` when ``ConcurrentMultiSpanProcessor`` is active (callers 101 should disable OTel ``context.attach`` and similar). Raises ``RuntimeError`` 102 for unknown processor layouts. 103 104 Uses OTel SDK-private ``_active_span_processor`` — re-validate on SDK upgrades. 105 """ 106 from opentelemetry.sdk.trace import ( 107 ConcurrentMultiSpanProcessor, 108 SynchronousMultiSpanProcessor, 109 ) 110 111 asp = getattr(provider, "_active_span_processor", None) 112 if asp is None: 113 raise RuntimeError( 114 "parlot: TracerProvider has no _active_span_processor; " 115 "cannot verify synchronous span processing" 116 ) 117 if isinstance(asp, ConcurrentMultiSpanProcessor): 118 logger.error( 119 "parlot: ConcurrentMultiSpanProcessor detected — " 120 "unsupported for Parlot span processors (use synchronous layout)" 121 ) 122 return False 123 if not isinstance(asp, SynchronousMultiSpanProcessor): 124 raise RuntimeError( 125 f"parlot: unsupported active span processor {type(asp).__name__!r}; " 126 "expected SynchronousMultiSpanProcessor" 127 ) 128 return True
Require SynchronousMultiSpanProcessor on the provider.
Parlot processors rely on synchronous on_start / on_end ordering and
often on contextvars copied at task creation. ConcurrentMultiSpanProcessor
runs handlers on worker threads and breaks that model.
Returns False when ConcurrentMultiSpanProcessor is active (callers
should disable OTel context.attach and similar). Raises RuntimeError
for unknown processor layouts.
Uses OTel SDK-private _active_span_processor — re-validate on SDK upgrades.
20@dataclass 21class SessionState: 22 """ 23 Per-trace accumulator for cross-span aggregates. 24 25 Framework-specific instrumentation packages subclass this to add their own 26 fields (e.g. LiveKit adds agent_chain, pending_handoff_end_ns). 27 """ 28 29 session_id: str = "" 30 room_name: str = "" 31 room_sid: str = "" 32 agent_label: str = "" 33 34 turn_count: int = 0 35 tool_call_count: int = 0 36 handoff_count: int = 0 37 38 total_input_tokens: int = 0 39 total_output_tokens: int = 0 40 total_cost_usd: float = 0.0 41 42 human_rep_participant_ids: set[str] = field(default_factory=set) 43 topology_agents: list[dict[str, Any]] = field(default_factory=list) 44 custom_metadata: dict[str, str] = field(default_factory=dict) 45 # Owning instrumentation package, e.g. "livekit" | "langgraph" 46 framework: str = ""
Per-trace accumulator for cross-span aggregates.
Framework-specific instrumentation packages subclass this to add their own fields (e.g. LiveKit adds agent_chain, pending_handoff_end_ns).
51@dataclass 52class SessionTopology: 53 """ 54 Accumulates intent segments and agent instructions during a live session. 55 56 Framework adapters call ``record_instructions`` / ``open_segment_after_handoff`` 57 during the session and ``stamp_session_span`` at close. 58 """ 59 60 prompts_by_agent: dict[str, str] = field(default_factory=dict) 61 intent_by_agent: dict[str, dict[str, str]] = field(default_factory=dict) 62 agents_seen: dict[str, dict[str, Any]] = field(default_factory=dict) 63 intent_segments: list[_MutableSegment] = field(default_factory=list) 64 active_segment: _MutableSegment | None = None 65 pending_segment_from_turn: int | None = None 66 first_agent_label: str = "" 67 agent_chain: list[str] = field(default_factory=list) 68 handoff_count: int = 0 69 _handoff_keys: set[str] = field(default_factory=set) 70 71 def upsert_agent(self, agent_id: str) -> None: 72 aid = agent_id.strip() 73 if not aid: 74 return 75 if not self.first_agent_label: 76 self.first_agent_label = aid 77 if aid not in self.agents_seen: 78 self.agents_seen[aid] = {"id": aid, "role": "agent"} 79 80 def record_instructions(self, agent_id: str, text: str) -> None: 81 aid = agent_id.strip() 82 if not aid or not text.strip(): 83 return 84 excerpt = _truncate(text.strip(), INSTRUCTIONS_PREVIEW) 85 prev = self.prompts_by_agent.get(aid, "") 86 if len(excerpt) >= len(prev): 87 self.prompts_by_agent[aid] = excerpt 88 self._sync_agent_intent_summary(aid) 89 self._refresh_open_segment_instructions(aid) 90 91 def push_agent_chain(self, agent_id: str) -> None: 92 aid = agent_id.strip() 93 if not aid: 94 return 95 if not self.agent_chain or self.agent_chain[-1] != aid: 96 self.agent_chain.append(aid) 97 98 def open_bootstrap_segment(self, agent_id: str, from_turn: int) -> None: 99 if self.active_segment is not None: 100 return 101 aid = agent_id.strip() or "unknown" 102 self.upsert_agent(aid) 103 self.active_segment = self._new_segment(aid, from_turn, 0) 104 105 def open_segment_after_handoff( 106 self, 107 agent_id: str, 108 *, 109 handoff_index: int | None = None, 110 turn_index: int, 111 from_agent: str = "", 112 ) -> None: 113 aid = agent_id.strip() 114 if not aid: 115 return 116 dedupe_key = f"{from_agent.strip()}->{aid}@{turn_index}" 117 if dedupe_key in self._handoff_keys: 118 return 119 self._handoff_keys.add(dedupe_key) 120 121 self.close_active_segment(turn_index) 122 self.upsert_agent(aid) 123 self.handoff_count += 1 124 effective_handoff = handoff_index if handoff_index is not None else self.handoff_count 125 if not self.intent_segments and self.active_segment is None: 126 effective_handoff = 0 127 self.pending_segment_from_turn = max(1, turn_index + 1) 128 self.active_segment = self._new_segment( 129 aid, 130 self.pending_segment_from_turn, 131 effective_handoff, 132 ) 133 134 def close_active_segment(self, to_turn: int) -> None: 135 if self.active_segment is None: 136 return 137 self.active_segment.to_turn = max(0, to_turn) 138 self.intent_segments.append(self.active_segment) 139 self.active_segment = None 140 141 def apply_pending_from_turn_on_emit(self, turn_index: int) -> None: 142 if self.active_segment is None: 143 return 144 if self.pending_segment_from_turn is not None: 145 self.active_segment.from_turn = self.pending_segment_from_turn 146 self.pending_segment_from_turn = None 147 elif self.active_segment.from_turn <= 0: 148 self.active_segment.from_turn = turn_index 149 150 def finalize_intent_sequence(self, final_turn: int) -> list[dict[str, Any]]: 151 if self.active_segment is not None: 152 self.close_active_segment(final_turn) 153 return self._collapse_segments( 154 [ 155 { 156 "segment_index": idx, 157 "agent_id": seg.agent_id, 158 "intent_key": seg.intent_key, 159 "intent_label": seg.intent_label, 160 "from_turn": seg.from_turn, 161 "handoff_index": seg.handoff_index, 162 "source": seg.source, 163 **( 164 {"to_turn": seg.to_turn} 165 if seg.to_turn is not None 166 else {} 167 ), 168 **( 169 {"instructions_excerpt": seg.instructions_excerpt} 170 if seg.instructions_excerpt 171 else {} 172 ), 173 } 174 for idx, seg in enumerate(self.intent_segments) 175 ] 176 ) 177 178 def bootstrap_instructions(self) -> str: 179 if self.first_agent_label: 180 text = self.prompts_by_agent.get(self.first_agent_label, "") 181 if text: 182 return text 183 if self.prompts_by_agent: 184 return max(self.prompts_by_agent.values(), key=len) 185 return "" 186 187 def agents_json(self) -> list[dict[str, Any]]: 188 out: list[dict[str, Any]] = [] 189 for node in self.agents_seen.values(): 190 summary = self.intent_by_agent.get(node["id"], {}) 191 entry = dict(node) 192 if summary: 193 entry["intent_key"] = summary.get("intent_key", node["id"]) 194 entry["intent_label"] = summary.get("intent_label", node["id"]) 195 excerpt = summary.get("instructions_excerpt", "") 196 if excerpt: 197 entry["instructions_excerpt"] = excerpt 198 out.append(entry) 199 return out 200 201 def stamp_session_span(self, span: Any, *, final_turn: int) -> None: 202 """Write topology attrs on the live ``parlot.session`` span.""" 203 if span is None or not hasattr(span, "set_attribute"): 204 return 205 sequence = self.finalize_intent_sequence(final_turn) 206 agents = self.agents_json() 207 bootstrap = self.bootstrap_instructions() 208 if agents: 209 span.set_attribute(ATTR_SESSION_TOPOLOGY_AGENTS, _json_dumps_cap(agents)) 210 if bootstrap: 211 span.set_attribute(ATTR_SESSION_TOPOLOGY_BOOTSTRAP_INSTRUCTIONS, bootstrap) 212 if sequence: 213 span.set_attribute(ATTR_SESSION_INTENT_SEQUENCE, _json_dumps_cap(sequence)) 214 if self.agent_chain: 215 span.set_attribute(ATTR_SESSION_AGENT_CHAIN, " → ".join(self.agent_chain)) 216 217 def _new_segment(self, agent_id: str, from_turn: int, handoff_index: int) -> _MutableSegment: 218 instructions = self.prompts_by_agent.get(agent_id, "") 219 derived = derive_intent(agent_id, instructions) 220 self._sync_agent_intent_summary(agent_id) 221 return _MutableSegment( 222 segment_index=0, 223 agent_id=agent_id, 224 intent_key=derived["intent_key"], 225 intent_label=derived["intent_label"], 226 from_turn=from_turn, 227 handoff_index=handoff_index, 228 instructions_excerpt=derived.get("instructions_excerpt", ""), 229 ) 230 231 def _sync_agent_intent_summary(self, agent_id: str) -> None: 232 instructions = self.prompts_by_agent.get(agent_id, "") 233 derived = derive_intent(agent_id, instructions) 234 self.intent_by_agent[agent_id] = derived 235 node = self.agents_seen.get(agent_id) 236 if node is not None: 237 node["intent_key"] = derived["intent_key"] 238 node["intent_label"] = derived["intent_label"] 239 if derived.get("instructions_excerpt"): 240 node["instructions_excerpt"] = derived["instructions_excerpt"] 241 242 def _refresh_open_segment_instructions(self, agent_id: str) -> None: 243 if self.active_segment is None or self.active_segment.agent_id != agent_id.strip(): 244 return 245 instructions = self.prompts_by_agent.get(agent_id, "") 246 derived = derive_intent(agent_id, instructions) 247 self.active_segment.intent_key = derived["intent_key"] 248 self.active_segment.intent_label = derived["intent_label"] 249 self.active_segment.instructions_excerpt = derived.get("instructions_excerpt", "") 250 self._sync_agent_intent_summary(agent_id) 251 252 @staticmethod 253 def _collapse_segments(segments: list[dict[str, Any]]) -> list[dict[str, Any]]: 254 if not segments: 255 return [] 256 merged: list[dict[str, Any]] = [] 257 for seg in segments: 258 if not merged: 259 merged.append(dict(seg)) 260 continue 261 prev = merged[-1] 262 if prev.get("agent_id") == seg.get("agent_id"): 263 prev["to_turn"] = seg.get("to_turn", prev.get("to_turn")) 264 if seg.get("instructions_excerpt") and not prev.get("instructions_excerpt"): 265 prev["instructions_excerpt"] = seg["instructions_excerpt"] 266 continue 267 merged.append(dict(seg)) 268 for idx, seg in enumerate(merged): 269 seg["segment_index"] = idx 270 return merged
Accumulates intent segments and agent instructions during a live session.
Framework adapters call record_instructions / open_segment_after_handoff
during the session and stamp_session_span at close.
80 def record_instructions(self, agent_id: str, text: str) -> None: 81 aid = agent_id.strip() 82 if not aid or not text.strip(): 83 return 84 excerpt = _truncate(text.strip(), INSTRUCTIONS_PREVIEW) 85 prev = self.prompts_by_agent.get(aid, "") 86 if len(excerpt) >= len(prev): 87 self.prompts_by_agent[aid] = excerpt 88 self._sync_agent_intent_summary(aid) 89 self._refresh_open_segment_instructions(aid)
105 def open_segment_after_handoff( 106 self, 107 agent_id: str, 108 *, 109 handoff_index: int | None = None, 110 turn_index: int, 111 from_agent: str = "", 112 ) -> None: 113 aid = agent_id.strip() 114 if not aid: 115 return 116 dedupe_key = f"{from_agent.strip()}->{aid}@{turn_index}" 117 if dedupe_key in self._handoff_keys: 118 return 119 self._handoff_keys.add(dedupe_key) 120 121 self.close_active_segment(turn_index) 122 self.upsert_agent(aid) 123 self.handoff_count += 1 124 effective_handoff = handoff_index if handoff_index is not None else self.handoff_count 125 if not self.intent_segments and self.active_segment is None: 126 effective_handoff = 0 127 self.pending_segment_from_turn = max(1, turn_index + 1) 128 self.active_segment = self._new_segment( 129 aid, 130 self.pending_segment_from_turn, 131 effective_handoff, 132 )
141 def apply_pending_from_turn_on_emit(self, turn_index: int) -> None: 142 if self.active_segment is None: 143 return 144 if self.pending_segment_from_turn is not None: 145 self.active_segment.from_turn = self.pending_segment_from_turn 146 self.pending_segment_from_turn = None 147 elif self.active_segment.from_turn <= 0: 148 self.active_segment.from_turn = turn_index
150 def finalize_intent_sequence(self, final_turn: int) -> list[dict[str, Any]]: 151 if self.active_segment is not None: 152 self.close_active_segment(final_turn) 153 return self._collapse_segments( 154 [ 155 { 156 "segment_index": idx, 157 "agent_id": seg.agent_id, 158 "intent_key": seg.intent_key, 159 "intent_label": seg.intent_label, 160 "from_turn": seg.from_turn, 161 "handoff_index": seg.handoff_index, 162 "source": seg.source, 163 **( 164 {"to_turn": seg.to_turn} 165 if seg.to_turn is not None 166 else {} 167 ), 168 **( 169 {"instructions_excerpt": seg.instructions_excerpt} 170 if seg.instructions_excerpt 171 else {} 172 ), 173 } 174 for idx, seg in enumerate(self.intent_segments) 175 ] 176 )
187 def agents_json(self) -> list[dict[str, Any]]: 188 out: list[dict[str, Any]] = [] 189 for node in self.agents_seen.values(): 190 summary = self.intent_by_agent.get(node["id"], {}) 191 entry = dict(node) 192 if summary: 193 entry["intent_key"] = summary.get("intent_key", node["id"]) 194 entry["intent_label"] = summary.get("intent_label", node["id"]) 195 excerpt = summary.get("instructions_excerpt", "") 196 if excerpt: 197 entry["instructions_excerpt"] = excerpt 198 out.append(entry) 199 return out
201 def stamp_session_span(self, span: Any, *, final_turn: int) -> None: 202 """Write topology attrs on the live ``parlot.session`` span.""" 203 if span is None or not hasattr(span, "set_attribute"): 204 return 205 sequence = self.finalize_intent_sequence(final_turn) 206 agents = self.agents_json() 207 bootstrap = self.bootstrap_instructions() 208 if agents: 209 span.set_attribute(ATTR_SESSION_TOPOLOGY_AGENTS, _json_dumps_cap(agents)) 210 if bootstrap: 211 span.set_attribute(ATTR_SESSION_TOPOLOGY_BOOTSTRAP_INSTRUCTIONS, bootstrap) 212 if sequence: 213 span.set_attribute(ATTR_SESSION_INTENT_SEQUENCE, _json_dumps_cap(sequence)) 214 if self.agent_chain: 215 span.set_attribute(ATTR_SESSION_AGENT_CHAIN, " → ".join(self.agent_chain))
Write topology attrs on the live parlot.session span.
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")
124def adopt_existing_tracer_provider() -> Any | None: 125 """Return the global tracer provider when already set by another Parlot package. 126 127 Returns ``None`` when only the default no-op / proxy provider is installed. 128 """ 129 from opentelemetry import trace 130 from opentelemetry.sdk.trace import TracerProvider as SdkTracerProvider 131 132 provider = trace.get_tracer_provider() 133 if isinstance(provider, SdkTracerProvider): 134 return provider 135 return None
Return the global tracer provider when already set by another Parlot package.
Returns None when only the default no-op / proxy provider is installed.
71def build_otlp_http_exporter( 72 *, 73 endpoint: str, 74 api_key: str = "", 75) -> SpanExporter: 76 """Build an OTLP/HTTP span exporter targeting ``{endpoint}/v1/traces``.""" 77 from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter 78 79 if not endpoint: 80 raise ValueError( 81 "No OTLP endpoint configured. Pass endpoint= or set the " 82 "PARLOT_ENDPOINT environment variable." 83 ) 84 headers = build_parlot_client_headers(api_key) 85 trace_endpoint = endpoint.rstrip("/") + "/v1/traces" 86 return OTLPSpanExporter(endpoint=trace_endpoint, headers=headers)
Build an OTLP/HTTP span exporter targeting {endpoint}/v1/traces.
56def build_parlot_client_headers(api_key: str = "") -> dict[str, str]: 57 """Build standard transport headers sent with Parlot telemetry exports.""" 58 from parlot.core.sdk_version import resolve_parlot_sdk_version 59 60 version = resolve_parlot_sdk_version() or "unknown" 61 headers: dict[str, str] = { 62 HEADER_SDK_NAME: "parlot-python", 63 HEADER_SDK_VERSION: version, 64 HEADER_INGESTION_VERSION: INGESTION_PROTOCOL_VERSION, 65 } 66 if api_key: 67 headers["Authorization"] = f"Bearer {api_key}" 68 return headers
Build standard transport headers sent with Parlot telemetry exports.
39def build_resource( 40 *, 41 service_name: Optional[str] = None, 42 service_version: Optional[str] = None, 43) -> Resource: 44 attrs = {service_attributes.SERVICE_NAME: service_name or "unknown"} 45 if service_version: 46 attrs[service_attributes.SERVICE_VERSION] = service_version 47 return Resource.create(attrs)
89def build_tracer_provider( 90 *, 91 endpoint: Optional[str] = None, 92 api_key: Optional[str] = None, 93 service_name: Optional[str] = None, 94 service_version: Optional[str] = None, 95 span_processors: Sequence[SpanProcessor] = (), 96 span_exporter: SpanExporter | None = None, 97) -> TracerProvider: 98 """Build a ``TracerProvider`` with optional processors and OTLP export. 99 100 When ``span_exporter`` is omitted, builds a default OTLP/HTTP exporter to 101 ``{endpoint}/v1/traces``. Callers typically wrap that exporter (filter, 102 sanitize, quiet) before passing it here. 103 """ 104 resolved_endpoint = resolve_endpoint(endpoint) 105 resolved_api_key = resolve_api_key(api_key) 106 resource = build_resource( 107 service_name=service_name, 108 service_version=service_version, 109 ) 110 exporter = span_exporter 111 if exporter is None: 112 exporter = build_otlp_http_exporter( 113 endpoint=resolved_endpoint, 114 api_key=resolved_api_key, 115 ) 116 117 provider = TracerProvider(resource=resource) 118 for processor in span_processors: 119 provider.add_span_processor(processor) 120 provider.add_span_processor(BatchSpanProcessor(exporter)) 121 return provider
Build a TracerProvider with optional processors and OTLP export.
When span_exporter is omitted, builds a default OTLP/HTTP exporter to
{endpoint}/v1/traces. Callers typically wrap that exporter (filter,
sanitize, quiet) before passing it here.
68def clear_active_session() -> None: 69 """Clear the active session bindings for this context.""" 70 set_active_session(None, None)
Clear the active session bindings for this context.
36def derive_intent(agent_id: str, instructions: str = "") -> dict[str, str]: 37 """Port of platform ingest ``deriveIntent``.""" 38 aid = agent_id.strip() 39 intent_key = _normalize_key(aid) if aid else "unknown" 40 excerpt = _truncate(instructions.strip(), INSTRUCTIONS_EXCERPT_MAX) if instructions.strip() else "" 41 42 if excerpt: 43 first = _first_line(excerpt) 44 intent_label = _truncate(first, INTENT_LABEL_MAX) if first else _title_from_id(intent_key) 45 else: 46 intent_label = _title_from_id(intent_key) 47 48 return { 49 "intent_key": intent_key, 50 "intent_label": intent_label, 51 "instructions_excerpt": excerpt, 52 }
Port of platform ingest deriveIntent.
15def fetch_telemetry_bootstrap( 16 endpoint: str, 17 api_key: str, 18 context: ParlotContext, 19 *, 20 timeout: float = 15.0, 21) -> Optional[dict[str, Any]]: 22 """GET ``/v1/telemetry/bootstrap`` and store runtime on ``context``, or None.""" 23 if not endpoint or not api_key: 24 return None 25 26 import httpx 27 28 from parlot.core.runtime import runtime_from_bootstrap 29 30 url = f"{endpoint.rstrip('/')}/v1/telemetry/bootstrap" 31 headers = {"Authorization": f"Bearer {api_key}"} 32 try: 33 with httpx.Client(timeout=timeout) as client: 34 resp = client.get(url, headers=headers) 35 if resp.status_code >= 400: 36 logger.error( 37 "parlot: telemetry bootstrap failed status=%s", 38 resp.status_code, 39 ) 40 return None 41 payload = resp.json() 42 if not isinstance(payload, dict): 43 logger.error("parlot: telemetry bootstrap returned non-object JSON") 44 return None 45 context.runtime = runtime_from_bootstrap(endpoint, api_key, payload) 46 return payload 47 except Exception as exc: 48 logger.error("parlot: telemetry bootstrap request failed — %s", exc) 49 logger.debug("parlot: telemetry bootstrap request failed", exc_info=True) 50 return None
GET /v1/telemetry/bootstrap and store runtime on context, or None.
49def get_active_session() -> SessionState | None: 50 """Return the session state bound to the current context, if any.""" 51 return _active_session_state.get()
Return the session state bound to the current context, if any.
54def get_active_session_span() -> Span | None: 55 """Return the active ``parlot.session`` span for the current context, if any.""" 56 return _active_session_span.get()
Return the active parlot.session span for the current context, if any.
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").
23def resolve_api_key(api_key: Optional[str] = None) -> str: 24 """Resolve API key from kwarg or ``PARLOT_API_KEY``.""" 25 return api_key or os.environ.get("PARLOT_API_KEY", "") or ""
Resolve API key from kwarg or PARLOT_API_KEY.
28def resolve_capture_genai_content(capture_genai_content: Optional[bool] = None) -> Optional[bool]: 29 """Return the explicit ``parlotize(capture_genai_content=)`` override, or None. 30 31 When None, callers should resolve via ``should_capture_genai_content`` against 32 bootstrap policy (default on). 33 """ 34 if capture_genai_content is not None: 35 return bool(capture_genai_content) 36 return None
Return the explicit parlotize(capture_genai_content=) override, or None.
When None, callers should resolve via should_capture_genai_content against
bootstrap policy (default on).
18def resolve_endpoint(endpoint: Optional[str] = None) -> str: 19 """Resolve OTLP base endpoint from kwarg or ``PARLOT_ENDPOINT``.""" 20 return (endpoint or os.environ.get("PARLOT_ENDPOINT", "") or "").rstrip("/")
Resolve OTLP base endpoint from kwarg or PARLOT_ENDPOINT.
26def resolve_parlot_sdk_version() -> str: 27 """Return installed parlot-core version, cached after first call.""" 28 global _cached_version 29 if _cached_version is not None: 30 return _cached_version 31 32 try: 33 _cached_version = _normalize_version(metadata.version(_PARLOT_CORE_DIST)) 34 except metadata.PackageNotFoundError: 35 _cached_version = "" 36 37 return _cached_version
Return installed parlot-core version, cached after first call.
136def runtime_from_bootstrap( 137 endpoint: str, api_key: str, payload: dict[str, Any] 138) -> ParlotRuntimeContext: 139 globs, agents = _parse_recording_policy(payload) 140 ( 141 logs_globs, 142 logs_agents, 143 logs_agent_min_levels, 144 logs_min_level, 145 logs_present, 146 ) = _parse_logs_policy(payload) 147 ( 148 capture_globs, 149 capture_agents, 150 capture_present, 151 ) = _parse_capture_genai_content_policy(payload) 152 return ParlotRuntimeContext( 153 endpoint=endpoint.rstrip("/"), 154 api_key=api_key, 155 tenant_id=str(payload.get("tenant_id") or ""), 156 content_bucket=str(payload.get("content_bucket") or ""), 157 r2_endpoint=str(payload.get("r2_endpoint") or ""), 158 recording_globs=globs, 159 recording_agents=agents, 160 logs_globs=logs_globs, 161 logs_agents=logs_agents, 162 logs_agent_min_levels=logs_agent_min_levels, 163 logs_min_level=logs_min_level, 164 logs_policy_present=logs_present, 165 capture_genai_content_globs=capture_globs, 166 capture_genai_content_agents=capture_agents, 167 capture_genai_content_policy_present=capture_present, 168 )
73def session_owned(*, framework: str | None = None) -> bool: 74 """True when an active Parlot session with a non-empty session_id is bound. 75 76 Args: 77 framework: If provided, also require ``state.framework`` to match 78 (e.g. ``session_owned(framework="livekit")`` for coexistence checks). 79 """ 80 state = get_active_session() 81 if state is None or not state.session_id: 82 return False 83 if framework is not None and state.framework != framework: 84 return False 85 return True
True when an active Parlot session with a non-empty session_id is bound.
Arguments:
- framework: If provided, also require
state.frameworkto match (e.g.session_owned(framework="livekit")for coexistence checks).
59def set_active_session( 60 span: Span | None, 61 state: SessionState | None, 62) -> None: 63 """Bind (or clear) the active session span + state for this context.""" 64 _active_session_span.set(span) 65 _active_session_state.set(state)
Bind (or clear) the active session span + state for this context.
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")
30def should_capture_genai_content( 31 agent_name: str, 32 *, 33 metadata_capture_genai_content: Optional[bool] = None, 34 capture_genai_content_config: Optional[bool] = None, 35 bootstrap_globs: Optional[Sequence[str]] = None, 36 bootstrap_agents: Optional[Mapping[str, bool]] = None, 37 bootstrap_present: bool = False, 38) -> bool: 39 """Return True when generative AI / tool content bodies should be captured. 40 41 Precedence: job metadata > ``parlotize(capture_genai_content=…)`` > bootstrap (UI). 42 When bootstrap is absent (``bootstrap_present=False``) or globs are unset, 43 default is **on**. Explicit empty globs means off. 44 """ 45 if metadata_capture_genai_content is False: 46 return False 47 if metadata_capture_genai_content is True: 48 return True 49 50 if capture_genai_content_config is not None: 51 return bool(capture_genai_content_config) 52 53 agents = bootstrap_agents or {} 54 if agent_name in agents: 55 return bool(agents[agent_name]) 56 57 if not bootstrap_present: 58 return True 59 60 if bootstrap_globs is None: 61 return True 62 63 globs = list(bootstrap_globs) 64 if not globs: 65 return False 66 return matches_allowlist(agent_name, globs)
Return True when generative AI / tool content bodies should be captured.
Precedence: job metadata > parlotize(capture_genai_content=…) > bootstrap (UI).
When bootstrap is absent (bootstrap_present=False) or globs are unset,
default is on. Explicit empty globs means off.
47def should_capture_logs( 48 agent_name: str, 49 *, 50 metadata_capture_logs: Optional[bool] = None, 51 capture_logs_config: CaptureLogsConfig = None, 52 bootstrap_globs: Optional[Sequence[str]] = None, 53 bootstrap_agents: Optional[Mapping[str, bool]] = None, 54 bootstrap_present: bool = False, 55) -> bool: 56 """Return True when application logs should be captured for this job. 57 58 Precedence: job metadata > ``parlotize(capture_logs=…)`` > bootstrap (UI). 59 When bootstrap is absent (``bootstrap_present=False``) or globs are unset, 60 default is **on**. Explicit empty globs means off. 61 """ 62 if metadata_capture_logs is False: 63 return False 64 if metadata_capture_logs is True: 65 return True 66 67 if capture_logs_config is not None: 68 if isinstance(capture_logs_config, bool): 69 return capture_logs_config 70 return matches_allowlist(agent_name, list(capture_logs_config)) 71 72 agents = bootstrap_agents or {} 73 if agent_name in agents: 74 return bool(agents[agent_name]) 75 76 if not bootstrap_present: 77 return True 78 79 if bootstrap_globs is None: 80 return True 81 82 globs = list(bootstrap_globs) 83 if not globs: 84 return False 85 return matches_allowlist(agent_name, globs)
Return True when application logs should be captured for this job.
Precedence: job metadata > parlotize(capture_logs=…) > bootstrap (UI).
When bootstrap is absent (bootstrap_present=False) or globs are unset,
default is on. Explicit empty globs means off.
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.
40def stamp_session_sdk_version(session_span: Any, state: Any | None = None) -> None: 41 """Stamp ``parlot.sdk.version`` and ``session.metadata.sdk_version`` on the session span and state.""" 42 version = resolve_parlot_sdk_version() 43 if not version: 44 return 45 46 if isinstance(session_span, dict): 47 session_span[ATTR_PARLOT_SDK_VERSION] = version 48 session_span[ATTR_SESSION_METADATA_SDK_VERSION] = version 49 elif session_span is not None and hasattr(session_span, "set_attribute"): 50 session_span.set_attribute(ATTR_PARLOT_SDK_VERSION, version) 51 session_span.set_attribute(ATTR_SESSION_METADATA_SDK_VERSION, version) 52 53 if state is None: 54 try: 55 from parlot.core.session import get_active_session 56 57 state = get_active_session() 58 except Exception: 59 state = None 60 61 if state is not None and hasattr(state, "custom_metadata") and isinstance(state.custom_metadata, dict): 62 state.custom_metadata[ATTR_SESSION_METADATA_SDK_VERSION] = version
Stamp parlot.sdk.version and session.metadata.sdk_version on the session span and state.
11def stamp_turn_utterance_text( 12 span: Any, 13 *, 14 participant_role: str, 15 utterance_text: str, 16) -> None: 17 """ 18 Stamp ``turn.user_text`` or ``turn.agent_text`` on a turn root span. 19 20 LiveKit and LangGraph adapters both call this after committing an utterance 21 so the OTLP contract stays consistent regardless of source. 22 """ 23 text = (utterance_text or "").strip() 24 if not text: 25 return 26 role = (participant_role or "").strip().lower() 27 if role == "user": 28 span.set_attribute(ATTR_TURN_USER_TEXT, text) 29 elif role == "agent": 30 span.set_attribute(ATTR_TURN_AGENT_TEXT, text)
Stamp turn.user_text or turn.agent_text on a turn root span.
LiveKit and LangGraph adapters both call this after committing an utterance so the OTLP contract stays consistent regardless of source.