Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
exporter.py498 linesDownload Raw Back to strands_agents
1"""Transform Strands OTEL spans into LangSmith-compatible formats.2 3This module wraps the standard OTLPSpanExporter and intercepts spans before export,4remapping attributes, span names, and structure to align with LangSmith's expected5OTEL ingest schema.6 7Strands emits GenAI message data as span events (gen_ai.user.message,8gen_ai.assistant.message, gen_ai.choice, etc.), which is non-standard — the OTEL9GenAI semantic conventions define these as Log Events, not span events.  This10exporter flattens those span events into span attributes (gen_ai.prompt,11gen_ai.completion) that LangSmith's server-side OTEL ingest can consume directly.12"""13 14import json15import logging16from collections.abc import Sequence17from typing import Any18 19from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter20from opentelemetry.sdk.trace import ReadableSpan21from opentelemetry.sdk.trace.export import (22    BatchSpanProcessor,23    ConsoleSpanExporter,24    SimpleSpanProcessor,25    SpanExporter,26    SpanExportResult,27)28from strands.telemetry import StrandsTelemetry29 30logger = logging.getLogger(__name__)31 32__all__ = [33    "LangSmithSpanExporter",34    "create_langsmith_exporter",35    "setup_langsmith_telemetry",36]37 38 39class LangSmithSpanExporter(SpanExporter):40    """Span exporter that reformats Strands OTEL spans for LangSmith compatibility.41 42    Wraps a delegate exporter (typically OTLPSpanExporter pointed at LangSmith's43    OTEL endpoint) and transforms each span before forwarding it.44 45    Args:46        delegate: The underlying SpanExporter to forward transformed spans to.47    """48 49    def __init__(self, delegate: SpanExporter) -> None:50        """Initialize the exporter.51 52        Args:53            delegate: The underlying SpanExporter to forward transformed spans to.54        """55        self._delegate = delegate56 57    def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult:58        """Transform spans and forward them to the delegate exporter.59 60        Args:61            spans: The batch of spans to export.62 63        Returns:64            The result from the delegate exporter.65        """66        transformed = []67        for span in spans:68            try:69                transformed.append(self._transform_span(span))70            except Exception:71                logger.warning(72                    "Failed to transform span %r, exporting original",73                    span.name,74                    exc_info=True,75                )76                transformed.append(span)77        return self._delegate.export(transformed)78 79    def shutdown(self) -> None:80        """Shut down the delegate exporter."""81        self._delegate.shutdown()82 83    def force_flush(self, timeout_millis: int = 30000) -> bool:84        """Force flush the delegate exporter.85 86        Args:87            timeout_millis: Maximum time to wait for flush to complete.88 89        Returns:90            True if flush succeeded.91        """92        return self._delegate.force_flush(timeout_millis)93 94    # -- gen_ai.operation.name → langsmith.span.kind mapping ---------------95 96    _OPERATION_TO_RUN_TYPE: dict[str, str] = {97        "chat": "llm",98        "invoke_agent": "chain",99        "execute_tool": "tool",100        "execute_event_loop_cycle": "chain",101    }102 103    # -- Span name fallback → langsmith.span.kind mapping ------------------104    # Some Strands spans (e.g. execute_event_loop_cycle) don't always carry105    # gen_ai.operation.name.  Fall back to the span *name* in that case.106    _SPAN_NAME_TO_RUN_TYPE: dict[str, str] = {107        "execute_event_loop_cycle": "chain",108    }109 110    # -- Event name → role mapping ----------------------------------------111 112    _EVENT_ROLE_MAP: dict[str, str] = {113        "gen_ai.user.message": "user",114        "gen_ai.assistant.message": "assistant",115        "gen_ai.system.message": "system",116        "gen_ai.tool.message": "tool",117        "gen_ai.choice": "assistant",118    }119 120    _MESSAGE_EVENTS: set[str] = {121        "gen_ai.user.message",122        "gen_ai.system.message",123        "gen_ai.tool.message",124        "gen_ai.assistant.message",125        "gen_ai.choice",126    }127 128    # ----------------------------------------------------------------------129 130    def _transform_span(self, span: ReadableSpan) -> ReadableSpan:131        """Flatten span events into prompt/completion attributes.132 133        Strands attaches GenAI message data as span events.  LangSmith expects134        them as JSON-serialized span attributes instead.  Each message object135        gets a ``role`` field injected when one can be inferred from the event136        name and isn't already present in the payload.137 138        All conversation history (user, system, tool, and intermediate assistant139        messages) is placed into ``gen_ai.prompt``.  Only the final140        ``gen_ai.choice`` event — the model's actual response — goes into141        ``gen_ai.completion``.142 143        Args:144            span: The original Strands span.145 146        Returns:147            A new ReadableSpan with the message attributes added.148        """149        span_attrs = dict(span.attributes) if span.attributes else {}150        operation_raw = span_attrs.get("gen_ai.operation.name", "")151        operation = operation_raw if isinstance(operation_raw, str) else ""152        # Strands puts tool metadata in span attributes (only on tool spans)153        tool_call_id_raw = span_attrs.get("gen_ai.tool.call.id", "")154        tool_call_id = tool_call_id_raw if isinstance(tool_call_id_raw, str) else ""155        tool_name_raw = span_attrs.get("gen_ai.tool.name", "")156        tool_name = tool_name_raw if isinstance(tool_name_raw, str) else ""157 158        # Maps toolUseId → tool name so tool-result messages can be labelled.159        # Seeded from span attrs on tool spans; extended inline as we encounter160        # assistant/choice events that contain toolUse blocks.161        tool_id_to_name: dict[str, str] = {}162        if tool_name and tool_call_id:163            tool_id_to_name[tool_call_id] = tool_name164 165        input_messages: list[dict[str, Any]] = []166        output_messages: list[dict[str, Any]] = []167        remaining_events: list[Any] = []168 169        for event in span.events:170            name = event.name171            attrs = dict(event.attributes) if event.attributes else {}172 173            if name == "gen_ai.choice":174                # The final model response is the only true output175                msg = self._event_to_message(176                    name, attrs, tool_id_to_name=tool_id_to_name177                )178                if tool_call_id:179                    msg["tool_call_id"] = tool_call_id180                output_messages.append(msg)181            elif name in self._MESSAGE_EVENTS:182                msg = self._event_to_message(183                    name, attrs, tool_id_to_name=tool_id_to_name184                )185                if tool_call_id:186                    msg["tool_call_id"] = tool_call_id187                input_messages.append(msg)188            else:189                # Preserve non-message events as-is190                remaining_events.append(event)191 192        # Merge new attributes with the originals.193        # LangSmith's server-side OTEL ingest expects inputs under "gen_ai.prompt"194        # and outputs under "gen_ai.completion" (matching the attribute names used195        # by LangSmith's own OTELExporter).196        new_attrs: dict[str, Any] = dict(span_attrs)197        if input_messages:198            new_attrs["gen_ai.prompt"] = json.dumps({"messages": input_messages})199        if output_messages:200            new_attrs["gen_ai.completion"] = json.dumps(output_messages[-1])201 202        # Map gen_ai.operation.name to langsmith.span.kind (run type).203        # Some Strands spans (e.g. execute_event_loop_cycle) don't carry204        # gen_ai.operation.name — fall back to the span name.205        run_type = self._OPERATION_TO_RUN_TYPE.get(operation)206        if not run_type:207            run_type = self._SPAN_NAME_TO_RUN_TYPE.get(span.name)208        if run_type:209            new_attrs["langsmith.span.kind"] = run_type210 211        # For LLM spans, set the provider metadata and display name.212        # Strands hardcodes gen_ai.system to "strands-agents" regardless of213        # backend, so we use langsmith.metadata.ls_provider to surface the214        # actual provider.215        # TODO: Detect non-Bedrock providers216        if run_type == "llm" and new_attrs.get("gen_ai.system"):217            new_attrs["langsmith.metadata.ls_provider"] = "amazon_bedrock"218            new_attrs["langsmith.metadata.ls_model_type"] = "chat"219 220        return ReadableSpan(221            name=span.name,222            context=span.context,223            parent=span.parent,224            resource=span.resource,225            attributes=new_attrs,226            events=remaining_events,227            links=span.links,228            kind=span.kind,229            status=span.status,230            start_time=span.start_time,231            end_time=span.end_time,232            instrumentation_scope=span.instrumentation_scope,233        )234 235    def _event_to_message(236        self,237        event_name: str,238        attrs: dict[str, Any],239        *,240        tool_id_to_name: dict[str, str] | None = None,241    ) -> dict[str, Any]:242        """Convert a span event into a message dict, injecting ``role`` if missing.243 244        For ``gen_ai.choice`` events the content lives under the ``message``245        key; for all other events it lives under ``content``.  In both cases246        the value may be a JSON-encoded string that we parse back so the final247        attribute is cleanly nested.248 249        Args:250            event_name: The OTEL event name (e.g. ``gen_ai.user.message``).251            attrs: The event's attribute dict.252            tool_id_to_name: Optional mapping of tool-call IDs to tool names.253 254        Returns:255            A message dict with at least ``role`` and ``content`` keys.256        """257        role = self._EVENT_ROLE_MAP.get(event_name, "unknown")258 259        # gen_ai.choice stores content in "message", others in "content"260        if event_name == "gen_ai.choice":261            raw = attrs.get("message", "[]")262        else:263            raw = attrs.get("content", "[]")264 265        # Parse the JSON string back into a list so the final serialized266        # attribute isn't double-encoded.267        try:268            content = json.loads(raw) if isinstance(raw, str) else raw269        except (json.JSONDecodeError, TypeError):270            content = raw271 272        # Tool messages in chat history contain Bedrock toolResult blocks.273        # Flatten them into the format LangSmith expects: a top-level274        # tool_call_id and plain-text content.275        if event_name == "gen_ai.tool.message" and isinstance(content, list):276            return self._flatten_tool_result_message(277                content,278                tool_id_to_name=tool_id_to_name or {},279            )280 281        # Convert Bedrock-shaped content blocks to LangSmith-shaped blocks.282        # While iterating, harvest toolUse names so later tool-result messages283        # can be labelled (assistant messages precede their tool results).284        if isinstance(content, list):285            converted = []286            for block in content:287                if (288                    tool_id_to_name is not None289                    and isinstance(block, dict)290                    and "toolUse" in block291                ):292                    tu = block["toolUse"]293                    tid, tname = tu.get("toolUseId", ""), tu.get("name", "")294                    if tid and tname:295                        tool_id_to_name[tid] = tname296                converted.append(self._convert_content_block(block))297            content = converted298 299        if event_name == "gen_ai.tool.message":300            content = self._stringify_tool_content(content)301 302        msg: dict[str, Any] = {"role": role, "content": content}303 304        # Carry over tool_call_id from event attributes (Strands stores it as "id")305        if "id" in attrs:306            msg["tool_call_id"] = attrs["id"]307 308        # Carry over finish_reason for choice events309        if "finish_reason" in attrs:310            msg["finish_reason"] = attrs["finish_reason"]311 312        return msg313 314    @staticmethod315    def _flatten_tool_result_message(316        content_blocks: list[Any],317        *,318        tool_id_to_name: dict[str, str] | None = None,319    ) -> dict[str, Any]:320        """Flatten Bedrock ``toolResult`` content blocks into a LangSmith tool message.321 322        Bedrock tool results arrive as::323 324            [325                {326                    "toolResult": {327                        "toolUseId": "x",328                        "status": "success",329                        "content": [{"text": "..."}],330                    }331                }332            ]333 334        LangSmith expects tool messages as::335 336            {"role": "tool", "name": "my_tool", "tool_call_id": "x", "content": "..."}337 338        If there are multiple toolResult blocks they are joined with newlines.339        Non-toolResult blocks are converted normally and appended.340 341        Args:342            content_blocks: The parsed content block list from the event.343            tool_id_to_name: Optional mapping of tool-call IDs to tool names.344 345        Returns:346            A flat tool message dict.347        """348        tool_call_id = ""349        text_parts: list[str] = []350        other_blocks: list[Any] = []351 352        for block in content_blocks:353            if isinstance(block, dict) and "toolResult" in block:354                tr = block["toolResult"]355                if not tool_call_id:356                    tool_call_id = tr.get("toolUseId", "")357                # Extract text from nested content blocks358                for nested in tr.get("content", []):359                    if isinstance(nested, dict) and "text" in nested:360                        text_parts.append(nested["text"])361                    else:362                        other_blocks.append(nested)363            else:364                other_blocks.append(block)365 366        if text_parts:367            flat_content: Any = "\n".join(text_parts)368        else:369            flat_content = other_blocks370 371        msg: dict[str, Any] = {372            "role": "tool",373            "content": LangSmithSpanExporter._stringify_tool_content(flat_content),374        }375        # Look up the tool name from the toolUseId → name mapping376        tool_name = (tool_id_to_name or {}).get(tool_call_id, "")377        if tool_name:378            msg["name"] = tool_name379        if tool_call_id:380            msg["tool_call_id"] = tool_call_id381        return msg382 383    @staticmethod384    def _stringify_tool_content(content: Any) -> str:385        """Ensure tool message content is a string.386 387        Strands emits tool-call inputs as JSON-serialized objects in388        ``gen_ai.tool.message`` events. After parsing event payloads for the389        rest of the exporter, convert those tool inputs back to strings so390        LangSmith receives tool message content as either a plain string or a391        stringified object.392        """393        if isinstance(content, str):394            return content395        try:396            return json.dumps(content, ensure_ascii=False)397        except (TypeError, ValueError):398            return str(content)399 400    @staticmethod401    def _convert_content_block(block: Any) -> Any:402        """Convert a single Bedrock/Converse content block to LangSmith format.403 404        Bedrock uses implicit typing (the key name *is* the type)::405 406            {"text": "hello"}407            {"toolUse": {"toolUseId": "x", "name": "f", "input": {...}}}408            {"toolResult": {"toolUseId": "x", "status": "success", "content": [...]}}409 410        LangSmith expects explicit ``type`` fields::411 412            {"type": "text", "text": "hello"}413            {"type": "tool_use", "id": "x", "name": "f", "input": {...}}414            {415                "type": "tool_result",416                "tool_use_id": "x",417                "status": "success",418                "content": [...],419            }420 421        Unrecognised blocks are returned as-is.422        """423        if not isinstance(block, dict):424            return block425 426        if "text" in block and len(block) == 1:427            return {"type": "text", "text": block["text"]}428 429        if "toolUse" in block:430            tu = block["toolUse"]431            return {432                "type": "tool_use",433                "id": tu.get("toolUseId", ""),434                "name": tu.get("name", ""),435                "input": tu.get("input", {}),436            }437 438        if "toolResult" in block:439            tr = block["toolResult"]440            converted: dict[str, Any] = {441                "type": "tool_result",442                "tool_use_id": tr.get("toolUseId", ""),443            }444            if "status" in tr:445                converted["status"] = tr["status"]446            if "content" in tr:447                # Recursively convert nested content blocks448                nested = tr["content"]449                if isinstance(nested, list):450                    nested = [451                        LangSmithSpanExporter._convert_content_block(b) for b in nested452                    ]453                converted["content"] = nested454            return converted455 456        # Unknown block shape — pass through unchanged457        return block458 459 460# ---------------------------------------------------------------------------461# Convenience wiring462# ---------------------------------------------------------------------------463 464 465def create_langsmith_exporter(**otlp_kwargs: Any) -> LangSmithSpanExporter:466    """Create a LangSmithSpanExporter wrapping a standard OTLPSpanExporter.467 468    Keyword arguments are forwarded to OTLPSpanExporter (endpoint, headers, etc.).469    If not provided, the exporter will fall back to the standard OTEL_EXPORTER_OTLP_*470    environment variables.471 472    Returns:473        A ready-to-use LangSmithSpanExporter instance.474    """475    delegate = OTLPSpanExporter(**otlp_kwargs)476    return LangSmithSpanExporter(delegate=delegate)477 478 479def setup_langsmith_telemetry(*, console: bool = False) -> None:480    """Wire up Strands telemetry with the LangSmith-compatible exporter.481 482    Call this instead of (or in addition to) the standard483    ``StrandsTelemetry().setup_otlp_exporter()`` flow.484 485    Args:486        console: If True, also add a ConsoleSpanExporter that prints transformed487                 spans to stdout (useful for debugging).488    """489    telemetry = StrandsTelemetry()490    exporter = create_langsmith_exporter()491    telemetry.tracer_provider.add_span_processor(BatchSpanProcessor(exporter))492 493    if console:494        console_exporter = LangSmithSpanExporter(delegate=ConsoleSpanExporter())495        telemetry.tracer_provider.add_span_processor(496            SimpleSpanProcessor(console_exporter)497        )498 
codekingpro/portable-devtools · Team Ai