mirror of
https://github.com/SigNoz/signoz.git
synced 2026-10-10 04:01:04 +01:00
Compare commits
4 Commits
issue_nerv
...
13065
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0ae98b05c5 | ||
|
|
7c76ba8d9e | ||
|
|
104758f3cf | ||
|
|
504be0dfd7 |
@@ -9878,19 +9878,6 @@ 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:
|
||||
@@ -10205,61 +10192,6 @@ 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:
|
||||
@@ -16077,103 +16009,6 @@ 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
|
||||
|
||||
@@ -11426,86 +11426,6 @@ 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;
|
||||
};
|
||||
@@ -13190,40 +13110,6 @@ 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
|
||||
|
||||
@@ -24,9 +24,6 @@ import type {
|
||||
GetTraceAggregationsPathParameters,
|
||||
GetTraceSummary200,
|
||||
GetTraceSummaryPathParameters,
|
||||
GetTraceThread200,
|
||||
GetTraceThreadParams,
|
||||
GetTraceThreadPathParameters,
|
||||
GetWaterfallV4200,
|
||||
GetWaterfallV4PathParameters,
|
||||
RenderErrorResponseDTO,
|
||||
@@ -260,121 +257,6 @@ 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
|
||||
|
||||
@@ -27,31 +27,14 @@ 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;
|
||||
isDataLoaded?: boolean;
|
||||
traceMetadata?: TraceMetadataForHeader;
|
||||
showTraceDetailsHeaderOptions?: boolean;
|
||||
}
|
||||
|
||||
const SKELETON_COUNT = 3;
|
||||
@@ -73,16 +56,15 @@ function DetailsLoader(): JSX.Element {
|
||||
}
|
||||
|
||||
function TraceDetailsHeader({
|
||||
filterMetadata,
|
||||
onFilteredSpansChange,
|
||||
isDataLoaded,
|
||||
traceMetadata,
|
||||
showTraceDetailsHeaderOptions,
|
||||
}: 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);
|
||||
|
||||
@@ -116,8 +98,11 @@ function TraceDetailsHeader({
|
||||
setShowTraceDetails((prev) => !prev);
|
||||
}, []);
|
||||
|
||||
const durationMs = traceMetadata
|
||||
? traceMetadata.endTimestampMillis - traceMetadata.startTimestampMillis
|
||||
const startTime = (traceSummary?.startTimestampMillis ?? 0) / 1e3;
|
||||
const endTime = (traceSummary?.endTimestampMillis ?? 0) / 1e3;
|
||||
|
||||
const durationMs = traceSummary
|
||||
? traceSummary.endTimestampMillis - traceSummary.startTimestampMillis
|
||||
: 0;
|
||||
|
||||
return (
|
||||
@@ -142,7 +127,7 @@ function TraceDetailsHeader({
|
||||
/>
|
||||
</div>
|
||||
)}
|
||||
{isDataLoaded && (
|
||||
{showTraceDetailsHeaderOptions && traceSummary && (
|
||||
<div
|
||||
className={cx(
|
||||
styles.filterSection,
|
||||
@@ -171,9 +156,9 @@ function TraceDetailsHeader({
|
||||
onToggleTraceDetails={handleToggleTraceDetails}
|
||||
onOpenPreviewFields={(): void => setIsPreviewFieldsOpen(true)}
|
||||
traceId={traceID || ''}
|
||||
startTime={filterMetadata.startTime}
|
||||
endTime={filterMetadata.endTime}
|
||||
totalSpansCount={traceMetadata?.totalSpansCount || 0}
|
||||
startTime={startTime}
|
||||
endTime={endTime}
|
||||
totalSpansCount={traceSummary.totalSpansCount}
|
||||
/>
|
||||
</div>
|
||||
</TooltipProvider>
|
||||
@@ -183,9 +168,9 @@ function TraceDetailsHeader({
|
||||
className={cx(styles.filter, isFilterExpanded && styles.isExpanded)}
|
||||
>
|
||||
<Filters
|
||||
startTime={filterMetadata.startTime}
|
||||
endTime={filterMetadata.endTime}
|
||||
traceID={filterMetadata.traceId}
|
||||
startTime={startTime}
|
||||
endTime={endTime}
|
||||
traceID={traceID || ''}
|
||||
onFilteredSpansChange={onFilteredSpansChange}
|
||||
isExpanded={isFilterExpanded}
|
||||
onExpand={(): void => setIsFilterExpanded(true)}
|
||||
@@ -198,18 +183,18 @@ function TraceDetailsHeader({
|
||||
|
||||
{showTraceDetails && (
|
||||
<div className={styles.subHeader}>
|
||||
{traceMetadata ? (
|
||||
{traceSummary ? (
|
||||
<EntityMetadataRow
|
||||
entity="trace"
|
||||
service={{
|
||||
name: traceMetadata.rootServiceName,
|
||||
entryPoint: traceMetadata.rootServiceEntryPoint,
|
||||
name: traceSummary.rootServiceName,
|
||||
entryPoint: traceSummary.rootServiceEntryPoint,
|
||||
}}
|
||||
durationMs={durationMs}
|
||||
timestamp={dayjs(traceMetadata.startTimestampMillis).format(
|
||||
timestamp={dayjs(traceSummary.startTimestampMillis).format(
|
||||
DATE_TIME_FORMATS.DD_MMM_YYYY_HH_MM_SS,
|
||||
)}
|
||||
statusCode={traceMetadata.rootSpanStatusCode}
|
||||
statusCode={traceSummary.rootSpanStatusCode}
|
||||
/>
|
||||
) : (
|
||||
<DetailsLoader />
|
||||
@@ -217,7 +202,7 @@ function TraceDetailsHeader({
|
||||
</div>
|
||||
)}
|
||||
|
||||
{traceMetadata?.hasMissingSpans && <MissingSpansBanner />}
|
||||
{traceSummary?.hasMissingSpans && <MissingSpansBanner />}
|
||||
|
||||
<FieldsSelector
|
||||
isOpen={isPreviewFieldsOpen}
|
||||
|
||||
@@ -5,6 +5,11 @@ 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();
|
||||
@@ -51,13 +56,19 @@ jest.mock('components/FieldsSelector', () => ({
|
||||
}));
|
||||
|
||||
const baseProps = {
|
||||
filterMetadata: {
|
||||
startTime: 0,
|
||||
endTime: 1,
|
||||
traceId: 'trace-123',
|
||||
},
|
||||
onFilteredSpansChange: jest.fn(),
|
||||
isDataLoaded: false,
|
||||
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,
|
||||
};
|
||||
|
||||
describe('TraceDetailsHeader – back button', () => {
|
||||
@@ -92,10 +103,32 @@ 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} isDataLoaded={false} />);
|
||||
render(
|
||||
<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions={false} />,
|
||||
);
|
||||
|
||||
expect(
|
||||
screen.queryByRole('button', { name: /^analytics$/i }),
|
||||
@@ -106,7 +139,7 @@ describe('TraceDetailsHeader – action cluster', () => {
|
||||
});
|
||||
|
||||
it('renders Analytics and Settings action buttons once data is loaded', () => {
|
||||
render(<TraceDetailsHeader {...baseProps} isDataLoaded />);
|
||||
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
|
||||
|
||||
expect(
|
||||
screen.getByRole('button', { name: /^analytics$/i }),
|
||||
@@ -117,7 +150,7 @@ describe('TraceDetailsHeader – action cluster', () => {
|
||||
});
|
||||
|
||||
it('toggles the AnalyticsPanel open state when the Analytics button is clicked', () => {
|
||||
render(<TraceDetailsHeader {...baseProps} isDataLoaded />);
|
||||
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
|
||||
|
||||
const panel = screen.getByTestId('analytics-panel');
|
||||
expect(panel).toHaveAttribute('data-open', 'false');
|
||||
@@ -133,7 +166,7 @@ describe('TraceDetailsHeader – action cluster', () => {
|
||||
});
|
||||
|
||||
describe('TraceDetailsHeader – trace metadata row', () => {
|
||||
// Plain prop, no API mock needed: traceMetadata is passed straight in.
|
||||
// useTraceSummary is mocked, so no API call is made.
|
||||
const traceMetadata = {
|
||||
startTimestampMillis: 1_700_000_000_000,
|
||||
endTimestampMillis: 1_700_000_120_000, // +120000ms = 2 min
|
||||
@@ -142,16 +175,20 @@ 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', () => {
|
||||
render(
|
||||
<TraceDetailsHeader
|
||||
{...baseProps}
|
||||
isDataLoaded
|
||||
traceMetadata={traceMetadata}
|
||||
/>,
|
||||
);
|
||||
mockSummary(traceMetadata);
|
||||
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
|
||||
|
||||
expect(screen.getByText(/inventory-frontend/)).toBeInTheDocument();
|
||||
expect(screen.getByText('large-trace-root')).toBeInTheDocument();
|
||||
@@ -166,13 +203,8 @@ 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 });
|
||||
render(
|
||||
<TraceDetailsHeader
|
||||
{...baseProps}
|
||||
isDataLoaded
|
||||
traceMetadata={traceMetadata}
|
||||
/>,
|
||||
);
|
||||
mockSummary(traceMetadata);
|
||||
render(<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />);
|
||||
|
||||
// Visible by default (showTraceDetails defaults to true).
|
||||
expect(screen.getByText(/inventory-frontend/)).toBeInTheDocument();
|
||||
@@ -192,9 +224,12 @@ describe('TraceDetailsHeader – trace metadata row', () => {
|
||||
expect(screen.getByText(/inventory-frontend/)).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('does not render the metadata row when traceMetadata is absent', () => {
|
||||
render(<TraceDetailsHeader {...baseProps} isDataLoaded />);
|
||||
it('shows skeletons instead of the metadata when the summary is absent', () => {
|
||||
const { container } = render(
|
||||
<TraceDetailsHeader {...baseProps} showTraceDetailsHeaderOptions />,
|
||||
);
|
||||
|
||||
expect(screen.queryByText(/inventory-frontend/)).not.toBeInTheDocument();
|
||||
expect(container.querySelectorAll('.ant-skeleton-input')).toHaveLength(3);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
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 };
|
||||
}
|
||||
@@ -28,7 +28,6 @@ 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';
|
||||
@@ -323,38 +322,6 @@ 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);
|
||||
@@ -393,10 +360,10 @@ function TraceDetailsV3(): JSX.Element {
|
||||
<TraceStoreSync availableColorByFields={availableColorByFields}>
|
||||
<div className={styles.root}>
|
||||
<TraceDetailsHeader
|
||||
filterMetadata={filterMetadata}
|
||||
onFilteredSpansChange={handleFilteredSpansChange}
|
||||
isDataLoaded={!!traceData?.payload?.spans?.length && !showNoData}
|
||||
traceMetadata={traceMetadataForHeader}
|
||||
showTraceDetailsHeaderOptions={
|
||||
!!traceData?.payload?.spans?.length && !showNoData
|
||||
}
|
||||
/>
|
||||
|
||||
{showNoData ? (
|
||||
|
||||
@@ -44,6 +44,7 @@ import {
|
||||
traceDetailFieldKeys,
|
||||
traceDetailFieldValues,
|
||||
traceFlamegraphResponse,
|
||||
traceSummaryResponse,
|
||||
traceWaterfallResponse,
|
||||
} from './__story_mockdata__/traceDetails';
|
||||
|
||||
@@ -149,6 +150,13 @@ 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)),
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
import type {
|
||||
GetFlamegraph200,
|
||||
GetTraceAggregations200,
|
||||
GetTraceSummary200,
|
||||
GetWaterfallV4200,
|
||||
SpantypesFlamegraphSpanDTO,
|
||||
SpantypesSpanAggregationDTO,
|
||||
@@ -303,6 +304,27 @@ 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,
|
||||
|
||||
@@ -84,23 +84,5 @@ 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
|
||||
}
|
||||
|
||||
@@ -93,31 +93,3 @@ 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)
|
||||
}
|
||||
|
||||
@@ -189,49 +189,6 @@ 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 {
|
||||
|
||||
@@ -267,7 +267,8 @@ 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",
|
||||
@@ -297,119 +298,6 @@ 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(
|
||||
|
||||
@@ -202,43 +202,3 @@ 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())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,7 +15,6 @@ 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.
|
||||
@@ -24,5 +23,4 @@ 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)
|
||||
}
|
||||
|
||||
@@ -566,6 +566,24 @@ 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.
|
||||
@@ -580,7 +598,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
|
||||
resStr, hasRes := data["resources_string"]
|
||||
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {
|
||||
attributes := make(map[string]any)
|
||||
attrJSON.FlattenInto("", attributes)
|
||||
flattenJSONPaths("", attrJSON, attributes)
|
||||
if m, ok := attrStr.(map[string]string); ok {
|
||||
for k, v := range m {
|
||||
attributes[k] = v
|
||||
|
||||
@@ -260,7 +260,6 @@ func NewSQLMigrationProviderFactories(
|
||||
sqlmigration.NewAddChannelSpecFactory(sqlschema),
|
||||
sqlmigration.NewAddUserTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddRuleViewFactory(sqlstore, sqlschema),
|
||||
sqlmigration.NewKeepRelevantLLMPricingRulesFactory(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -1,116 +0,0 @@
|
||||
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
|
||||
}
|
||||
@@ -37,8 +37,6 @@ 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)
|
||||
|
||||
@@ -1,202 +0,0 @@
|
||||
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()
|
||||
}
|
||||
@@ -1,75 +0,0 @@
|
||||
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))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
@@ -94,36 +93,35 @@ 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"`
|
||||
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"`
|
||||
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"`
|
||||
}
|
||||
|
||||
// MinimalSpan with only the fields needed to build the parent-child tree.
|
||||
|
||||
@@ -35,21 +35,3 @@ 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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,327 +0,0 @@
|
||||
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
|
||||
Reference in New Issue
Block a user