Compare commits

...

25 Commits

Author SHA1 Message Date
nityanandagohain
7b51c47e6d fix: refactor 2026-10-09 21:54:18 +05:30
nityanandagohain
b9565c4866 fix: more cleanup 2026-10-09 18:33:33 +05:30
nityanandagohain
1096a9f70c fix: update tests 2026-10-09 18:18:37 +05:30
nityanandagohain
b19ba11d98 fix: format file 2026-10-09 18:01:16 +05:30
nityanandagohain
51663240e0 fix: more fixes 2026-10-09 17:48:08 +05:30
nityanandagohain
143c3a64f1 fix: more changes 2026-10-09 17:16:07 +05:30
nityanandagohain
b4da88d5ec fix: more updates 2026-10-09 16:42:03 +05:30
nityanandagohain
277124320a fix: update struct 2026-10-09 16:21:26 +05:30
nityanandagohain
9a11ebf350 Merge remote-tracking branch 'origin/main' into feat/ai-trace-thread-normalize 2026-10-09 13:10:11 +05:30
nityanandagohain
a4cb6ca847 Merge remote-tracking branch 'origin/main' into feat/ai-trace-thread-normalize 2026-10-07 17:54:46 +05:30
nityanandagohain
f50d786155 Merge branch 'feat/ai-trace-thread' into feat/ai-trace-thread-normalize 2026-10-07 15:45:01 +05:30
nityanandagohain
bbd9fd1ad5 fix: address comments 2026-10-07 15:32:06 +05:30
nityanandagohain
b0b3409827 Merge remote-tracking branch 'origin/main' into feat/ai-trace-thread 2026-10-07 15:15:35 +05:30
nityanandagohain
37932c7e84 fix: address comments 2026-10-07 12:10:55 +05:30
nityanandagohain
7313b02455 Merge branch 'feat/ai-trace-thread' into feat/ai-trace-thread-normalize 2026-10-05 16:08:42 +05:30
nityanandagohain
a1a38fc7da Merge remote-tracking branch 'origin/main' into feat/ai-trace-thread 2026-10-05 16:03:40 +05:30
nityanandagohain
f76e3e3085 fix: update comments and tests 2026-10-05 16:02:31 +05:30
nityanandagohain
5b9d854518 Merge branch 'feat/ai-trace-thread' into feat/ai-trace-thread-normalize
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-01 20:05:05 +05:30
nityanandagohain
0a9920583b fix: api changes 2026-10-01 19:59:16 +05:30
nityanandagohain
ae0d8cb1a4 fix: cleanup 2026-09-29 16:13:20 +05:30
nityanandagohain
abb242b1f9 feat(tracedetail): normalise gen_ai messages in thread spans 2026-09-29 12:47:25 +05:30
nityanandagohain
61f6370f45 fix: remove changes from waterfall 2026-09-29 12:42:30 +05:30
nityanandagohain
f8b22c0feb refactor(tracedetail): split message normalisation out of the thread api 2026-09-29 11:20:56 +05:30
nityanandagohain
9e3a6bb35d Merge remote-tracking branch 'origin/main' into feat/ai-trace-thread 2026-09-29 10:25:57 +05:30
nityanandagohain
f5f019f61b feat: trace detail thread endpoint 2026-09-28 17:32:55 +05:30
19 changed files with 1883 additions and 93 deletions

View File

@@ -5099,6 +5099,220 @@ components:
required:
- config
type: object
GenaiBlobPart:
properties:
content:
type: string
mime_type:
nullable: true
type: string
modality:
type: string
type:
type: string
required:
- type
- modality
- content
type: object
GenaiChatMessage:
properties:
name:
nullable: true
type: string
parts:
$ref: '#/components/schemas/GenaiParts'
role:
type: string
required:
- role
- parts
type: object
GenaiCompactionPart:
properties:
content:
nullable: true
type: string
id:
nullable: true
type: string
type:
type: string
required:
- type
type: object
GenaiFilePart:
properties:
file_id:
type: string
mime_type:
nullable: true
type: string
modality:
type: string
type:
type: string
required:
- type
- modality
- file_id
type: object
GenaiGenericPart:
additionalProperties: {}
type: object
GenaiGenericServerToolCall:
additionalProperties: {}
type: object
GenaiGenericServerToolCallResponse:
additionalProperties: {}
type: object
GenaiInputMessages:
items:
$ref: '#/components/schemas/GenaiChatMessage'
type: array
GenaiOutputMessage:
properties:
finish_reason:
nullable: true
type: string
name:
nullable: true
type: string
parts:
$ref: '#/components/schemas/GenaiParts'
role:
type: string
required:
- role
- parts
type: object
GenaiOutputMessages:
items:
$ref: '#/components/schemas/GenaiOutputMessage'
type: array
GenaiPart:
discriminator:
mapping:
blob: '#/components/schemas/GenaiBlobPart'
compaction: '#/components/schemas/GenaiCompactionPart'
file: '#/components/schemas/GenaiFilePart'
reasoning: '#/components/schemas/GenaiReasoningPart'
server_tool_call: '#/components/schemas/GenaiServerToolCallPart'
server_tool_call_response: '#/components/schemas/GenaiServerToolCallResponsePart'
text: '#/components/schemas/GenaiTextPart'
tool_call: '#/components/schemas/GenaiToolCallRequestPart'
tool_call_response: '#/components/schemas/GenaiToolCallResponsePart'
uri: '#/components/schemas/GenaiUriPart'
propertyName: type
oneOf:
- $ref: '#/components/schemas/GenaiTextPart'
- $ref: '#/components/schemas/GenaiToolCallRequestPart'
- $ref: '#/components/schemas/GenaiToolCallResponsePart'
- $ref: '#/components/schemas/GenaiServerToolCallPart'
- $ref: '#/components/schemas/GenaiServerToolCallResponsePart'
- $ref: '#/components/schemas/GenaiBlobPart'
- $ref: '#/components/schemas/GenaiFilePart'
- $ref: '#/components/schemas/GenaiUriPart'
- $ref: '#/components/schemas/GenaiReasoningPart'
- $ref: '#/components/schemas/GenaiCompactionPart'
- $ref: '#/components/schemas/GenaiGenericPart'
type: object
GenaiParts:
items:
$ref: '#/components/schemas/GenaiPart'
type: array
GenaiReasoningPart:
properties:
content:
type: string
type:
type: string
required:
- type
- content
type: object
GenaiServerToolCallPart:
properties:
id:
nullable: true
type: string
name:
type: string
server_tool_call:
$ref: '#/components/schemas/GenaiGenericServerToolCall'
type:
type: string
required:
- type
- name
- server_tool_call
type: object
GenaiServerToolCallResponsePart:
properties:
id:
nullable: true
type: string
server_tool_call_response:
$ref: '#/components/schemas/GenaiGenericServerToolCallResponse'
type:
type: string
required:
- type
- server_tool_call_response
type: object
GenaiTextPart:
properties:
content:
type: string
type:
type: string
required:
- type
- content
type: object
GenaiToolCallRequestPart:
properties:
arguments: {}
id:
nullable: true
type: string
name:
type: string
type:
type: string
required:
- type
- name
type: object
GenaiToolCallResponsePart:
properties:
id:
nullable: true
type: string
response:
nullable: true
type:
type: string
required:
- type
- response
type: object
GenaiUriPart:
properties:
mime_type:
nullable: true
type: string
modality:
type: string
type:
type: string
uri:
type: string
required:
- type
- modality
- uri
type: object
GlobaltypesAPIKeyConfig:
properties:
enabled:
@@ -10217,6 +10431,16 @@ components:
items:
$ref: '#/components/schemas/SpantypesEvent'
type: array
formatted_input:
$ref: '#/components/schemas/GenaiInputMessages'
formatted_output:
$ref: '#/components/schemas/GenaiOutputMessages'
formatter:
type: string
formatter_warnings:
items:
type: string
type: array
has_error:
type: boolean
kind_string:
@@ -10259,6 +10483,10 @@ components:
- attributes
- events
- references
- formatted_input
- formatted_output
- formatter
- formatter_warnings
type: object
SpantypesTraceAISummary:
properties:
@@ -16080,9 +16308,12 @@ paths:
/api/v1/traces/{traceID}/thread:
get:
deprecated: false
description: Returns the spans carrying gen_ai input or output messages in timestamp
order. Pass nextCursor as after or prevCursor as before to page, or spanId
to open the page around a span.
description: Returns the spans carrying gen_ai input or output messages, or
a tool execution, in timestamp order. Messages already in the OTel GenAI shape
are decoded into formatted_input and formatted_output, with formatter naming
the converter and formatter_warnings what it could not resolve. Pass nextCursor
as after or prevCursor as before to page, or spanId to open the page around
a span.
operationId: GetTraceThread
parameters:
- description: Page size, at most 100. 0 means 20.

View File

@@ -6675,6 +6675,260 @@ export interface GatewaytypesUpdatableIngestionKeyLimitDTO {
tags?: string[] | null;
}
export enum GenaiBlobPartDTOType {
blob = 'blob',
}
export interface GenaiBlobPartDTO {
/**
* @type string
*/
content: string;
/**
* @type string,null
*/
mime_type?: string | null;
/**
* @type string
*/
modality: string;
/**
* @type string
* @enum blob
*/
type: GenaiBlobPartDTOType;
}
export enum GenaiTextPartDTOType {
text = 'text',
}
export interface GenaiTextPartDTO {
/**
* @type string
*/
content: string;
/**
* @type string
* @enum text
*/
type: GenaiTextPartDTOType;
}
export enum GenaiToolCallRequestPartDTOType {
tool_call = 'tool_call',
}
export interface GenaiToolCallRequestPartDTO {
arguments?: unknown;
/**
* @type string,null
*/
id?: string | null;
/**
* @type string
*/
name: string;
/**
* @type string
* @enum tool_call
*/
type: GenaiToolCallRequestPartDTOType;
}
export enum GenaiToolCallResponsePartDTOType {
tool_call_response = 'tool_call_response',
}
export type GenaiToolCallResponsePartDTOResponse = unknown | null;
export interface GenaiToolCallResponsePartDTO {
/**
* @type string,null
*/
id?: string | null;
/**
* @nullable true
*/
response: GenaiToolCallResponsePartDTOResponse;
/**
* @type string
* @enum tool_call_response
*/
type: GenaiToolCallResponsePartDTOType;
}
export interface GenaiGenericServerToolCallDTO {
[key: string]: unknown;
}
export enum GenaiServerToolCallPartDTOType {
server_tool_call = 'server_tool_call',
}
export interface GenaiServerToolCallPartDTO {
/**
* @type string,null
*/
id?: string | null;
/**
* @type string
*/
name: string;
server_tool_call: GenaiGenericServerToolCallDTO;
/**
* @type string
* @enum server_tool_call
*/
type: GenaiServerToolCallPartDTOType;
}
export interface GenaiGenericServerToolCallResponseDTO {
[key: string]: unknown;
}
export enum GenaiServerToolCallResponsePartDTOType {
server_tool_call_response = 'server_tool_call_response',
}
export interface GenaiServerToolCallResponsePartDTO {
/**
* @type string,null
*/
id?: string | null;
server_tool_call_response: GenaiGenericServerToolCallResponseDTO;
/**
* @type string
* @enum server_tool_call_response
*/
type: GenaiServerToolCallResponsePartDTOType;
}
export enum GenaiFilePartDTOType {
file = 'file',
}
export interface GenaiFilePartDTO {
/**
* @type string
*/
file_id: string;
/**
* @type string,null
*/
mime_type?: string | null;
/**
* @type string
*/
modality: string;
/**
* @type string
* @enum file
*/
type: GenaiFilePartDTOType;
}
export enum GenaiUriPartDTOType {
uri = 'uri',
}
export interface GenaiUriPartDTO {
/**
* @type string,null
*/
mime_type?: string | null;
/**
* @type string
*/
modality: string;
/**
* @type string
* @enum uri
*/
type: GenaiUriPartDTOType;
/**
* @type string
*/
uri: string;
}
export enum GenaiReasoningPartDTOType {
reasoning = 'reasoning',
}
export interface GenaiReasoningPartDTO {
/**
* @type string
*/
content: string;
/**
* @type string
* @enum reasoning
*/
type: GenaiReasoningPartDTOType;
}
export enum GenaiCompactionPartDTOType {
compaction = 'compaction',
}
export interface GenaiCompactionPartDTO {
/**
* @type string,null
*/
content?: string | null;
/**
* @type string,null
*/
id?: string | null;
/**
* @type string
* @enum compaction
*/
type: GenaiCompactionPartDTOType;
}
export interface GenaiGenericPartDTO {
[key: string]: unknown;
}
export type GenaiPartDTO =
| GenaiTextPartDTO
| GenaiToolCallRequestPartDTO
| GenaiToolCallResponsePartDTO
| GenaiServerToolCallPartDTO
| GenaiServerToolCallResponsePartDTO
| GenaiBlobPartDTO
| GenaiFilePartDTO
| GenaiUriPartDTO
| GenaiReasoningPartDTO
| GenaiCompactionPartDTO
| GenaiGenericPartDTO;
export type GenaiPartsDTO = GenaiPartDTO[];
export interface GenaiChatMessageDTO {
/**
* @type string,null
*/
name?: string | null;
parts: GenaiPartsDTO;
/**
* @type string
*/
role: string;
}
export type GenaiInputMessagesDTO = GenaiChatMessageDTO[];
export interface GenaiOutputMessageDTO {
/**
* @type string,null
*/
finish_reason?: string | null;
/**
* @type string,null
*/
name?: string | null;
parts: GenaiPartsDTO;
/**
* @type string
*/
role: string;
}
export type GenaiOutputMessagesDTO = GenaiOutputMessageDTO[];
export interface GlobaltypesAPIKeyConfigDTO {
/**
* @type boolean
@@ -11444,6 +11698,16 @@ export interface SpantypesThreadSpanDTO {
* @type array
*/
events: SpantypesEventDTO[];
formatted_input: GenaiInputMessagesDTO;
formatted_output: GenaiOutputMessagesDTO;
/**
* @type string
*/
formatter: string;
/**
* @type array
*/
formatter_warnings: string[];
/**
* @type boolean
*/

View File

@@ -261,7 +261,7 @@ export const invalidateGetTraceSummary = async (
};
/**
* Returns the spans carrying gen_ai input or output messages in timestamp order. Pass nextCursor as after or prevCursor as before to page, or spanId to open the page around a span.
* Returns the spans carrying gen_ai input or output messages, or a tool execution, in timestamp order. Messages already in the OTel GenAI shape are decoded into formatted_input and formatted_output, with formatter naming the converter and formatter_warnings what it could not resolve. Pass nextCursor as after or prevCursor as before to page, or spanId to open the page around a span.
* @summary Get thread view for a trace
*/
export const getTraceThread = (

View File

@@ -90,7 +90,7 @@ func (provider *provider) addTraceDetailRoutes(router *mux.Router) error {
ID: "GetTraceThread",
Tags: []string{"tracedetail"},
Summary: "Get thread view for a trace",
Description: "Returns the spans carrying gen_ai input or output messages in timestamp order. Pass nextCursor as after or prevCursor as before to page, or spanId to open the page around a span.",
Description: "Returns the spans carrying gen_ai input or output messages, or a tool execution, in timestamp order. Messages already in the OTel GenAI shape are decoded into formatted_input and formatted_output, with formatter naming the converter and formatter_warnings what it could not resolve. Pass nextCursor as after or prevCursor as before to page, or spanId to open the page around a span.",
RequestQuery: new(spantypes.GetTraceThreadParams),
Response: new(spantypes.GettableTraceThread),
ResponseContentType: "application/json",

View File

@@ -158,10 +158,6 @@ func (s *traceStore) GetTraceStats(ctx context.Context, orgID valuer.UUID, trace
// genAISpanColumns returns the gen_ai columns aggregated per span, resolved across attribute evolutions.
func (s *traceStore) genAISpanColumns(ctx context.Context, orgID valuer.UUID, bounds *spantypes.TraceBounds, sb *sqlbuilder.SelectBuilder) ([]string, error) {
attributeKey := func(name string, dataType telemetrytypes.FieldDataType) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{Name: name, Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: dataType}
}
values := []struct{ key, alias string }{
{aiobservabilitytypes.GenAIUsageInputTokens, "input_tokens_value"},
{aiobservabilitytypes.GenAIUsageOutputTokens, "output_tokens_value"},
@@ -175,26 +171,17 @@ func (s *traceStore) genAISpanColumns(ctx context.Context, orgID valuer.UUID, bo
for _, value := range values {
names = append(names, value.key)
}
selectors := make([]*telemetrytypes.FieldKeySelector, 0, len(names))
for _, name := range names {
selectors = append(selectors, &telemetrytypes.FieldKeySelector{Name: name, Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact})
}
keys, _, err := s.metadataStore.GetKeysMulti(ctx, orgID, querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, s.flagger, selectors))
keys, err := s.genAIFieldKeys(ctx, orgID, names)
if err != nil {
return nil, err
}
q := querybuilder.NewQueryInfo(ctx, orgID, s.flagger, telemetrytypes.SignalTraces, nil, uint64(bounds.Start.UnixNano()), uint64(bounds.End.UnixNano()))
gate := make([]string, 0, len(aiobservabilitytypes.GenAISpanGateKeys))
for _, name := range aiobservabilitytypes.GenAISpanGateKeys {
conds, _, err := querybuilder.Conditions(ctx, q, s.storage, attributeKey(name, telemetrytypes.FieldDataTypeString), qbtypes.FilterOperatorExists, nil, keys, false, sb)
if err != nil {
return nil, err
}
gate = append(gate, conds...)
isGenAI, err := s.anyExistsCondition(ctx, q, aiobservabilitytypes.GenAISpanGateKeys, keys, sb)
if err != nil {
return nil, err
}
columns := []string{sb.Or(gate...) + " AS is_gen_ai"}
columns := []string{isGenAI + " AS is_gen_ai"}
for _, value := range values {
// lookup by number, the type metadata stores numeric attributes under; float64 is only the output cast
@@ -310,7 +297,11 @@ func (s *traceStore) GetThreadSpans(ctx context.Context, orgID valuer.UUID, trac
sb.SelectMore("attributes")
}
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
hasMessages, err := s.messagesExistCondition(ctx, q, orgID, bounds, sb)
keys, err := s.genAIFieldKeys(ctx, orgID, aiobservabilitytypes.GenAIThreadKeys)
if err != nil {
return nil, err
}
isThreadSpan, err := s.anyExistsCondition(ctx, q, aiobservabilitytypes.GenAIThreadKeys, keys, sb)
if err != nil {
return nil, err
}
@@ -318,7 +309,7 @@ func (s *traceStore) GetThreadSpans(ctx context.Context, orgID valuer.UUID, trac
sb.E("trace_id", traceID),
sb.GE("ts_bucket_start", bounds.Start.Unix()-1800),
sb.LE("ts_bucket_start", bounds.End.Unix()),
hasMessages,
isThreadSpan,
)
if cursor := page.Cursor; cursor != nil {
// ClickHouse can't use an index for a tuple comparison, so the separate timestamp and
@@ -355,38 +346,6 @@ func (s *traceStore) GetThreadSpans(ctx context.Context, orgID valuer.UUID, trac
return spans, nil
}
// messagesExistCondition resolves the gen_ai message keys through the attribute evolution metadata
// and the use_trace_attributes_json flag, so the filter reads the same columns the query builder does.
func (s *traceStore) messagesExistCondition(ctx context.Context, q qbtypes.QueryInfo, orgID valuer.UUID, bounds *spantypes.TraceBounds, sb *sqlbuilder.SelectBuilder) (string, error) {
names := []string{aiobservabilitytypes.GenAIInputMessages, aiobservabilitytypes.GenAIOutputMessages}
selectors := make([]*telemetrytypes.FieldKeySelector, len(names))
for i, name := range names {
selectors[i] = &telemetrytypes.FieldKeySelector{
StartUnixMilli: bounds.Start.UnixMilli(),
EndUnixMilli: bounds.End.UnixMilli(),
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
Name: name,
SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact,
}
}
fieldKeys, _, err := s.metadataStore.GetKeysMulti(ctx, orgID, selectors)
if err != nil {
return "", errors.WrapInternalf(err, errors.CodeInternal, "error fetching thread field keys")
}
conds := make([]string, 0, len(names))
for _, name := range names {
key := &telemetrytypes.TelemetryFieldKey{Name: name, Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute}
keyConds, _, err := querybuilder.Conditions(ctx, q, s.storage, key, qbtypes.FilterOperatorExists, nil, fieldKeys, false, sb)
if err != nil {
return "", err
}
conds = append(conds, keyConds...)
}
return sb.Or(conds...), nil
}
func (s *traceStore) GetThreadCursor(ctx context.Context, traceID string, bounds *spantypes.TraceBounds, spanID string) (*spantypes.ThreadCursor, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("toUnixTimestamp64Nano(timestamp)")
@@ -532,3 +491,34 @@ func (s *traceStore) GetSpanDurationByField(ctx context.Context, traceID string,
}
return result, nil
}
func attributeKey(name string, dataType telemetrytypes.FieldDataType) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{Name: name, Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: dataType}
}
// genAIFieldKeys fetches the metadata keys for the attribute names, expanded across semconv families.
func (s *traceStore) genAIFieldKeys(ctx context.Context, orgID valuer.UUID, names []string) (map[string][]*telemetrytypes.TelemetryFieldKey, error) {
selectors := make([]*telemetrytypes.FieldKeySelector, 0, len(names))
for _, name := range names {
selectors = append(selectors, &telemetrytypes.FieldKeySelector{Name: name, Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextAttribute, SelectorMatchType: telemetrytypes.FieldSelectorMatchTypeExact})
}
keys, _, err := s.metadataStore.GetKeysMulti(ctx, orgID, querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, s.flagger, selectors))
if err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error fetching gen_ai field keys")
}
return keys, nil
}
// anyExistsCondition ORs an EXISTS test per attribute name, resolved across attribute evolutions
// and the use_trace_attributes_json flag so the filter reads the same columns the query builder does.
func (s *traceStore) anyExistsCondition(ctx context.Context, q qbtypes.QueryInfo, names []string, keys map[string][]*telemetrytypes.TelemetryFieldKey, sb *sqlbuilder.SelectBuilder) (string, error) {
conds := make([]string, 0, len(names))
for _, name := range names {
keyConds, _, err := querybuilder.Conditions(ctx, q, s.storage, attributeKey(name, telemetrytypes.FieldDataTypeString), qbtypes.FilterOperatorExists, nil, keys, false, sb)
if err != nil {
return "", err
}
conds = append(conds, keyConds...)
}
return sb.Or(conds...), nil
}

View File

@@ -220,14 +220,14 @@ func TestGetThreadSpans(t *testing.T) {
name: "FlagOff_ReadsAndFiltersLegacyMaps",
jsonRelease: &jsonInsideTrace,
selectSQL: selectSQL,
whereSQL: "(mapContains(attributes_string, 'gen_ai.input.messages') OR mapContains(attributes_string, 'gen_ai.output.messages'))",
whereSQL: "(mapContains(attributes_string, 'gen_ai.input.messages') OR mapContains(attributes_string, 'gen_ai.output.messages') OR mapContains(attributes_string, 'gen_ai.tool.name'))",
},
{
name: "FlagOn_ReleasedDuringTrace_ReadsJSONFiltersJSONThenMaps",
jsonOn: true,
jsonRelease: &jsonInsideTrace,
selectSQL: selectSQL + ", attributes",
whereSQL: "(multiIf(attributes.`gen_ai.input.messages` IS NOT NULL, attributes.`gen_ai.input.messages`::String, mapContains(attributes_string, 'gen_ai.input.messages'), attributes_string['gen_ai.input.messages'], NULL) IS NOT NULL OR multiIf(attributes.`gen_ai.output.messages` IS NOT NULL, attributes.`gen_ai.output.messages`::String, mapContains(attributes_string, 'gen_ai.output.messages'), attributes_string['gen_ai.output.messages'], NULL) IS NOT NULL)",
whereSQL: "(multiIf(attributes.`gen_ai.input.messages` IS NOT NULL, attributes.`gen_ai.input.messages`::String, mapContains(attributes_string, 'gen_ai.input.messages'), attributes_string['gen_ai.input.messages'], NULL) IS NOT NULL OR multiIf(attributes.`gen_ai.output.messages` IS NOT NULL, attributes.`gen_ai.output.messages`::String, mapContains(attributes_string, 'gen_ai.output.messages'), attributes_string['gen_ai.output.messages'], NULL) IS NOT NULL OR multiIf(attributes.`gen_ai.tool.name` IS NOT NULL, attributes.`gen_ai.tool.name`::String, mapContains(attributes_string, 'gen_ai.tool.name'), attributes_string['gen_ai.tool.name'], NULL) IS NOT NULL)",
},
}

View File

@@ -0,0 +1,57 @@
// Package genai mirrors the OpenTelemetry GenAI message schemas, one type per schema definition:
// open-telemetry/semantic-conventions-genai@06ec68e, model/gen-ai/gen-ai-{input,output}-messages.json.
// Keys the schema does not name are dropped; the raw attribute on the span still has them.
package genai
import "github.com/SigNoz/signoz/pkg/valuer"
var (
RoleSystem = Role{valuer.NewString("system")}
RoleUser = Role{valuer.NewString("user")}
RoleAssistant = Role{valuer.NewString("assistant")}
RoleTool = Role{valuer.NewString("tool")}
)
var (
FinishReasonStop = FinishReason{valuer.NewString("stop")}
FinishReasonLength = FinishReason{valuer.NewString("length")}
FinishReasonContentFilter = FinishReason{valuer.NewString("content_filter")}
FinishReasonToolCall = FinishReason{valuer.NewString("tool_call")}
FinishReasonCompaction = FinishReason{valuer.NewString("compaction")}
FinishReasonError = FinishReason{valuer.NewString("error")}
)
var (
ModalityImage = Modality{valuer.NewString("image")}
ModalityVideo = Modality{valuer.NewString("video")}
ModalityAudio = Modality{valuer.NewString("audio")}
ModalityDocument = Modality{valuer.NewString("document")}
)
// Role, FinishReason and Modality are open enums: the variables are the values the schema
// names, any other string is kept as sent, so none of them lists an Enum for the spec.
type Role struct{ valuer.String }
type FinishReason struct{ valuer.String }
type Modality struct{ valuer.String }
type InputMessages []ChatMessage
// OutputMessages holds one message per choice.
type OutputMessages []OutputMessage
type ChatMessage struct {
Role Role `json:"role" required:"true"`
Parts Parts `json:"parts" required:"true" nullable:"false"`
Name *string `json:"name,omitempty"`
}
// OutputMessage is one choice. The schema deprecates FinishReason in favour of
// gen_ai.response.finish_reasons; it stays so each choice carries its own.
type OutputMessage struct {
Role Role `json:"role" required:"true"`
Parts Parts `json:"parts" required:"true" nullable:"false"`
Name *string `json:"name,omitempty"`
FinishReason *FinishReason `json:"finish_reason,omitempty"`
}

View File

@@ -0,0 +1,194 @@
package genai
import (
"encoding/json"
"github.com/swaggest/jsonschema-go"
"github.com/SigNoz/signoz/pkg/valuer"
)
var (
PartTypeText = PartType{valuer.NewString("text")}
PartTypeToolCall = PartType{valuer.NewString("tool_call")}
PartTypeToolCallResponse = PartType{valuer.NewString("tool_call_response")}
PartTypeServerToolCall = PartType{valuer.NewString("server_tool_call")}
PartTypeServerToolCallResponse = PartType{valuer.NewString("server_tool_call_response")}
PartTypeBlob = PartType{valuer.NewString("blob")}
PartTypeFile = PartType{valuer.NewString("file")}
PartTypeURI = PartType{valuer.NewString("uri")}
PartTypeReasoning = PartType{valuer.NewString("reasoning")}
PartTypeCompaction = PartType{valuer.NewString("compaction")}
)
var (
_ jsonschema.OneOfExposer = Part{}
_ jsonschema.Preparer = Part{}
)
// partVariants are the part types the schema names; any other type decodes as GenericPart.
// schemaRef is the component in the generated OpenAPI spec, for the discriminator mapping.
var partVariants = []partVariant{
{typ: PartTypeText, decode: decodePart[TextPart], schema: TextPart{}, schemaRef: "#/components/schemas/GenaiTextPart"},
{typ: PartTypeToolCall, decode: decodePart[ToolCallRequestPart], schema: ToolCallRequestPart{}, schemaRef: "#/components/schemas/GenaiToolCallRequestPart"},
{typ: PartTypeToolCallResponse, decode: decodePart[ToolCallResponsePart], schema: ToolCallResponsePart{}, schemaRef: "#/components/schemas/GenaiToolCallResponsePart"},
{typ: PartTypeServerToolCall, decode: decodePart[ServerToolCallPart], schema: ServerToolCallPart{}, schemaRef: "#/components/schemas/GenaiServerToolCallPart"},
{typ: PartTypeServerToolCallResponse, decode: decodePart[ServerToolCallResponsePart], schema: ServerToolCallResponsePart{}, schemaRef: "#/components/schemas/GenaiServerToolCallResponsePart"},
{typ: PartTypeBlob, decode: decodePart[BlobPart], schema: BlobPart{}, schemaRef: "#/components/schemas/GenaiBlobPart"},
{typ: PartTypeFile, decode: decodePart[FilePart], schema: FilePart{}, schemaRef: "#/components/schemas/GenaiFilePart"},
{typ: PartTypeURI, decode: decodePart[UriPart], schema: UriPart{}, schemaRef: "#/components/schemas/GenaiUriPart"},
{typ: PartTypeReasoning, decode: decodePart[ReasoningPart], schema: ReasoningPart{}, schemaRef: "#/components/schemas/GenaiReasoningPart"},
{typ: PartTypeCompaction, decode: decodePart[CompactionPart], schema: CompactionPart{}, schemaRef: "#/components/schemas/GenaiCompactionPart"},
}
type PartType struct{ valuer.String }
// Part is one message part, the schema's anyOf discriminated on "type". Value holds the part
// struct for that type and is encoded as the part itself.
type Part struct {
Value any `json:"-"`
}
type Parts []Part
type TextPart struct {
Type PartType `json:"type" required:"true"`
Content string `json:"content" required:"true"`
}
type ToolCallRequestPart struct {
Type PartType `json:"type" required:"true"`
Name string `json:"name" required:"true"`
ID *string `json:"id,omitempty"`
Arguments any `json:"arguments,omitempty"`
}
// ToolCallResponsePart carries a client tool result, or a built-in tool outcome.
type ToolCallResponsePart struct {
Type PartType `json:"type" required:"true"`
Response any `json:"response" required:"true" nullable:"true"`
ID *string `json:"id,omitempty"`
}
// ServerToolCallPart is a tool the provider ran itself, such as code_interpreter or web_search.
type ServerToolCallPart struct {
Type PartType `json:"type" required:"true"`
Name string `json:"name" required:"true"`
ID *string `json:"id,omitempty"`
ServerToolCall GenericServerToolCall `json:"server_tool_call" required:"true" nullable:"false"`
}
// GenericServerToolCall is the provider's own call object, {type, ...} with any other keys.
type GenericServerToolCall map[string]any
type ServerToolCallResponsePart struct {
Type PartType `json:"type" required:"true"`
ID *string `json:"id,omitempty"`
ServerToolCallResponse GenericServerToolCallResponse `json:"server_tool_call_response" required:"true" nullable:"false"`
}
// GenericServerToolCallResponse is the provider's own response object, {type, ...} with any other keys.
type GenericServerToolCallResponse map[string]any
// BlobPart is inline data; Content is base64.
type BlobPart struct {
Type PartType `json:"type" required:"true"`
Modality Modality `json:"modality" required:"true"`
Content string `json:"content" required:"true"`
MimeType *string `json:"mime_type,omitempty"`
}
// FilePart references a file uploaded to the provider.
type FilePart struct {
Type PartType `json:"type" required:"true"`
Modality Modality `json:"modality" required:"true"`
FileID string `json:"file_id" required:"true"`
MimeType *string `json:"mime_type,omitempty"`
}
// UriPart references external data; a data: URL belongs in BlobPart instead.
type UriPart struct {
Type PartType `json:"type" required:"true"`
Modality Modality `json:"modality" required:"true"`
URI string `json:"uri" required:"true"`
MimeType *string `json:"mime_type,omitempty"`
}
type ReasoningPart struct {
Type PartType `json:"type" required:"true"`
Content string `json:"content" required:"true"`
}
// CompactionPart is compacted conversation state; Content is the summary when it is not encrypted.
type CompactionPart struct {
Type PartType `json:"type" required:"true"`
ID *string `json:"id,omitempty"`
Content *string `json:"content,omitempty"`
}
// GenericPart keeps a part of any other type as sent.
type GenericPart map[string]any
func (Part) JSONSchemaOneOf() []any {
oneOf := make([]any, 0, len(partVariants)+1)
for _, variant := range partVariants {
oneOf = append(oneOf, variant.schema)
}
return append(oneOf, GenericPart{})
}
func (Part) PrepareJSONSchema(schema *jsonschema.Schema) error {
if schema.ExtraProperties == nil {
schema.ExtraProperties = map[string]any{}
}
mapping := make(map[string]string, len(partVariants))
for _, variant := range partVariants {
mapping[variant.typ.StringValue()] = variant.schemaRef
}
schema.ExtraProperties["x-signoz-discriminator"] = map[string]any{
"propertyName": "type",
"mapping": mapping,
}
return nil
}
func (p Part) MarshalJSON() ([]byte, error) {
return json.Marshal(p.Value)
}
func (p *Part) UnmarshalJSON(data []byte) error {
var head struct {
Type PartType `json:"type"`
}
if err := json.Unmarshal(data, &head); err != nil {
return err
}
decode := decodePart[GenericPart]
for _, variant := range partVariants {
if variant.typ == head.Type {
decode = variant.decode
break
}
}
value, err := decode(data)
if err != nil {
return err
}
p.Value = value
return nil
}
type partVariant struct {
typ PartType
decode func(data []byte) (any, error)
schema any
schemaRef string
}
func decodePart[T any](data []byte) (any, error) {
var value T
if err := json.Unmarshal(data, &value); err != nil {
return nil, err
}
return value, nil
}

View File

@@ -0,0 +1,74 @@
package genai
import (
"encoding/json"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestParts_RoundTrip(t *testing.T) {
id := "call_1"
mime := "image/png"
summary := "earlier turns"
testCases := []struct {
name string
json string
want Parts
}{
{
name: "EveryNamedType_DecodesToItsStruct",
json: `[{"type":"text","content":"hi"},` +
`{"type":"tool_call","name":"get_weather","id":"call_1","arguments":{"city":"Paris"}},` +
`{"type":"tool_call_response","id":"call_1","response":"rainy"},` +
`{"type":"server_tool_call","name":"code_interpreter","id":"call_1","server_tool_call":{"type":"code_interpreter","code":"1+1"}},` +
`{"type":"server_tool_call_response","id":"call_1","server_tool_call_response":{"type":"code_interpreter","outputs":[]}},` +
`{"type":"blob","modality":"image","content":"aGk=","mime_type":"image/png"},` +
`{"type":"file","modality":"image","file_id":"file_1"},` +
`{"type":"uri","modality":"image","uri":"gs://b/x.png","mime_type":"image/png"},` +
`{"type":"reasoning","content":"thinking"},` +
`{"type":"compaction","id":"call_1","content":"earlier turns"}]`,
want: Parts{
{Value: TextPart{Type: PartTypeText, Content: "hi"}},
{Value: ToolCallRequestPart{Type: PartTypeToolCall, Name: "get_weather", ID: &id, Arguments: map[string]any{"city": "Paris"}}},
{Value: ToolCallResponsePart{Type: PartTypeToolCallResponse, ID: &id, Response: "rainy"}},
{Value: ServerToolCallPart{Type: PartTypeServerToolCall, Name: "code_interpreter", ID: &id, ServerToolCall: GenericServerToolCall{"type": "code_interpreter", "code": "1+1"}}},
{Value: ServerToolCallResponsePart{Type: PartTypeServerToolCallResponse, ID: &id, ServerToolCallResponse: GenericServerToolCallResponse{"type": "code_interpreter", "outputs": []any{}}}},
{Value: BlobPart{Type: PartTypeBlob, Modality: ModalityImage, Content: "aGk=", MimeType: &mime}},
{Value: FilePart{Type: PartTypeFile, Modality: ModalityImage, FileID: "file_1"}},
{Value: UriPart{Type: PartTypeURI, Modality: ModalityImage, URI: "gs://b/x.png", MimeType: &mime}},
{Value: ReasoningPart{Type: PartTypeReasoning, Content: "thinking"}},
{Value: CompactionPart{Type: PartTypeCompaction, ID: &id, Content: &summary}},
},
},
{
name: "UnknownKeys_Dropped",
json: `[{"type":"reasoning","content":"","signature":"sig"}]`,
want: Parts{{Value: ReasoningPart{Type: PartTypeReasoning}}},
},
{
name: "UnknownType_DecodesToGenericMap",
json: `[{"type":"refusal","refusal":"no"}]`,
want: Parts{{Value: GenericPart{"type": "refusal", "refusal": "no"}}},
},
{
name: "MissingType_DecodesToGenericMap",
json: `[{"text":"bare"}]`,
want: Parts{{Value: GenericPart{"text": "bare"}}},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
var parts Parts
require.NoError(t, json.Unmarshal([]byte(testCase.json), &parts))
assert.Equal(t, testCase.want, parts)
encoded, err := json.Marshal(parts)
require.NoError(t, err)
var again Parts
require.NoError(t, json.Unmarshal(encoded, &again))
assert.Equal(t, parts, again)
})
}
}

View File

@@ -0,0 +1,252 @@
// Package genaiformatter builds the OTel GenAI view of a span's messages from the formats SDKs
// write into gen_ai.input.messages and gen_ai.output.messages.
package genaiformatter
import (
"encoding/json"
"fmt"
"regexp"
"strings"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes/genai"
"github.com/SigNoz/signoz/pkg/valuer"
)
const (
FormatterSemconv = "semconv"
FormatterOpenAIChat = "openai.chat"
FormatterText = "text"
FormatterGeneric = "generic"
)
var dataURL = regexp.MustCompile(`^data:([^;,]+);base64,(.+)$`)
// Provider spellings of the finish reasons; the enum values themselves pass through.
var finishReasonAliases = map[string]genai.FinishReason{
"tool_calls": genai.FinishReasonToolCall, "function_call": genai.FinishReasonToolCall,
"end_turn": genai.FinishReasonStop, "stop_sequence": genai.FinishReasonStop,
"max_tokens": genai.FinishReasonLength,
}
// Formatted is the OTel GenAI view of one span's messages. Formatter names the converter that
// produced it, Warnings what the converter could not resolve. The lists are never nil.
type Formatted struct {
Input genai.InputMessages
Output genai.OutputMessages
Formatter string
Warnings []string
}
// Format converts the two attribute values, each a JSON string or a structured value; nil means
// the attribute is absent.
func Format(input, output any) Formatted {
f := &formatting{out: Formatted{Input: genai.InputMessages{}, Output: genai.OutputMessages{}, Warnings: []string{}}}
var inputLabel, outputLabel string
if input != nil {
var msgs genai.OutputMessages
msgs, inputLabel = f.side(decode(input), genai.RoleUser)
for _, m := range msgs {
f.out.Input = append(f.out.Input, genai.ChatMessage{Role: m.Role, Parts: m.Parts, Name: m.Name})
}
}
if output != nil {
f.out.Output, outputLabel = f.side(decode(output), genai.RoleAssistant)
}
f.out.Formatter = inputLabel
if inputLabel == "" {
f.out.Formatter = outputLabel
} else if outputLabel != "" && outputLabel != inputLabel {
f.warn("output formatted as %s, input as %s", outputLabel, inputLabel)
}
return f.out
}
// object is a decoded JSON object whose accessors return zero values for missing keys.
type object map[string]any
type formatting struct {
out Formatted
}
func (o object) str(key string) string {
return stringOf(o[key])
}
func (o object) text(keys ...string) string {
for _, key := range keys {
if s, ok := o[key].(string); ok && s != "" {
return s
}
}
return ""
}
func (o object) obj(key string) (object, bool) {
m, ok := o[key].(map[string]any)
return m, ok
}
func (o object) list(key string) ([]any, bool) {
l, ok := o[key].([]any)
return l, ok
}
func (f *formatting) warn(format string, args ...any) {
f.out.Warnings = append(f.out.Warnings, fmt.Sprintf(format, args...))
}
// side converts one attribute value by its shape; role is what a bare string on this side is
// spoken as.
func (f *formatting) side(value any, role genai.Role) (genai.OutputMessages, string) {
switch v := value.(type) {
case []any:
if first, ok := firstObject(v); ok {
if _, ok := first["parts"]; ok {
return f.semconv(v), FormatterSemconv
}
if _, ok := first["role"]; ok {
return f.openAIChatMessages(v), FormatterOpenAIChat
}
}
case map[string]any:
if choices, ok := object(v).list("choices"); ok {
return f.openAIChatResponse(choices), FormatterOpenAIChat
}
if messages, ok := object(v).list("messages"); ok {
return f.openAIChatMessages(messages), FormatterOpenAIChat
}
case string:
if v == "" {
return genai.OutputMessages{}, FormatterText
}
f.warn("bare %s text, role assumed", role.StringValue())
return genai.OutputMessages{{Role: role, Parts: genai.Parts{part(textPart(v))}}}, FormatterText
}
f.warn("message format not recognised, kept as generic")
return genai.OutputMessages{genericMessage(value)}, FormatterGeneric
}
func firstObject(list []any) (object, bool) {
if len(list) == 0 {
return nil, false
}
first, ok := list[0].(map[string]any)
return first, ok
}
func genericMessage(value any) genai.OutputMessage {
return genai.OutputMessage{Parts: genai.Parts{part(genai.GenericPart{"type": FormatterGeneric, "content": stringOf(value)})}}
}
// decode parses a JSON-encoded attribute value, twice when an SDK encoded the string again.
// Anything that is not JSON is returned as is.
func decode(value any) any {
s, ok := value.(string)
if !ok {
return value
}
trimmed := strings.TrimSpace(s)
if trimmed == "" || !strings.ContainsAny(trimmed[:1], `[{"`) {
return s
}
var decoded any
if err := json.Unmarshal([]byte(trimmed), &decoded); err != nil {
return s
}
if inner, ok := decoded.(string); ok {
return decode(inner)
}
return decoded
}
func role(value string) genai.Role {
lower := strings.ToLower(strings.TrimSpace(value))
switch lower {
case "human":
return genai.RoleUser
case "ai", "model":
return genai.RoleAssistant
case "function":
return genai.RoleTool
}
return genai.Role{String: valuer.NewString(lower)}
}
func finishReason(value string) *genai.FinishReason {
lower := strings.ToLower(strings.TrimSpace(value))
if lower == "" {
return nil
}
fr, ok := finishReasonAliases[lower]
if !ok {
fr = genai.FinishReason{String: valuer.NewString(lower)}
}
return &fr
}
func optionalString(value any) *string {
s := stringOf(value)
if s == "" {
return nil
}
return &s
}
func part(value any) genai.Part {
return genai.Part{Value: value}
}
func textPart(content string) genai.TextPart {
return genai.TextPart{Type: genai.PartTypeText, Content: content}
}
// contentPart converts one OpenAI content block, which semconv parts may also carry; a block of
// any other type is kept as sent.
func contentPart(item any) genai.Part {
p, ok := item.(map[string]any)
if !ok {
if s, ok := item.(string); ok {
return part(textPart(s))
}
return part(genai.GenericPart{"type": FormatterGeneric, "content": stringOf(item)})
}
block := object(p)
switch block.str("type") {
case "text":
return part(textPart(block.str("text")))
case "refusal":
return part(textPart(block.str("refusal")))
case "image_url":
url := block.str("image_url")
if ref, ok := block.obj("image_url"); ok {
url = ref.str("url")
}
if url != "" {
return part(mediaPart(url, genai.ModalityImage))
}
}
return part(genai.GenericPart(p))
}
// mediaPart places a data URL in a BlobPart, any other URL in a UriPart.
func mediaPart(url string, modality genai.Modality) any {
if m := dataURL.FindStringSubmatch(url); m != nil {
return genai.BlobPart{Type: genai.PartTypeBlob, Modality: modality, MimeType: &m[1], Content: m[2]}
}
return genai.UriPart{Type: genai.PartTypeURI, Modality: modality, URI: url}
}
// stringOf renders nil as "" and non-strings as compact JSON.
func stringOf(value any) string {
switch v := value.(type) {
case nil:
return ""
case string:
return v
}
data, err := json.Marshal(value)
if err != nil {
return ""
}
return string(data)
}

View File

@@ -0,0 +1,101 @@
package genaiformatter
import (
"encoding/json"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
type formatCase struct {
name string
input any
output any
formatter string
warnings []string
wantInput string
wantOutput string
}
func runFormatCases(t *testing.T, testCases []formatCase) {
t.Helper()
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
formatted := Format(testCase.input, testCase.output)
assert.Equal(t, testCase.formatter, formatted.Formatter)
assert.Equal(t, append([]string{}, testCase.warnings...), formatted.Warnings)
input, err := json.Marshal(formatted.Input)
require.NoError(t, err)
assert.JSONEq(t, testCase.wantInput, string(input))
output, err := json.Marshal(formatted.Output)
require.NoError(t, err)
assert.JSONEq(t, testCase.wantOutput, string(output))
})
}
}
func TestFormat(t *testing.T) {
runFormatCases(t, []formatCase{
{
// bifrost-gateway capture, HTTP span
name: "BareStrings_RolesAssumed_WarnsPerSide",
input: `A bat and a ball cost $1.10. The bat costs $1 more than the ball. How much is the ball?`,
output: `Let the cost of the ball be x dollars. Then the bat costs x + $1.00. According to the problem:
  x + (x + 1.00) = 1.10
Combine like terms:
  2x + 1.00 = 1.10
Subtract 1.00 from both sides:
  2x = 0.10
Divide by 2:
  x = 0.05
So, the ball costs 5 cents.`,
formatter: FormatterText,
warnings: []string{"bare user text, role assumed", "bare assistant text, role assumed"},
wantInput: `[
{
"role": "user",
"parts": [
{
"type": "text",
"content": "A bat and a ball cost $1.10. The bat costs $1 more than the ball. How much is the ball?"
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "text",
"content": "Let the cost of the ball be x dollars. Then the bat costs x + $1.00. According to the problem:\n\n  x + (x + 1.00) = 1.10\n\nCombine like terms:\n\n  2x + 1.00 = 1.10\n\nSubtract 1.00 from both sides:\n\n  2x = 0.10\n\nDivide by 2:\n\n  x = 0.05\n\nSo, the ball costs 5 cents."
}
]
}
]`,
},
{
name: "UnrecognisedObject_GenericPart_Warns",
input: `{"city": "Paris", "temp_c": 21}`,
formatter: FormatterGeneric,
warnings: []string{"message format not recognised, kept as generic"},
wantInput: `[{"role": "", "parts": [{"type": "generic", "content": "{\"city\":\"Paris\",\"temp_c\":21}"}]}]`,
wantOutput: `[]`,
},
{
name: "NoAttributes_EmptyListsNoFormatter",
wantInput: `[]`,
wantOutput: `[]`,
},
})
}

View File

@@ -0,0 +1,122 @@
package genaiformatter
import (
"strings"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes/genai"
)
// Format: OpenAI Chat Completions, a request {messages}, a response {choices}, or a bare message list.
// Written by: bifrost, openrouter requests, the raw OpenAI calls OpenInference records.
// Handled: string and block content, tool_calls nested or flattened, tool messages, refusal,
// reasoning_content, finish reasons.
// Not yet: streaming delta chunks, legacy completions choices[].text, audio output.
func (f *formatting) openAIChatMessages(list []any) genai.OutputMessages {
out := make(genai.OutputMessages, 0, len(list))
for _, item := range list {
m, ok := item.(map[string]any)
if !ok {
out = append(out, genericMessage(item))
continue
}
out = append(out, f.openAIChatMessage(m))
}
return out
}
// openAIChatResponse converts the choices of an OpenAI chat response, one message each with its
// finish reason.
func (f *formatting) openAIChatResponse(choices []any) genai.OutputMessages {
out := make(genai.OutputMessages, 0, len(choices))
for _, item := range choices {
choice, ok := item.(map[string]any)
if !ok {
continue
}
inner, _ := object(choice).obj("message")
msg := f.openAIChatMessage(inner)
if msg.Role.IsZero() {
msg.Role = genai.RoleAssistant
}
if fr := finishReason(object(choice).str("finish_reason")); fr != nil {
msg.FinishReason = fr
}
out = append(out, msg)
}
return out
}
// openAIChatMessage converts one {role, content, tool_calls, ...} message in the OpenAI chat shape.
func (f *formatting) openAIChatMessage(m object) genai.OutputMessage {
msg := genai.OutputMessage{Role: role(m.str("role")), Parts: genai.Parts{}, Name: optionalString(m["name"]), FinishReason: finishReason(m.str("finish_reason"))}
if id := m.str("tool_call_id"); msg.Role == genai.RoleTool && id != "" {
msg.Parts = genai.Parts{part(genai.ToolCallResponsePart{Type: genai.PartTypeToolCallResponse, ID: &id, Response: toolResponse(m["content"])})}
return msg
}
switch content := m["content"].(type) {
case string:
if content != "" {
msg.Parts = append(msg.Parts, part(textPart(content)))
}
case []any:
for _, item := range content {
msg.Parts = append(msg.Parts, contentPart(item))
}
}
if refusal := m.str("refusal"); refusal != "" {
msg.Parts = append(msg.Parts, part(textPart(refusal)))
}
if reasoning := m.text("reasoning_content", "reasoning"); reasoning != "" {
msg.Parts = append(msg.Parts, part(genai.ReasoningPart{Type: genai.PartTypeReasoning, Content: reasoning}))
}
if calls, ok := m.list("tool_calls"); ok {
for _, call := range calls {
if c, ok := call.(map[string]any); ok {
msg.Parts = append(msg.Parts, toolCallPart(c))
}
}
}
if call, ok := m.obj("function_call"); ok {
msg.Parts = append(msg.Parts, part(genai.ToolCallRequestPart{Type: genai.PartTypeToolCall, Name: call.str("name"), Arguments: decode(call["arguments"])}))
}
return msg
}
// toolCallPart reads {id, function: {name, arguments}}; a gateway may flatten it to {id, name, arguments|args}.
func toolCallPart(c object) genai.Part {
call := genai.ToolCallRequestPart{Type: genai.PartTypeToolCall, ID: optionalString(c["id"]), Name: c.str("name")}
args := c["arguments"]
if args == nil {
args = c["args"]
}
if fn, ok := c.obj("function"); ok {
call.Name = fn.str("name")
args = fn["arguments"]
}
if args != nil {
call.Arguments = decode(args)
}
return part(call)
}
// toolResponse parses a JSON string and joins text blocks; anything else stays as sent.
func toolResponse(value any) any {
switch v := value.(type) {
case string:
return decode(v)
case []any:
texts := make([]string, 0, len(v))
for _, item := range v {
block, ok := item.(map[string]any)
if !ok || object(block).str("type") != "text" {
return v
}
texts = append(texts, object(block).str("text"))
}
return strings.Join(texts, "\n")
}
return value
}

View File

@@ -0,0 +1,173 @@
package genaiformatter
import "testing"
// Captured cases are gen_ai.input.messages and gen_ai.output.messages as the collector stored
// them for the SigNoz/scripts static-telemetry-generator captures.
func TestFormat_OpenAIChat(t *testing.T) {
runFormatCases(t, []formatCase{
{
// openinference capture, ChatCompletion span: the raw OpenAI request and response
name: "ResponseWithToolCalls_ToolCallParts_FinishReasonToolCall",
input: `{"messages": [{"role": "user", "content": "What's the weather in Paris and in London? Use the tool for each."}], "model": "gpt-4o-mini", "max_tokens": 120, "tool_choice": "auto", "tools": [{"type": "function", "function": {"name": "get_current_weather", "description": "Get the current weather for a city.", "parameters": {"type": "object", "properties": {"city": {"type": "string"}, "unit": {"type": "string", "enum": ["c", "f"]}}, "required": ["city"]}}}]}`,
output: `{"id":"chatcmpl-DksBDptnQ3aTu18YrqSdYQgiPOuRV","choices":[{"finish_reason":"tool_calls","index":0,"logprobs":null,"message":{"content":null,"refusal":null,"role":"assistant","annotations":[],"tool_calls":[{"id":"call_SW4hst7nzysgvWSAz4OdyK9D","function":{"arguments":"{\"city\": \"Paris\"}","name":"get_current_weather"},"type":"function"},{"id":"call_si8wiETYGpooMB7gMXSgMQpk","function":{"arguments":"{\"city\": \"London\"}","name":"get_current_weather"},"type":"function"}]}}],"created":1780063727,"model":"gpt-4o-mini-2024-07-18","object":"chat.completion","service_tier":"default","system_fingerprint":"fp_6e71a9f378","usage":{"completion_tokens":46,"prompt_tokens":71,"total_tokens":117,"completion_tokens_details":{"accepted_prediction_tokens":0,"audio_tokens":0,"reasoning_tokens":0,"rejected_prediction_tokens":0},"prompt_tokens_details":{"audio_tokens":0,"cached_tokens":0}}}`,
formatter: FormatterOpenAIChat,
wantInput: `[
{
"role": "user",
"parts": [
{
"type": "text",
"content": "What's the weather in Paris and in London? Use the tool for each."
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "tool_call",
"name": "get_current_weather",
"id": "call_SW4hst7nzysgvWSAz4OdyK9D",
"arguments": {
"city": "Paris"
}
},
{
"type": "tool_call",
"name": "get_current_weather",
"id": "call_si8wiETYGpooMB7gMXSgMQpk",
"arguments": {
"city": "London"
}
}
],
"finish_reason": "tool_call"
}
]`,
},
{
// openinference capture, the follow-up ChatCompletion span
name: "ToolRoleMessages_BecomeToolCallResponseParts",
input: `{"messages": [{"role": "user", "content": "What's the weather in Paris and in London? Use the tool for each."}, {"content": null, "refusal": null, "role": "assistant", "annotations": [], "audio": null, "function_call": null, "tool_calls": [{"id": "call_SW4hst7nzysgvWSAz4OdyK9D", "function": {"arguments": "{\"city\": \"Paris\"}", "name": "get_current_weather"}, "type": "function"}, {"id": "call_si8wiETYGpooMB7gMXSgMQpk", "function": {"arguments": "{\"city\": \"London\"}", "name": "get_current_weather"}, "type": "function"}]}, {"role": "tool", "tool_call_id": "call_SW4hst7nzysgvWSAz4OdyK9D", "content": "{\"city\": \"Paris\", \"temp_c\": 18, \"summary\": \"Clear\"}"}, {"role": "tool", "tool_call_id": "call_si8wiETYGpooMB7gMXSgMQpk", "content": "{\"city\": \"London\", \"temp_c\": 18, \"summary\": \"Clear\"}"}], "model": "gpt-4o-mini", "max_tokens": 80}`,
output: `{"id":"chatcmpl-DksBF2QuaoMMMFg9UAXX9Rc28XVJT","choices":[{"finish_reason":"stop","index":0,"logprobs":null,"message":{"content":"The current weather is as follows:\n\n- **Paris**: 18°C, Clear\n- **London**: 18°C, Clear\n\nBoth cities are enjoying clear skies with the same temperature.","refusal":null,"role":"assistant","annotations":[]}}],"created":1780063729,"model":"gpt-4o-mini-2024-07-18","object":"chat.completion","service_tier":"default","system_fingerprint":"fp_e2d886d409","usage":{"completion_tokens":40,"prompt_tokens":118,"total_tokens":158,"completion_tokens_details":{"accepted_prediction_tokens":0,"audio_tokens":0,"reasoning_tokens":0,"rejected_prediction_tokens":0},"prompt_tokens_details":{"audio_tokens":0,"cached_tokens":0}}}`,
formatter: FormatterOpenAIChat,
wantInput: `[
{
"role": "user",
"parts": [
{
"type": "text",
"content": "What's the weather in Paris and in London? Use the tool for each."
}
]
},
{
"role": "assistant",
"parts": [
{
"type": "tool_call",
"name": "get_current_weather",
"id": "call_SW4hst7nzysgvWSAz4OdyK9D",
"arguments": {
"city": "Paris"
}
},
{
"type": "tool_call",
"name": "get_current_weather",
"id": "call_si8wiETYGpooMB7gMXSgMQpk",
"arguments": {
"city": "London"
}
}
]
},
{
"role": "tool",
"parts": [
{
"type": "tool_call_response",
"response": {
"city": "Paris",
"summary": "Clear",
"temp_c": 18
},
"id": "call_SW4hst7nzysgvWSAz4OdyK9D"
}
]
},
{
"role": "tool",
"parts": [
{
"type": "tool_call_response",
"response": {
"city": "London",
"summary": "Clear",
"temp_c": 18
},
"id": "call_si8wiETYGpooMB7gMXSgMQpk"
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "text",
"content": "The current weather is as follows:\n\n- **Paris**: 18°C, Clear\n- **London**: 18°C, Clear\n\nBoth cities are enjoying clear skies with the same temperature."
}
],
"finish_reason": "stop"
}
]`,
},
{
// bifrost-gateway capture, HTTP span
name: "GatewayFlattenedToolCalls_ToolCallParts_BareInputWarns",
input: `What's the weather in Paris and in London? Use the tool for each.`,
output: `[{"role":"assistant","content":"","tool_calls":[{"id":"call_xP33RT5exgBrNsvsWQZwiY9D","type":"function","name":"get_current_weather","args":"{\"city\": \"Paris\"}"},{"id":"call_lC6GgT8JW2ygJqWEM2wcohIz","type":"function","name":"get_current_weather","args":"{\"city\": \"London\"}"}]}]`,
formatter: FormatterText,
warnings: []string{"bare user text, role assumed", "output formatted as openai.chat, input as text"},
wantInput: `[
{
"role": "user",
"parts": [
{
"type": "text",
"content": "What's the weather in Paris and in London? Use the tool for each."
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "tool_call",
"name": "get_current_weather",
"id": "call_xP33RT5exgBrNsvsWQZwiY9D",
"arguments": {
"city": "Paris"
}
},
{
"type": "tool_call",
"name": "get_current_weather",
"id": "call_lC6GgT8JW2ygJqWEM2wcohIz",
"arguments": {
"city": "London"
}
}
]
}
]`,
},
})
}

View File

@@ -0,0 +1,70 @@
package genaiformatter
import (
"encoding/json"
"maps"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes/genai"
)
// Format: the OTel shape itself, [{role, parts}].
// Written by: litellm, langchain, openllmetry, portkey, traceloop LLM spans.
// Decoding is the genai types. Repairs: "text" used for "content" on text parts; OpenAI content
// blocks left inside parts.
// Not yet: "thinking" for "reasoning"; "result" or "output" for "response" on tool call responses.
// semconv decodes messages that already follow the schema through the genai types. A part of a
// type the schema does not name is a provider content block, which contentPart converts.
func (f *formatting) semconv(list []any) genai.OutputMessages {
data, err := json.Marshal(semconvTextKey(list))
if err != nil {
f.warn("messages could not be encoded: %s", err)
return genai.OutputMessages{genericMessage(list)}
}
var msgs genai.OutputMessages
if err := json.Unmarshal(data, &msgs); err != nil {
f.warn("messages do not follow the schema: %s", err)
return genai.OutputMessages{genericMessage(list)}
}
for i := range msgs {
if msgs[i].Parts == nil {
msgs[i].Parts = genai.Parts{}
}
for j, p := range msgs[i].Parts {
if block, ok := p.Value.(genai.GenericPart); ok {
msgs[i].Parts[j] = contentPart(map[string]any(block))
}
}
}
return msgs
}
// semconvTextKey copies the text and reasoning parts some SDKs write with "text" in place of the
// schema's "content". The input is left untouched, it is also the span's raw attribute.
func semconvTextKey(list []any) []any {
out := make([]any, 0, len(list))
for _, item := range list {
m, ok := item.(map[string]any)
parts, hasParts := object(m).list("parts")
if !ok || !hasParts {
out = append(out, item)
continue
}
fixedParts := make([]any, 0, len(parts))
for _, p := range parts {
part, ok := p.(map[string]any)
typ := object(part).str("type")
if ok && (typ == "text" || typ == "reasoning") && part["content"] == nil && part["text"] != nil {
fixed := maps.Clone(part)
fixed["content"] = fixed["text"]
delete(fixed, "text")
p = fixed
}
fixedParts = append(fixedParts, p)
}
msg := maps.Clone(m)
msg["parts"] = fixedParts
out = append(out, msg)
}
return out
}

View File

@@ -0,0 +1,163 @@
package genaiformatter
import "testing"
// Captured cases are gen_ai.input.messages and gen_ai.output.messages as the collector stored
// them for the SigNoz/scripts static-telemetry-generator captures.
func TestFormat_Semconv(t *testing.T) {
runFormatCases(t, []formatCase{
{
// litellm capture, litellm_request span
name: "SemconvMessages_DecodedUnchanged",
input: `[{"role": "system", "parts": [{"type": "text", "content": "You are a concise assistant."}]}, {"role": "user", "parts": [{"type": "text", "content": "Give me a one-line definition of observability."}]}]`,
output: `[{"role": "assistant", "parts": [{"type": "text", "content": "Observability is the ability to understand the internal state of a system based on external outputs, enabling effective monitoring and troubleshooting."}], "finish_reason": "stop"}]`,
formatter: FormatterSemconv,
wantInput: `[
{
"role": "system",
"parts": [
{
"type": "text",
"content": "You are a concise assistant."
}
]
},
{
"role": "user",
"parts": [
{
"type": "text",
"content": "Give me a one-line definition of observability."
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "text",
"content": "Observability is the ability to understand the internal state of a system based on external outputs, enabling effective monitoring and troubleshooting."
}
],
"finish_reason": "stop"
}
]`,
},
{
// litellm capture, vision call; base64 trimmed
name: "SemconvTextSpelledAsTextKey_ContentFilled_ImageBlockBecomesBlob",
input: `[{"role": "user", "parts": [{"type": "text", "text": "What animal is in this image? One sentence."}, {"type": "image_url", "image_url": {"url": "data:image/jpeg;base64,/9j/4AAQSkZJRgABAQAAAQABAAD/2wBD"}}]}]`,
output: `[{"role": "assistant", "parts": [{"type": "text", "content": "The animal in the image is a cartoon duck."}], "finish_reason": "stop"}]`,
formatter: FormatterSemconv,
wantInput: `[
{
"role": "user",
"parts": [
{
"type": "text",
"content": "What animal is in this image? One sentence."
},
{
"type": "blob",
"modality": "image",
"content": "/9j/4AAQSkZJRgABAQAAAQABAAD/2wBD",
"mime_type": "image/jpeg"
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "text",
"content": "The animal in the image is a cartoon duck."
}
],
"finish_reason": "stop"
}
]`,
},
{
// langchain capture, ChatOpenAI.chat span
name: "SemconvToolCallTurn_DecodedUnchanged",
input: `[{"role": "user", "parts": [{"type": "text", "content": "Search for SigNoz AI observability docs, then tell me the weather in Bengaluru."}]}, {"role": "assistant", "parts": [{"type": "tool_call", "id": "call_jL9UhSdUC9wWn2qfWwOjpzFd", "name": "search_web", "arguments": {"query": "SigNoz AI observability documentation"}}, {"type": "tool_call", "id": "call_nkVdZ8xX8d17HtmVddbUgJUc", "name": "get_weather", "arguments": {"city": "Bengaluru"}}]}, {"role": "tool", "parts": [{"type": "tool_call_response", "id": "call_jL9UhSdUC9wWn2qfWwOjpzFd", "response": "{\"query\": \"SigNoz AI observability documentation\", \"results\": [{\"title\": \"SigNoz docs: AI observability\", \"snippet\": \"...\"}, {\"title\": \"OpenTelemetry GenAI conventions\", \"snippet\": \"...\"}]}"}]}, {"role": "tool", "parts": [{"type": "tool_call_response", "id": "call_nkVdZ8xX8d17HtmVddbUgJUc", "response": "{\"city\": \"Bengaluru\", \"temp_c\": 18, \"summary\": \"Clear\"}"}]}]`,
output: `[{"role": "assistant", "parts": [{"type": "text", "content": "I found some results for the SigNoz AI observability documentation. Here's one of the titles you might be interested in:\n\n- **SigNoz docs: AI observability**\n\nUnfortunately, a snippet was not provided, but you can look it up directly for more details.\n\nAs for the weather in Bengaluru, it's currently 18\u00b0C and clear."}], "finish_reason": "stop"}]`,
formatter: FormatterSemconv,
wantInput: `[
{
"role": "user",
"parts": [
{
"type": "text",
"content": "Search for SigNoz AI observability docs, then tell me the weather in Bengaluru."
}
]
},
{
"role": "assistant",
"parts": [
{
"type": "tool_call",
"name": "search_web",
"id": "call_jL9UhSdUC9wWn2qfWwOjpzFd",
"arguments": {
"query": "SigNoz AI observability documentation"
}
},
{
"type": "tool_call",
"name": "get_weather",
"id": "call_nkVdZ8xX8d17HtmVddbUgJUc",
"arguments": {
"city": "Bengaluru"
}
}
]
},
{
"role": "tool",
"parts": [
{
"type": "tool_call_response",
"response": "{\"query\": \"SigNoz AI observability documentation\", \"results\": [{\"title\": \"SigNoz docs: AI observability\", \"snippet\": \"...\"}, {\"title\": \"OpenTelemetry GenAI conventions\", \"snippet\": \"...\"}]}",
"id": "call_jL9UhSdUC9wWn2qfWwOjpzFd"
}
]
},
{
"role": "tool",
"parts": [
{
"type": "tool_call_response",
"response": "{\"city\": \"Bengaluru\", \"temp_c\": 18, \"summary\": \"Clear\"}",
"id": "call_nkVdZ8xX8d17HtmVddbUgJUc"
}
]
}
]`,
wantOutput: `[
{
"role": "assistant",
"parts": [
{
"type": "text",
"content": "I found some results for the SigNoz AI observability documentation. Here's one of the titles you might be interested in:\n\n- **SigNoz docs: AI observability**\n\nUnfortunately, a snippet was not provided, but you can look it up directly for more details.\n\nAs for the weather in Bengaluru, it's currently 18°C and clear."
}
],
"finish_reason": "stop"
}
]`,
},
{
name: "StructuredSemconvValue_DecodedWithoutJSONString",
input: []any{map[string]any{"role": "user", "parts": []any{map[string]any{"type": "text", "content": "hi"}}}},
formatter: FormatterSemconv,
wantInput: `[{"role": "user", "parts": [{"type": "text", "content": "hi"}]}]`,
wantOutput: `[]`,
},
})
}

View File

@@ -35,6 +35,9 @@ const (
// agent span. A trace belongs to the AI explorer when any span carries one.
var GenAISpanGateKeys = []string{GenAIRequestModel, GenAIToolName, GenAIAgentName}
// GenAIThreadKeys mark a span as a thread row: a model call with its messages, or a tool execution.
var GenAIThreadKeys = []string{GenAIInputMessages, GenAIOutputMessages, GenAIToolName}
// GenAISpanFilterExpression renders the gate as a query-builder filter
// expression: each gate key ORed on EXISTS.
func GenAISpanFilterExpression() string {

View File

@@ -4,8 +4,12 @@ import (
"encoding/base64"
"encoding/json"
"maps"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes/genai"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes/genaiformatter"
)
const (
@@ -76,6 +80,12 @@ type ThreadSpan struct {
Attributes map[string]any `json:"attributes" required:"true" nullable:"false"`
Events []Event `json:"events" required:"true" nullable:"false"`
References []OtelSpanRef `json:"references" required:"true" nullable:"false"`
// The formatted fields hold the span's messages in the OTel GenAI shape; Formatter names the
// converter that produced them and FormatterWarnings what it could not resolve.
FormattedInput genai.InputMessages `json:"formatted_input" required:"true" nullable:"false"`
FormattedOutput genai.OutputMessages `json:"formatted_output" required:"true" nullable:"false"`
Formatter string `json:"formatter" required:"true"`
FormatterWarnings []string `json:"formatter_warnings" required:"true" nullable:"false"`
timeUnixNano uint64
}
@@ -171,32 +181,55 @@ func newThreadSpan(traceID string, storable *StorableSpan) *ThreadSpan {
resources := make(map[string]string, len(storable.ResourcesString))
maps.Copy(resources, storable.ResourcesString)
timeUnixNano := uint64(storable.StartTime.UnixNano())
attributes := threadAttributes(storable)
formatted := genaiformatter.Format(rawAttribute(storable, aiobservabilitytypes.GenAIInputMessages), rawAttribute(storable, aiobservabilitytypes.GenAIOutputMessages))
return &ThreadSpan{
SpanID: storable.SpanID,
TraceID: traceID,
ParentSpanID: storable.ParentSpanID,
Name: storable.Name,
KindString: storable.SpanKind,
TimeUnix: timeUnixNano / 1_000_000, // client expects millis, as in the waterfall
DurationNano: storable.DurationNano,
HasError: storable.HasError,
StatusCodeString: storable.StatusCodeString,
StatusMessage: storable.StatusMessage,
Resource: resources,
Attributes: threadAttributes(storable),
Events: storable.UnmarshalledEvents(),
References: storable.UnmarshalledRefs(),
timeUnixNano: timeUnixNano,
SpanID: storable.SpanID,
TraceID: traceID,
ParentSpanID: storable.ParentSpanID,
Name: storable.Name,
KindString: storable.SpanKind,
TimeUnix: timeUnixNano / 1_000_000, // client expects millis, as in the waterfall
DurationNano: storable.DurationNano,
HasError: storable.HasError,
StatusCodeString: storable.StatusCodeString,
StatusMessage: storable.StatusMessage,
Resource: resources,
Attributes: attributes,
Events: storable.UnmarshalledEvents(),
References: storable.UnmarshalledRefs(),
FormattedInput: formatted.Input,
FormattedOutput: formatted.Output,
Formatter: formatted.Formatter,
FormatterWarnings: formatted.Warnings,
timeUnixNano: timeUnixNano,
}
}
// threadAttributes reads the JSON column and falls back to the legacy maps for spans written
// before the JSON rollout.
// threadAttributes flattens the JSON column into dotted keys, as the querier does for list
// responses, and falls back to the legacy maps for spans written before the JSON rollout.
func threadAttributes(storable *StorableSpan) map[string]any {
if len(storable.AttributesJSON) > 0 {
attributes := make(map[string]any, len(storable.AttributesJSON))
storable.AttributesJSON.FlattenInto("", attributes)
return attributes
if len(storable.AttributesJSON) == 0 {
return storable.Attributes()
}
return storable.Attributes()
attributes := make(map[string]any, len(storable.AttributesJSON))
storable.AttributesJSON.FlattenInto("", attributes)
return attributes
}
// rawAttribute reads one attribute for decoding from the JSON document, where an object value is
// still whole. The legacy maps hold objects split into one key per field, so spans from before
// the JSON rollout are not decoded. Not handled: a list of JSON strings, which formats as generic,
// and a scalar at a prefix of the path, such as gen_ai.input beside gen_ai.input.messages, which
// ends the walk early.
func rawAttribute(storable *StorableSpan, key string) any {
var current any = map[string]any(storable.AttributesJSON)
for segment := range strings.SplitSeq(key, ".") {
object, ok := current.(map[string]any)
if !ok {
return nil
}
current = object[segment]
}
return current
}

View File

@@ -73,3 +73,35 @@ func TestThreadAttributes(t *testing.T) {
})
}
}
func TestRawAttribute(t *testing.T) {
arguments := map[string]any{"city": "Paris", "days": []any{1.0, 2.0}}
testCases := []struct {
name string
storable StorableSpan
key string
want any
}{
{
name: "JSONColumn_ObjectReturnedWhole",
storable: StorableSpan{AttributesJSON: telemetrystoretypes.JSONValue{"gen_ai": map[string]any{"tool": map[string]any{"call": map[string]any{"arguments": arguments}}}}},
key: "gen_ai.tool.call.arguments",
want: arguments,
},
{
name: "JSONColumn_MissingPath_Nil",
storable: StorableSpan{AttributesJSON: telemetrystoretypes.JSONValue{"gen_ai": map[string]any{"tool": map[string]any{"name": "get_weather"}}}},
key: "gen_ai.tool.call.arguments",
},
{
name: "LegacyMapsOnly_Nil",
storable: StorableSpan{AttributesString: map[string]string{"gen_ai.output.messages": "sunny", "gen_ai.tool.call.arguments.city": "Paris"}},
key: "gen_ai.output.messages",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.want, rawAttribute(&testCase.storable, testCase.key))
})
}
}

View File

@@ -21,7 +21,7 @@ def test_thread_returns_message_spans_in_order(
seed_attribute_evolution("traces", ATTRIBUTE_JSON_ROLLOUT_TIME)
now = datetime.now(tz=UTC).replace(microsecond=0)
trace_id = TraceIdGenerator.trace_id()
root_id, first_llm_id, tool_id, second_llm_id, third_llm_id = (TraceIdGenerator.span_id() for _ in range(5))
root_id, first_llm_id, tool_id, agent_id, second_llm_id, third_llm_id = (TraceIdGenerator.span_id() for _ in range(6))
resources = {"service.name": "tracedetail-thread"}
first_input = json.dumps([{"role": "user", "parts": [{"type": "text", "content": "weather in Bangalore?"}]}])
first_output = json.dumps([{"role": "assistant", "parts": [{"type": "tool_call", "id": "call_1", "name": "get_weather", "arguments": {"city": "Bangalore"}}], "finish_reason": "tool_call"}])
@@ -33,7 +33,10 @@ def test_thread_returns_message_spans_in_order(
Traces(
timestamp=now - timedelta(seconds=8), trace_id=trace_id, span_id=first_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.input.messages": first_input, "gen_ai.output.messages": first_output}, attribute_write_mode="json_only"
),
Traces(timestamp=now - timedelta(seconds=6), trace_id=trace_id, span_id=tool_id, parent_span_id=root_id, name="execute_tool get_weather", resources=resources, attributes={"gen_ai.tool.name": "get_weather"}, attribute_write_mode="json_only"),
Traces(
timestamp=now - timedelta(seconds=6), trace_id=trace_id, span_id=tool_id, parent_span_id=root_id, name="execute_tool get_weather", resources=resources, attributes={"gen_ai.tool.name": "get_weather", "gen_ai.tool.call.id": "call_1", "gen_ai.tool.call.arguments": json.dumps({"city": "Bangalore"}), "gen_ai.tool.call.result": "sunny"}, attribute_write_mode="json_only"
),
Traces(timestamp=now - timedelta(seconds=5), trace_id=trace_id, span_id=agent_id, parent_span_id=root_id, name="invoke_agent planner", resources=resources, attributes={"gen_ai.agent.name": "planner"}, attribute_write_mode="json_only"),
Traces(timestamp=now - timedelta(seconds=4), trace_id=trace_id, span_id=second_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.input.messages": second_input}, attribute_write_mode="json_only"),
Traces(timestamp=now - timedelta(seconds=2), trace_id=trace_id, span_id=third_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.output.messages": "It is sunny in Bangalore."}, attribute_write_mode="json_only"),
]
@@ -44,10 +47,12 @@ def test_thread_returns_message_spans_in_order(
assert response.status_code == HTTPStatus.OK, response.text
thread = response.json()["data"]
assert [span["span_id"] for span in thread["spans"]] == [first_llm_id, second_llm_id, third_llm_id]
assert [span["span_id"] for span in thread["spans"]] == [first_llm_id, tool_id, second_llm_id, third_llm_id]
assert "nextCursor" not in thread
first, input_only, output_only = thread["spans"]
first, tool, input_only, output_only = thread["spans"]
assert tool["attributes"]["gen_ai.tool.call.arguments"] == json.dumps({"city": "Bangalore"})
assert tool["attributes"]["gen_ai.tool.call.result"] == "sunny"
assert first["time_unix"] == int((now - timedelta(seconds=8)).timestamp() * 1000)
assert first["attributes"]["gen_ai.input.messages"] == first_input
assert first["attributes"]["gen_ai.request.model"] == "gpt-4o"
@@ -57,6 +62,23 @@ def test_thread_returns_message_spans_in_order(
assert "gen_ai.input.messages" not in output_only["attributes"]
assert output_only["attributes"]["gen_ai.output.messages"] == "It is sunny in Bangalore."
tool_call = {"type": "tool_call", "name": "get_weather", "id": "call_1", "arguments": {"city": "Bangalore"}}
assert first["formatter"] == "semconv"
assert first["formatted_input"] == [{"role": "user", "parts": [{"type": "text", "content": "weather in Bangalore?"}]}]
assert first["formatted_output"] == [{"role": "assistant", "parts": [tool_call], "finish_reason": "tool_call"}]
assert first["formatter_warnings"] == []
# tool spans carry no messages until their converter lands
assert tool["formatter"] == ""
assert tool["formatted_input"] == []
assert tool["formatted_output"] == []
assert input_only["formatter"] == "openai.chat"
assert input_only["formatted_input"] == [{"role": "tool", "parts": [{"type": "tool_call_response", "id": "call_1", "response": "sunny"}]}]
assert input_only["formatted_output"] == []
assert output_only["formatter"] == "text"
assert output_only["formatted_input"] == []
assert output_only["formatted_output"] == [{"role": "assistant", "parts": [{"type": "text", "content": "It is sunny in Bangalore."}]}]
assert output_only["formatter_warnings"] == ["bare assistant text, role assumed"]
def test_thread_paginates_with_cursors(
signoz: types.SigNoz,
@@ -180,8 +202,8 @@ def test_thread_opens_around_span(
resources = {"service.name": "tracedetail-thread-anchor"}
root_id = TraceIdGenerator.span_id()
llm_ids = [TraceIdGenerator.span_id() for _ in range(5)]
tool_id = TraceIdGenerator.span_id()
# tool span sits between the third and fourth llm spans
agent_id = TraceIdGenerator.span_id()
# agent span, not a thread row, sits between the third and fourth llm spans
insert_traces(
[
Traces(timestamp=now - timedelta(seconds=20), duration=timedelta(seconds=19), trace_id=trace_id, span_id=root_id, name="POST /chat", kind=TracesKind.SPAN_KIND_SERVER, resources=resources, attribute_write_mode="json_only"),
@@ -189,7 +211,7 @@ def test_thread_opens_around_span(
Traces(timestamp=now - timedelta(seconds=18 - 3 * i), trace_id=trace_id, span_id=span_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": span_id}])}, attribute_write_mode="json_only")
for i, span_id in enumerate(llm_ids)
),
Traces(timestamp=now - timedelta(seconds=11), trace_id=trace_id, span_id=tool_id, parent_span_id=root_id, name="execute_tool get_weather", resources=resources, attributes={"gen_ai.tool.name": "get_weather"}, attribute_write_mode="json_only"),
Traces(timestamp=now - timedelta(seconds=11), trace_id=trace_id, span_id=agent_id, parent_span_id=root_id, name="invoke_agent planner", resources=resources, attributes={"gen_ai.agent.name": "planner"}, attribute_write_mode="json_only"),
]
)
@@ -211,10 +233,10 @@ def test_thread_opens_around_span(
assert [span["span_id"] for span in get_page({"limit": 3, "after": around["nextCursor"]})["spans"]] == llm_ids[4:]
# span without messages: only its neighbours
around_tool = get_page({"limit": 2, "spanId": tool_id})
assert [span["span_id"] for span in around_tool["spans"]] == llm_ids[2:4]
assert around_tool["prevCursor"]
assert around_tool["nextCursor"]
around_agent = get_page({"limit": 2, "spanId": agent_id})
assert [span["span_id"] for span in around_agent["spans"]] == llm_ids[2:4]
assert around_agent["prevCursor"]
assert around_agent["nextCursor"]
# near the start, the short side gives its room to the other
at_start = get_page({"limit": 3, "spanId": llm_ids[0]})
@@ -259,7 +281,7 @@ def test_thread_reads_spans_across_json_rollout(
insert_traces(
[
Traces(timestamp=rollout - timedelta(minutes=10), trace_id=before_trace_id, span_id=before_ids[0], name="chat gpt-4o", resources=resources, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": "first"}])}, attribute_write_mode="legacy_only"),
Traces(timestamp=rollout - timedelta(minutes=8), trace_id=before_trace_id, span_id=TraceIdGenerator.span_id(), name="execute_tool get_weather", resources=resources, attributes={"gen_ai.tool.name": "get_weather"}, attribute_write_mode="legacy_only"),
Traces(timestamp=rollout - timedelta(minutes=8), trace_id=before_trace_id, span_id=TraceIdGenerator.span_id(), name="invoke_agent planner", resources=resources, attributes={"gen_ai.agent.name": "planner"}, attribute_write_mode="legacy_only"),
Traces(timestamp=rollout - timedelta(minutes=5), trace_id=before_trace_id, span_id=before_ids[1], name="chat gpt-4o", resources=resources, attributes={"gen_ai.output.messages": json.dumps([{"role": "assistant", "content": "second"}])}, attribute_write_mode="legacy_only"),
Traces(timestamp=rollout - timedelta(minutes=5), trace_id=straddle_trace_id, span_id=legacy_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": "legacy"}])}, attribute_write_mode="legacy_only"),
Traces(timestamp=rollout + timedelta(minutes=5), trace_id=straddle_trace_id, span_id=json_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": "json"}])}, attribute_write_mode="json_only"),
@@ -275,6 +297,11 @@ def test_thread_reads_spans_across_json_rollout(
assert [span["span_id"] for span in before_spans] == before_ids
assert before_spans[0]["attributes"]["gen_ai.input.messages"] == json.dumps([{"role": "user", "content": "first"}])
assert before_spans[1]["attributes"]["gen_ai.output.messages"] == json.dumps([{"role": "assistant", "content": "second"}])
# messages are decoded from the JSON column only, never from the legacy maps
for span in before_spans:
assert span["formatter"] == ""
assert span["formatted_input"] == []
assert span["formatted_output"] == []
straddle = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{straddle_trace_id}/thread"), headers=headers, timeout=10)
assert straddle.status_code == HTTPStatus.OK, straddle.text
@@ -282,6 +309,10 @@ def test_thread_reads_spans_across_json_rollout(
assert [span["span_id"] for span in straddle_spans] == [legacy_id, json_id]
assert straddle_spans[0]["attributes"]["gen_ai.input.messages"] == json.dumps([{"role": "user", "content": "legacy"}])
assert straddle_spans[1]["attributes"]["gen_ai.input.messages"] == json.dumps([{"role": "user", "content": "json"}])
assert straddle_spans[0]["formatter"] == ""
assert straddle_spans[0]["formatted_input"] == []
assert straddle_spans[1]["formatter"] == "openai.chat"
assert straddle_spans[1]["formatted_input"] == [{"role": "user", "parts": [{"type": "text", "content": "json"}]}]
def test_thread_without_messages_is_empty(