mirror of
https://github.com/SigNoz/signoz.git
synced 2026-10-09 19:50:43 +01:00
Compare commits
25 Commits
feat/trace
...
feat/ai-tr
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7b51c47e6d | ||
|
|
b9565c4866 | ||
|
|
1096a9f70c | ||
|
|
b19ba11d98 | ||
|
|
51663240e0 | ||
|
|
143c3a64f1 | ||
|
|
b4da88d5ec | ||
|
|
277124320a | ||
|
|
9a11ebf350 | ||
|
|
a4cb6ca847 | ||
|
|
f50d786155 | ||
|
|
bbd9fd1ad5 | ||
|
|
b0b3409827 | ||
|
|
37932c7e84 | ||
|
|
7313b02455 | ||
|
|
a1a38fc7da | ||
|
|
f76e3e3085 | ||
|
|
5b9d854518 | ||
|
|
0a9920583b | ||
|
|
ae0d8cb1a4 | ||
|
|
abb242b1f9 | ||
|
|
61f6370f45 | ||
|
|
f8b22c0feb | ||
|
|
9e3a6bb35d | ||
|
|
f5f019f61b |
@@ -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.
|
||||
|
||||
@@ -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
|
||||
*/
|
||||
|
||||
@@ -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 = (
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)",
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
57
pkg/types/aiobservabilitytypes/genai/message.go
Normal file
57
pkg/types/aiobservabilitytypes/genai/message.go
Normal 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"`
|
||||
}
|
||||
194
pkg/types/aiobservabilitytypes/genai/part.go
Normal file
194
pkg/types/aiobservabilitytypes/genai/part.go
Normal 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
|
||||
}
|
||||
74
pkg/types/aiobservabilitytypes/genai/part_test.go
Normal file
74
pkg/types/aiobservabilitytypes/genai/part_test.go
Normal 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)
|
||||
})
|
||||
}
|
||||
}
|
||||
252
pkg/types/aiobservabilitytypes/genaiformatter/formatter.go
Normal file
252
pkg/types/aiobservabilitytypes/genaiformatter/formatter.go
Normal 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)
|
||||
}
|
||||
101
pkg/types/aiobservabilitytypes/genaiformatter/formatter_test.go
Normal file
101
pkg/types/aiobservabilitytypes/genaiformatter/formatter_test.go
Normal 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: `[]`,
|
||||
},
|
||||
})
|
||||
}
|
||||
122
pkg/types/aiobservabilitytypes/genaiformatter/openai_chat.go
Normal file
122
pkg/types/aiobservabilitytypes/genaiformatter/openai_chat.go
Normal 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
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
]`,
|
||||
},
|
||||
})
|
||||
}
|
||||
70
pkg/types/aiobservabilitytypes/genaiformatter/semconv.go
Normal file
70
pkg/types/aiobservabilitytypes/genaiformatter/semconv.go
Normal 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
|
||||
}
|
||||
163
pkg/types/aiobservabilitytypes/genaiformatter/semconv_test.go
Normal file
163
pkg/types/aiobservabilitytypes/genaiformatter/semconv_test.go
Normal 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: `[]`,
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
Reference in New Issue
Block a user