Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/strands_evals/mappers/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
"""Converters for transforming telemetry data to Session format."""

from .adk_otel_session_mapper import ADKOtelSessionMapper
from .cloudwatch_parser import CloudWatchLogsParser, parse_cloudwatch_logs
from .cloudwatch_session_mapper import CloudWatchSessionMapper
from .langchain_otel_session_mapper import LangChainOtelSessionMapper
Expand All @@ -10,6 +11,7 @@
from .utils import detect_otel_mapper, get_scope_name, readable_spans_to_dicts

__all__ = [
"ADKOtelSessionMapper",
"CloudWatchLogsParser",
"CloudWatchSessionMapper",
"GenAIConventionVersion",
Expand Down
637 changes: 637 additions & 0 deletions src/strands_evals/mappers/adk_otel_session_mapper.py

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions src/strands_evals/mappers/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
SCOPE_LANGCHAIN_OTEL = "opentelemetry.instrumentation.langchain"
SCOPE_OPENINFERENCE = "openinference.instrumentation.langchain"
SCOPE_OPENINFERENCE_SMOLAGENTS = "openinference.instrumentation.smolagents"
SCOPE_ADK = "gcp.vertex.agent"
Comment thread
poshinchen marked this conversation as resolved.
SCOPE_STRANDS = "strands.telemetry.tracer"

# All scopes that should route to OpenInferenceSessionMapper
Expand Down
63 changes: 16 additions & 47 deletions src/strands_evals/mappers/langchain_otel_session_mapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,6 @@
import json
import logging
from collections import defaultdict
from datetime import datetime, timezone
from typing import Any

from ..types.trace import (
Expand Down Expand Up @@ -47,6 +46,7 @@
SCOPE_LANGCHAIN_OTEL,
)
from .session_mapper import SessionMapper
from .utils import safe_json_parse

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -248,7 +248,7 @@ def _parse_adot_body(self, span: dict) -> Any:
result = None
else:
in_content = input_messages[0].get("content", "")
result = self._safe_json_parse(in_content) if isinstance(in_content, str) else in_content
result = safe_json_parse(in_content) if isinstance(in_content, str) else in_content

if span_id:
self._adot_body_cache[span_id] = result
Expand Down Expand Up @@ -327,7 +327,7 @@ def _convert_tool_execution_span(self, span: dict, session_id: str) -> ToolExecu
# Direct inputs dict: {"inputs": {"a": 1, "b": 2}}
tool_parameters = inputs
if tool_parameters is None and ADOT_INPUT_STR_KEY in in_parsed:
params_parsed = self._safe_json_parse(in_parsed.get(ADOT_INPUT_STR_KEY, ""))
params_parsed = safe_json_parse(in_parsed.get(ADOT_INPUT_STR_KEY, ""))
if isinstance(params_parsed, dict):
tool_parameters = params_parsed

Expand All @@ -351,17 +351,17 @@ def _convert_tool_execution_span(self, span: dict, session_id: str) -> ToolExecu
entity_output = attrs.get(ATTR_TRACELOOP_ENTITY_OUTPUT, "")

if entity_input:
parsed = self._safe_json_parse(entity_input)
parsed = safe_json_parse(entity_input)
if isinstance(parsed, dict):
if "inputs" in parsed and isinstance(parsed.get("inputs"), dict):
tool_parameters = parsed.get("inputs")
elif ADOT_INPUT_STR_KEY in parsed:
params_parsed = self._safe_json_parse(parsed.get(ADOT_INPUT_STR_KEY, ""))
params_parsed = safe_json_parse(parsed.get(ADOT_INPUT_STR_KEY, ""))
if isinstance(params_parsed, dict):
tool_parameters = params_parsed

if entity_output:
parsed = self._safe_json_parse(entity_output)
parsed = safe_json_parse(entity_output)
lc_kwargs = self._extract_lc_kwargs(parsed, "output") if isinstance(parsed, dict) else None
if lc_kwargs:
tool_output_content = str(lc_kwargs.get("content", ""))
Expand Down Expand Up @@ -410,14 +410,14 @@ def _convert_agent_invocation_span(
entity_output = attrs.get(ATTR_TRACELOOP_ENTITY_OUTPUT, "")

if entity_input:
parsed = self._safe_json_parse(entity_input)
parsed = safe_json_parse(entity_input)
if isinstance(parsed, dict) and "inputs" in parsed:
inputs = parsed["inputs"]
if isinstance(inputs, dict) and "messages" in inputs:
user_query = self._get_last_message_text(inputs["messages"])

if entity_output:
parsed = self._safe_json_parse(entity_output)
parsed = safe_json_parse(entity_output)
if isinstance(parsed, dict) and "outputs" in parsed:
outputs = parsed["outputs"]
if isinstance(outputs, dict) and "messages" in outputs:
Expand Down Expand Up @@ -451,8 +451,8 @@ def _get_scope_name(self, span: dict) -> str:

def _create_span_info(self, span: dict, session_id: str) -> SpanInfo:
"""Create SpanInfo from span dict."""
start_time = self._parse_timestamp(span.get("start_time"))
end_time = self._parse_timestamp(span.get("end_time"))
start_time = self.parse_timestamp(span.get("start_time"))
end_time = self.parse_timestamp(span.get("end_time"))

return SpanInfo(
trace_id=span.get("trace_id"),
Expand All @@ -463,44 +463,13 @@ def _create_span_info(self, span: dict, session_id: str) -> SpanInfo:
end_time=end_time,
)

def _parse_timestamp(self, value: Any) -> datetime:
"""Parse timestamp from various formats."""
if value is None:
return datetime.now(timezone.utc)
if isinstance(value, datetime):
return value
if isinstance(value, str):
try:
if value.endswith("Z"):
value = value[:-1] + "+00:00"
return datetime.fromisoformat(value)
except ValueError:
return datetime.now(timezone.utc)
if isinstance(value, (int, float)):
# Handle nanoseconds
if value > 1e12:
value = value / 1e9
return datetime.fromtimestamp(value, tz=timezone.utc)
return datetime.now(timezone.utc)

def _safe_json_parse(self, content: Any) -> Any:
"""Safely parse JSON content."""
if isinstance(content, dict):
return content
if isinstance(content, str):
try:
return json.loads(content)
except json.JSONDecodeError:
return content
return content

def _parse_adot_tool_content(self, input_messages: list[dict], output_messages: list[dict]) -> tuple[Any, Any]:
"""Parse and return (input_parsed, output_parsed) from ADOT body messages."""
in_content = input_messages[-1].get("content", "")
in_parsed = self._safe_json_parse(in_content) if isinstance(in_content, str) else in_content
in_parsed = safe_json_parse(in_content) if isinstance(in_content, str) else in_content

out_content = output_messages[-1].get("content", "")
out_parsed = self._safe_json_parse(out_content) if isinstance(out_content, str) else out_content
out_parsed = safe_json_parse(out_content) if isinstance(out_content, str) else out_content

return in_parsed, out_parsed

Expand Down Expand Up @@ -639,7 +608,7 @@ def _extract_user_message(self, message: dict, index: int, attrs: dict) -> UserM
msg_content = message.get("content", "")
if isinstance(msg_content, str):
# ADOT double-encodes strings — decode outer JSON quotes if present
text = self._safe_json_parse(msg_content) if msg_content.startswith('"') else msg_content
text = safe_json_parse(msg_content) if msg_content.startswith('"') else msg_content
if not isinstance(text, str):
text = msg_content
if is_tool_msg:
Expand All @@ -658,7 +627,7 @@ def _extract_assistant_message(self, message: dict, index: int, attrs: dict) ->
msg_content = message.get("content", "")
if isinstance(msg_content, str):
# ADOT double-encodes empty strings as '""' — decode and skip if empty
text = self._safe_json_parse(msg_content) if msg_content.startswith('"') else msg_content
text = safe_json_parse(msg_content) if msg_content.startswith('"') else msg_content
if not isinstance(text, str):
text = msg_content
if text:
Expand Down Expand Up @@ -716,7 +685,7 @@ def _extract_user_prompt_from_input(self, input_messages: list[dict]) -> str | N
msg = input_messages[-1]
content = msg.get("content", "")
if isinstance(content, str):
parsed = self._safe_json_parse(content)
parsed = safe_json_parse(content)
if isinstance(parsed, dict) and "inputs" in parsed:
inputs = parsed["inputs"]
if isinstance(inputs, dict) and "messages" in inputs:
Expand All @@ -741,7 +710,7 @@ def _extract_agent_response_from_output(self, output_messages: list[dict]) -> st
if not isinstance(content, str):
return None

parsed = self._safe_json_parse(content)
parsed = safe_json_parse(content)
if not isinstance(parsed, dict):
return None

Expand Down
42 changes: 6 additions & 36 deletions src/strands_evals/mappers/openinference_session_mapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@
import json
import logging
from collections import defaultdict
from datetime import datetime, timezone
from typing import Any

from ..types.trace import (
Expand All @@ -37,6 +36,7 @@
)
from .constants import SCOPE_OPENINFERENCE_SMOLAGENTS, SCOPES_OPENINFERENCE_FAMILY
from .session_mapper import SessionMapper
from .utils import safe_json_parse

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -315,7 +315,7 @@ def _is_tool_execution_span(self, span: dict) -> bool:
input_messages, _ = self._get_messages_from_span_events(span)
if input_messages:
in_content = input_messages[0].get("content", "")
in_parsed = self._safe_json_parse(in_content) if isinstance(in_content, str) else in_content
in_parsed = safe_json_parse(in_content) if isinstance(in_content, str) else in_content
if isinstance(in_parsed, dict) and in_parsed.get("__type") == "tool_call_with_context":
return False
return True
Expand Down Expand Up @@ -361,7 +361,7 @@ def _is_agent_invocation_span(self, span: dict) -> bool:
input_messages, _ = self._get_messages_from_span_events(span)
if input_messages:
in_content = input_messages[0].get("content", "")
in_parsed = self._safe_json_parse(in_content) if isinstance(in_content, str) else in_content
in_parsed = safe_json_parse(in_content) if isinstance(in_content, str) else in_content
if isinstance(in_parsed, dict) and "messages" in in_parsed and "remaining_steps" not in in_parsed:
out_parsed = self._parse_adot_output(span)
if isinstance(out_parsed, dict) and "messages" in out_parsed:
Expand Down Expand Up @@ -599,8 +599,8 @@ def _get_scope_name(self, span: dict) -> str:

def _create_span_info(self, span: dict, session_id: str) -> SpanInfo:
"""Create SpanInfo from span dict."""
start_time = self._parse_timestamp(span.get("start_time"))
end_time = self._parse_timestamp(span.get("end_time"))
start_time = self.parse_timestamp(span.get("start_time"))
end_time = self.parse_timestamp(span.get("end_time"))

return SpanInfo(
trace_id=span.get("trace_id"),
Expand All @@ -611,36 +611,6 @@ def _create_span_info(self, span: dict, session_id: str) -> SpanInfo:
end_time=end_time,
)

def _parse_timestamp(self, value: Any) -> datetime:
"""Parse timestamp from various formats."""
if value is None:
return datetime.now(timezone.utc)
if isinstance(value, datetime):
return value
if isinstance(value, str):
try:
if value.endswith("Z"):
value = value[:-1] + "+00:00"
return datetime.fromisoformat(value)
except ValueError:
return datetime.now(timezone.utc)
if isinstance(value, (int, float)):
if value > 1e12:
value = value / 1e9
return datetime.fromtimestamp(value, tz=timezone.utc)
return datetime.now(timezone.utc)

def _safe_json_parse(self, content: Any) -> Any:
"""Safely parse JSON content."""
if isinstance(content, dict):
return content
if isinstance(content, str):
try:
return json.loads(content)
except json.JSONDecodeError:
return content
return content

def _parse_adot_output(self, span: dict) -> Any:
"""Parse the output content from the first ADOT body message.
Expand All @@ -656,7 +626,7 @@ def _parse_adot_output(self, span: dict) -> Any:
result = None
else:
out_content = output_messages[0].get("content", "")
result = self._safe_json_parse(out_content) if isinstance(out_content, str) else out_content
result = safe_json_parse(out_content) if isinstance(out_content, str) else out_content

if span_id:
self._adot_output_cache[span_id] = result
Expand Down
37 changes: 37 additions & 0 deletions src/strands_evals/mappers/session_mapper.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
"""

from abc import ABC, abstractmethod
from datetime import datetime, timezone

from typing_extensions import Any

Expand Down Expand Up @@ -72,3 +73,39 @@ def _normalize_to_flat_spans(self, data: Any) -> list[dict]:

# Fallback for unexpected types
return []

def parse_timestamp(self, value: Any) -> datetime:
"""Parse timestamp from various formats.

Handles:
- None → current UTC time
- datetime → passthrough
- ISO 8601 string (with optional trailing Z) → parsed datetime
- Numeric (int/float) nanosecond epoch → datetime
- String-encoded nanosecond epoch → datetime

Args:
value: Raw timestamp value from a span dict.

Returns:
Timezone-aware datetime in UTC.
"""
if value is None:
return datetime.now(timezone.utc)
if isinstance(value, datetime):
return value
if isinstance(value, str):
if value.isdigit():
return datetime.fromtimestamp(int(value) / 1e9, tz=timezone.utc)
try:
if value.endswith("Z"):
value = value[:-1] + "+00:00"
return datetime.fromisoformat(value)
except ValueError:
return datetime.now(timezone.utc)
if isinstance(value, (int, float)):
# Handle nanoseconds
if value > 1e12:
value = value / 1e9
return datetime.fromtimestamp(value, tz=timezone.utc)
return datetime.now(timezone.utc)
29 changes: 28 additions & 1 deletion src/strands_evals/mappers/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,35 @@
import logging
from typing import Any

from .constants import SCOPE_LANGCHAIN_OTEL, SCOPE_STRANDS, SCOPES_OPENINFERENCE_FAMILY
from .constants import SCOPE_ADK, SCOPE_LANGCHAIN_OTEL, SCOPE_STRANDS, SCOPES_OPENINFERENCE_FAMILY
from .session_mapper import SessionMapper

logger = logging.getLogger(__name__)


def safe_json_parse(content: Any) -> Any:
"""Safely parse JSON content, returning the original value on failure.

If content is already a dict, returns it as-is. If it's a string, attempts
JSON parsing and falls back to returning the raw string on decode error.
For all other types, returns the value unchanged.

Args:
content: Value to parse — typically a str or dict from span attributes.

Returns:
Parsed dict/list on success, or the original value if parsing fails or is unnecessary.
"""
if isinstance(content, dict):
return content
if isinstance(content, str):
try:
return json.loads(content)
except json.JSONDecodeError:
return content
return content


def join_tool_result_content(content: Any) -> str:
"""Join all blocks in a Bedrock-style toolResult content list into one string.

Expand Down Expand Up @@ -89,6 +112,7 @@ def detect_otel_mapper(spans: list[Any]) -> SessionMapper:
>>> session = mapper.map_to_session(spans, "session-123")
"""
# Import here to avoid circular imports
from .adk_otel_session_mapper import ADKOtelSessionMapper
from .cloudwatch_session_mapper import CloudWatchSessionMapper
from .langchain_otel_session_mapper import LangChainOtelSessionMapper
from .openinference_session_mapper import OpenInferenceSessionMapper
Expand All @@ -107,6 +131,9 @@ def detect_otel_mapper(spans: list[Any]) -> SessionMapper:
if scope_name in SCOPES_OPENINFERENCE_FAMILY:
return OpenInferenceSessionMapper()
Comment thread
liramon2 marked this conversation as resolved.

if scope_name == SCOPE_ADK:
return ADKOtelSessionMapper()

if scope_name == SCOPE_STRANDS:
# CloudWatch split format puts body on a separate entry from
# the scoped metadata entry. Break here and let the fallback
Expand Down
Loading