codekingpro/portable-devtools
114k
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 