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
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,9 @@
from opentelemetry.util.genai.invocation import InferenceInvocation
from opentelemetry.util.genai.types import (
InputMessage,
MessagePart,
OutputMessage,
SystemInstructionPart,
TextPart,
)
from opentelemetry.util.types import AttributeValue

Expand Down Expand Up @@ -125,10 +126,16 @@ def get_input_messages(

def get_system_instruction(
system: str | Iterable[TextBlockParam] | None,
) -> list[MessagePart]:
) -> list[SystemInstructionPart]:
if system is None:
return []
return convert_content_to_parts(system)
if isinstance(system, str):
return [TextPart(content=system)] if system else []
return [
TextPart(content=block["text"])
for block in system
if block.get("text")
]


def get_output_messages_from_message(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
from __future__ import annotations

import json
from collections.abc import Mapping, Sequence
from typing import Any, TypeGuard
from urllib.parse import urlparse

Expand All @@ -19,6 +20,7 @@
OutputMessage,
ReasoningPart,
Role,
SystemInstructionPart,
TextPart,
ToolCallRequestPart,
ToolCallResponsePart,
Expand Down Expand Up @@ -235,7 +237,7 @@ def extract_content_block(block: dict[str, Any]) -> MessagePart | None:
"toolRemoval",
):
if key in block:
return GenericPart(type=key, value=None)
return GenericPart(type=key)

return None

Expand All @@ -256,6 +258,30 @@ def _extract_parts(content: Any) -> list[MessagePart]:
return parts


def _extract_system_parts(
content: str | Sequence[Mapping[str, Any] | str] | None,
) -> list[SystemInstructionPart]:
if not content:
return []
if isinstance(content, str):
return [TextPart(content=content)] if content else []
parts: list[SystemInstructionPart] = []
for item in content:
if isinstance(item, str):
if item:
parts.append(TextPart(content=item))
elif _is_dict(item):
text = item.get("text")
if isinstance(text, str) and text:
parts.append(TextPart(content=text))
else:
for key in item:
if key != "text":
parts.append(GenericPart(type=key))
break
return parts


def _extract_guardrail_id(
params: dict[str, Any], invocation: InferenceInvocation
) -> None:
Expand Down Expand Up @@ -324,7 +350,7 @@ def extract_converse_request(
# system instruction
raw_system = kwargs.get("system")
if capture_content and raw_system:
system_parts = _extract_parts(raw_system)
system_parts = _extract_system_parts(raw_system)
if system_parts:
invocation.system_instruction = system_parts

Expand Down Expand Up @@ -533,7 +559,7 @@ def extract_invoke_model_request(
# System instruction (e.g. Anthropic / Nova)
raw_system = body.get("system")
if raw_system:
system_parts = _extract_parts(raw_system)
system_parts = _extract_system_parts(raw_system)
if system_parts:
invocation.system_instruction = system_parts

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
)
from opentelemetry.trace import StatusCode
from opentelemetry.util.genai.handler import TelemetryHandler
from opentelemetry.util.genai.types import GenericPart, TextPart


@pytest.mark.vcr
Expand Down Expand Up @@ -596,3 +597,47 @@ def test_extract_converse_request_prompt_variables_no_content(
== "sgi5gkybzqak"
)
assert "gen_ai.prompt.variable.user_name" not in invocation.attributes


def test_extract_converse_request_system_instruction(tracer_provider) -> None:
handler = TelemetryHandler(tracer_provider=tracer_provider)
invocation = handler.inference(provider="aws.bedrock")

extract_converse_request(
{
"system": [
{"text": "Be concise"},
{"text": "Answer politely"},
],
},
invocation,
)

assert invocation.system_instruction == [
TextPart(content="Be concise"),
TextPart(content="Answer politely"),
]


def test_extract_converse_request_system_instruction_generic(
tracer_provider,
) -> None:
handler = TelemetryHandler(tracer_provider=tracer_provider)
invocation = handler.inference(provider="aws.bedrock")

extract_converse_request(
{
"system": [
{"text": "Be concise"},
{"guardContent": {"guardrailIdentifier": "gr-123"}},
{"cachePoint": {"type": "default"}},
],
},
invocation,
)

assert invocation.system_instruction == [
TextPart(content="Be concise"),
GenericPart(type="guardContent"),
GenericPart(type="cachePoint"),
]

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
prepare_tool_definitions,
resolve_response_model_and_id,
response_fields_from_generation,
split_system_and_input_messages,
to_input_messages,
)
from opentelemetry.util.genai.handler import TelemetryHandler
from opentelemetry.util.genai.invocation import (
Expand Down Expand Up @@ -327,24 +327,21 @@ def on_chat_model_start(
if "ls_max_tokens" in metadata:
max_tokens = metadata.get("ls_max_tokens")

# Flatten ``list[list[BaseMessage]]`` (one inner list per generation
# request) before splitting into system / input.
# ``messages`` from on_chat_model_start is ``list[list[BaseMessage]]``
# (one inner list per generation request). Flatten and let
# :func:`to_input_messages` produce spec-conformant ``InputMessage`` s
# with proper roles, tool-call requests, tool results, and reasoning.
flattened: list[BaseMessage] = [msg for sub in messages for msg in sub]
system_instruction: list[MessagePart] = []
input_messages: list[InputMessage] = []
if self._telemetry_handler.should_capture_content():
system_instruction, input_messages = (
split_system_and_input_messages(flattened)
)
input_messages = to_input_messages(flattened)
Comment thread
lmolkova marked this conversation as resolved.

llm_invocation = self._telemetry_handler.inference(
provider,
request_model=request_model,
)
llm_invocation.conversation_id = _conversation_id(metadata)
llm_invocation.input_messages = input_messages
if system_instruction:
llm_invocation.system_instruction = system_instruction
llm_invocation.top_p = top_p
llm_invocation.frequency_penalty = frequency_penalty
llm_invocation.presence_penalty = presence_penalty
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -367,41 +367,6 @@ def to_input_messages(
return result


def split_system_and_input_messages(
messages: Iterable[Any],
) -> tuple[list[MessagePart], list[InputMessage]]:
"""Split ``messages`` into ``system_instruction`` parts and ``InputMessage`` s.

Called only when content capture is enabled
(``TelemetryHandler.should_capture_content()``).
"""
materialized = list(messages)
try:
normalized: Iterable[BaseMessage] = convert_to_messages(materialized)
except Exception: # pylint: disable=broad-except
normalized = [m for m in materialized if isinstance(m, BaseMessage)]

system_parts: list[MessagePart] = []
input_messages: list[InputMessage] = []

for message in normalized:
if isinstance(message, SystemMessage):
system_parts.extend(_content_to_parts(message.content))
else:
parts = _message_parts(message)
if not parts and not _has_content(message):
continue
input_messages.append(
InputMessage(
role=_normalize_role(message) or Role.USER.value,
parts=parts,
name=_message_name(message),
)
)

return system_parts, input_messages


def to_output_messages(
messages: Iterable[BaseMessage],
*,
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -16,15 +16,12 @@
from opentelemetry.sdk._logs import LoggerProvider
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.test.weaver_live_check import LiveCheckReport
from opentelemetry.test_util_genai.conformance import (
ExpectedViolation,
Scenario,
)
from opentelemetry.test_util_genai.instrumentor import instrument

from ._shared import span_attribute_values


class InferenceScenario(Scenario):
expected_spans = {"chat": 1}
Expand All @@ -40,20 +37,6 @@ class InferenceScenario(Scenario):
),
)

def validate(self, report: LiveCheckReport) -> None:
super().validate(report)
system_instructions = span_attribute_values(
report, "gen_ai.system_instructions"
)
assert len(system_instructions) == 1, (
"chat span with a SystemMessage input should set "
f"gen_ai.system_instructions once; saw {system_instructions}"
)
assert "You are a helpful assistant!" in system_instructions[0], (
"gen_ai.system_instructions should carry the SystemMessage "
f"content; got {system_instructions[0]}"
)

def run(
self,
*,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,6 @@
)
from opentelemetry.test_util_genai.instrumentor import instrument

from ._shared import span_attribute_values


class InferenceStreamingScenario(Scenario):
expected_spans = {"chat": 1}
Expand All @@ -55,22 +53,17 @@ class InferenceStreamingScenario(Scenario):

def validate(self, report: LiveCheckReport) -> None:
super().validate(report)
stream_values = span_attribute_values(report, "gen_ai.request.stream")
stream_values = [
attr["value"]
for entry in report["samples"]
if "span" in entry
for attr in entry["span"]["attributes"]
if attr["name"] == "gen_ai.request.stream"
]
assert stream_values == [True], (
"streaming chat should set gen_ai.request.stream=true on the chat "
f"span; saw {stream_values}"
)
system_instructions = span_attribute_values(
report, "gen_ai.system_instructions"
)
assert len(system_instructions) == 1, (
"streaming chat span with a SystemMessage input should set "
f"gen_ai.system_instructions once; saw {system_instructions}"
)
assert "You are a helpful assistant!" in system_instructions[0], (
"gen_ai.system_instructions should carry the SystemMessage "
f"content; got {system_instructions[0]}"
)

def run(
self,
Expand Down
Loading