Compare commits

...

4 Commits

Author SHA1 Message Date
nityanandagohain
61f6370f45 fix: remove changes from waterfall 2026-09-29 12:42:30 +05:30
nityanandagohain
f8b22c0feb refactor(tracedetail): split message normalisation out of the thread api 2026-09-29 11:20:56 +05:30
nityanandagohain
9e3a6bb35d Merge remote-tracking branch 'origin/main' into feat/ai-trace-thread 2026-09-29 10:25:57 +05:30
nityanandagohain
f5f019f61b feat: trace detail thread endpoint 2026-09-28 17:32:55 +05:30
17 changed files with 1045 additions and 60 deletions

View File

@@ -68,6 +68,7 @@ jobs:
- semconvfamilies
- serviceaccount
- spanmapper
- tracedetail
- querier_json_body
- querier_skip_resource_fingerprint
- ttl

View File

@@ -9647,6 +9647,17 @@ components:
required:
- aggregations
type: object
SpantypesGettableTraceThread:
properties:
nextCursor:
type: string
spans:
items:
$ref: '#/components/schemas/SpantypesThreadSpan'
type: array
required:
- spans
type: object
SpantypesGettableWaterfallTrace:
properties:
endTimestampMillis:
@@ -9961,6 +9972,84 @@ components:
nullable: true
type: object
type: object
SpantypesThreadSpan:
properties:
attributes:
additionalProperties: {}
nullable: true
type: object
db_name:
type: string
db_operation:
type: string
duration_nano:
minimum: 0
type: integer
events:
items:
$ref: '#/components/schemas/SpantypesEvent'
nullable: true
type: array
external_http_method:
type: string
external_http_url:
type: string
flags:
minimum: 0
type: integer
has_children:
type: boolean
has_error:
type: boolean
http_host:
type: string
http_method:
type: string
http_url:
type: string
is_remote:
type: string
kind_string:
type: string
level:
minimum: 0
type: integer
name:
type: string
parent_span_id:
type: string
references:
items:
$ref: '#/components/schemas/SpantypesOtelSpanRef'
type: array
resource:
additionalProperties:
type: string
nullable: true
type: object
response_status_code:
type: string
span_id:
type: string
status_code:
type: integer
status_code_string:
type: string
status_message:
type: string
sub_tree_node_count:
minimum: 0
type: integer
time_unix:
minimum: 0
type: integer
trace_id:
type: string
trace_state:
type: string
required:
- references
type: object
SpantypesUpdatableSpanMapper:
properties:
config:
@@ -15685,6 +15774,80 @@ paths:
tags:
- tracedetail
x-signoz-stability: alpha
/api/v1/traces/{traceID}/thread:
get:
deprecated: false
description: Returns the spans carrying gen_ai input or output messages in timestamp
order. Pages are fetched with the returned nextCursor.
operationId: GetTraceThread
parameters:
- in: query
name: limit
schema:
type: integer
- in: query
name: cursor
schema:
type: string
- in: path
name: traceID
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/SpantypesGettableTraceThread'
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: Get thread view for a trace
tags:
- tracedetail
x-signoz-stability: alpha
/api/v1/user/me:
get:
deprecated: true

View File

@@ -11142,6 +11142,157 @@ export interface SpantypesOtelSpanRefDTO {
traceId?: string;
}
export type SpantypesThreadSpanDTOAttributesAnyOf = { [key: string]: unknown };
/**
* @nullable
*/
export type SpantypesThreadSpanDTOAttributes =
SpantypesThreadSpanDTOAttributesAnyOf | null;
export type SpantypesThreadSpanDTOResourceAnyOf = { [key: string]: string };
/**
* @nullable
*/
export type SpantypesThreadSpanDTOResource =
SpantypesThreadSpanDTOResourceAnyOf | null;
export interface SpantypesThreadSpanDTO {
/**
* @type object,null
*/
attributes?: SpantypesThreadSpanDTOAttributes;
/**
* @type string
*/
db_name?: string;
/**
* @type string
*/
db_operation?: string;
/**
* @type integer
* @minimum 0
*/
duration_nano?: number;
/**
* @type array,null
*/
events?: SpantypesEventDTO[] | null;
/**
* @type string
*/
external_http_method?: string;
/**
* @type string
*/
external_http_url?: string;
/**
* @type integer
* @minimum 0
*/
flags?: number;
/**
* @type boolean
*/
has_children?: boolean;
/**
* @type boolean
*/
has_error?: boolean;
/**
* @type string
*/
http_host?: string;
/**
* @type string
*/
http_method?: string;
/**
* @type string
*/
http_url?: string;
/**
* @type string
*/
is_remote?: string;
/**
* @type string
*/
kind_string?: string;
/**
* @type integer
* @minimum 0
*/
level?: number;
/**
* @type string
*/
name?: string;
/**
* @type string
*/
parent_span_id?: string;
/**
* @type array
*/
references: SpantypesOtelSpanRefDTO[];
/**
* @type object,null
*/
resource?: SpantypesThreadSpanDTOResource;
/**
* @type string
*/
response_status_code?: string;
/**
* @type string
*/
span_id?: string;
/**
* @type integer
*/
status_code?: number;
/**
* @type string
*/
status_code_string?: string;
/**
* @type string
*/
status_message?: string;
/**
* @type integer
* @minimum 0
*/
sub_tree_node_count?: number;
/**
* @type integer
* @minimum 0
*/
time_unix?: number;
/**
* @type string
*/
trace_id?: string;
/**
* @type string
*/
trace_state?: string;
}
export interface SpantypesGettableTraceThreadDTO {
/**
* @type string
*/
nextCursor?: string;
/**
* @type array
*/
spans: SpantypesThreadSpanDTO[];
}
export type SpantypesWaterfallSpanDTOAttributesAnyOf = {
[key: string]: unknown;
};
@@ -12815,6 +12966,30 @@ export type GetTraceAggregations200 = {
status: string;
};
export type GetTraceThreadPathParameters = {
traceID: string;
};
export type GetTraceThreadParams = {
/**
* @type integer
* @description undefined
*/
limit?: number;
/**
* @type string
* @description undefined
*/
cursor?: string;
};
export type GetTraceThread200 = {
data: SpantypesGettableTraceThreadDTO;
/**
* @type string
*/
status: string;
};
export type ListUserPreferences200 = {
/**
* @type array

View File

@@ -4,11 +4,17 @@
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
import { useMutation } from 'react-query';
import { useMutation, useQuery } from 'react-query';
import type {
InvalidateOptions,
MutationFunction,
QueryClient,
QueryFunction,
QueryKey,
UseMutationOptions,
UseMutationResult,
UseQueryOptions,
UseQueryResult,
} from 'react-query';
import type {
@@ -16,6 +22,9 @@ import type {
GetFlamegraphPathParameters,
GetTraceAggregations200,
GetTraceAggregationsPathParameters,
GetTraceThread200,
GetTraceThreadParams,
GetTraceThreadPathParameters,
GetWaterfallV4200,
GetWaterfallV4PathParameters,
RenderErrorResponseDTO,
@@ -27,6 +36,26 @@ import type {
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
const withQueryKey = <T extends object, K>(
query: T,
queryKey: K,
): T & { queryKey: K } => {
const result = { queryKey } as T & { queryKey: K };
for (const key of Object.keys(query)) {
// The explicit queryKey always wins, matching the previous
// `{ ...query, queryKey }` spread where it was set last.
if (key === 'queryKey') {
continue;
}
Object.defineProperty(result, key, {
enumerable: true,
configurable: true,
get: () => (query as Record<string, unknown>)[key],
});
}
return result;
};
/**
* Computes span aggregations grouped by requested field.
* @summary Get aggregations for a trace
@@ -127,6 +156,121 @@ export const useGetTraceAggregations = <
> => {
return useMutation(getGetTraceAggregationsMutationOptions(options));
};
/**
* Returns the spans carrying gen_ai input or output messages in timestamp order. Pages are fetched with the returned nextCursor.
* @summary Get thread view for a trace
*/
export const getTraceThread = (
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
signal?: AbortSignal,
) => {
return GeneratedAPIInstance<GetTraceThread200>({
url: `/api/v1/traces/${traceID}/thread`,
method: 'GET',
params,
signal,
});
};
export const getGetTraceThreadQueryKey = (
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
) => {
return [
`/api/v1/traces/${traceID}/thread`,
...(params ? [params] : []),
] as const;
};
export const getGetTraceThreadQueryOptions = <
TData = Awaited<ReturnType<typeof getTraceThread>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getTraceThread>>,
TError,
TData
>;
},
) => {
const { query: queryOptions } = options ?? {};
const queryKey =
queryOptions?.queryKey ?? getGetTraceThreadQueryKey({ traceID }, params);
const queryFn: QueryFunction<Awaited<ReturnType<typeof getTraceThread>>> = ({
signal,
}) => getTraceThread({ traceID }, params, signal);
return {
queryKey,
queryFn,
enabled: traceID !== null && traceID !== undefined,
...queryOptions,
} as UseQueryOptions<
Awaited<ReturnType<typeof getTraceThread>>,
TError,
TData
> & { queryKey: QueryKey };
};
export type GetTraceThreadQueryResult = NonNullable<
Awaited<ReturnType<typeof getTraceThread>>
>;
export type GetTraceThreadQueryError = ErrorType<RenderErrorResponseDTO>;
/**
* @summary Get thread view for a trace
*/
export function useGetTraceThread<
TData = Awaited<ReturnType<typeof getTraceThread>>,
TError = ErrorType<RenderErrorResponseDTO>,
>(
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
options?: {
query?: UseQueryOptions<
Awaited<ReturnType<typeof getTraceThread>>,
TError,
TData
>;
},
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
const queryOptions = getGetTraceThreadQueryOptions(
{ traceID },
params,
options,
);
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
queryKey: QueryKey;
};
return withQueryKey(query, queryOptions.queryKey);
}
/**
* @summary Get thread view for a trace
*/
export const invalidateGetTraceThread = async (
queryClient: QueryClient,
{ traceID }: GetTraceThreadPathParameters,
params?: GetTraceThreadParams,
options?: InvalidateOptions,
): Promise<QueryClient> => {
await queryClient.invalidateQueries(
{ queryKey: getGetTraceThreadQueryKey({ traceID }, params) },
options,
);
return queryClient;
};
/**
* Returns the flamegraph view of spans for a given trace ID.
* @summary Get flamegraph view for a trace

View File

@@ -67,5 +67,23 @@ func (provider *provider) addTraceDetailRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v1/traces/{traceID}/thread", handler.New(
provider.authzMiddleware.ViewAccess(provider.traceDetailHandler.GetThread),
handler.OpenAPIDef{
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. Pages are fetched with the returned nextCursor.",
RequestQuery: new(spantypes.QueryableThread),
Response: new(spantypes.GettableTraceThread),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
},
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
return nil
}

View File

@@ -75,3 +75,25 @@ func (h *handler) GetFlamegraph(rw http.ResponseWriter, r *http.Request) {
render.Success(rw, http.StatusOK, result)
}
func (h *handler) GetThread(rw http.ResponseWriter, r *http.Request) {
req := new(spantypes.QueryableThread)
if err := binding.Query.BindQuery(r.URL.Query(), req); err != nil {
render.Error(rw, err)
return
}
query, err := spantypes.NewThreadQuery(req)
if err != nil {
render.Error(rw, err)
return
}
result, err := h.module.GetThread(r.Context(), mux.Vars(r)["traceID"], query)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, result)
}

View File

@@ -173,6 +173,19 @@ func (m *module) getWindowedWaterfall(ctx context.Context, traceID, selectedSpan
), nil
}
func (m *module) GetThread(ctx context.Context, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error) {
summary, err := m.store.GetTraceSummary(ctx, traceID)
if err != nil {
return nil, err
}
spans, err := m.store.GetThreadSpans(ctx, traceID, summary, query.Cursor, query.Limit+1)
if err != nil {
return nil, err
}
return spantypes.NewGettableTraceThread(traceID, spans, query.Limit), nil
}
func (m *module) getFullFlamegraph(ctx context.Context, traceID string, summary *spantypes.TraceSummary, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error) {
fullSpans, err := m.store.GetFlamegraphSpans(ctx, traceID, summary.Start, summary.End, nil)
if err != nil {

View File

@@ -11,12 +11,23 @@ import (
"github.com/SigNoz/signoz/pkg/clickhousesql"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
const colServiceName = `resource_string_service$$$$name` // $ gets escaped so $$$$ converts to $$.
var fullSpanColumns = []string{
"duration_nano", "span_id", "has_error", "kind",
colServiceName, "name",
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
"flags", "is_remote", "trace_state", "status_code",
"db_name", "db_operation", "http_method", "http_url", "http_host",
"external_http_method", "external_http_url", "response_status_code", "links as references",
}
func buildFieldExpr(fieldKey telemetrytypes.TelemetryFieldKey) (string, error) {
switch fieldKey.FieldContext {
case telemetrytypes.FieldContextResource:
@@ -123,16 +134,8 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return []spantypes.StorableSpan{}, nil
}
sb := sqlbuilder.NewSelectBuilder()
sb.Select(
"DISTINCT ON (span_id) timestamp",
"duration_nano", "span_id", "has_error", "kind",
colServiceName, "name",
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
"flags", "is_remote", "trace_state", "status_code",
"db_name", "db_operation", "http_method", "http_url", "http_host",
"external_http_method", "external_http_url", "response_status_code", "links as references",
)
sb.Select("DISTINCT ON (span_id) timestamp")
sb.SelectMore(fullSpanColumns...)
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
ids := make([]any, len(spanIDs))
for i, id := range spanIDs {
@@ -155,6 +158,44 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return spans, nil
}
func (s *traceStore) GetThreadSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary, cursor *spantypes.ThreadCursor, limit int) ([]spantypes.StorableSpan, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("DISTINCT ON (span_id) timestamp")
sb.SelectMore(fullSpanColumns...)
sb.SelectMore("attributes")
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
sb.Where(
sb.E("trace_id", traceID),
sb.GE("ts_bucket_start", summary.Start.Unix()-1800),
sb.LE("ts_bucket_start", summary.End.Unix()),
// Reads only the JSON column; spans with messages only in the legacy maps are skipped.
// todo(nitya): pick the column from the attribute evolution metadata.
sb.Or(
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIInputMessages))),
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIOutputMessages))),
),
)
if cursor != nil {
// ClickHouse can't use an index for a tuple comparison, so the separate timestamp and
// ts_bucket_start bounds are what skip the data before the cursor.
sb.Where(
sb.GE("ts_bucket_start", int64(cursor.TimeUnixNano/uint64(time.Second))-1800),
sb.GE("timestamp", fmt.Sprintf("%d", cursor.TimeUnixNano)),
sb.GT("(toUnixTimestamp64Nano(timestamp), span_id)", sqlbuilder.Tuple(cursor.TimeUnixNano, cursor.SpanID)),
)
}
sb.OrderByAsc("timestamp")
sb.OrderByAsc("span_id")
sb.Limit(limit)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
var spans []spantypes.StorableSpan
if err := s.telemetryStore.ClickhouseDB().Select(ctx, &spans, query, args...); err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error querying thread spans")
}
return spans, nil
}
func (s *traceStore) GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]spantypes.StorableSpan, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select(

View File

@@ -13,6 +13,7 @@ type Handler interface {
GetWaterfallV4(http.ResponseWriter, *http.Request)
GetTraceAggregations(http.ResponseWriter, *http.Request)
GetFlamegraph(http.ResponseWriter, *http.Request)
GetThread(http.ResponseWriter, *http.Request)
}
// Module defines the business logic for trace detail operations.
@@ -20,4 +21,5 @@ type Module interface {
GetWaterfallV4(ctx context.Context, traceID string, selectedSpanID string, uncollapsedSpans []string) (*spantypes.GettableWaterfallTrace, error)
GetTraceAggregations(ctx context.Context, traceID string, req *spantypes.PostableTraceAggregations) (*spantypes.GettableTraceAggregations, error)
GetFlamegraph(ctx context.Context, traceID string, selectedSpanID string, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error)
GetThread(ctx context.Context, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error)
}

View File

@@ -566,24 +566,6 @@ func readAsRaw(rows driver.Rows, queryName string) (*qbtypes.RawData, error) {
}, nil
}
// flattenJSONPaths flattens a decoded JSON document into dotted keys, overwriting existing keys in out.
func flattenJSONPaths(prefix string, m map[string]any, out map[string]any) {
for k, v := range m {
key := k
if prefix != "" {
key = prefix + "." + k
}
switch child := v.(type) {
case map[string]any:
flattenJSONPaths(key, child, out)
case telemetrystoretypes.JSONValue:
flattenJSONPaths(key, child, out)
default:
out[key] = v
}
}
}
// mergeSpanAttributeColumns merges (attributes_string, attributes_number, attributes_bool, resources_string) into
// unified "attributes" and "resource" keys, and parses the stringified `events`
// and `links` columns into structured slices. Raw DB columns are removed.
@@ -598,7 +580,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
resStr, hasRes := data["resources_string"]
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {
attributes := make(map[string]any)
flattenJSONPaths("", attrJSON, attributes)
attrJSON.FlattenInto("", attributes)
if m, ok := attrStr.(map[string]string); ok {
for k, v := range m {
attributes[k] = v

View File

@@ -36,6 +36,7 @@ type TraceStore interface {
GetMinimalSpans(ctx context.Context, traceID string, start, end time.Time) ([]MinimalSpan, error)
GetTraceSpansByIDs(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
GetThreadSpans(ctx context.Context, traceID string, summary *TraceSummary, cursor *ThreadCursor, limit int) ([]StorableSpan, error)
GetSpanCountByField(ctx context.Context, traceID string, summary *TraceSummary, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)
GetSpanDurationByField(ctx context.Context, traceID string, summary *TraceSummary, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)

View File

@@ -0,0 +1,113 @@
package spantypes
import (
"encoding/base64"
"encoding/json"
"github.com/SigNoz/signoz/pkg/errors"
)
const (
threadDefaultLimit = 100
threadMaxLimit = 1000
)
var (
ErrCodeThreadInvalidLimit = errors.MustNewCode("trace_thread_invalid_limit")
ErrCodeThreadInvalidCursor = errors.MustNewCode("trace_thread_invalid_cursor")
)
type QueryableThread struct {
// Limit is the page size; 0 means 100.
Limit int `query:"limit"`
// Cursor is the nextCursor of the previous page; empty for the first page.
Cursor string `query:"cursor"`
}
type ThreadQuery struct {
Limit int
Cursor *ThreadCursor
}
func NewThreadQuery(queryable *QueryableThread) (*ThreadQuery, error) {
query := &ThreadQuery{Limit: queryable.Limit}
if query.Limit < 0 {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidLimit, "limit cannot be negative, got %d", query.Limit)
}
if query.Limit == 0 {
query.Limit = threadDefaultLimit
}
if query.Limit > threadMaxLimit {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidLimit, "limit cannot exceed %d, got %d", threadMaxLimit, query.Limit)
}
if queryable.Cursor != "" {
cursor, err := DecodeThreadCursor(queryable.Cursor)
if err != nil {
return nil, err
}
query.Cursor = cursor
}
return query, nil
}
// ThreadCursor is the (TimeUnixNano, SpanID) of the last span of a page.
type ThreadCursor struct {
TimeUnixNano uint64 `json:"t"`
SpanID string `json:"s"`
}
func (c ThreadCursor) Encode() string {
data, _ := json.Marshal(c)
return base64.RawURLEncoding.EncodeToString(data)
}
func DecodeThreadCursor(cursor string) (*ThreadCursor, error) {
data, err := base64.RawURLEncoding.DecodeString(cursor)
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
}
c := new(ThreadCursor)
if err := json.Unmarshal(data, c); err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
}
if c.SpanID == "" {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidCursor, "invalid cursor: missing span id")
}
return c, nil
}
type GettableTraceThread struct {
Spans []*ThreadSpan `json:"spans" required:"true" nullable:"false"`
NextCursor string `json:"nextCursor,omitempty"`
}
type ThreadSpan struct {
WaterfallSpan
}
// NewGettableTraceThread expects limit+1 spans; the extra one only signals a next page.
func NewGettableTraceThread(traceID string, spans []StorableSpan, limit int) *GettableTraceThread {
hasMore := len(spans) > limit
if hasMore {
spans = spans[:limit]
}
out := make([]*ThreadSpan, len(spans))
for i := range spans {
out[i] = newThreadSpan(traceID, &spans[i])
}
thread := &GettableTraceThread{Spans: out}
if hasMore {
last := spans[len(spans)-1]
thread.NextCursor = ThreadCursor{TimeUnixNano: uint64(last.StartTime.UnixNano()), SpanID: last.SpanID}.Encode()
}
return thread
}
func newThreadSpan(traceID string, storable *StorableSpan) *ThreadSpan {
span := &ThreadSpan{WaterfallSpan: *storable.ToWaterfallSpan(traceID)}
// client expects millis, as in the waterfall
span.TimeUnix = span.TimeUnix / 1_000_000
return span
}

View File

@@ -0,0 +1,102 @@
package spantypes
import (
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestNewThreadQuery(t *testing.T) {
cursor := ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}
testCases := []struct {
name string
queryable QueryableThread
want *ThreadQuery
wantErr bool
}{
{name: "ZeroLimit_UsesDefault", queryable: QueryableThread{}, want: &ThreadQuery{Limit: threadDefaultLimit}},
{name: "PositiveLimit_Kept", queryable: QueryableThread{Limit: 25}, want: &ThreadQuery{Limit: 25}},
{name: "MaxLimit_Kept", queryable: QueryableThread{Limit: threadMaxLimit}, want: &ThreadQuery{Limit: threadMaxLimit}},
{name: "AboveMaxLimit_Rejected", queryable: QueryableThread{Limit: threadMaxLimit + 1}, wantErr: true},
{name: "NegativeLimit_Rejected", queryable: QueryableThread{Limit: -1}, wantErr: true},
{name: "Cursor_Decoded", queryable: QueryableThread{Limit: 10, Cursor: cursor.Encode()}, want: &ThreadQuery{Limit: 10, Cursor: &cursor}},
{name: "InvalidCursor_Rejected", queryable: QueryableThread{Cursor: "not base64!"}, wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
got, err := NewThreadQuery(&testCase.queryable)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, got)
})
}
}
func TestDecodeThreadCursor(t *testing.T) {
testCases := []struct {
name string
cursor string
want *ThreadCursor
wantErr bool
}{
{name: "EncodedCursor_RoundTrips", cursor: ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}.Encode(), want: &ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}},
{name: "NotBase64_Rejected", cursor: "not base64!", wantErr: true},
{name: "NotJSON_Rejected", cursor: "bm90IGpzb24", wantErr: true},
{name: "MissingSpanID_Rejected", cursor: "eyJ0IjogMX0", wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
got, err := DecodeThreadCursor(testCase.cursor)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, got)
})
}
}
func TestNewGettableTraceThread(t *testing.T) {
spans := []StorableSpan{
{SpanID: "a", StartTime: time.Unix(1, 500_000_000)},
{SpanID: "b", StartTime: time.Unix(2, 0)},
{SpanID: "c", StartTime: time.Unix(3, 0)},
}
testCases := []struct {
name string
spans []StorableSpan
limit int
wantSpanIDs []string
wantTimeUnix []uint64
wantNextCursor string
}{
{name: "MoreThanLimit_TrimsAndSetsCursor", spans: spans, limit: 2, wantSpanIDs: []string{"a", "b"}, wantTimeUnix: []uint64{1500, 2000}, wantNextCursor: ThreadCursor{TimeUnixNano: 2_000_000_000, SpanID: "b"}.Encode()},
{name: "WithinLimit_NoCursor", spans: spans, limit: 3, wantSpanIDs: []string{"a", "b", "c"}, wantTimeUnix: []uint64{1500, 2000, 3000}},
{name: "NoSpans_EmptyList", limit: 3, wantSpanIDs: []string{}, wantTimeUnix: []uint64{}},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
thread := NewGettableTraceThread("trace-1", testCase.spans, testCase.limit)
require.NotNil(t, thread.Spans)
spanIDs := make([]string, len(thread.Spans))
timeUnix := make([]uint64, len(thread.Spans))
for i, span := range thread.Spans {
spanIDs[i] = span.SpanID
timeUnix[i] = span.TimeUnix
assert.Equal(t, "trace-1", span.TraceID)
}
assert.Equal(t, testCase.wantSpanIDs, spanIDs)
assert.Equal(t, testCase.wantTimeUnix, timeUnix)
assert.Equal(t, testCase.wantNextCursor, thread.NextCursor)
})
}
}

View File

@@ -8,6 +8,7 @@ import (
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
@@ -93,35 +94,36 @@ type WaterfallSpan struct {
// StorableSpan is the ClickHouse scan struct for the v3 waterfall query.
type StorableSpan struct {
StartTime time.Time `ch:"timestamp"`
DurationNano uint64 `ch:"duration_nano"`
SpanID string `ch:"span_id"`
HasError bool `ch:"has_error"`
Kind int8 `ch:"kind"`
ServiceName string `ch:"resource_string_service$$name"`
Name string `ch:"name"`
AttributesString map[string]string `ch:"attributes_string"`
AttributesNumber map[string]float64 `ch:"attributes_number"`
AttributesBool map[string]bool `ch:"attributes_bool"`
ResourcesString map[string]string `ch:"resources_string"`
Events []string `ch:"events"`
StatusMessage string `ch:"status_message"`
StatusCodeString string `ch:"status_code_string"`
SpanKind string `ch:"kind_string"`
ParentSpanID string `ch:"parent_span_id"`
Flags uint32 `ch:"flags"`
IsRemote string `ch:"is_remote"`
TraceState string `ch:"trace_state"`
StatusCode int16 `ch:"status_code"`
DBName string `ch:"db_name"`
DBOperation string `ch:"db_operation"`
HTTPMethod string `ch:"http_method"`
HTTPURL string `ch:"http_url"`
HTTPHost string `ch:"http_host"`
ExternalHTTPMethod string `ch:"external_http_method"`
ExternalHTTPURL string `ch:"external_http_url"`
ResponseStatusCode string `ch:"response_status_code"`
References string `ch:"references"`
StartTime time.Time `ch:"timestamp"`
DurationNano uint64 `ch:"duration_nano"`
SpanID string `ch:"span_id"`
HasError bool `ch:"has_error"`
Kind int8 `ch:"kind"`
ServiceName string `ch:"resource_string_service$$name"`
Name string `ch:"name"`
AttributesString map[string]string `ch:"attributes_string"`
AttributesNumber map[string]float64 `ch:"attributes_number"`
AttributesBool map[string]bool `ch:"attributes_bool"`
AttributesJSON telemetrystoretypes.JSONValue `ch:"attributes"`
ResourcesString map[string]string `ch:"resources_string"`
Events []string `ch:"events"`
StatusMessage string `ch:"status_message"`
StatusCodeString string `ch:"status_code_string"`
SpanKind string `ch:"kind_string"`
ParentSpanID string `ch:"parent_span_id"`
Flags uint32 `ch:"flags"`
IsRemote string `ch:"is_remote"`
TraceState string `ch:"trace_state"`
StatusCode int16 `ch:"status_code"`
DBName string `ch:"db_name"`
DBOperation string `ch:"db_operation"`
HTTPMethod string `ch:"http_method"`
HTTPURL string `ch:"http_url"`
HTTPHost string `ch:"http_host"`
ExternalHTTPMethod string `ch:"external_http_method"`
ExternalHTTPURL string `ch:"external_http_url"`
ResponseStatusCode string `ch:"response_status_code"`
References string `ch:"references"`
}
// MinimalSpan with only the fields needed to build the parent-child tree.
@@ -277,8 +279,10 @@ func (item *StorableSpan) AttributeValue(name string) any {
return nil
}
// Attributes flattens the JSON column first, so the legacy maps win on collision.
func (item *StorableSpan) Attributes() map[string]any {
attributes := make(map[string]any, len(item.AttributesString)+len(item.AttributesNumber)+len(item.AttributesBool))
attributes := make(map[string]any, len(item.AttributesString)+len(item.AttributesNumber)+len(item.AttributesBool)+len(item.AttributesJSON))
item.AttributesJSON.FlattenInto("", attributes)
for k, v := range item.AttributesString {
attributes[k] = v
}

View File

@@ -0,0 +1,59 @@
package spantypes
import (
"testing"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/stretchr/testify/assert"
)
func TestStorableSpanAttributes(t *testing.T) {
testCases := []struct {
name string
span StorableSpan
wantAttrs map[string]any
}{
{
name: "LegacyMapOnly_Kept",
span: StorableSpan{AttributesString: map[string]string{
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
}},
wantAttrs: map[string]any{
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
},
},
{
name: "JSONColumn_FlattenedToDottedKeys",
span: StorableSpan{AttributesJSON: telemetrystoretypes.JSONValue{
"gen_ai": map[string]any{
"input": map[string]any{"messages": `[{"role":"user","content":"hi"}]`},
"request": map[string]any{"model": "gpt-4o"},
},
}},
wantAttrs: map[string]any{
"gen_ai.input.messages": `[{"role":"user","content":"hi"}]`,
"gen_ai.request.model": "gpt-4o",
},
},
{
name: "LegacyMapWinsOverJSONColumn",
span: StorableSpan{
AttributesJSON: telemetrystoretypes.JSONValue{"gen_ai": map[string]any{"request": map[string]any{"model": "json"}}},
AttributesString: map[string]string{"gen_ai.request.model": "map"},
},
wantAttrs: map[string]any{"gen_ai.request.model": "map"},
},
{
name: "NoMessages_AttributesKept",
span: StorableSpan{AttributesString: map[string]string{"http.method": "GET"}},
wantAttrs: map[string]any{"http.method": "GET"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.wantAttrs, testCase.span.Attributes())
})
}
}

View File

@@ -35,3 +35,21 @@ func (v *JSONValue) Scan(src any) error {
*v = decoded
return nil
}
// FlattenInto writes v into out under dotted keys, overwriting existing keys.
func (v JSONValue) FlattenInto(prefix string, out map[string]any) {
for k, value := range v {
key := k
if prefix != "" {
key = prefix + "." + k
}
switch child := value.(type) {
case map[string]any:
JSONValue(child).FlattenInto(key, out)
case JSONValue:
child.FlattenInto(key, out)
default:
out[key] = value
}
}
}

View File

@@ -0,0 +1,127 @@
import json
from collections.abc import Callable
from datetime import UTC, datetime, timedelta
from http import HTTPStatus
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
from fixtures.traces import TraceIdGenerator, Traces, TracesKind
def test_thread_returns_message_spans_in_order(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
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))
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"}])
second_input = json.dumps([{"role": "tool", "content": "sunny", "tool_call_id": "call_1"}])
insert_traces(
[
Traces(timestamp=now - timedelta(seconds=10), duration=timedelta(seconds=9), trace_id=trace_id, span_id=root_id, name="POST /chat", kind=TracesKind.SPAN_KIND_SERVER, resources=resources, attribute_write_mode="json_only"),
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=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"),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread"), headers={"Authorization": f"Bearer {token}"}, timeout=10)
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 "nextCursor" not in thread
first, input_only, output_only = thread["spans"]
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"
assert first["attributes"]["gen_ai.output.messages"] == first_output
assert input_only["attributes"]["gen_ai.input.messages"] == second_input
assert "gen_ai.output.messages" not in input_only["attributes"]
assert "gen_ai.input.messages" not in output_only["attributes"]
assert output_only["attributes"]["gen_ai.output.messages"] == "It is sunny in Bangalore."
def test_thread_paginates_with_cursor(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
now = datetime.now(tz=UTC).replace(microsecond=0)
trace_id = TraceIdGenerator.trace_id()
span_ids = [TraceIdGenerator.span_id() for _ in range(3)]
# identical timestamps on the last two exercise the span_id tie-break
timestamps = [now - timedelta(seconds=6), now - timedelta(seconds=3), now - timedelta(seconds=3)]
insert_traces(
[
Traces(timestamp=timestamp, trace_id=trace_id, span_id=span_id, name="chat gpt-4o", resources={"service.name": "tracedetail-thread-pages"}, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": span_id}])}, attribute_write_mode="json_only")
for span_id, timestamp in zip(span_ids, timestamps, strict=True)
]
)
expected_order = [span_ids[0], *sorted(span_ids[1:])]
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
url = signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread")
headers = {"Authorization": f"Bearer {token}"}
first_page = requests.get(url, params={"limit": 2}, headers=headers, timeout=10)
assert first_page.status_code == HTTPStatus.OK, first_page.text
first = first_page.json()["data"]
assert [span["span_id"] for span in first["spans"]] == expected_order[:2]
assert first["nextCursor"]
second_page = requests.get(url, params={"limit": 2, "cursor": first["nextCursor"]}, headers=headers, timeout=10)
assert second_page.status_code == HTTPStatus.OK, second_page.text
second = second_page.json()["data"]
assert [span["span_id"] for span in second["spans"]] == expected_order[2:]
assert "nextCursor" not in second
def test_thread_without_messages_is_empty(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
trace_id = TraceIdGenerator.trace_id()
insert_traces([Traces(timestamp=datetime.now(tz=UTC) - timedelta(seconds=5), trace_id=trace_id, span_id=TraceIdGenerator.span_id(), name="GET /health", resources={"service.name": "tracedetail-thread-empty"}, attribute_write_mode="json_only")])
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
response = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread"), headers={"Authorization": f"Bearer {token}"}, timeout=10)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"] == {"spans": []}
def test_thread_rejects_invalid_requests(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
trace_id = TraceIdGenerator.trace_id()
insert_traces([Traces(timestamp=datetime.now(tz=UTC) - timedelta(seconds=5), trace_id=trace_id, span_id=TraceIdGenerator.span_id(), name="chat gpt-4o", resources={"service.name": "tracedetail-thread-invalid"}, attributes={"gen_ai.input.messages": "hi"}, attribute_write_mode="json_only")])
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
headers = {"Authorization": f"Bearer {token}"}
url = signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread")
for params in ({"limit": -1}, {"limit": 1001}, {"cursor": "not-a-cursor"}):
response = requests.get(url, params=params, headers=headers, timeout=10)
assert response.status_code == HTTPStatus.BAD_REQUEST, f"{params}: {response.text}"
missing = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{TraceIdGenerator.trace_id()}/thread"), headers=headers, timeout=10)
assert missing.status_code == HTTPStatus.NOT_FOUND, missing.text