Compare commits

..

3 Commits

Author SHA1 Message Date
nityanandagohain
dbae3f0115 fix: update name 2026-10-08 10:11:02 +05:30
nityanandagohain
ffcdf2aa67 chore: add migration to remove unwanted llm models 2026-10-08 09:45:25 +05:30
Nityananda Gohain
bb62c9c060 feat: trace detail thread endpoint (#12999)
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description
- Adds `GET /api/v1/traces/{traceID}/thread`, which returns the spans
that carry `gen_ai.input.messages` or `gen_ai.output.messages`, in
timestamp order.
- Pages with `after` / `before` cursors (`nextCursor` / `prevCursor`),
or `spanId` to open the page around a span (404 if absent); the three
are exclusive.
- Message normalisation (`formatted_input` / `formatted_output`) comes
in a stacked follow-up PR.

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR
Part of https://github.com/SigNoz/nerve-pod/issues/266

<!--Anything reviewers should keep in mind while reviewing -->
#### Additional Information
- The thread reads only the `attributes` JSON column, so spans that have
messages only in the old attribute maps don't appear.
  - The waterfall API is unchanged.
2026-10-07 12:08:32 +00:00
24 changed files with 1514 additions and 182 deletions

View File

@@ -9878,6 +9878,19 @@ components:
- totalErrorSpansCount
- hasMissingSpans
type: object
SpantypesGettableTraceThread:
properties:
nextCursor:
type: string
prevCursor:
type: string
spans:
items:
$ref: '#/components/schemas/SpantypesThreadSpan'
type: array
required:
- spans
type: object
SpantypesGettableWaterfallTrace:
properties:
endTimestampMillis:
@@ -10192,6 +10205,61 @@ components:
nullable: true
type: object
type: object
SpantypesThreadSpan:
properties:
attributes:
additionalProperties: {}
type: object
duration_nano:
minimum: 0
type: integer
events:
items:
$ref: '#/components/schemas/SpantypesEvent'
type: array
has_error:
type: boolean
kind_string:
type: string
name:
type: string
parent_span_id:
type: string
references:
items:
$ref: '#/components/schemas/SpantypesOtelSpanRef'
type: array
resource:
additionalProperties:
type: string
type: object
span_id:
type: string
status_code_string:
type: string
status_message:
type: string
time_unix:
minimum: 0
type: integer
trace_id:
type: string
required:
- span_id
- trace_id
- parent_span_id
- name
- kind_string
- time_unix
- duration_nano
- has_error
- status_code_string
- status_message
- resource
- attributes
- events
- references
type: object
SpantypesTraceAISummary:
properties:
tokens:
@@ -16009,6 +16077,103 @@ 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. 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.
in: query
name: limit
schema:
description: Page size, at most 100. 0 means 20.
type: integer
- description: The nextCursor of a page; returns the spans after it. Set only
one of after, before and spanId.
in: query
name: after
schema:
description: The nextCursor of a page; returns the spans after it. Set only
one of after, before and spanId.
type: string
- description: The prevCursor of a page; returns the spans before it. Set only
one of after, before and spanId.
in: query
name: before
schema:
description: The prevCursor of a page; returns the spans before it. Set
only one of after, before and spanId.
type: string
- description: Returns the page around this span. Set only one of after, before
and spanId.
in: query
name: spanId
schema:
description: Returns the page around this span. Set only one of after, before
and spanId.
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

@@ -11426,6 +11426,86 @@ export interface SpantypesOtelSpanRefDTO {
traceId?: string;
}
export type SpantypesThreadSpanDTOAttributes = { [key: string]: unknown };
export type SpantypesThreadSpanDTOResource = { [key: string]: string };
export interface SpantypesThreadSpanDTO {
/**
* @type object
*/
attributes: SpantypesThreadSpanDTOAttributes;
/**
* @type integer
* @minimum 0
*/
duration_nano: number;
/**
* @type array
*/
events: SpantypesEventDTO[];
/**
* @type boolean
*/
has_error: boolean;
/**
* @type string
*/
kind_string: string;
/**
* @type string
*/
name: string;
/**
* @type string
*/
parent_span_id: string;
/**
* @type array
*/
references: SpantypesOtelSpanRefDTO[];
/**
* @type object
*/
resource: SpantypesThreadSpanDTOResource;
/**
* @type string
*/
span_id: string;
/**
* @type string
*/
status_code_string: string;
/**
* @type string
*/
status_message: string;
/**
* @type integer
* @minimum 0
*/
time_unix: number;
/**
* @type string
*/
trace_id: string;
}
export interface SpantypesGettableTraceThreadDTO {
/**
* @type string
*/
nextCursor?: string;
/**
* @type string
*/
prevCursor?: string;
/**
* @type array
*/
spans: SpantypesThreadSpanDTO[];
}
export type SpantypesWaterfallSpanDTOAttributesAnyOf = {
[key: string]: unknown;
};
@@ -13110,6 +13190,40 @@ export type GetTraceSummary200 = {
status: string;
};
export type GetTraceThreadPathParameters = {
traceID: string;
};
export type GetTraceThreadParams = {
/**
* @type integer
* @description Page size, at most 100. 0 means 20.
*/
limit?: number;
/**
* @type string
* @description The nextCursor of a page; returns the spans after it. Set only one of after, before and spanId.
*/
after?: string;
/**
* @type string
* @description The prevCursor of a page; returns the spans before it. Set only one of after, before and spanId.
*/
before?: string;
/**
* @type string
* @description Returns the page around this span. Set only one of after, before and spanId.
*/
spanId?: string;
};
export type GetTraceThread200 = {
data: SpantypesGettableTraceThreadDTO;
/**
* @type string
*/
status: string;
};
export type ListUserPreferences200 = {
/**
* @type array

View File

@@ -24,6 +24,9 @@ import type {
GetTraceAggregationsPathParameters,
GetTraceSummary200,
GetTraceSummaryPathParameters,
GetTraceThread200,
GetTraceThreadParams,
GetTraceThreadPathParameters,
GetWaterfallV4200,
GetWaterfallV4PathParameters,
RenderErrorResponseDTO,
@@ -257,6 +260,121 @@ export const invalidateGetTraceSummary = async (
return queryClient;
};
/**
* 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.
* @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

@@ -27,14 +27,31 @@ import AnalyticsPanel from '../SpanDetailsPanel/AnalyticsPanel/AnalyticsPanel';
import Filters from '../TraceWaterfall/TraceWaterfallStates/Success/Filters/Filters';
import MissingSpansBanner from './MissingSpansBanner';
import TraceOptionsMenu from './TraceOptionsMenu';
import { useTraceSummary } from './useTraceSummary';
import styles from './TraceDetailsHeader.module.scss';
import { DATE_TIME_FORMATS } from 'constants/dateTimeFormats';
interface FilterMetadata {
startTime: number;
endTime: number;
traceId: string;
}
export interface TraceMetadataForHeader {
startTimestampMillis: number;
endTimestampMillis: number;
rootServiceName: string;
rootServiceEntryPoint: string;
rootSpanStatusCode: string;
hasMissingSpans: boolean;
totalSpansCount: number;
}
interface TraceDetailsHeaderProps {
filterMetadata: FilterMetadata;
onFilteredSpansChange: (spanIds: string[], isFilterActive: boolean) => void;
showTraceDetailsHeaderOptions?: boolean;
isDataLoaded?: boolean;
traceMetadata?: TraceMetadataForHeader;
}
const SKELETON_COUNT = 3;
@@ -56,15 +73,16 @@ function DetailsLoader(): JSX.Element {
}
function TraceDetailsHeader({
filterMetadata,
onFilteredSpansChange,
showTraceDetailsHeaderOptions,
isDataLoaded,
traceMetadata,
}: TraceDetailsHeaderProps): JSX.Element {
const { id: traceID } = useParams<TraceDetailV3URLProps>();
const [showTraceDetails, setShowTraceDetails] = useState(true);
const [isFilterExpanded, setIsFilterExpanded] = useState(false);
const [isPreviewFieldsOpen, setIsPreviewFieldsOpen] = useState(false);
const [isAnalyticsOpen, setIsAnalyticsOpen] = useState(false);
const { data: traceSummary } = useTraceSummary(traceID || '');
const previewFields = useTraceStore((s) => s.previewFields);
const setPreviewFields = useTraceStore((s) => s.setPreviewFields);
@@ -98,11 +116,8 @@ function TraceDetailsHeader({
setShowTraceDetails((prev) => !prev);
}, []);
const startTime = (traceSummary?.startTimestampMillis ?? 0) / 1e3;
const endTime = (traceSummary?.endTimestampMillis ?? 0) / 1e3;
const durationMs = traceSummary
? traceSummary.endTimestampMillis - traceSummary.startTimestampMillis
const durationMs = traceMetadata
? traceMetadata.endTimestampMillis - traceMetadata.startTimestampMillis
: 0;
return (
@@ -127,7 +142,7 @@ function TraceDetailsHeader({
/>
</div>
)}
{showTraceDetailsHeaderOptions && traceSummary && (
{isDataLoaded && (
<div
className={cx(
styles.filterSection,
@@ -156,9 +171,9 @@ function TraceDetailsHeader({
onToggleTraceDetails={handleToggleTraceDetails}
onOpenPreviewFields={(): void => setIsPreviewFieldsOpen(true)}
traceId={traceID || ''}
startTime={startTime}
endTime={endTime}
totalSpansCount={traceSummary.totalSpansCount}
startTime={filterMetadata.startTime}
endTime={filterMetadata.endTime}
totalSpansCount={traceMetadata?.totalSpansCount || 0}
/>
</div>
</TooltipProvider>
@@ -168,9 +183,9 @@ function TraceDetailsHeader({
className={cx(styles.filter, isFilterExpanded && styles.isExpanded)}
>
<Filters
startTime={startTime}
endTime={endTime}
traceID={traceID || ''}
startTime={filterMetadata.startTime}
endTime={filterMetadata.endTime}
traceID={filterMetadata.traceId}
onFilteredSpansChange={onFilteredSpansChange}
isExpanded={isFilterExpanded}
onExpand={(): void => setIsFilterExpanded(true)}
@@ -183,18 +198,18 @@ function TraceDetailsHeader({
{showTraceDetails && (
<div className={styles.subHeader}>
{traceSummary ? (
{traceMetadata ? (
<EntityMetadataRow
entity="trace"
service={{
name: traceSummary.rootServiceName,
entryPoint: traceSummary.rootServiceEntryPoint,
name: traceMetadata.rootServiceName,
entryPoint: traceMetadata.rootServiceEntryPoint,
}}
durationMs={durationMs}
timestamp={dayjs(traceSummary.startTimestampMillis).format(
timestamp={dayjs(traceMetadata.startTimestampMillis).format(
DATE_TIME_FORMATS.DD_MMM_YYYY_HH_MM_SS,
)}
statusCode={traceSummary.rootSpanStatusCode}
statusCode={traceMetadata.rootSpanStatusCode}
/>
) : (
<DetailsLoader />
@@ -202,7 +217,7 @@ function TraceDetailsHeader({
</div>
)}
{traceSummary?.hasMissingSpans && <MissingSpansBanner />}
{traceMetadata?.hasMissingSpans && <MissingSpansBanner />}
<FieldsSelector
isOpen={isPreviewFieldsOpen}

View File

@@ -5,11 +5,6 @@ import ROUTES from 'constants/routes';
import { render } from 'tests/test-utils';
import TraceDetailsHeader from '../TraceDetailsHeader';
import { useTraceSummary } from '../useTraceSummary';
jest.mock('../useTraceSummary', () => ({
useTraceSummary: jest.fn(() => ({ data: undefined, isLoading: false })),
}));
const mockGoBack = jest.fn();
const mockPush = jest.fn();
@@ -56,19 +51,13 @@ jest.mock('components/FieldsSelector', () => ({
}));
const baseProps = {
filterMetadata: {
startTime: 0,
endTime: 1,
traceId: 'trace-123',
},
onFilteredSpansChange: jest.fn(),
showTraceDetailsHeaderOptions: false,
};
const SUMMARY = {
startTimestampMillis: 1_700_000_000_000,
endTimestampMillis: 1_700_000_120_000,
rootServiceName: 'frontend',
rootServiceEntryPoint: 'GET /checkout',
rootSpanStatusCode: '200',
hasMissingSpans: false,
totalSpansCount: 3,
totalErrorSpansCount: 0,
isDataLoaded: false,
};
describe('TraceDetailsHeader – back button', () => {
@@ -103,32 +92,10 @@ describe('TraceDetailsHeader – back button', () => {
describe('TraceDetailsHeader – action cluster', () => {
beforeEach(() => {
mockReplace.mockClear();
jest
.mocked(useTraceSummary)
.mockReturnValue({ data: SUMMARY, isLoading: false });
});
afterEach(() => {
jest
.mocked(useTraceSummary)
.mockReturnValue({ data: undefined, isLoading: false });
});
it('does not render the action buttons until the summary loads', () => {
jest
.mocked(useTraceSummary)
.mockReturnValue({ data: undefined, isLoading: true });
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
expect(
screen.queryByRole('button', { name: /^analytics$/i }),
).not.toBeInTheDocument();
});
it('does not render the action buttons while data is still loading', () => {
render(
<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions={false} />,
);
render(<TraceDetailsHeader {...baseProps} isDataLoaded={false} />);
expect(
screen.queryByRole('button', { name: /^analytics$/i }),
@@ -139,7 +106,7 @@ describe('TraceDetailsHeader – action cluster', () => {
});
it('renders Analytics and Settings action buttons once data is loaded', () => {
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
render(<TraceDetailsHeader {...baseProps} isDataLoaded />);
expect(
screen.getByRole('button', { name: /^analytics$/i }),
@@ -150,7 +117,7 @@ describe('TraceDetailsHeader – action cluster', () => {
});
it('toggles the AnalyticsPanel open state when the Analytics button is clicked', () => {
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
render(<TraceDetailsHeader {...baseProps} isDataLoaded />);
const panel = screen.getByTestId('analytics-panel');
expect(panel).toHaveAttribute('data-open', 'false');
@@ -166,7 +133,7 @@ describe('TraceDetailsHeader – action cluster', () => {
});
describe('TraceDetailsHeader – trace metadata row', () => {
// useTraceSummary is mocked, so no API call is made.
// Plain prop, no API mock needed: traceMetadata is passed straight in.
const traceMetadata = {
startTimestampMillis: 1_700_000_000_000,
endTimestampMillis: 1_700_000_120_000, // +120000ms = 2 min
@@ -175,20 +142,16 @@ describe('TraceDetailsHeader – trace metadata row', () => {
rootSpanStatusCode: '404',
hasMissingSpans: false,
totalSpansCount: 42,
totalErrorSpansCount: 0,
};
const mockSummary = (data?: typeof traceMetadata): void => {
jest.mocked(useTraceSummary).mockReturnValue({ data, isLoading: false });
};
afterEach(() => {
mockSummary(undefined);
});
it('renders the metadata (service, entry point, duration, status) when provided', () => {
mockSummary(traceMetadata);
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
render(
<TraceDetailsHeader
{...baseProps}
isDataLoaded
traceMetadata={traceMetadata}
/>,
);
expect(screen.getByText(/inventory-frontend/)).toBeInTheDocument();
expect(screen.getByText('large-trace-root')).toBeInTheDocument();
@@ -203,8 +166,13 @@ describe('TraceDetailsHeader – trace metadata row', () => {
it('is shown by default and can be hidden / shown again via the Trace options menu', async () => {
const user = userEvent.setup({ delay: null });
mockSummary(traceMetadata);
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
render(
<TraceDetailsHeader
{...baseProps}
isDataLoaded
traceMetadata={traceMetadata}
/>,
);
// Visible by default (showTraceDetails defaults to true).
expect(screen.getByText(/inventory-frontend/)).toBeInTheDocument();
@@ -224,12 +192,9 @@ describe('TraceDetailsHeader – trace metadata row', () => {
expect(screen.getByText(/inventory-frontend/)).toBeInTheDocument();
});
it('shows skeletons instead of the metadata when the summary is absent', () => {
const { container } = render(
<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />,
);
it('does not render the metadata row when traceMetadata is absent', () => {
render(<TraceDetailsHeader {...baseProps} isDataLoaded />);
expect(screen.queryByText(/inventory-frontend/)).not.toBeInTheDocument();
expect(container.querySelectorAll('.ant-skeleton-input')).toHaveLength(3);
});
});

View File

@@ -1,16 +0,0 @@
import { useGetTraceSummary } from 'api/generated/services/tracedetail';
import type { SpantypesGettableTraceSummaryDTO } from 'api/generated/services/sigNoz.schemas';
interface UseTraceSummaryResult {
data: SpantypesGettableTraceSummaryDTO | undefined;
isLoading: boolean;
}
export function useTraceSummary(traceId: string): UseTraceSummaryResult {
const { data, isLoading } = useGetTraceSummary(
{ traceID: traceId },
{ query: { enabled: !!traceId } },
);
return { data: data?.data, isLoading };
}

View File

@@ -28,6 +28,7 @@ import TraceStoreSync from './stores/TraceStoreSync';
import { useTraceStore } from './stores/traceStore';
import { SpanDetailVariant } from './SpanDetailsPanel/constants';
import SpanDetailsPanel from './SpanDetailsPanel/SpanDetailsPanel';
import type { TraceMetadataForHeader } from './TraceDetailsHeader/TraceDetailsHeader';
import TraceDetailsHeader from './TraceDetailsHeader/TraceDetailsHeader';
import { FLAMEGRAPH_SPAN_LIMIT } from './TraceFlamegraph/constants';
import TraceFlamegraph from './TraceFlamegraph/TraceFlamegraph';
@@ -322,6 +323,38 @@ function TraceDetailsV3(): JSX.Element {
[],
);
const filterMetadata = useMemo(
() => ({
startTime: (traceData?.payload?.startTimestampMillis || 0) / 1e3,
endTime: (traceData?.payload?.endTimestampMillis || 0) / 1e3,
traceId: traceId || '',
}),
[
traceData?.payload?.startTimestampMillis,
traceData?.payload?.endTimestampMillis,
traceId,
],
);
const traceMetadataForHeader = useMemo(():
| TraceMetadataForHeader
| undefined => {
const payload = traceData?.payload;
if (!payload) {
return undefined;
}
const rootSpan = payload.spans?.find((s) => s.level === 0);
return {
startTimestampMillis: payload.startTimestampMillis,
endTimestampMillis: payload.endTimestampMillis,
rootServiceName: payload.rootServiceName,
rootServiceEntryPoint: payload.rootServiceEntryPoint,
rootSpanStatusCode: rootSpan?.response_status_code || '',
hasMissingSpans: payload.hasMissingSpans || false,
totalSpansCount: payload.totalSpansCount || 0,
};
}, [traceData?.payload]);
const showNoData =
!isFetchingTraceData &&
(!!errorFetchingTraceData || !traceData?.payload?.spans?.length);
@@ -360,10 +393,10 @@ function TraceDetailsV3(): JSX.Element {
<TraceStoreSync availableColorByFields={availableColorByFields}>
<div className={styles.root}>
<TraceDetailsHeader
filterMetadata={filterMetadata}
onFilteredSpansChange={handleFilteredSpansChange}
showTraceDetailsHeaderOptions={
!!traceData?.payload?.spans?.length && !showNoData
}
isDataLoaded={!!traceData?.payload?.spans?.length && !showNoData}
traceMetadata={traceMetadataForHeader}
/>
{showNoData ? (

View File

@@ -44,7 +44,6 @@ import {
traceDetailFieldKeys,
traceDetailFieldValues,
traceFlamegraphResponse,
traceSummaryResponse,
traceWaterfallResponse,
} from './__story_mockdata__/traceDetails';
@@ -150,13 +149,6 @@ export const traceDetailsMocks = defineStoryMocks({
),
),
rest.get(
'http://localhost/api/v1/traces/:traceId/summary',
response.json(() =>
traceSummaryResponse({ ...trace, missingSpans: values.missingSpans }),
),
),
rest.post(
'http://localhost/api/v3/traces/:traceId/flamegraph',
response.json(() => traceFlamegraphResponse(trace)),

View File

@@ -6,7 +6,6 @@
import type {
GetFlamegraph200,
GetTraceAggregations200,
GetTraceSummary200,
GetWaterfallV4200,
SpantypesFlamegraphSpanDTO,
SpantypesSpanAggregationDTO,
@@ -304,27 +303,6 @@ export const traceWaterfallResponse = (
};
};
export const traceSummaryResponse = (
options: TraceOptions & { missingSpans: boolean },
): GetTraceSummary200 => {
const spans = buildSpans(options);
const root = spans[0];
return {
status: 'success',
data: {
startTimestampMillis: Math.round(options.traceStart),
endTimestampMillis: Math.round(options.traceStart + ROOT_DURATION_MS),
rootServiceName: root?.template.service ?? '',
rootServiceEntryPoint: root?.template.name ?? '',
rootSpanStatusCode: root?.hasError ? '503' : '200',
totalSpansCount: spans.length,
totalErrorSpansCount: spans.filter(({ hasError }) => hasError).length,
hasMissingSpans: options.missingSpans,
},
};
};
const flamegraphSpan = (span: BuiltSpan): SpantypesFlamegraphSpanDTO => ({
spanId: span.spanId,
parentSpanId: span.parentSpanId,

View File

@@ -84,5 +84,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. 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",
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

@@ -93,3 +93,31 @@ 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) {
claims, err := authtypes.ClaimsFromContext(r.Context())
if err != nil {
render.Error(rw, err)
return
}
req := new(spantypes.GetTraceThreadParams)
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(), valuer.MustNewUUID(claims.OrgID), mux.Vars(r)["traceID"], query)
if err != nil {
render.Error(rw, err)
return
}
render.Success(rw, http.StatusOK, result)
}

View File

@@ -189,6 +189,49 @@ func (m *module) getWindowedWaterfall(ctx context.Context, traceID, selectedSpan
), nil
}
func (m *module) GetThread(ctx context.Context, orgID valuer.UUID, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error) {
bounds, err := m.store.GetTraceBounds(ctx, traceID)
if err != nil {
return nil, err
}
// One extra row per side signals a next/prev page; NewGettableTraceThread trims the response back to query.Limit.
page := spantypes.ThreadPage{Limit: query.Limit + 1}
switch {
case query.SpanID != "":
// Window centred on the span: fetch both directions, NewGettableTraceThread splits the limit.
anchor, err := m.store.GetThreadCursor(ctx, traceID, bounds, query.SpanID)
if err != nil {
return nil, err
}
before, err := m.store.GetThreadSpans(ctx, orgID, traceID, bounds, spantypes.ThreadPage{Cursor: anchor, From: spantypes.ThreadBefore, Limit: page.Limit})
if err != nil {
return nil, err
}
// ThreadAt keeps the anchor when it carries messages; otherwise the filter drops it.
after, err := m.store.GetThreadSpans(ctx, orgID, traceID, bounds, spantypes.ThreadPage{Cursor: anchor, From: spantypes.ThreadAt, Limit: page.Limit})
if err != nil {
return nil, err
}
return spantypes.NewGettableTraceThread(traceID, query, before, after), nil
case query.Before != nil:
page.Cursor, page.From = query.Before, spantypes.ThreadBefore
before, err := m.store.GetThreadSpans(ctx, orgID, traceID, bounds, page)
if err != nil {
return nil, err
}
return spantypes.NewGettableTraceThread(traceID, query, before, nil), nil
default:
// nil for the first page.
page.Cursor = query.After
after, err := m.store.GetThreadSpans(ctx, orgID, traceID, bounds, page)
if err != nil {
return nil, err
}
return spantypes.NewGettableTraceThread(traceID, query, nil, after), nil
}
}
func (m *module) getFullFlamegraph(ctx context.Context, traceID string, bounds *spantypes.TraceBounds, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error) {
fullSpans, err := m.store.GetFlamegraphSpans(ctx, traceID, bounds.Start, bounds.End, nil)
if err != nil {

View File

@@ -267,8 +267,7 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
}
sb := sqlbuilder.NewSelectBuilder()
sb.Select(
"DISTINCT ON (span_id) timestamp",
"duration_nano", "span_id", "has_error", "kind",
"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",
@@ -298,6 +297,119 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return spans, nil
}
func (s *traceStore) GetThreadSpans(ctx context.Context, orgID valuer.UUID, traceID string, bounds *spantypes.TraceBounds, page spantypes.ThreadPage) ([]spantypes.StorableSpan, error) {
q := querybuilder.NewQueryInfo(ctx, orgID, s.flagger, telemetrytypes.SignalTraces, nil, uint64(bounds.Start.UnixNano()), uint64(bounds.End.UnixNano()))
sb := sqlbuilder.NewSelectBuilder()
sb.Select(
"DISTINCT ON (span_id) timestamp", "duration_nano", "span_id", "parent_span_id", "has_error", "name", "kind_string",
"status_code_string", "status_message", "resources_string",
"attributes_string", "attributes_number", "attributes_bool",
"events", "links as references",
)
if q.TraceAttrsJSONOn {
sb.SelectMore("attributes")
}
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
hasMessages, err := s.messagesExistCondition(ctx, q, orgID, bounds, sb)
if err != nil {
return nil, err
}
sb.Where(
sb.E("trace_id", traceID),
sb.GE("ts_bucket_start", bounds.Start.Unix()-1800),
sb.LE("ts_bucket_start", bounds.End.Unix()),
hasMessages,
)
if cursor := page.Cursor; 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 on the far side of the cursor.
key := "(toUnixTimestamp64Nano(timestamp), span_id)"
bucket := int64(cursor.TimeUnixNano / uint64(time.Second))
timestamp := fmt.Sprintf("fromUnixTimestamp64Nano(toInt64(%s))", sb.Var(cursor.TimeUnixNano))
tuple := sqlbuilder.Tuple(cursor.TimeUnixNano, cursor.SpanID)
switch page.From {
case spantypes.ThreadBefore:
sb.Where(sb.LE("ts_bucket_start", bucket), "timestamp <= "+timestamp, sb.LT(key, tuple))
case spantypes.ThreadAt:
sb.Where(sb.GE("ts_bucket_start", bucket-1800), "timestamp >= "+timestamp, sb.GE(key, tuple))
default:
sb.Where(sb.GE("ts_bucket_start", bucket-1800), "timestamp >= "+timestamp, sb.GT(key, tuple))
}
}
// span_id breaks timestamp ties so the order matches the cursor key; otherwise tied spans
// can be skipped or repeated across pages.
if page.From == spantypes.ThreadBefore {
sb.OrderByDesc("timestamp")
sb.OrderByDesc("span_id")
} else {
sb.OrderByAsc("timestamp")
sb.OrderByAsc("span_id")
}
sb.Limit(page.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
}
// 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)")
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
sb.Where(
sb.E("trace_id", traceID),
sb.GE("ts_bucket_start", bounds.Start.Unix()-1800),
sb.LE("ts_bucket_start", bounds.End.Unix()),
sb.E("span_id", spanID),
)
sb.Limit(1)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
var timeUnixNano int64
if err := s.telemetryStore.ClickhouseDB().QueryRow(ctx, query, args...).Scan(&timeUnixNano); err != nil {
if errors.Is(err, sql.ErrNoRows) {
return nil, errors.NewNotFoundf(spantypes.ErrCodeThreadSpanNotFound, "span %s not found in trace %s", spanID, traceID)
}
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error querying thread span")
}
return &spantypes.ThreadCursor{TimeUnixNano: uint64(timeUnixNano), SpanID: spanID}, 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

@@ -202,3 +202,43 @@ func TestGetSpanDurationByField(t *testing.T) {
})
}
}
func TestGetThreadSpans(t *testing.T) {
selectSQL := "SELECT DISTINCT ON (span_id) timestamp, duration_nano, span_id, parent_span_id, has_error, name, kind_string, status_code_string, status_message, resources_string, attributes_string, attributes_number, attributes_bool, events, links as references"
fromSQL := " FROM signoz_traces.distributed_signoz_index_v3 WHERE trace_id = ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND "
orderSQL := " ORDER BY timestamp ASC, span_id ASC LIMIT ?"
jsonInsideTrace := testStart.Add(500 * time.Second)
testCases := []struct {
name string
jsonOn bool
jsonRelease *time.Time
selectSQL string
whereSQL string
}{
{
name: "FlagOff_ReadsAndFiltersLegacyMaps",
jsonRelease: &jsonInsideTrace,
selectSQL: selectSQL,
whereSQL: "(mapContains(attributes_string, 'gen_ai.input.messages') OR mapContains(attributes_string, 'gen_ai.output.messages'))",
},
{
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)",
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{flagger.FeatureUseTraceAttributesJSON.String(): testCase.jsonOn})
s := newTestStoreWithMetadata(sqlmock.QueryMatcherRegexp, genAIMetadataStore(testCase.jsonRelease), fl)
s.Mock().ExpectSelect(regexp.QuoteMeta(testCase.selectSQL + fromSQL + testCase.whereSQL + orderSQL)).
WillReturnRows(cmock.NewRows(nil, nil))
_, _ = s.Store().GetThreadSpans(context.Background(), valuer.GenerateUUID(), testTraceID, testBounds, spantypes.ThreadPage{Limit: 3})
assert.NoError(t, s.Mock().ExpectationsWereMet())
})
}
}

View File

@@ -15,6 +15,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.
@@ -23,4 +24,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, orgID valuer.UUID, 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

@@ -260,6 +260,7 @@ func NewSQLMigrationProviderFactories(
sqlmigration.NewAddChannelSpecFactory(sqlschema),
sqlmigration.NewAddUserTuplesFactory(sqlstore),
sqlmigration.NewAddRuleViewFactory(sqlstore, sqlschema),
sqlmigration.NewKeepRelevantLLMPricingRulesFactory(),
)
}

View File

@@ -0,0 +1,116 @@
package sqlmigration
import (
"context"
"encoding/json"
"log/slog"
"regexp"
"strings"
"github.com/uptrace/bun"
"github.com/uptrace/bun/migrate"
"github.com/SigNoz/signoz/pkg/factory"
)
// Mirrors the zeus LLMPriceFilter and its default model authors.
// https://github.com/SigNoz/zeus/pull/591/changes
var (
llmPricingRuleProviders = map[string]struct{}{
"openai": {}, "anthropic": {}, "google": {}, "mistralai": {}, "deepseek": {},
"qwen": {}, "x-ai": {}, "meta-llama": {}, "cohere": {}, "amazon": {},
}
qwenHostedModelPattern = regexp.MustCompile(`max|plus|flash|turbo`)
)
type llmPricingRuleSyncedRow struct {
bun.BaseModel `bun:"table:llm_pricing_rule"`
ID string `bun:"id,pk"`
Provider string `bun:"provider"`
Model string `bun:"model"`
Pricing string `bun:"pricing"`
}
type llmPricingRulePrices struct {
Input float64 `json:"input"`
Output float64 `json:"output"`
}
type keepRelevantLLMPricingRules struct {
settings factory.ProviderSettings
}
func NewKeepRelevantLLMPricingRulesFactory() factory.ProviderFactory[SQLMigration, Config] {
return factory.NewProviderFactory(factory.MustNewName("keep_relevant_llm_pricing_rules"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
return &keepRelevantLLMPricingRules{settings: ps}, nil
})
}
func (migration *keepRelevantLLMPricingRules) Register(migrations *migrate.Migrations) error {
return migrations.Register(migration.Up, migration.Down)
}
func (migration *keepRelevantLLMPricingRules) Up(ctx context.Context, db *bun.DB) error {
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
var rows []*llmPricingRuleSyncedRow
if err := tx.NewSelect().
Model(&rows).
Where("source_id IS NOT NULL").
Where("NOT is_override").
Scan(ctx); err != nil {
return err
}
ids := make([]string, 0)
for _, row := range rows {
var prices llmPricingRulePrices
if err := json.Unmarshal([]byte(row.Pricing), &prices); err != nil {
migration.settings.Logger.WarnContext(ctx, "llm pricing rule has unparseable pricing, leaving it untouched", slog.String("rule_id", row.ID), slog.String("raw_pricing", row.Pricing))
continue
}
if !llmPricingRuleRelevant(row.Provider, row.Model, prices) {
ids = append(ids, row.ID)
}
}
if len(ids) > 0 {
if _, err := tx.NewDelete().
Model((*llmPricingRuleSyncedRow)(nil)).
Where("id IN (?)", bun.In(ids)).
Exec(ctx); err != nil {
return err
}
}
migration.settings.Logger.InfoContext(ctx, "deleted irrelevant llm pricing rules", slog.Int("total", len(rows)), slog.Int("deleted", len(ids)))
return tx.Commit()
}
func (migration *keepRelevantLLMPricingRules) Down(context.Context, *bun.DB) error {
return nil
}
// The provider allowlist also rejects "~" alias ids, whose provider segment
// starts with "~", and ids without a "/" that were stored as "unknown".
func llmPricingRuleRelevant(provider, model string, prices llmPricingRulePrices) bool {
if _, ok := llmPricingRuleProviders[provider]; !ok {
return false
}
if strings.Contains(model, ":") {
return false
}
if prices.Input <= 0 || prices.Output <= 0 {
return false
}
if provider == "qwen" && !qwenHostedModelPattern.MatchString(model) {
return false
}
return true
}

View File

@@ -37,6 +37,8 @@ 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, orgID valuer.UUID, traceID string, bounds *TraceBounds, page ThreadPage) ([]StorableSpan, error)
GetThreadCursor(ctx context.Context, traceID string, bounds *TraceBounds, spanID string) (*ThreadCursor, error)
GetSpanCountByField(ctx context.Context, traceID string, bounds *TraceBounds, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)
GetSpanDurationByField(ctx context.Context, traceID string, bounds *TraceBounds, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)

View File

@@ -0,0 +1,202 @@
package spantypes
import (
"encoding/base64"
"encoding/json"
"maps"
"github.com/SigNoz/signoz/pkg/errors"
)
const (
threadDefaultLimit = 20
threadMaxLimit = 100
)
const (
ThreadAfter ThreadFrom = iota // > cursor, ascending
ThreadAt // >= cursor, ascending
ThreadBefore // < cursor, descending
)
var (
ErrCodeThreadInvalidLimit = errors.MustNewCode("trace_thread_invalid_limit")
ErrCodeThreadInvalidCursor = errors.MustNewCode("trace_thread_invalid_cursor")
ErrCodeThreadInvalidAnchor = errors.MustNewCode("trace_thread_invalid_anchor")
ErrCodeThreadSpanNotFound = errors.MustNewCode("trace_thread_span_not_found")
)
type GetTraceThreadParams struct {
Limit int `query:"limit" description:"Page size, at most 100. 0 means 20."`
After string `query:"after" description:"The nextCursor of a page; returns the spans after it. Set only one of after, before and spanId."`
Before string `query:"before" description:"The prevCursor of a page; returns the spans before it. Set only one of after, before and spanId."`
SpanID string `query:"spanId" description:"Returns the page around this span. Set only one of after, before and spanId."`
}
type ThreadQuery struct {
Limit int
After *ThreadCursor
Before *ThreadCursor
SpanID string
}
// ThreadCursor is the (TimeUnixNano, SpanID) key of a span.
type ThreadCursor struct {
TimeUnixNano uint64 `json:"timeUnixNano"`
SpanID string `json:"spanId"`
}
type ThreadFrom int
type ThreadPage struct {
Cursor *ThreadCursor
From ThreadFrom
Limit int
}
type GettableTraceThread struct {
Spans []*ThreadSpan `json:"spans" required:"true" nullable:"false"`
PrevCursor string `json:"prevCursor,omitempty"`
NextCursor string `json:"nextCursor,omitempty"`
}
// ThreadSpan carries the fields the span details pane reads; snake_case keys match WaterfallSpan.
type ThreadSpan struct {
SpanID string `json:"span_id" required:"true"`
TraceID string `json:"trace_id" required:"true"`
ParentSpanID string `json:"parent_span_id" required:"true"`
Name string `json:"name" required:"true"`
KindString string `json:"kind_string" required:"true"`
TimeUnix uint64 `json:"time_unix" required:"true"`
DurationNano uint64 `json:"duration_nano" required:"true"`
HasError bool `json:"has_error" required:"true"`
StatusCodeString string `json:"status_code_string" required:"true"`
StatusMessage string `json:"status_message" required:"true"`
Resource map[string]string `json:"resource" required:"true" nullable:"false"`
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"`
timeUnixNano uint64
}
func NewThreadQuery(params *GetTraceThreadParams) (*ThreadQuery, error) {
query := &ThreadQuery{Limit: params.Limit, SpanID: params.SpanID}
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)
}
anchors := 0
for _, value := range []string{params.After, params.Before, params.SpanID} {
if value != "" {
anchors++
}
}
if anchors > 1 {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidAnchor, "only one of after, before and spanId can be set")
}
encoded := params.After
if encoded == "" {
encoded = params.Before
}
if encoded == "" {
return query, nil
}
data, err := base64.RawURLEncoding.DecodeString(encoded)
if err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
}
cursor := new(ThreadCursor)
if err := json.Unmarshal(data, cursor); err != nil {
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
}
if cursor.SpanID == "" {
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidCursor, "invalid cursor: missing span id")
}
if params.After != "" {
query.After = cursor
} else {
query.Before = cursor
}
return query, nil
}
func (c ThreadCursor) Encode() string {
data, _ := json.Marshal(c)
return base64.RawURLEncoding.EncodeToString(data)
}
// NewGettableTraceThread takes up to limit+1 spans on each side of the anchor: before in
// descending order, after in ascending order. Half the page goes to before, the rest to after,
// and a short side gives its room to the other. An extra span on a side sets that side's cursor.
func NewGettableTraceThread(traceID string, query *ThreadQuery, before, after []StorableSpan) *GettableTraceThread {
nAfter := min(len(after), query.Limit-min(len(before), query.Limit/2))
nBefore := min(len(before), query.Limit-nAfter)
hasPrev := len(before) > nBefore || query.After != nil
hasNext := len(after) > nAfter || query.Before != nil
spans := make([]*ThreadSpan, 0, nBefore+nAfter)
for i := nBefore - 1; i >= 0; i-- {
spans = append(spans, newThreadSpan(traceID, &before[i]))
}
for i := range nAfter {
spans = append(spans, newThreadSpan(traceID, &after[i]))
}
thread := &GettableTraceThread{Spans: spans}
if len(spans) == 0 {
return thread
}
if hasPrev {
thread.PrevCursor = spans[0].cursor().Encode()
}
if hasNext {
thread.NextCursor = spans[len(spans)-1].cursor().Encode()
}
return thread
}
func (s *ThreadSpan) cursor() ThreadCursor {
return ThreadCursor{TimeUnixNano: s.timeUnixNano, SpanID: s.SpanID}
}
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())
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,
}
}
// threadAttributes reads the JSON column 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
}
return storable.Attributes()
}

View File

@@ -0,0 +1,75 @@
package spantypes
import (
"testing"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestNewThreadQuery_Cursor(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: "eyJ0aW1lVW5peE5hbm8iOiAxfQ", wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
after, err := NewThreadQuery(&GetTraceThreadParams{After: testCase.cursor})
before, errBefore := NewThreadQuery(&GetTraceThreadParams{Before: testCase.cursor})
if testCase.wantErr {
assert.Error(t, err)
assert.Error(t, errBefore)
return
}
require.NoError(t, err)
require.NoError(t, errBefore)
assert.Equal(t, testCase.want, after.After)
assert.Nil(t, after.Before)
assert.Equal(t, testCase.want, before.Before)
assert.Nil(t, before.After)
})
}
}
func TestThreadAttributes(t *testing.T) {
testCases := []struct {
name string
span StorableSpan
wantAttrs map[string]any
}{
{
name: "JSONColumnPresent_LegacyMapsIgnored",
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", "http.method": "GET"},
},
wantAttrs: map[string]any{"gen_ai.request.model": "json"},
},
{
name: "NoJSONColumn_FallsBackToLegacyMaps",
span: StorableSpan{
AttributesString: map[string]string{"gen_ai.output.messages": `[{"role":"assistant","content":"hello"}]`},
AttributesNumber: map[string]float64{"gen_ai.usage.input_tokens": 12},
AttributesBool: map[string]bool{"gen_ai.stream": true},
},
wantAttrs: map[string]any{
"gen_ai.output.messages": `[{"role":"assistant","content":"hello"}]`,
"gen_ai.usage.input_tokens": float64(12),
"gen_ai.stream": true,
},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.wantAttrs, threadAttributes(&testCase.span))
})
}
}

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.

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,327 @@
import base64
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 ATTRIBUTE_JSON_ROLLOUT_TIME, 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],
seed_attribute_evolution: Callable[[str, datetime], None],
) -> None:
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))
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_cursors(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
seed_attribute_evolution: Callable[[str, datetime], None],
) -> None:
seed_attribute_evolution("traces", ATTRIBUTE_JSON_ROLLOUT_TIME)
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}"}
def get_page(params: dict) -> dict:
response = requests.get(url, params=params, headers=headers, timeout=10)
assert response.status_code == HTTPStatus.OK, f"{params}: {response.text}"
return response.json()["data"]
first = get_page({"limit": 2})
assert [span["span_id"] for span in first["spans"]] == expected_order[:2]
assert "prevCursor" not in first
assert first["nextCursor"]
last = get_page({"limit": 2, "after": first["nextCursor"]})
assert [span["span_id"] for span in last["spans"]] == expected_order[2:]
assert last["prevCursor"]
assert "nextCursor" not in last
previous = get_page({"limit": 2, "before": last["prevCursor"]})
assert [span["span_id"] for span in previous["spans"]] == expected_order[:2]
assert "prevCursor" not in previous
assert previous["nextCursor"] == first["nextCursor"]
middle = get_page({"limit": 1, "before": last["prevCursor"]})
assert [span["span_id"] for span in middle["spans"]] == expected_order[1:2]
assert middle["prevCursor"]
assert middle["nextCursor"]
start = get_page({"limit": 2, "before": middle["prevCursor"]})
assert [span["span_id"] for span in start["spans"]] == expected_order[:1]
assert "prevCursor" not in start
def test_thread_paginates_across_buckets(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
seed_attribute_evolution: Callable[[str, datetime], None],
) -> None:
seed_attribute_evolution("traces", ATTRIBUTE_JSON_ROLLOUT_TIME)
now = datetime.now(tz=UTC).replace(microsecond=0)
bucket = now.replace(minute=0 if now.minute < 30 else 30, second=0)
trace_id = TraceIdGenerator.trace_id()
# neighbours within one 30-minute ts_bucket_start and across bucket boundaries
timestamps = [
bucket - timedelta(minutes=59, seconds=59),
bucket - timedelta(minutes=30, seconds=1),
bucket - timedelta(minutes=30),
bucket - timedelta(seconds=1),
bucket,
]
span_ids = [TraceIdGenerator.span_id() for _ in timestamps]
insert_traces(
[
Traces(timestamp=timestamp, trace_id=trace_id, span_id=span_id, name="chat gpt-4o", resources={"service.name": "tracedetail-thread-buckets"}, 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)
]
)
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}"}
def get_page(params: dict) -> dict:
response = requests.get(url, params=params, headers=headers, timeout=10)
assert response.status_code == HTTPStatus.OK, f"{params}: {response.text}"
return response.json()["data"]
page = get_page({"limit": 1})
forward = [span["span_id"] for span in page["spans"]]
while "nextCursor" in page:
page = get_page({"limit": 1, "after": page["nextCursor"]})
forward += [span["span_id"] for span in page["spans"]]
assert forward == span_ids
backward = [span["span_id"] for span in page["spans"]]
while "prevCursor" in page:
page = get_page({"limit": 1, "before": page["prevCursor"]})
backward = [span["span_id"] for span in page["spans"]] + backward
assert backward == span_ids
for index in range(1, len(span_ids)):
around = get_page({"limit": 2, "spanId": span_ids[index]})
assert [span["span_id"] for span in around["spans"]] == span_ids[index - 1 : index + 1], index
def test_thread_opens_around_span(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
seed_attribute_evolution: Callable[[str, datetime], None],
) -> None:
seed_attribute_evolution("traces", ATTRIBUTE_JSON_ROLLOUT_TIME)
now = datetime.now(tz=UTC).replace(microsecond=0)
trace_id = TraceIdGenerator.trace_id()
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
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"),
*(
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"),
]
)
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}"}
def get_page(params: dict) -> dict:
response = requests.get(url, params=params, headers=headers, timeout=10)
assert response.status_code == HTTPStatus.OK, f"{params}: {response.text}"
return response.json()["data"]
# span with messages: included with its neighbours
around = get_page({"limit": 3, "spanId": llm_ids[2]})
assert [span["span_id"] for span in around["spans"]] == llm_ids[1:4]
assert around["prevCursor"]
assert around["nextCursor"]
assert [span["span_id"] for span in get_page({"limit": 3, "before": around["prevCursor"]})["spans"]] == llm_ids[:1]
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"]
# near the start, the short side gives its room to the other
at_start = get_page({"limit": 3, "spanId": llm_ids[0]})
assert [span["span_id"] for span in at_start["spans"]] == llm_ids[:3]
assert "prevCursor" not in at_start
assert at_start["nextCursor"]
# near the end
at_end = get_page({"limit": 3, "spanId": llm_ids[4]})
assert [span["span_id"] for span in at_end["spans"]] == llm_ids[2:]
assert at_end["prevCursor"]
assert "nextCursor" not in at_end
# page covers the whole thread
whole = get_page({"limit": 10, "spanId": root_id})
assert [span["span_id"] for span in whole["spans"]] == llm_ids
assert "prevCursor" not in whole
assert "nextCursor" not in whole
missing = requests.get(url, params={"spanId": TraceIdGenerator.span_id()}, headers=headers, timeout=10)
assert missing.status_code == HTTPStatus.NOT_FOUND, missing.text
def test_thread_reads_spans_across_json_rollout(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
seed_attribute_evolution: Callable[[str, datetime], None],
) -> None:
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
rollout = now - timedelta(minutes=30)
seed_attribute_evolution("traces", rollout)
resources = {"service.name": "tracedetail-thread-rollout"}
# trace entirely before the rollout: messages live only in the legacy maps
before_trace_id = TraceIdGenerator.trace_id()
before_ids = [TraceIdGenerator.span_id() for _ in range(2)]
# trace straddling the rollout: one span in the maps, one in the JSON column
straddle_trace_id = TraceIdGenerator.trace_id()
legacy_id, json_id = TraceIdGenerator.span_id(), TraceIdGenerator.span_id()
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=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"),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
headers = {"Authorization": f"Bearer {token}"}
before = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{before_trace_id}/thread"), headers=headers, timeout=10)
assert before.status_code == HTTPStatus.OK, before.text
before_spans = before.json()["data"]["spans"]
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"}])
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
straddle_spans = straddle.json()["data"]["spans"]
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"}])
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],
seed_attribute_evolution: Callable[[str, datetime], None],
) -> None:
seed_attribute_evolution("traces", ATTRIBUTE_JSON_ROLLOUT_TIME)
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],
) -> None:
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/{TraceIdGenerator.trace_id()}/thread")
cursor = base64.urlsafe_b64encode(json.dumps({"timeUnixNano": 1, "spanId": "f1fa1bc863e94dd0"}).encode()).decode().rstrip("=")
for params in (
{"limit": -1},
{"limit": 101},
{"after": "not-a-cursor"},
{"before": "not-a-cursor"},
{"after": cursor, "before": cursor},
{"after": cursor, "spanId": "f1fa1bc863e94dd0"},
{"before": cursor, "spanId": "f1fa1bc863e94dd0"},
):
response = requests.get(url, params=params, headers=headers, timeout=10)
assert response.status_code == HTTPStatus.BAD_REQUEST, f"{params}: {response.text}"
missing = requests.get(url, headers=headers, timeout=10)
assert missing.status_code == HTTPStatus.NOT_FOUND, missing.text