mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-29 14:50:41 +01:00
Compare commits
4 Commits
v0.144.0
...
feat/ai-tr
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
61f6370f45 | ||
|
|
f8b22c0feb | ||
|
|
9e3a6bb35d | ||
|
|
f5f019f61b |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -68,6 +68,7 @@ jobs:
|
||||
- semconvfamilies
|
||||
- serviceaccount
|
||||
- spanmapper
|
||||
- tracedetail
|
||||
- querier_json_body
|
||||
- querier_skip_resource_fingerprint
|
||||
- ttl
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
113
pkg/types/spantypes/thread.go
Normal file
113
pkg/types/spantypes/thread.go
Normal 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
|
||||
}
|
||||
102
pkg/types/spantypes/thread_test.go
Normal file
102
pkg/types/spantypes/thread_test.go
Normal 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)
|
||||
})
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
59
pkg/types/spantypes/waterfall_span_test.go
Normal file
59
pkg/types/spantypes/waterfall_span_test.go
Normal 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())
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
127
tests/integration/tests/tracedetail/01_thread.py
Normal file
127
tests/integration/tests/tracedetail/01_thread.py
Normal 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
|
||||
Reference in New Issue
Block a user