Skip to content

Commit 064f7dd

Browse files
add input output message format
1 parent ebc3e9a commit 064f7dd

19 files changed

Lines changed: 1231 additions & 92 deletions

libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/__init__.py

Lines changed: 47 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,13 +19,37 @@
1919
unregister_span_enricher,
2020
)
2121
from .exporters.spectra_exporter_options import SpectraExporterOptions
22-
from .inference_call_details import InferenceCallDetails, ServiceEndpoint
22+
from .inference_call_details import InferenceCallDetails
23+
from .models.service_endpoint import ServiceEndpoint
2324
from .inference_operation_type import InferenceOperationType
2425
from .inference_scope import InferenceScope
2526
from .invoke_agent_details import InvokeAgentScopeDetails
2627
from .invoke_agent_scope import InvokeAgentScope
2728
from .middleware.baggage_builder import BaggageBuilder
2829
from .models.caller_details import CallerDetails
30+
from .models.messages import (
31+
BlobPart,
32+
ChatMessage,
33+
FilePart,
34+
FinishReason,
35+
GenericPart,
36+
InputMessages,
37+
InputMessagesParam,
38+
MessagePart,
39+
MessageRole,
40+
Modality,
41+
OutputMessage,
42+
OutputMessages,
43+
OutputMessagesParam,
44+
ReasoningPart,
45+
ServerToolCallPart,
46+
ServerToolCallResponsePart,
47+
TextPart,
48+
ToolCallRequestPart,
49+
ToolCallResponsePart,
50+
UriPart,
51+
)
52+
from .models.response import Response
2953
from .models.user_details import UserDetails
3054
from .opentelemetry_scope import OpenTelemetryScope
3155
from .request import Request
@@ -70,12 +94,34 @@
7094
"ToolCallDetails",
7195
"Channel",
7296
"Request",
97+
"Response",
7398
"SpanDetails",
7499
"InferenceCallDetails",
75100
"ServiceEndpoint",
76101
# Enums
77102
"InferenceOperationType",
78103
"ToolType",
104+
# OTEL gen-ai message format types
105+
"MessageRole",
106+
"FinishReason",
107+
"Modality",
108+
"TextPart",
109+
"ToolCallRequestPart",
110+
"ToolCallResponsePart",
111+
"ReasoningPart",
112+
"BlobPart",
113+
"FilePart",
114+
"UriPart",
115+
"ServerToolCallPart",
116+
"ServerToolCallResponsePart",
117+
"GenericPart",
118+
"MessagePart",
119+
"ChatMessage",
120+
"OutputMessage",
121+
"InputMessages",
122+
"OutputMessages",
123+
"InputMessagesParam",
124+
"OutputMessagesParam",
79125
# Utility functions
80126
"extract_context_from_headers",
81127
"get_traceparent",

libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/execution_type.py

Lines changed: 0 additions & 14 deletions
This file was deleted.

libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/inference_call_details.py

Lines changed: 1 addition & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -4,17 +4,7 @@
44
from dataclasses import dataclass
55

66
from .inference_operation_type import InferenceOperationType
7-
8-
9-
@dataclass
10-
class ServiceEndpoint:
11-
"""Represents a service endpoint with hostname and optional port."""
12-
13-
hostname: str
14-
"""The hostname of the service endpoint."""
15-
16-
port: int | None = None
17-
"""The port of the service endpoint."""
7+
from .models.service_endpoint import ServiceEndpoint
188

199

2010
@dataclass

libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/inference_scope.py

Lines changed: 23 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,12 @@
2727
GEN_AI_CALLER_CLIENT_IP_KEY,
2828
)
2929
from .inference_call_details import InferenceCallDetails
30+
from .message_utils import (
31+
normalize_input_messages,
32+
normalize_output_messages,
33+
serialize_messages,
34+
)
35+
from .models.messages import InputMessagesParam, OutputMessagesParam
3036
from .models.user_details import UserDetails
3137
from .opentelemetry_scope import OpenTelemetryScope
3238
from .request import Request
@@ -97,7 +103,9 @@ def __init__(
97103
)
98104

99105
if request.content:
100-
self.set_tag_maybe(GEN_AI_INPUT_MESSAGES_KEY, request.content)
106+
# Wrap bare string into list for backward compatibility
107+
content = [request.content] if isinstance(request.content, str) else request.content
108+
self.record_input_messages(content)
101109
self.set_tag_maybe(GEN_AI_CONVERSATION_ID_KEY, request.conversation_id)
102110

103111
self.set_tag_maybe(GEN_AI_OPERATION_NAME_KEY, details.operationName.value)
@@ -138,21 +146,29 @@ def __init__(
138146
validate_and_normalize_ip(user_details.user_client_ip),
139147
)
140148

141-
def record_input_messages(self, messages: List[str]) -> None:
149+
def record_input_messages(self, messages: InputMessagesParam) -> None:
142150
"""Records the input messages for telemetry tracking.
143151
152+
Accepts plain strings (auto-wrapped as OTEL ChatMessage with role ``user``)
153+
or a versioned ``InputMessages`` wrapper.
154+
144155
Args:
145-
messages: List of input messages
156+
messages: List of input message strings or an InputMessages wrapper
146157
"""
147-
self.set_tag_maybe(GEN_AI_INPUT_MESSAGES_KEY, safe_json_dumps(messages))
158+
wrapper = normalize_input_messages(messages)
159+
self.set_tag_maybe(GEN_AI_INPUT_MESSAGES_KEY, serialize_messages(wrapper))
148160

149-
def record_output_messages(self, messages: List[str]) -> None:
161+
def record_output_messages(self, messages: OutputMessagesParam) -> None:
150162
"""Records the output messages for telemetry tracking.
151163
164+
Accepts plain strings (auto-wrapped as OTEL OutputMessage with role ``assistant``)
165+
or a versioned ``OutputMessages`` wrapper.
166+
152167
Args:
153-
messages: List of output messages
168+
messages: List of output message strings or an OutputMessages wrapper
154169
"""
155-
self.set_tag_maybe(GEN_AI_OUTPUT_MESSAGES_KEY, safe_json_dumps(messages))
170+
wrapper = normalize_output_messages(messages)
171+
self.set_tag_maybe(GEN_AI_OUTPUT_MESSAGES_KEY, serialize_messages(wrapper))
156172

157173
def record_input_tokens(self, input_tokens: int) -> None:
158174
"""Records the number of input tokens for telemetry tracking.

libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/invoke_agent_details.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,12 @@
44
# Data class for invoke agent scope details.
55

66
from dataclasses import dataclass
7-
from urllib.parse import ParseResult
7+
8+
from .models.service_endpoint import ServiceEndpoint
89

910

1011
@dataclass
1112
class InvokeAgentScopeDetails:
1213
"""Scope-level configuration for agent invocation tracing."""
1314

14-
endpoint: ParseResult | None = None
15+
endpoint: ServiceEndpoint | None = None

libraries/microsoft-agents-a365-observability-core/microsoft_agents_a365/observability/core/invoke_agent_scope.py

Lines changed: 24 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -31,11 +31,17 @@
3131
USER_NAME_KEY,
3232
)
3333
from .invoke_agent_details import InvokeAgentScopeDetails
34+
from .message_utils import (
35+
normalize_input_messages,
36+
normalize_output_messages,
37+
serialize_messages,
38+
)
3439
from .models.caller_details import CallerDetails
40+
from .models.messages import InputMessagesParam, OutputMessagesParam
3541
from .opentelemetry_scope import OpenTelemetryScope
3642
from .request import Request
3743
from .span_details import SpanDetails
38-
from .utils import safe_json_dumps, validate_and_normalize_ip
44+
from .utils import validate_and_normalize_ip
3945

4046
logger = logging.getLogger(__name__)
4147

@@ -130,7 +136,9 @@ def __init__(
130136
self.set_tag_maybe(CHANNEL_NAME_KEY, request.channel.name)
131137
self.set_tag_maybe(CHANNEL_LINK_KEY, request.channel.link)
132138
if request.content:
133-
self.set_tag_maybe(GEN_AI_INPUT_MESSAGES_KEY, safe_json_dumps([request.content]))
139+
# Wrap bare string into list for backward compatibility
140+
content = [request.content] if isinstance(request.content, str) else request.content
141+
self.record_input_messages(content)
134142

135143
# Set caller details tags
136144
if caller_details:
@@ -176,18 +184,26 @@ def record_response(self, response: str) -> None:
176184
"""
177185
self.record_output_messages([response])
178186

179-
def record_input_messages(self, messages: list[str]) -> None:
187+
def record_input_messages(self, messages: InputMessagesParam) -> None:
180188
"""Record the input messages for telemetry tracking.
181189
190+
Accepts plain strings (auto-wrapped as OTEL ChatMessage with role ``user``)
191+
or a versioned ``InputMessages`` wrapper.
192+
182193
Args:
183-
messages: List of input messages to record
194+
messages: List of input message strings or an InputMessages wrapper
184195
"""
185-
self.set_tag_maybe(GEN_AI_INPUT_MESSAGES_KEY, safe_json_dumps(messages))
196+
wrapper = normalize_input_messages(messages)
197+
self.set_tag_maybe(GEN_AI_INPUT_MESSAGES_KEY, serialize_messages(wrapper))
186198

187-
def record_output_messages(self, messages: list[str]) -> None:
199+
def record_output_messages(self, messages: OutputMessagesParam) -> None:
188200
"""Record the output messages for telemetry tracking.
189201
202+
Accepts plain strings (auto-wrapped as OTEL OutputMessage with role ``assistant``)
203+
or a versioned ``OutputMessages`` wrapper.
204+
190205
Args:
191-
messages: List of output messages to record
206+
messages: List of output message strings or an OutputMessages wrapper
192207
"""
193-
self.set_tag_maybe(GEN_AI_OUTPUT_MESSAGES_KEY, safe_json_dumps(messages))
208+
wrapper = normalize_output_messages(messages)
209+
self.set_tag_maybe(GEN_AI_OUTPUT_MESSAGES_KEY, serialize_messages(wrapper))
Lines changed: 137 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,137 @@
1+
# Copyright (c) Microsoft Corporation.
2+
# Licensed under the MIT License.
3+
4+
"""Conversion and serialization helpers for OTEL gen-ai message format.
5+
6+
Provides normalization from plain ``list[str]`` (backward compat) to the
7+
versioned wrapper format, and a non-throwing ``serialize_messages`` function.
8+
"""
9+
10+
from __future__ import annotations
11+
12+
import json
13+
import logging
14+
from dataclasses import asdict
15+
from typing import Union
16+
17+
from .models.messages import (
18+
A365_MESSAGE_SCHEMA_VERSION,
19+
ChatMessage,
20+
InputMessages,
21+
InputMessagesParam,
22+
MessageRole,
23+
OutputMessage,
24+
OutputMessages,
25+
OutputMessagesParam,
26+
TextPart,
27+
)
28+
29+
logger = logging.getLogger(__name__)
30+
31+
32+
def is_string_list(
33+
param: Union[InputMessagesParam, OutputMessagesParam],
34+
) -> bool:
35+
"""Return ``True`` when *param* is a plain ``list[str]``."""
36+
return isinstance(param, list) and all(isinstance(item, str) for item in param)
37+
38+
39+
def is_wrapped_messages(
40+
param: Union[InputMessagesParam, OutputMessagesParam],
41+
) -> bool:
42+
"""Return ``True`` when *param* is a versioned wrapper (``InputMessages`` or ``OutputMessages``)."""
43+
return isinstance(param, (InputMessages, OutputMessages))
44+
45+
46+
# ---------------------------------------------------------------------------
47+
# Plain-string → structured conversion
48+
# ---------------------------------------------------------------------------
49+
50+
51+
def to_input_messages(messages: list[str]) -> list[ChatMessage]:
52+
"""Convert plain input strings into OTEL ``ChatMessage`` objects."""
53+
return [
54+
ChatMessage(role=MessageRole.USER.value, parts=[TextPart(content=content)])
55+
for content in messages
56+
]
57+
58+
59+
def to_output_messages(messages: list[str]) -> list[OutputMessage]:
60+
"""Convert plain output strings into OTEL ``OutputMessage`` objects."""
61+
return [
62+
OutputMessage(role=MessageRole.ASSISTANT.value, parts=[TextPart(content=content)])
63+
for content in messages
64+
]
65+
66+
67+
# ---------------------------------------------------------------------------
68+
# Normalization (union → versioned wrapper)
69+
# ---------------------------------------------------------------------------
70+
71+
72+
def normalize_input_messages(param: InputMessagesParam) -> InputMessages:
73+
"""Normalize an ``InputMessagesParam`` to a versioned ``InputMessages`` wrapper.
74+
75+
- ``list[str]`` → converted to ``ChatMessage`` list and wrapped.
76+
- ``InputMessages`` → returned as-is.
77+
"""
78+
if is_string_list(param):
79+
return InputMessages(messages=to_input_messages(param)) # type: ignore[arg-type]
80+
return param # type: ignore[return-value]
81+
82+
83+
def normalize_output_messages(param: OutputMessagesParam) -> OutputMessages:
84+
"""Normalize an ``OutputMessagesParam`` to a versioned ``OutputMessages`` wrapper.
85+
86+
- ``list[str]`` → converted to ``OutputMessage`` list and wrapped.
87+
- ``OutputMessages`` → returned as-is.
88+
"""
89+
if is_string_list(param):
90+
return OutputMessages(messages=to_output_messages(param)) # type: ignore[arg-type]
91+
return param # type: ignore[return-value]
92+
93+
94+
# ---------------------------------------------------------------------------
95+
# Serialization
96+
# ---------------------------------------------------------------------------
97+
98+
99+
def _message_dict_factory(items: list[tuple[str, object]]) -> dict[str, object]:
100+
"""Custom dict factory for ``dataclasses.asdict`` that drops ``None`` values."""
101+
return {k: v for k, v in items if v is not None}
102+
103+
104+
def serialize_messages(wrapper: Union[InputMessages, OutputMessages]) -> str:
105+
"""Serialize a versioned message wrapper to JSON.
106+
107+
The output is the full wrapper object:
108+
``{"version":"0.1.0","messages":[...]}``.
109+
110+
The try/except ensures telemetry recording is non-throwing even when
111+
message parts contain non-JSON-serializable values.
112+
"""
113+
try:
114+
return json.dumps(
115+
asdict(wrapper, dict_factory=_message_dict_factory),
116+
default=str,
117+
ensure_ascii=False,
118+
)
119+
except Exception:
120+
logger.warning("Failed to serialize messages; using fallback.", exc_info=True)
121+
count = len(wrapper.messages)
122+
noun = "message" if count == 1 else "messages"
123+
fallback = {
124+
"version": A365_MESSAGE_SCHEMA_VERSION,
125+
"messages": [
126+
{
127+
"role": MessageRole.SYSTEM.value,
128+
"parts": [
129+
{
130+
"type": "text",
131+
"content": f"[serialization failed: {count} {noun}]",
132+
}
133+
],
134+
}
135+
],
136+
}
137+
return json.dumps(fallback, ensure_ascii=False)

0 commit comments

Comments
 (0)