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]
@dataclass
class BaseParlotizeResult:
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.

BaseParlotizeResult( context: ParlotContext, endpoint: str, api_key: str, tracer_provider: Any, agent_id: str, agent_version: str, capture_genai_content: Optional[bool], capture_logs: bool | list[str] | None, log_level: Optional[str])
context: ParlotContext
endpoint: str
api_key: str
tracer_provider: Any
agent_id: str
agent_version: str
capture_genai_content: Optional[bool]
capture_logs: bool | list[str] | None
log_level: Optional[str]
@runtime_checkable
class ParlotizeProtocol(typing.Protocol):
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=).

ParlotizeProtocol(*args, **kwargs)
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)
@dataclass
class ParlotContext:
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.

ParlotContext( runtime: Optional[ParlotRuntimeContext] = None, session_logs: SessionLogsCollector = <factory>)
runtime: Optional[ParlotRuntimeContext] = None
session_logs: SessionLogsCollector
def shutdown(self) -> None:
23    def shutdown(self) -> None:
24        self.session_logs.shutdown()
25        self.runtime = None
@dataclass(frozen=True)
class ParlotRuntimeContext:
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)
ParlotRuntimeContext( endpoint: str, api_key: str, tenant_id: str, content_bucket: str, r2_endpoint: str, recording_globs: tuple[str, ...] = (), recording_agents: tuple[tuple[str, bool], ...] = (), logs_globs: tuple[str, ...] = (), logs_agents: tuple[tuple[str, bool], ...] = (), logs_agent_min_levels: tuple[tuple[str, str], ...] = (), logs_min_level: str = 'INFO', logs_policy_present: bool = False, capture_genai_content_globs: tuple[str, ...] = (), capture_genai_content_agents: tuple[tuple[str, bool], ...] = (), capture_genai_content_policy_present: bool = False)
endpoint: str
api_key: str
tenant_id: str
content_bucket: str
r2_endpoint: str
recording_globs: tuple[str, ...] = ()
recording_agents: tuple[tuple[str, bool], ...] = ()
logs_globs: tuple[str, ...] = ()
logs_agents: tuple[tuple[str, bool], ...] = ()
logs_agent_min_levels: tuple[tuple[str, str], ...] = ()
logs_min_level: str = 'INFO'
logs_policy_present: bool = False
capture_genai_content_globs: tuple[str, ...] = ()
capture_genai_content_agents: tuple[tuple[str, bool], ...] = ()
capture_genai_content_policy_present: bool = False
def recording_agents_map(self) -> dict[str, bool]:
32    def recording_agents_map(self) -> dict[str, bool]:
33        return dict(self.recording_agents)
def logs_agents_map(self) -> dict[str, bool]:
35    def logs_agents_map(self) -> dict[str, bool]:
36        return dict(self.logs_agents)
def logs_agent_min_levels_map(self) -> dict[str, str]:
38    def logs_agent_min_levels_map(self) -> dict[str, str]:
39        return dict(self.logs_agent_min_levels)
def capture_genai_content_agents_map(self) -> dict[str, bool]:
41    def capture_genai_content_agents_map(self) -> dict[str, bool]:
42        return dict(self.capture_genai_content_agents)
class SessionLogsCollector:
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.

def bind_context(self, context: ParlotContext) -> None:
264    def bind_context(self, context: ParlotContext) -> None:
265        self._context_ref = weakref.ref(context)
def set_capture_logs_config( self, capture_logs: Union[bool, Sequence[str], NoneType] = None, *, log_level: Optional[str] = None) -> None:
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)
def set_resolvers( self, *, session_resolver: Optional[Callable[[], Optional[dict[str, Any]]]] = None, agent_id_resolver: Optional[Callable[[], str]] = None, metadata_resolver: Optional[Callable[[], Optional[bool]]] = None) -> None:
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).

def capture_logs_enabled_for_agent(self, agent_name: str = '') -> bool:
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        )
def handle_record(self, handler: logging.Handler, record: logging.LogRecord) -> None:
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)
def init(self, *, endpoint: str = '', api_key: str = '') -> None:
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.

def drain(self) -> list[parlot.core.session_logs.SessionLogRecord]:
405    def drain(self) -> list[SessionLogRecord]:
406        with self._lock:
407            events = list(self._buffer)
408            self._buffer.clear()
409            return events
def snapshot(self) -> list[parlot.core.session_logs.SessionLogRecord]:
411    def snapshot(self) -> list[SessionLogRecord]:
412        with self._lock:
413            return list(self._buffer)
def shutdown(self) -> None:
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.

def base_parlotize( agent_id: str, *, endpoint: Optional[str] = None, api_key: Optional[str] = None, capture_genai_content: Optional[bool] = None, tracer_provider: opentelemetry.trace.TracerProvider | None = None, version: Optional[str] = None, capture_logs: bool | list[str] | None = None, log_level: Optional[str] = None) -> BaseParlotizeResult:
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 TracerProvider to 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.

def configure_parlot_logging() -> None:
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.

class ExportFilterSpanExporter(opentelemetry.sdk.trace.export.SpanExporter):
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.

ExportFilterSpanExporter(exporter: opentelemetry.sdk.trace.export.SpanExporter)
17    def __init__(self, exporter: SpanExporter) -> None:
18        self._exporter = exporter
def export( self, spans: Sequence[opentelemetry.sdk.trace.ReadableSpan]) -> opentelemetry.sdk.trace.export.SpanExportResult:
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.Span objects to be exported
Returns:

The result of the export

def shutdown(self) -> None:
26    def shutdown(self) -> None:
27        self._exporter.shutdown()

Shuts down the exporter.

Called when the SDK is shut down.

def force_flush(self, timeout_millis: int = 30000) -> bool:
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.

class ParlotBaseProcessor(opentelemetry.sdk.trace.SpanProcessor):
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))
def on_start(self, span, parent_context=None) -> None:
36    def on_start(self, span, parent_context=None) -> None:
37        log_span_event("on_start", span)

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.Span that just started.
  • parent_context: The parent context of the span that just started.
def on_end(self, span: opentelemetry.sdk.trace.ReadableSpan) -> None:
39    def on_end(self, span: ReadableSpan) -> None:
40        log_span_event("on_end", span)

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.Span that just ended.
def shutdown(self) -> None:
42    def shutdown(self) -> None:
43        pass

Called when a opentelemetry.sdk.trace.TracerProvider is shutdown.

def force_flush(self, timeout_millis: int = 30000) -> bool:
45    def force_flush(self, timeout_millis: int = 30_000) -> bool:
46        return True

Export all ended spans to the configured Exporter that have not yet been exported.

Arguments:
  • timeout_millis: The maximum amount of time to wait for spans to be exported.
Returns:

False if the timeout is exceeded, True otherwise.

def assert_sync_span_processors(provider) -> bool:
 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.

@dataclass
class SessionState:
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).

SessionState( session_id: str = '', room_name: str = '', room_sid: str = '', agent_label: str = '', turn_count: int = 0, tool_call_count: int = 0, handoff_count: int = 0, total_input_tokens: int = 0, total_output_tokens: int = 0, total_cost_usd: float = 0.0, human_rep_participant_ids: set[str] = <factory>, topology_agents: list[dict[str, typing.Any]] = <factory>, custom_metadata: dict[str, str] = <factory>, framework: str = '')
session_id: str = ''
room_name: str = ''
room_sid: str = ''
agent_label: str = ''
turn_count: int = 0
tool_call_count: int = 0
handoff_count: int = 0
total_input_tokens: int = 0
total_output_tokens: int = 0
total_cost_usd: float = 0.0
human_rep_participant_ids: set[str]
topology_agents: list[dict[str, typing.Any]]
custom_metadata: dict[str, str]
framework: str = ''
@dataclass
class SessionTopology:
 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.

SessionTopology( prompts_by_agent: dict[str, str] = <factory>, intent_by_agent: dict[str, dict[str, str]] = <factory>, agents_seen: dict[str, dict[str, typing.Any]] = <factory>, intent_segments: list[parlot.core.topology._MutableSegment] = <factory>, active_segment: parlot.core.topology._MutableSegment | None = None, pending_segment_from_turn: int | None = None, first_agent_label: str = '', agent_chain: list[str] = <factory>, handoff_count: int = 0, _handoff_keys: set[str] = <factory>)
prompts_by_agent: dict[str, str]
intent_by_agent: dict[str, dict[str, str]]
agents_seen: dict[str, dict[str, typing.Any]]
intent_segments: list[parlot.core.topology._MutableSegment]
active_segment: parlot.core.topology._MutableSegment | None = None
pending_segment_from_turn: int | None = None
first_agent_label: str = ''
agent_chain: list[str]
handoff_count: int = 0
def upsert_agent(self, agent_id: str) -> None:
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"}
def record_instructions(self, agent_id: str, text: str) -> None:
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)
def push_agent_chain(self, agent_id: str) -> None:
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)
def open_bootstrap_segment(self, agent_id: str, from_turn: int) -> None:
 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)
def open_segment_after_handoff( self, agent_id: str, *, handoff_index: int | None = None, turn_index: int, from_agent: str = '') -> None:
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        )
def close_active_segment(self, to_turn: int) -> None:
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
def apply_pending_from_turn_on_emit(self, turn_index: int) -> None:
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
def finalize_intent_sequence(self, final_turn: int) -> list[dict[str, typing.Any]]:
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        )
def bootstrap_instructions(self) -> str:
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 ""
def agents_json(self) -> list[dict[str, typing.Any]]:
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
def stamp_session_span(self, span: Any, *, final_turn: int) -> None:
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.

def add_platform_ref(kind: str, value: str, *, framework: str = 'custom') -> None:
 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")
def adopt_existing_tracer_provider() -> typing.Any | None:
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.

def build_otlp_http_exporter( *, endpoint: str, api_key: str = '') -> opentelemetry.sdk.trace.export.SpanExporter:
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.

def build_parlot_client_headers(api_key: str = '') -> dict[str, str]:
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.

def build_resource( *, service_name: Optional[str] = None, service_version: Optional[str] = None) -> opentelemetry.sdk.resources.Resource:
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)
def build_tracer_provider( *, endpoint: Optional[str] = None, api_key: Optional[str] = None, service_name: Optional[str] = None, service_version: Optional[str] = None, span_processors: Sequence[opentelemetry.sdk.trace.SpanProcessor] = (), span_exporter: opentelemetry.sdk.trace.export.SpanExporter | None = None) -> opentelemetry.sdk.trace.TracerProvider:
 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.

HEADER_SDK_NAME = 'x-parlot-sdk-name'
HEADER_SDK_VERSION = 'x-parlot-sdk-version'
HEADER_INGESTION_VERSION = 'x-parlot-ingestion-version'
INGESTION_PROTOCOL_VERSION = '1'
def clear_active_session() -> None:
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.

def derive_intent(agent_id: str, instructions: str = '') -> dict[str, str]:
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.

def fetch_telemetry_bootstrap( endpoint: str, api_key: str, context: ParlotContext, *, timeout: float = 15.0) -> Optional[dict[str, Any]]:
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.

def get_active_session() -> SessionState | 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.

def get_active_session_span() -> opentelemetry.trace.span.Span | None:
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.

@contextmanager
def human_escalation(label: str | None = None) -> Iterator[NoneType]:
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.
def record_human_rep(participant_id: str, *, label: str | None = None) -> None:
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").
def resolve_api_key(api_key: Optional[str] = None) -> str:
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.

def resolve_capture_genai_content(capture_genai_content: Optional[bool] = None) -> Optional[bool]:
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).

def resolve_endpoint(endpoint: Optional[str] = None) -> str:
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.

def resolve_parlot_sdk_version() -> str:
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.

def runtime_from_bootstrap( endpoint: str, api_key: str, payload: dict[str, typing.Any]) -> ParlotRuntimeContext:
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    )
def session_owned(*, framework: str | None = None) -> bool:
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.framework to match (e.g. session_owned(framework="livekit") for coexistence checks).
def set_active_session( span: opentelemetry.trace.span.Span | None, state: SessionState | None) -> None:
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.

def set_session_attribute(key: str, value: str | int | float | bool) -> None:
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, or bool).
def set_session_metadata(**pairs: str | int | float | bool) -> None:
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, or bool).

Example::

from parlot.instrumentation.livekit import set_session_metadata

set_session_metadata(order_id="12345", crm_ticket="TKT-9")
def should_capture_genai_content( agent_name: str, *, metadata_capture_genai_content: Optional[bool] = None, capture_genai_content_config: Optional[bool] = None, bootstrap_globs: Optional[Sequence[str]] = None, bootstrap_agents: Optional[Mapping[str, bool]] = None, bootstrap_present: bool = False) -> bool:
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.

def should_capture_logs( agent_name: str, *, metadata_capture_logs: Optional[bool] = None, capture_logs_config: Union[bool, Sequence[str], NoneType] = None, bootstrap_globs: Optional[Sequence[str]] = None, bootstrap_agents: Optional[Mapping[str, bool]] = None, bootstrap_present: bool = False) -> bool:
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.

def stamp_platform_refs(span: Any, refs: list[tuple[str, str, str]]) -> None:
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 a ReadableSpan with a mutable _attributes dict.
  • refs: Non-empty list of reference triples.
def stamp_session_sdk_version(session_span: Any, state: typing.Any | None = None) -> None:
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.

def stamp_turn_utterance_text(span: Any, *, participant_role: str, utterance_text: str) -> None:
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.