Compare commits

..

1 Commits

Author SHA1 Message Date
nityanandagohain
f5f019f61b feat: trace detail thread endpoint 2026-09-28 17:32:55 +05:30
45 changed files with 2666 additions and 1402 deletions

View File

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

View File

@@ -1,5 +1,48 @@
components:
schemas:
AiobservabilitytypesMessage:
properties:
content:
items:
$ref: '#/components/schemas/AiobservabilitytypesPart'
type: array
finishReason:
type: string
role:
type: string
required:
- content
type: object
AiobservabilitytypesPart:
properties:
arguments: {}
content:
type: string
id:
type: string
isError:
type: boolean
name:
type: string
redacted:
type: boolean
server:
type: boolean
toolCallId:
type: string
type:
$ref: '#/components/schemas/AiobservabilitytypesPartType'
required:
- type
type: object
AiobservabilitytypesPartType:
enum:
- text
- thinking
- tool_call
- tool_result
- generic
type: string
AlertmanagertypesChannel:
properties:
createdAt:
@@ -9647,6 +9690,17 @@ components:
required:
- aggregations
type: object
SpantypesGettableTraceThread:
properties:
nextCursor:
type: string
spans:
items:
$ref: '#/components/schemas/SpantypesThreadSpan'
type: array
required:
- spans
type: object
SpantypesGettableWaterfallTrace:
properties:
endTimestampMillis:
@@ -9961,6 +10015,92 @@ components:
nullable: true
type: object
type: object
SpantypesThreadSpan:
properties:
attributes:
additionalProperties: {}
nullable: true
type: object
db_name:
type: string
db_operation:
type: string
duration_nano:
minimum: 0
type: integer
events:
items:
$ref: '#/components/schemas/SpantypesEvent'
nullable: true
type: array
external_http_method:
type: string
external_http_url:
type: string
flags:
minimum: 0
type: integer
formatted_input:
items:
$ref: '#/components/schemas/AiobservabilitytypesMessage'
type: array
formatted_output:
items:
$ref: '#/components/schemas/AiobservabilitytypesMessage'
type: array
has_children:
type: boolean
has_error:
type: boolean
http_host:
type: string
http_method:
type: string
http_url:
type: string
is_remote:
type: string
kind_string:
type: string
level:
minimum: 0
type: integer
name:
type: string
parent_span_id:
type: string
references:
items:
$ref: '#/components/schemas/SpantypesOtelSpanRef'
type: array
resource:
additionalProperties:
type: string
nullable: true
type: object
response_status_code:
type: string
span_id:
type: string
status_code:
type: integer
status_code_string:
type: string
status_message:
type: string
sub_tree_node_count:
minimum: 0
type: integer
time_unix:
minimum: 0
type: integer
trace_id:
type: string
trace_state:
type: string
required:
- references
type: object
SpantypesUpdatableSpanMapper:
properties:
config:
@@ -15685,6 +15825,81 @@ 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, each with the messages normalised into formatted_input and formatted_output.
Pages are fetched with the returned nextCursor.
operationId: GetTraceThread
parameters:
- in: query
name: limit
schema:
type: integer
- in: query
name: cursor
schema:
type: string
- in: path
name: traceID
required: true
schema:
type: string
responses:
"200":
content:
application/json:
schema:
properties:
data:
$ref: '#/components/schemas/SpantypesGettableTraceThread'
status:
type: string
required:
- status
- data
type: object
description: OK
"400":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Bad Request
"401":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Unauthorized
"403":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Forbidden
"404":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Not Found
"500":
content:
application/json:
schema:
$ref: '#/components/schemas/RenderErrorResponse'
description: Internal Server Error
security:
- api_key:
- VIEWER
- tokenizer:
- VIEWER
summary: Get thread view for a trace
tags:
- tracedetail
x-signoz-stability: alpha
/api/v1/user/me:
get:
deprecated: true

View File

@@ -23,10 +23,6 @@ const IGNORED_MESSAGES = [
// (YouTube embeds, the docs pane) so they hit the real network instead of
// an unanswered msw request; the block is the point, not a bug.
/violates the following Content Security Policy directive/,
// The filter editor's ANTLR parser reports every syntax error through
// `console.error` (`line 1:14 missing ...`), so each partial expression
// typed into it logs one; the editor shows the same errors on screen.
/^line \d+:\d+ /,
];
interface CapturedMessage {

View File

@@ -4,6 +4,61 @@
* * regenerate with 'pnpm generate:api'
* SigNoz
*/
export enum AiobservabilitytypesPartTypeDTO {
text = 'text',
thinking = 'thinking',
tool_call = 'tool_call',
tool_result = 'tool_result',
generic = 'generic',
}
export interface AiobservabilitytypesPartDTO {
arguments?: unknown;
/**
* @type string
*/
content?: string;
/**
* @type string
*/
id?: string;
/**
* @type boolean
*/
isError?: boolean;
/**
* @type string
*/
name?: string;
/**
* @type boolean
*/
redacted?: boolean;
/**
* @type boolean
*/
server?: boolean;
/**
* @type string
*/
toolCallId?: string;
type: AiobservabilitytypesPartTypeDTO;
}
export interface AiobservabilitytypesMessageDTO {
/**
* @type array
*/
content: AiobservabilitytypesPartDTO[];
/**
* @type string
*/
finishReason?: string;
/**
* @type string
*/
role?: string;
}
export interface AlertmanagertypesChannelDTO {
/**
* @type string
@@ -11142,6 +11197,165 @@ export interface SpantypesOtelSpanRefDTO {
traceId?: string;
}
export type SpantypesThreadSpanDTOAttributesAnyOf = { [key: string]: unknown };
/**
* @nullable
*/
export type SpantypesThreadSpanDTOAttributes =
SpantypesThreadSpanDTOAttributesAnyOf | null;
export type SpantypesThreadSpanDTOResourceAnyOf = { [key: string]: string };
/**
* @nullable
*/
export type SpantypesThreadSpanDTOResource =
SpantypesThreadSpanDTOResourceAnyOf | null;
export interface SpantypesThreadSpanDTO {
/**
* @type object,null
*/
attributes?: SpantypesThreadSpanDTOAttributes;
/**
* @type string
*/
db_name?: string;
/**
* @type string
*/
db_operation?: string;
/**
* @type integer
* @minimum 0
*/
duration_nano?: number;
/**
* @type array,null
*/
events?: SpantypesEventDTO[] | null;
/**
* @type string
*/
external_http_method?: string;
/**
* @type string
*/
external_http_url?: string;
/**
* @type integer
* @minimum 0
*/
flags?: number;
/**
* @type array
*/
formatted_input?: AiobservabilitytypesMessageDTO[];
/**
* @type array
*/
formatted_output?: AiobservabilitytypesMessageDTO[];
/**
* @type boolean
*/
has_children?: boolean;
/**
* @type boolean
*/
has_error?: boolean;
/**
* @type string
*/
http_host?: string;
/**
* @type string
*/
http_method?: string;
/**
* @type string
*/
http_url?: string;
/**
* @type string
*/
is_remote?: string;
/**
* @type string
*/
kind_string?: string;
/**
* @type integer
* @minimum 0
*/
level?: number;
/**
* @type string
*/
name?: string;
/**
* @type string
*/
parent_span_id?: string;
/**
* @type array
*/
references: SpantypesOtelSpanRefDTO[];
/**
* @type object,null
*/
resource?: SpantypesThreadSpanDTOResource;
/**
* @type string
*/
response_status_code?: string;
/**
* @type string
*/
span_id?: string;
/**
* @type integer
*/
status_code?: number;
/**
* @type string
*/
status_code_string?: string;
/**
* @type string
*/
status_message?: string;
/**
* @type integer
* @minimum 0
*/
sub_tree_node_count?: number;
/**
* @type integer
* @minimum 0
*/
time_unix?: number;
/**
* @type string
*/
trace_id?: string;
/**
* @type string
*/
trace_state?: string;
}
export interface SpantypesGettableTraceThreadDTO {
/**
* @type string
*/
nextCursor?: string;
/**
* @type array
*/
spans: SpantypesThreadSpanDTO[];
}
export type SpantypesWaterfallSpanDTOAttributesAnyOf = {
[key: string]: unknown;
};
@@ -12815,6 +13029,30 @@ export type GetTraceAggregations200 = {
status: string;
};
export type GetTraceThreadPathParameters = {
traceID: string;
};
export type GetTraceThreadParams = {
/**
* @type integer
* @description undefined
*/
limit?: number;
/**
* @type string
* @description undefined
*/
cursor?: string;
};
export type GetTraceThread200 = {
data: SpantypesGettableTraceThreadDTO;
/**
* @type string
*/
status: string;
};
export type ListUserPreferences200 = {
/**
* @type array

View File

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

View File

@@ -1,241 +0,0 @@
import { EditorView } from '@codemirror/view';
import { userEvent, waitFor, within } from 'storybook/test';
/** Suggestions wait on a 300ms debounce and a fetch, past the 1s default. */
const untilLoaded = { timeout: 15_000 };
/**
* The editor is controlled: each change round-trips through React state before
* the next one is applied on top of it. People type slower than this.
*/
const KEYSTROKE_MS = 50;
/** Throws until `found` holds something, which is what `waitFor` retries on. */
const present = <TValue>(
found: TValue | null | undefined,
what: string,
): TValue => {
if (found === null || found === undefined) {
throw new Error(`${what} not found`);
}
return found;
};
const pause = (ms: number): Promise<void> =>
new Promise((resolve) => {
setTimeout(resolve, ms);
});
const suggestionList = (canvasElement: HTMLElement): HTMLElement | null =>
canvasElement.querySelector<HTMLElement>('.cm-tooltip-autocomplete');
const suggestionRow = (
canvasElement: HTMLElement,
text: string,
): HTMLElement | undefined => {
const list = suggestionList(canvasElement);
return list
? within(list)
.queryAllByRole('option')
.find((option) => option.textContent?.includes(text))
: undefined;
};
/** Ctrl+Space, the editor's own shortcut for asking for suggestions. */
const requestSuggestions = (editor: HTMLElement): void => {
editor.dispatchEvent(
new KeyboardEvent('keydown', {
key: ' ',
code: 'Space',
ctrlKey: true,
bubbles: true,
}),
);
};
const currentView = (
canvasElement: HTMLElement,
): { editor: HTMLElement; view: EditorView } => {
// An explorer renders one editor per query; the first is the one on screen.
const editor = present(
canvasElement.querySelector<HTMLElement>(
'.code-mirror-where-clause .cm-content',
),
'filter editor',
);
return {
editor,
view: present(EditorView.findFromDOM(editor), 'editor view'),
};
};
/**
* The explorers wrap the filter in `OverlayScrollbar`, which initialises when
* the browser is idle. Initialising moves the content, the editor with it, and
* focuses the editor again through the DOM, which puts the caret back at the
* start and swaps the suggestions for the key list. Throws until every wrapper
* around the filter has initialised.
*/
const assertScrollbarsReady = (editor: HTMLElement): void => {
for (
let wrapper = editor.closest('.overlay-scrollbar');
wrapper;
wrapper = wrapper.parentElement?.closest('.overlay-scrollbar') ?? null
) {
if (!wrapper.hasAttribute('data-overlayscrollbars')) {
throw new Error('scrollbars around the filter still initialising');
}
}
};
/**
* Waits until `text` shows in the suggestion list, asking for suggestions
* whenever the list is shut. Focus and typing only open it once the keys have
* loaded, and moving the caret never does.
*/
const waitForSuggestion = (
canvasElement: HTMLElement,
text: string,
): Promise<HTMLElement> =>
waitFor(
() => {
const { editor } = currentView(canvasElement);
if (!suggestionList(canvasElement)) {
requestSuggestions(editor);
}
return present(suggestionRow(canvasElement, text), `suggestion "${text}"`);
},
{ ...untilLoaded, interval: 250 },
);
/**
* Focuses the filter once the scrollbars around it have initialised, and waits
* for its suggestion list.
*/
const focusFilter = async (canvasElement: HTMLElement): Promise<EditorView> => {
await waitFor(
() => {
const { editor, view } = currentView(canvasElement);
assertScrollbarsReady(editor);
if (!view.hasFocus) {
view.focus();
}
if (!suggestionList(canvasElement)) {
requestSuggestions(editor);
}
return present(suggestionList(canvasElement), 'suggestion list');
},
{ ...untilLoaded, interval: 250 },
);
return currentView(canvasElement).view;
};
/** Waits for a row of the suggestion list. */
export const findSuggestion = (
canvasElement: HTMLElement,
text: string,
): Promise<HTMLElement> => waitForSuggestion(canvasElement, text);
/** Focuses the empty filter: every key, with any recent filters above them. */
export const openKeySuggestions = async (
canvasElement: HTMLElement,
row: string,
): Promise<void> => {
await focusFilter(canvasElement);
await waitForSuggestion(canvasElement, row);
};
/**
* Focuses the filter and types onto the end of it one character at a time,
* each as the transaction a keystroke makes, leaving the caret at the end so
* the suggestion list follows what was typed. Quotes and brackets are not
* closed for it: type both.
*
* `userEvent.type` cannot be used: CodeMirror redraws the line as tokens are
* highlighted, which strands the caret `userEvent` tracks.
*/
export const typeFilter = async (
canvasElement: HTMLElement,
text: string,
): Promise<void> => {
const view = await focusFilter(canvasElement);
for (const character of text) {
const at = view.state.doc.length;
view.dispatch({
changes: { from: at, insert: character },
selection: { anchor: at + character.length },
userEvent: 'input.type',
});
await pause(KEYSTROKE_MS);
}
};
/**
* Types an expression, then steps the caret back inside it, before a closing
* bracket or parenthesis, where the suggestions are about what goes in there.
*/
export const typeFilterWithCaretBack = async (
canvasElement: HTMLElement,
text: string,
stepsBack: number,
row: string,
): Promise<void> => {
await typeFilter(canvasElement, text);
const { view } = currentView(canvasElement);
view.dispatch({
selection: { anchor: view.state.doc.length - stepsBack },
userEvent: 'select',
});
await waitForSuggestion(canvasElement, row);
};
/**
* Moves focus off the filter, which is when the expression is validated and
* the error marker can show.
*/
export const blurFilter = async (canvasElement: HTMLElement): Promise<void> => {
await userEvent.keyboard('{Escape}');
await userEvent.click(canvasElement.ownerDocument.body);
};
/** Types an expression, leaves the filter and opens its validation errors. */
export const showFilterErrors = async (
canvasElement: HTMLElement,
text: string,
): Promise<void> => {
await typeFilter(canvasElement, text);
await blurFilter(canvasElement);
const marker = await waitFor(
() =>
present(
canvasElement.querySelector<HTMLElement>('.query-status-container button'),
'error marker',
),
untilLoaded,
);
await userEvent.hover(marker);
await waitFor(
() =>
present(
canvasElement.ownerDocument.querySelector('.query-validation-error'),
'validation error',
),
untilLoaded,
);
};

View File

@@ -8,7 +8,6 @@ import { QueryParams } from 'constants/query';
import ROUTES from 'constants/routes';
import { encode } from 'js-base64';
import type { Tags } from 'hooks/useResourceAttribute/types';
import { fireEvent, userEvent, waitFor, within } from 'storybook/test';
import {
choiceControl,
@@ -32,9 +31,6 @@ import {
type ServiceHealth,
} from './__story_mockdata__/serviceMap';
/** The keys are only fetched once the select opens, past the 1s default. */
const untilLoaded = { timeout: 15_000 };
const GRAPH = 'Service map · graph';
const FILTERS = 'Service map · filters';
@@ -133,80 +129,3 @@ export const serviceMapMocks = defineStoryMocks({
],
config: (values) => ({ route: serviceMapRoute(values.filters) }),
});
/** Opens the select under a test id; it closes again after every pick. */
export const openSelect = async (
canvasElement: HTMLElement,
testId: string,
): Promise<void> => {
const select = await within(canvasElement).findByTestId(
testId,
undefined,
untilLoaded,
);
await userEvent.click(within(select).getByRole('combobox'));
};
export const OPEN_DROPDOWN =
'.ant-select-dropdown:not(.ant-select-dropdown-hidden)';
/** antd keeps a hidden copy of each label for screen readers; the title skips it. */
const visibleOption = (title: string): HTMLElement | null =>
document.querySelector<HTMLElement>(
`${OPEN_DROPDOWN} .ant-select-item-option[title="${title}"]`,
);
/**
* Picks the option titled `title` in the open dropdown. `userEvent.click`
* moves focus off the select on the way, which closes it before the option
* takes the click.
*/
export const pickOption = async (title: string): Promise<void> => {
const option = await waitFor(() => {
const match = visibleOption(title);
if (!match) {
throw new Error(`option "${title}" not found`);
}
return match;
}, untilLoaded);
await fireEvent.click(option);
};
/**
* Opens the attribute filter on its next step. Each single-choice pick closes
* the dropdown and swaps the select for the next step's, and a click that lands
* before the swap opens nothing, so it opens again until `title` shows.
*/
export const openAttributeFilterOn = (
canvasElement: HTMLElement,
title: string,
): Promise<void> =>
waitFor(
async () => {
if (visibleOption(title)) {
return;
}
if (!document.querySelector(OPEN_DROPDOWN)) {
await openSelect(canvasElement, 'resource-attributes-filter');
}
throw new Error(`option "${title}" not shown`);
},
{ ...untilLoaded, interval: 500 },
);
/** Stages `k8s.cluster.name IN` and leaves the filter open on its values. */
export const stageClusterIn = async (
canvasElement: HTMLElement,
): Promise<void> => {
await openAttributeFilterOn(canvasElement, 'k8s.cluster.name');
await pickOption('k8s.cluster.name');
await openAttributeFilterOn(canvasElement, 'IN');
await pickOption('IN');
await openAttributeFilterOn(canvasElement, 'staging-eu');
};

View File

@@ -1,17 +1,10 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import { expect, screen, userEvent, waitFor, within } from 'storybook/test';
import { screen, userEvent, within } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
import type { PageStoryArgs } from '@/storybook/runtime/resolveStory';
import {
OPEN_DROPDOWN,
openAttributeFilterOn,
openSelect,
pickOption,
serviceMapMocks,
stageClusterIn,
} from './ServiceMap.stories.mocks';
import { serviceMapMocks } from './ServiceMap.stories.mocks';
import ServiceMapContainer from '../index';
@@ -104,49 +97,3 @@ export const FilterAttributes: Story = {
await screen.findByText('k8s.cluster.name', undefined, untilLoaded);
},
};
/** A key staged as a chip, the filter open again on how to match it. */
export const FilterOperators: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openAttributeFilterOn(canvasElement, 'k8s.cluster.name');
await pickOption('k8s.cluster.name');
await openAttributeFilterOn(canvasElement, 'Not IN');
},
};
/** A key and `IN` staged, the filter open on the values the key holds. */
export const FilterValues: Story = {
play: async ({ canvasElement }): Promise<void> => {
await stageClusterIn(canvasElement);
},
};
/** Two values ticked before the filter is left, which is what applies it. */
export const FilterValuesSelected: Story = {
play: async ({ canvasElement }): Promise<void> => {
await stageClusterIn(canvasElement);
await pickOption('prod-us-east');
await pickOption('prod-eu-west');
await waitFor(
() =>
expect(
document.querySelectorAll('.ant-select-item-option-selected'),
).toHaveLength(2),
untilLoaded,
);
},
};
/** The environment selector open on the environments the calls came from. */
export const EnvironmentOptions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openSelect(canvasElement, 'resource-environment-filter');
await waitFor(
() =>
expect(
document.querySelector(`${OPEN_DROPDOWN} .ant-select-item-option`),
).not.toBeNull(),
untilLoaded,
);
},
};

View File

@@ -103,18 +103,6 @@ export const DeleteDowntimeConfirm: Story = {
},
};
/** The new-downtime form set to repeat weekly, which adds the days and duration. */
export const NewDowntimeRecurring: Story = {
play: async ({ canvasElement }): Promise<void> => {
await NewDowntime.play?.({ canvasElement } as never);
await userEvent.click(
await screen.findByRole('combobox', { name: 'Repeats every' }),
);
await userEvent.click(await screen.findByText('Weekly'));
await screen.findByText('Duration');
},
};
/** A client-side search with no matching downtime schedule. */
export const SearchNoResults: Story = {
play: async ({ canvasElement }): Promise<void> => {

View File

@@ -6,19 +6,9 @@
import { rest } from 'msw';
import set from 'api/browser/localstorage/set';
import { LOCALSTORAGE } from 'constants/localStorage';
import { screen, userEvent, waitFor, within } from 'storybook/test';
import {
choiceControl,
countControl,
toggleControl,
} from '@/storybook/controls/controls';
import { countControl, toggleControl } from '@/storybook/controls/controls';
import { defineStoryMocks } from '@/storybook/controls/defineStoryMocks';
import {
RESPONSE_STATES,
type ResponseState,
respondWith,
} from '@/storybook/runtime/responseState';
import { fieldValuesResponse } from '@/storybook/msw/__story_mockdata__/fields';
import {
@@ -33,9 +23,6 @@ import {
type ListErrorsBody,
} from './__story_mockdata__/exceptions';
/** The page fetches before it renders a row, which outlasts the 1s default. */
const untilLoaded = { timeout: 15_000 };
const LIST = 'Exceptions · list';
const FILTERS = 'Exceptions · filters';
@@ -55,13 +42,6 @@ export const exceptionsMocks = defineStoryMocks({
value: 6,
max: EXCEPTION_QUICK_FILTER_CAP,
}),
filterKeys: choiceControl<ResponseState>('Filter keys', {
group: FILTERS,
description:
'How `/autocomplete/attribute_keys` answers when the resource filter opens, apart from the page-wide Data control.',
options: RESPONSE_STATES,
value: 'loaded',
}),
filterPanel: toggleControl('Quick filters panel', {
group: FILTERS,
description:
@@ -112,7 +92,7 @@ export const exceptionsMocks = defineStoryMocks({
rest.get(
'http://localhost/api/v3/autocomplete/attribute_keys',
respondWith(values.filterKeys, (req) =>
response.json((req) =>
exceptionAttributeKeysResponse(req.url.searchParams.get('searchText')),
),
),
@@ -131,61 +111,3 @@ export const exceptionsMocks = defineStoryMocks({
set(LOCALSTORAGE.SHOW_EXCEPTIONS_QUICK_FILTERS, String(values.filterPanel));
},
});
export const openResourceFilter = async (
canvasElement: HTMLElement,
): Promise<HTMLElement> => {
const filter = await within(canvasElement).findByTestId(
'qb-search-select',
undefined,
untilLoaded,
);
await userEvent.click(within(filter).getByRole('combobox'));
return filter;
};
/**
* Clicks the visible row whose label is `text`. The dropdown renders in the
* body, a key row carries its type beside the label, and antd keeps a hidden
* copy of each label for screen readers that takes no clicks.
*
* The key list renders twice after the filter opens, since a second key
* request empties it until it answers. The click happens in the same task as
* the lookup: `fireEvent` from `storybook/test` dispatches a tick later, which
* can land on a row already removed and never reach React.
*/
export const pickSuggestion = async (text: string): Promise<void> => {
await waitFor(() => {
const match = Array.from(
document.querySelectorAll<HTMLElement>(
'.query-builder-search.ant-select-dropdown:not(.ant-select-dropdown-hidden) .ant-select-item-option',
),
).find((option) =>
Array.from(option.querySelectorAll('*')).some(
(node) => node.children.length === 0 && node.textContent === text,
),
);
if (!match) {
throw new Error(`suggestion "${text}" not found`);
}
// `userEvent.click` moves focus off the search input on the way, which closes
// the dropdown before the row takes the click.
match.click();
}, untilLoaded);
};
export const commitFilter = async (
key: string,
operator: string,
value: string,
): Promise<void> => {
await pickSuggestion(key);
await screen.findByText('Operator for', { exact: false }, untilLoaded);
await pickSuggestion(operator);
await screen.findByText('Value(s) for', { exact: false }, untilLoaded);
await pickSuggestion(value);
};

View File

@@ -5,12 +5,7 @@ import { expect, screen, userEvent, waitFor, within } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
import type { PageStoryArgs } from '@/storybook/runtime/resolveStory';
import {
commitFilter,
exceptionsMocks,
openResourceFilter,
pickSuggestion,
} from './AllErrors.stories.mocks';
import { exceptionsMocks } from './AllErrors.stories.mocks';
import AllErrors from '../index';
type AllErrorsArgs = PageStoryArgs<typeof exceptionsMocks>;
@@ -130,108 +125,3 @@ export const QuickFiltersSettingsWithBanner: Story = {
args: { banner: 'trial-expiry' },
play: dirtyQuickFiltersSettings,
};
/** The resource filter opened: every key the exceptions can be narrowed by. */
export const FilterKeySuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await screen.findByText('Suggested Filters', undefined, untilLoaded);
},
};
/** The key list grown past its first rows with the Show all shortcut. */
export const FilterAllKeys: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await screen.findByText('Show all filter items', undefined, untilLoaded);
await userEvent.keyboard('{Control>}/{/Control}');
await screen.findByText('cloud.region', undefined, untilLoaded);
},
};
/** A partial key: the typed text as a free search, then the keys that match. */
export const FilterPartialKey: Story = {
play: async ({ canvasElement }): Promise<void> => {
const filter = await openResourceFilter(canvasElement);
await userEvent.type(within(filter).getByRole('combobox'), 'serv');
await screen.findByText('service.namespace', undefined, untilLoaded);
},
};
/** A key picked: the operators the exceptions page allows for it. */
export const FilterOperatorSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await pickSuggestion('service.name');
await screen.findByText('Operator for', { exact: false }, untilLoaded);
},
};
/** A key and an operator picked: the values the key holds. */
export const FilterValueSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await pickSuggestion('service.name');
await pickSuggestion('=');
await screen.findByText('Value(s) for', { exact: false }, untilLoaded);
},
};
/** Two conditions committed as chips, with the dropdown closed again. */
export const FilterChips: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await commitFilter('service.name', '=', 'checkout');
await commitFilter('deployment.environment', '!=', 'staging');
await userEvent.click(canvasElement.ownerDocument.body);
await within(canvasElement).findByText(
'deployment.environment != staging',
undefined,
untilLoaded,
);
},
};
/**
* A committed chip clicked to change it: its text goes back into the input,
* with the dropdown shut until the input is typed into.
*/
export const FilterEditChip: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await commitFilter('service.name', '=', 'checkout');
await userEvent.click(
await within(canvasElement).findByText(
'service.name = checkout',
undefined,
untilLoaded,
),
);
// The select remounts whenever its chips change, so it is looked up again.
await waitFor(
() =>
expect(
within(within(canvasElement).getByTestId('qb-search-select')).getByRole(
'combobox',
),
).toHaveValue('service.name = checkout'),
untilLoaded,
);
},
};
/** The key list while its request is still in flight. */
export const FilterKeysLoading: Story = {
args: { filterKeys: 'loading' },
play: async ({ canvasElement }): Promise<void> => {
await openResourceFilter(canvasElement);
await waitFor(
() =>
expect(
document.querySelector('.query-builder-search .ant-spin'),
).not.toBeNull(),
untilLoaded,
);
},
};

View File

@@ -15,10 +15,7 @@ import {
toggleControl,
} from '@/storybook/controls/controls';
import { defineStoryMocks } from '@/storybook/controls/defineStoryMocks';
import {
fieldKeysResponse,
fieldValuesResponse,
} from '@/storybook/msw/__story_mockdata__/fields';
import { fieldValuesResponse } from '@/storybook/msw/__story_mockdata__/fields';
import { queryRangeV5ScalarResponse } from '@/storybook/msw/__story_mockdata__/queryRange';
import {
@@ -34,7 +31,6 @@ import {
import {
emptyPanelResponse,
NAMESPACE_VALUES,
VARIABLE_ATTRIBUTES,
panelResponse,
serviceVariableValues,
} from './__story_mockdata__/panelData';
@@ -242,12 +238,6 @@ export const dashboardMocks = defineStoryMocks({
response.json(() => fieldValuesResponse(NAMESPACE_VALUES)),
),
// The dynamic variable editor lists the attributes a variable can read.
rest.get(
'http://localhost/api/v1/fields/keys',
response.json(() => fieldKeysResponse(VARIABLE_ATTRIBUTES)),
),
// The header reads the public link on every load, so it answers even while
// the panels are held in the loading or failed state.
rest.get('http://localhost/api/v1/dashboards/:id/public', (_req, res, ctx) =>

View File

@@ -2,7 +2,7 @@ import type { ComponentType } from 'react';
import type { Meta, StoryObj } from '@storybook/react-vite';
import { Route } from 'react-router-dom';
import ROUTES from 'constants/routes';
import { expect, screen, userEvent, waitFor, within } from 'storybook/test';
import { screen, userEvent, within } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
import type { PageStoryArgs } from '@/storybook/runtime/resolveStory';
@@ -285,67 +285,6 @@ export const SectionActionsMenu: Story = {
},
};
/** The panel menu's move-to-section submenu, open on the sections it can go to. */
export const PanelMoveToSectionSubmenu: Story = {
play: async (context) => {
await PanelActionsMenu.play?.(context);
await userEvent.hover(await screen.findByText('Move to section'));
// The submenu lists the sections the panel is not already in.
await waitFor(() => expect(screen.getAllByRole('menu')).toHaveLength(2), {
timeout: 10000,
});
},
};
/** Dashboard settings on the Publish tab, where the public link is managed. */
export const SettingsPublicDashboard: Story = {
play: async ({ canvasElement }) => {
await userEvent.click(
await within(canvasElement).findByRole(
'button',
{ name: 'Configure' },
{ timeout: 10000 },
),
);
await userEvent.click(await screen.findByRole('tab', { name: 'Publish' }));
await screen.findByText('Default time range');
},
};
/** The Publish tab with its default time range select open. */
export const SettingsPublicDashboardSelectOpen: Story = {
play: async (context) => {
await SettingsPublicDashboard.play?.(context);
await userEvent.click(
within(screen.getByRole('tabpanel')).getByRole('combobox'),
);
await screen.findByRole('listbox');
},
};
/** The Variables tab of dashboard settings with a new variable's form open. */
export const SettingsVariablesNew: Story = {
play: async ({ canvasElement }) => {
await userEvent.click(
await within(canvasElement).findByRole(
'button',
{ name: 'Configure' },
{ timeout: 10000 },
),
);
await userEvent.click(await screen.findByRole('tab', { name: 'Variables' }));
const add = await within(await screen.findByRole('tabpanel')).findByRole(
'button',
{ name: 'Add variable' },
);
// The button stays disabled until its permission check resolves.
await waitFor(() => expect(add).toBeEnabled(), { timeout: 10000 });
await userEvent.click(add);
await screen.findByText('Variable Type');
},
};
/**
* A dashboard id nobody has, which is what a deleted or mistyped link opens on.
*

View File

@@ -168,15 +168,6 @@ export const serviceVariableValues = (count: number): string[] =>
);
/** Values the dynamic `namespace` variable resolves from the fields endpoint. */
/** Attributes the dynamic variable editor offers a variable to read. */
export const VARIABLE_ATTRIBUTES = [
'k8s.namespace.name',
'k8s.cluster.name',
'service.name',
'deployment.environment',
'host.name',
];
export const NAMESPACE_VALUES = [
'checkout-prod',
'payments-prod',

View File

@@ -6,7 +6,6 @@
import { rest } from 'msw';
import type { GetDashboardV2200 } from 'api/generated/services/sigNoz.schemas';
import ROUTES from 'constants/routes';
import { screen, userEvent, waitFor, within } from 'storybook/test';
import {
choiceControl,
@@ -343,30 +342,3 @@ export const dashboardsListMocks = defineStoryMocks({
});
},
});
/** Opens the actions menu of the row at `index`. */
export const openRowActions = async (
canvasElement: HTMLElement,
index: number,
): Promise<void> => {
// The icon-only trigger carries no accessible name.
const triggers = await within(canvasElement).findAllByTestId(
'dashboard-action-icon',
{},
{ timeout: 10000 },
);
await userEvent.click(triggers[index]);
await screen.findByText('Rename');
};
/** Picks a row action, retrying while its permission check still disables it. */
export const pickRowAction = async (label: string | RegExp): Promise<void> => {
await waitFor(
async () => {
await userEvent.click(screen.getByText(label));
await screen.findByRole('dialog', {}, { timeout: 500 });
},
{ timeout: 10000 },
);
};

View File

@@ -6,9 +6,7 @@ import type { PageStoryArgs } from '@/storybook/runtime/resolveStory';
import {
dashboardsListMocks,
openRowActions,
overflowingRows,
pickRowAction,
} from './DashboardsListPage.stories.mocks';
import { BuiltinViewId } from '../types';
@@ -116,45 +114,6 @@ export const Tooltips: Story = {
parameters: { msw: { handlers: [overflowingRows] } },
};
/** The first row's actions menu, open over the list. */
export const RowActionsMenu: Story = {
play: async ({ canvasElement }) => {
await openRowActions(canvasElement, 0);
},
};
/** The rename dialog, opened from the menu of the second row (the first is locked). */
export const RenameDashboardDialog: Story = {
play: async ({ canvasElement }) => {
await openRowActions(canvasElement, 1);
await pickRowAction('Rename');
await screen.findByRole('dialog', { name: 'Rename dashboard' });
},
};
/** The tags dialog, opened from the menu of the second row (the first is locked). */
export const EditTagsDialog: Story = {
play: async ({ canvasElement }) => {
await openRowActions(canvasElement, 1);
await pickRowAction(/^(Edit|Add) Tags$/);
await screen.findByRole('dialog', { name: /^(Edit|Add) tags$/ });
},
};
/** The popover that names the current filters as a new saved view. */
export const SaveViewPopover: Story = {
play: async ({ canvasElement }) => {
await userEvent.click(
await within(canvasElement).findByRole(
'button',
{ name: 'Save current filters as a view' },
{ timeout: 10000 },
),
);
await screen.findByText('Save as view');
},
};
/**
* The query the backend refused: the parse error it returned replaces the
* generic failure copy, and there is nothing to retry.

View File

@@ -1,9 +1,4 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import {
findSuggestion,
openKeySuggestions,
typeFilter,
} from 'components/QueryBuilderV2/QueryV2/QuerySearch/stories/__story_mockdata__/querySearch.play';
import { screen, userEvent, within } from 'storybook/test';
import { VIEWS } from 'container/InfraMonitoringK8sV2/constants';
@@ -46,21 +41,6 @@ export const PodDetailsEvents: StoryObj<PodsArgs> = {
args: { drawer: true, drawerTab: VIEWS.EVENTS },
};
/** The selected pod's details drawer, switched to its logs tab. */
export const DetailsDrawerLogsTab: StoryObj<PodsArgs> = {
args: { drawer: true },
play: async () => {
const drawer = within(
await screen.findByRole('dialog', {}, { timeout: 10000 }),
);
await userEvent.click(
await drawer.findByText('Logs', {}, { timeout: 10000 }),
);
await drawer.findAllByText(/handled request in/, {}, { timeout: 10000 });
},
};
/**
* Every tooltip the pod list carries, held open: Collapse Filters beside the
* quick filters, Options above the table, the Pod Name, Status, Age and Restarts
@@ -106,27 +86,3 @@ export const TooltipsInOptionsPanel: StoryObj<PodsArgs> = {
await screen.findByText('Columns');
},
};
/**
* The pod list re-renders the filter when the viewport grows, and each render
* reconfigures the editor, which closes its suggestions. Shot at the height
* the page opened at.
*/
const heldViewport = { sbshot: { viewport: { height: 1200 } } };
/** The pod filter focused: the Kubernetes keys pods can be narrowed by. */
export const FilterKeySuggestions: StoryObj<PodsArgs> = {
parameters: heldViewport,
play: async ({ canvasElement }): Promise<void> => {
await openKeySuggestions(canvasElement, 'k8s.node.name');
},
};
/** The pod filter on a namespace: the namespaces the pods run in. */
export const FilterValueSuggestions: StoryObj<PodsArgs> = {
parameters: heldViewport,
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'k8s.namespace.name = ');
await findSuggestion(canvasElement, 'kube-system');
},
};

View File

@@ -69,12 +69,3 @@ export const GroupActionsMenu: Story = {
await screen.findByRole('menu');
},
};
/** The first mapping group's edit drawer, opened from its menu. */
export const GroupFormDrawer: Story = {
play: async (context): Promise<void> => {
await GroupActionsMenu.play?.(context);
await userEvent.click(await screen.findByText('Edit'));
await screen.findByText('Edit group');
},
};

View File

@@ -54,24 +54,3 @@ export const ModelCostActionsMenu: Story = {
await screen.findByRole('menu');
},
};
/** The first pricing rule's drawer, opened from its row menu. */
export const ModelCostDrawer: Story = {
play: async (context): Promise<void> => {
await ModelCostActionsMenu.play?.(context);
await userEvent.click(await screen.findByText('Edit'));
await screen.findByText('Edit model cost');
},
};
/** The drawer with its cache mode select open. */
export const ModelCostDrawerCacheModeOpen: Story = {
play: async (context): Promise<void> => {
await ModelCostDrawer.play?.(context);
await userEvent.click(
await screen.findByRole('combobox', { name: 'Cache mode' }),
);
await screen.findByRole('listbox');
},
};

View File

@@ -25,11 +25,6 @@ import {
toggleControl,
} from '@/storybook/controls/controls';
import { defineStoryMocks } from '@/storybook/controls/defineStoryMocks';
import {
RESPONSE_STATES,
type ResponseState,
respondWith,
} from '@/storybook/runtime/responseState';
import {
logsSavedViewsResponse,
@@ -49,8 +44,6 @@ import {
logRowsResponse,
QUICK_FILTER_MAX,
logsQuickFiltersResponse,
RECENT_FILTER_MAX,
recentFiltersStorage,
RELATIVE_TIME,
timeRangeState,
} from './__story_mockdata__/logs';
@@ -165,20 +158,6 @@ export const logsMocks = defineStoryMocks({
value: QUICK_FILTER_MAX,
max: QUICK_FILTER_MAX,
}),
filterValues: choiceControl<ResponseState>('Filter values', {
group: FILTERS,
description:
'How `/fields/values` answers once a key and an operator are typed in the filter, apart from the page-wide Data control.',
options: RESPONSE_STATES,
value: 'loaded',
}),
recentFilters: countControl('Recent filters', {
group: FILTERS,
description:
'Filters run before in this browser, which the filter offers above its key suggestions.',
value: 0,
max: RECENT_FILTER_MAX,
}),
savedViews: countControl('Saved views', {
group: FILTERS,
description: 'The views the view picker above the query builder lists.',
@@ -238,7 +217,7 @@ export const logsMocks = defineStoryMocks({
rest.get(
'http://localhost/api/v1/fields/values',
respondWith(values.filterValues, (req) =>
response.json((req) =>
logFieldValuesResponse(
req.url.searchParams.get('name') ?? '',
req.url.searchParams.get('searchText') ?? '',
@@ -285,16 +264,8 @@ export const logsMocks = defineStoryMocks({
}
: {},
}),
effect: ({
frequencyChart,
filtersPanel,
format,
maxLines,
fontSize,
recentFilters,
}) => {
effect: ({ frequencyChart, filtersPanel, format, maxLines, fontSize }) => {
setLocalStorage(LOCALSTORAGE.SHOW_FREQUENCY_CHART, String(frequencyChart));
setLocalStorage(...recentFiltersStorage(recentFilters));
setLocalStorage(LOCALSTORAGE.SHOW_LOGS_QUICK_FILTERS, String(filtersPanel));
// The preferences loader reads localStorage ahead of the URL, so this is

View File

@@ -1,12 +1,4 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import {
blurFilter,
openKeySuggestions,
showFilterErrors,
typeFilter,
typeFilterWithCaretBack,
findSuggestion,
} from 'components/QueryBuilderV2/QueryV2/QuerySearch/stories/__story_mockdata__/querySearch.play';
import { expect, screen, userEvent, waitFor, within } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
@@ -217,138 +209,3 @@ export const Tooltips: Story = {
);
},
};
/** The filter focused before anything is typed: every key the logs carry. */
export const FilterKeySuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openKeySuggestions(canvasElement, 'severity_text');
},
};
/** A partial key, with the keys that still match and the typed part marked. */
export const FilterPartialKey: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'serv');
await findSuggestion(canvasElement, 'service.name');
},
};
/** A string key followed by a space: the operators a string compares with. */
export const FilterOperatorSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name ');
await findSuggestion(canvasElement, 'CONTAINS');
},
};
/** A number key puts the range comparisons first. */
export const FilterNumberOperatorSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'http.status_code ');
await findSuggestion(canvasElement, 'BETWEEN');
},
};
/** `NOT` after a key narrows the list to the operators it can negate. */
export const FilterNegatedOperatorSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name NOT ');
await findSuggestion(canvasElement, 'IN');
},
};
/** A key and an operator: the values the key holds, fetched for it. */
export const FilterValueSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name = ');
await findSuggestion(canvasElement, 'checkout');
},
};
/** Values still being fetched for the key. */
export const FilterValuesLoading: Story = {
args: { filterValues: 'loading' },
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name = ');
await findSuggestion(canvasElement, 'Loading suggestions');
},
};
/** A key the backend holds no values for, such as the free-text body. */
export const FilterNoValueSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'body = ');
await findSuggestion(canvasElement, 'No suggestions available');
},
};
/** The values request failed. */
export const FilterValuesError: Story = {
args: { filterValues: 'error' },
// The values request deliberately fails.
parameters: { allowConsoleErrors: true },
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name = ');
await findSuggestion(canvasElement, 'Error loading suggestions');
},
};
/** Inside an `IN` list, after the first value: the rest of the values. */
export const FilterInList: Story = {
// The editor logs a TypeError while the list is open, which it survives.
parameters: { allowConsoleErrors: true },
play: async ({ canvasElement }): Promise<void> => {
await typeFilterWithCaretBack(
canvasElement,
"service.name IN ['auth', ]",
1,
'checkout',
);
},
};
/** A complete condition: the conjunctions that start the next one. */
export const FilterConjunctionSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, "service.name = 'checkout' ");
await findSuggestion(canvasElement, 'OR');
},
};
/** Inside an opened group: keys, another group and `NOT`. */
export const FilterNestedGroup: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilterWithCaretBack(canvasElement, '()', 1, 'NOT');
},
};
/** A long valid expression mixing operators, left for the next run. */
export const FilterComplete: Story = {
// The editor logs a TypeError while the `IN` list is typed, which it survives.
parameters: { allowConsoleErrors: true },
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(
canvasElement,
"service.name IN ['checkout', 'payments'] AND severity_text = 'ERROR' AND http.status_code >= 500 AND body CONTAINS 'timeout'",
);
await blurFilter(canvasElement);
},
};
/** An incomplete expression after focus left: the marker and its errors. */
export const FilterSyntaxError: Story = {
play: async ({ canvasElement }): Promise<void> => {
await showFilterErrors(canvasElement, 'service.name = ');
},
};
/** Filters run before, offered above the key suggestions. */
export const FilterRecentSearches: Story = {
args: { recentFilters: 3 },
play: async ({ canvasElement }): Promise<void> => {
await openKeySuggestions(
canvasElement,
"k8s.namespace.name = 'observability'",
);
},
};

View File

@@ -15,9 +15,6 @@ import {
} from 'api/generated/services/sigNoz.schemas';
import { defaultLogsSelectedColumns } from 'container/OptionsMenu/constants';
import type { OptionsQuery } from 'container/OptionsMenu/types';
import { STORAGE_VERSION } from 'lib/recentQueries/constants';
import type { RecentQueriesStoreShape } from 'lib/recentQueries/types';
import { makeId, storageKeyFor } from 'lib/recentQueries/utils';
import type { Time } from 'container/TopNav/DateTimeSelectionV2/types';
import { quickFiltersListResponse } from 'mocks-server/__mockdata__/customQuickFilters';
import type { AppState } from 'store/reducers';
@@ -473,34 +470,3 @@ const DASHBOARD_NAMES = [
export const dashboardsResponse = (): ListDashboardsForUserV2200 =>
dashboardsForUserResponse(DASHBOARD_NAMES);
const RECENT_FILTERS = [
"service.name = 'checkout' AND severity_text = 'ERROR'",
"k8s.namespace.name = 'observability'",
"body CONTAINS 'timeout'",
'http.status_code >= 500',
"deployment.environment IN ['production', 'staging']",
];
/** The filter lists this many recent entries at most. */
export const RECENT_FILTER_MAX = RECENT_FILTERS.length;
/**
* The explorer's recent filters as the store keeps them in localStorage,
* newest first, five minutes apart. Dated off `new Date()`, which the story
* clock freezes, so the "5 minutes ago" labels hold still.
*/
export const recentFiltersStorage = (count: number): [string, string] => {
const store: RecentQueriesStoreShape = {
version: STORAGE_VERSION,
entries: RECENT_FILTERS.slice(0, count).map((expression, index) => ({
id: makeId('logs', '', expression),
signal: 'logs',
source: '',
filter: { expression },
lastUsedAt: new Date().getTime() - (index + 1) * 5 * 60_000,
})),
};
return [storageKeyFor('logs', ''), JSON.stringify(store)];
};

View File

@@ -1,9 +1,4 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import {
findSuggestion,
openKeySuggestions,
typeFilter,
} from 'components/QueryBuilderV2/QueryV2/QuerySearch/stories/__story_mockdata__/querySearch.play';
import { screen, userEvent } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
@@ -119,18 +114,3 @@ export const MetricDetailsDashboardsMenu: Story = {
await screen.findByRole('menu');
},
};
/** The summary's metric search focused: the attributes metrics can be found by. */
export const FilterKeySuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openKeySuggestions(canvasElement, 'k8s.cluster.name');
},
};
/** The metric search on an attribute and an operator: the values it holds. */
export const FilterValueSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name = ');
await findSuggestion(canvasElement, 'checkout');
},
};

View File

@@ -15,14 +15,10 @@ import {
toggleControl,
} from '@/storybook/controls/controls';
import { defineStoryMocks } from '@/storybook/controls/defineStoryMocks';
import type { MockResolver } from '@/storybook/msw/types';
import {
CREATE_OUTCOMES,
type CreateOutcome,
EXPIRIES,
type Expiry,
ingestionKeyCreateError,
ingestionKeysResponse,
KEYS_PER_PAGE,
legacyIngestionResponse,
@@ -33,11 +29,6 @@ import {
const KEYS = 'Ingestion · keys';
const LIMITS = 'Ingestion · limits';
const rejectCreate: MockResolver = (_req, res, ctx) =>
res(ctx.status(409), ctx.json(ingestionKeyCreateError()));
const holdCreate: MockResolver = (_req, res, ctx) => res(ctx.delay('infinite'));
export const ingestionMocks = defineStoryMocks({
controls: {
gateway: toggleControl('Gateway', {
@@ -59,13 +50,6 @@ export const ingestionMocks = defineStoryMocks({
options: EXPIRIES,
value: 'none',
}),
create: choiceControl<CreateOutcome>('Creating a key', {
group: KEYS,
description:
'What the create behind the new key form answers. `hangs` holds the submit button in its loading state.',
options: CREATE_OUTCOMES,
value: 'succeeds',
}),
limits: multiChoiceControl<LimitSignal>('Signals with a limit', {
group: LIMITS,
description:
@@ -101,14 +85,10 @@ export const ingestionMocks = defineStoryMocks({
rest.post(
'http://localhost/api/v2/gateway/ingestion_keys',
{
succeeds: response.json(() => ({
status: 'success',
data: { id: 'ingestion-key-new', value: 'sk_new' },
})),
fails: rejectCreate,
hangs: holdCreate,
}[values.create],
response.json(() => ({
status: 'success',
data: { id: 'ingestion-key-new', value: 'sk_new' },
})),
),
rest.patch(

View File

@@ -1,6 +1,5 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import dayjs from 'dayjs';
import { expect, screen, userEvent, waitFor, within } from 'storybook/test';
import { screen, userEvent, within } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
import type { PageStoryArgs } from '@/storybook/runtime/resolveStory';
@@ -77,149 +76,16 @@ export const KeyLimits: Story = {
},
};
async function openCreateKey(
canvasElement: HTMLElement,
): Promise<ReturnType<typeof within>> {
await userEvent.click(
await within(canvasElement).findByText(
'New Ingestion key',
undefined,
untilLoaded,
),
);
return within(
await screen.findByRole(
'dialog',
{ name: 'Create new ingestion key' },
untilLoaded,
),
);
}
async function addTag(
dialog: ReturnType<typeof within>,
tag: string,
): Promise<void> {
await userEvent.click(await dialog.findByRole('button', { name: /New Tag/ }));
await userEvent.keyboard(`${tag}{Enter}`);
}
async function submitCreateKey(
dialog: ReturnType<typeof within>,
): Promise<void> {
await userEvent.type(dialog.getByLabelText('Name'), 'otel-collectors');
await userEvent.click(dialog.getByLabelText('Expiration'));
await userEvent.click(
await screen.findByTitle(dayjs().add(1, 'day').format('YYYY-MM-DD')),
);
await userEvent.click(
dialog.getByRole('button', { name: 'Create new Ingestion key' }),
);
}
/** The form a new key is named and dated in. */
export const CreateKey: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openCreateKey(canvasElement);
},
};
/** A tag typed on the new key and not confirmed yet. */
export const CreateKeyAddingTag: Story = {
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
await userEvent.click(dialog.getByRole('button', { name: /New Tag/ }));
await userEvent.keyboard('team-payments');
},
};
/** The new key with its tags confirmed, each one removable. */
export const CreateKeyTagsAdded: Story = {
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
for (const tag of ['team-payments', 'env:production']) {
await addTag(dialog, tag);
await dialog.findByText(tag);
}
},
};
/**
* A tag typed again while the key already has it: Enter keeps the input open
* and adds nothing, with no message saying why.
*/
export const CreateKeyDuplicateTag: Story = {
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
await addTag(dialog, 'team-payments');
await addTag(dialog, 'team-payments');
await dialog.findByDisplayValue('team-payments');
},
};
/** A tag longer than its pill, cut short with an ellipsis. */
export const CreateKeyLongTag: Story = {
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
const tag =
'team-payments-platform-observability-production-us-east-1-canary-collectors';
await addTag(dialog, tag);
await dialog.findByText(tag);
},
};
/** A name with a space and no expiration, submitted: each field names its rule. */
export const CreateKeyInvalid: Story = {
// The page logs the rejected validation through `console.error`.
parameters: { allowConsoleErrors: true },
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
await userEvent.type(dialog.getByLabelText('Name'), 'otel collectors');
await userEvent.click(
dialog.getByRole('button', { name: 'Create new Ingestion key' }),
);
await dialog.findByText(/should only contain letters/);
},
};
/** The expiration calendar, where today and every day before it are disabled. */
export const CreateKeyExpirationOpen: Story = {
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
await userEvent.click(dialog.getByLabelText('Expiration'));
await screen.findByTitle(dayjs().format('YYYY-MM-DD'));
},
};
/** The new key sent, with the create still in flight. */
export const CreateKeySubmitting: Story = {
args: { create: 'hangs' },
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
await submitCreateKey(dialog);
await waitFor(() =>
expect(
dialog.getByRole('button', { name: 'Create new Ingestion key' }),
).toHaveAttribute('aria-busy', 'true'),
);
},
};
/**
* A create the gateway rejects: the notification names the conflict and the
* form keeps what was typed.
*/
export const CreateKeyFailed: Story = {
args: { create: 'fails' },
// The deliberate 409 is the state under test.
parameters: { allowConsoleErrors: true },
play: async ({ canvasElement }): Promise<void> => {
const dialog = await openCreateKey(canvasElement);
await submitCreateKey(dialog);
await screen.findByText(
'An ingestion key with this name already exists.',
undefined,
untilLoaded,
await within(canvasElement).findByText(
'New Ingestion key',
undefined,
untilLoaded,
),
);
await screen.findByText('Create new ingestion key', undefined, untilLoaded);
},
};

View File

@@ -7,7 +7,6 @@ import type {
GatewaytypesIngestionKeyDTO,
GatewaytypesLimitDTO,
GetIngestionKeys200,
RenderErrorResponseDTO,
} from 'api/generated/services/sigNoz.schemas';
import type { IngestionInfo } from 'types/api/settings/ingestion';
@@ -22,10 +21,6 @@ export const EXPIRIES = ['none', 'soon', 'expired'] as const;
export type Expiry = (typeof EXPIRIES)[number];
export const CREATE_OUTCOMES = ['succeeds', 'fails', 'hangs'] as const;
export type CreateOutcome = (typeof CREATE_OUTCOMES)[number];
const KEY_NAMES = [
'production-us-east',
'production-eu-west',
@@ -122,18 +117,6 @@ export const ingestionKeysResponse = (
},
});
export const ingestionKeyCreateError = (): RenderErrorResponseDTO => ({
status: 'error',
error: {
code: 'already_exists',
type: 'already_exists',
message: 'An ingestion key with this name already exists.',
url: '',
errors: [],
suggestions: [],
},
});
/**
* The pre-gateway endpoint, which answers with one record and no limits: the tab
* falls back to it when the gateway feature is off.

View File

@@ -5,7 +5,6 @@
import ROUTES from 'constants/routes';
import { rest } from 'msw';
import { userEvent, within } from 'storybook/test';
import { RoleType } from 'types/roles';
import { choiceControl } from '@/storybook/controls/controls';
@@ -85,18 +84,3 @@ export const roleEditorMocks = defineStoryMocks({
],
config: (values) => ({ route: routeFor(values.mode, values.editor) }),
});
/** Opens the Logs card and returns it. */
export const openLogsCard = async (
canvasElement: HTMLElement,
): Promise<HTMLElement> => {
const header = await within(canvasElement).findByRole(
'button',
{ name: /^Logs:/ },
{ timeout: 15_000 },
);
await userEvent.click(header);
return header.parentElement as HTMLElement;
};

View File

@@ -1,10 +1,10 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import { screen, userEvent, within } from 'storybook/test';
import { userEvent, within } from 'storybook/test';
import { storyMocks } from '@/storybook/controls/defineStoryMocks';
import type { PageStoryArgs } from '@/storybook/runtime/resolveStory';
import { openLogsCard, roleEditorMocks } from './RoleEditor.stories.mocks';
import { roleEditorMocks } from './RoleEditor.stories.mocks';
import SettingsPage from '../../../Settings';
@@ -75,27 +75,3 @@ export const Tooltips: Story = {
export const TooltipsInJsonEditor: Story = {
args: { tooltipsOpen: true, mode: 'edit', editor: 'json' },
};
/** The Logs card's first verb granted over everything instead of nothing. */
export const PermissionEditorScope: Story = {
play: async ({ canvasElement }): Promise<void> => {
const card = await openLogsCard(canvasElement);
const [all] = await within(card).findAllByText('All');
await userEvent.click(all);
},
};
/** The selector wizard, opened from a Logs verb scoped to named objects. */
export const TelemetrySelectorWizard: Story = {
play: async ({ canvasElement }): Promise<void> => {
const card = await openLogsCard(canvasElement);
const [onlySelected] = await within(card).findAllByText('Only selected');
await userEvent.click(onlySelected);
await userEvent.click(
await within(card).findByRole('button', { name: 'Wizard' }),
);
await screen.findByRole('dialog', { name: 'Selector Wizard' });
},
};

View File

@@ -172,31 +172,3 @@ export const TraceOptionsMenu: Story = {
await screen.findByText('Preview fields');
},
};
/** The span panel moved from the right edge to the bottom through its dock toggle. */
export const DockModeSwitched: Story = {
play: async ({ canvasElement }): Promise<void> => {
// The dock options are icon-only and carry no accessible name.
const option = await within(canvasElement).findByTestId(
'dock-mode-docked',
undefined,
untilLoaded,
);
await userEvent.click(option.querySelector('button') ?? option);
},
};
/** The span's percentile badge expanded into the distribution it was ranked in. */
export const SpanPercentileOpen: Story = {
play: async ({ canvasElement }): Promise<void> => {
await userEvent.click(
await within(canvasElement).findByRole(
'button',
{ name: /^p\d+/ },
untilLoaded,
),
);
await screen.findByText(/This span duration is/);
},
};

View File

@@ -1,10 +1,4 @@
import type { Meta, StoryObj } from '@storybook/react-vite';
import {
findSuggestion,
openKeySuggestions,
showFilterErrors,
typeFilter,
} from 'components/QueryBuilderV2/QueryV2/QuerySearch/stories/__story_mockdata__/querySearch.play';
import { ExplorerViews } from 'pages/LogsExplorer/utils';
import { expect, screen, userEvent, waitFor } from 'storybook/test';
@@ -148,41 +142,3 @@ export const QuickFiltersSettingsWithBanner: Story = {
args: { banner: 'trial-expiry' },
play: dirtyQuickFiltersSettings,
};
/** The filter focused before anything is typed: span and resource keys together. */
export const FilterKeySuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await openKeySuggestions(canvasElement, 'status_code_string');
},
};
/** A context prefix: only the keys that live on the resource. */
export const FilterResourceKeys: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'resource.');
await findSuggestion(canvasElement, 'resource.service.name');
},
};
/** A key and an operator: the services the spans came from. */
export const FilterValueSuggestions: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'service.name = ');
await findSuggestion(canvasElement, 'checkout');
},
};
/** A boolean key offers its two values. */
export const FilterBooleanValues: Story = {
play: async ({ canvasElement }): Promise<void> => {
await typeFilter(canvasElement, 'has_error = ');
await findSuggestion(canvasElement, 'false');
},
};
/** A dangling conjunction after focus left: the marker and its errors. */
export const FilterSyntaxError: Story = {
play: async ({ canvasElement }): Promise<void> => {
await showFilterErrors(canvasElement, 'has_error = true AND');
},
};

View File

@@ -67,5 +67,23 @@ func (provider *provider) addTraceDetailRoutes(router *mux.Router) error {
return err
}
if err := router.Handle("/api/v1/traces/{traceID}/thread", handler.New(
provider.authzMiddleware.ViewAccess(provider.traceDetailHandler.GetThread),
handler.OpenAPIDef{
ID: "GetTraceThread",
Tags: []string{"tracedetail"},
Summary: "Get thread view for a trace",
Description: "Returns the spans carrying gen_ai input or output messages in timestamp order, each with the messages normalised into formatted_input and formatted_output. Pages are fetched with the returned nextCursor.",
RequestQuery: new(spantypes.PostableThreadQuery),
Response: new(spantypes.GettableTraceThread),
ResponseContentType: "application/json",
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
},
)).Methods(http.MethodGet).GetError(); err != nil {
return err
}
return nil
}

View File

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

View File

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

View File

@@ -4,6 +4,7 @@ import (
"context"
"database/sql"
"fmt"
"strings"
"time"
sqlbuilder "github.com/huandu/go-sqlbuilder"
@@ -11,12 +12,23 @@ import (
"github.com/SigNoz/signoz/pkg/clickhousesql"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
const colServiceName = `resource_string_service$$$$name` // $ gets escaped so $$$$ converts to $$.
var fullSpanColumns = []string{
"duration_nano", "span_id", "has_error", "kind",
colServiceName, "name",
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
"flags", "is_remote", "trace_state", "status_code",
"db_name", "db_operation", "http_method", "http_url", "http_host",
"external_http_method", "external_http_url", "response_status_code", "links as references",
}
func buildFieldExpr(fieldKey telemetrytypes.TelemetryFieldKey) (string, error) {
switch fieldKey.FieldContext {
case telemetrytypes.FieldContextResource:
@@ -68,18 +80,11 @@ func (s *traceStore) GetTraceSummary(ctx context.Context, traceID string) (*span
func (s *traceStore) GetTraceSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary) ([]spantypes.StorableSpan, error) {
// DISTINCT ON (span_id) is ClickHouse-specific syntax not supported by sqlbuilder
query := fmt.Sprintf(`
SELECT DISTINCT ON (span_id)
timestamp, duration_nano, span_id, has_error, kind,
resource_string_service$$name, name,
attributes_string, attributes_number, attributes_bool, resources_string,
events, status_message, status_code_string, kind_string, parent_span_id,
flags, is_remote, trace_state, status_code,
db_name, db_operation, http_method, http_url, http_host,
external_http_method, external_http_url, response_status_code, links as references
SELECT DISTINCT ON (span_id) timestamp, %s
FROM %s.%s
WHERE trace_id=? AND ts_bucket_start>=? AND ts_bucket_start<=?
ORDER BY timestamp ASC, name ASC`,
spantypes.TraceDB, spantypes.TraceTable,
strings.Join(fullSpanColumns, ", "), spantypes.TraceDB, spantypes.TraceTable,
)
var spanItems []spantypes.StorableSpan
err := s.telemetryStore.ClickhouseDB().Select(
@@ -123,16 +128,8 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return []spantypes.StorableSpan{}, nil
}
sb := sqlbuilder.NewSelectBuilder()
sb.Select(
"DISTINCT ON (span_id) timestamp",
"duration_nano", "span_id", "has_error", "kind",
colServiceName, "name",
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
"flags", "is_remote", "trace_state", "status_code",
"db_name", "db_operation", "http_method", "http_url", "http_host",
"external_http_method", "external_http_url", "response_status_code", "links as references",
)
sb.Select("DISTINCT ON (span_id) timestamp")
sb.SelectMore(fullSpanColumns...)
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
ids := make([]any, len(spanIDs))
for i, id := range spanIDs {
@@ -155,6 +152,36 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
return spans, nil
}
func (s *traceStore) GetThreadSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary, cursor *spantypes.ThreadCursor, limit int) ([]spantypes.StorableSpan, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("DISTINCT ON (span_id) timestamp")
sb.SelectMore(fullSpanColumns...)
sb.SelectMore("attributes")
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
sb.Where(
sb.E("trace_id", traceID),
sb.GE("ts_bucket_start", summary.Start.Unix()-1800),
sb.LE("ts_bucket_start", summary.End.Unix()),
sb.Or(
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIInputMessages))),
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIOutputMessages))),
),
)
if cursor != nil {
sb.Where(sb.GT("(toUnixTimestamp64Nano(timestamp), span_id)", sqlbuilder.Tuple(cursor.TimeUnixNano, cursor.SpanID)))
}
sb.OrderByAsc("timestamp")
sb.OrderByAsc("span_id")
sb.Limit(limit)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
var spans []spantypes.StorableSpan
if err := s.telemetryStore.ClickhouseDB().Select(ctx, &spans, query, args...); err != nil {
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error querying thread spans")
}
return spans, nil
}
func (s *traceStore) GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]spantypes.StorableSpan, error) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select(

View File

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

View File

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

View File

@@ -0,0 +1,102 @@
package aiobservabilitytypes
import "strings"
const (
MessageRoleSystem MessageRole = "system"
MessageRoleUser MessageRole = "user"
MessageRoleAssistant MessageRole = "assistant"
MessageRoleTool MessageRole = "tool"
)
const (
FinishReasonStop FinishReason = "stop"
FinishReasonToolCall FinishReason = "tool_call"
FinishReasonLength FinishReason = "length"
FinishReasonContentFilter FinishReason = "content_filter"
FinishReasonError FinishReason = "error"
)
const (
PartTypeText PartType = "text"
PartTypeThinking PartType = "thinking"
PartTypeToolCall PartType = "tool_call"
PartTypeToolResult PartType = "tool_result"
PartTypeGeneric PartType = "generic"
)
type MessageRole string
type FinishReason string
type PartType string
// Part is one piece of a message. Which fields are set depends on Type:
//
// text Content
// thinking Content, Redacted
// tool_call ID, Name, Arguments, Server
// tool_result ToolCallID, Name, Content, IsError, Server
// generic Content (the original value, always a string)
type Part struct {
Type PartType `json:"type" required:"true"`
Content string `json:"content,omitempty"`
Redacted bool `json:"redacted,omitempty"`
ID string `json:"id,omitempty"`
Name string `json:"name,omitempty"`
Arguments any `json:"arguments,omitempty"`
Server bool `json:"server,omitempty"`
ToolCallID string `json:"toolCallId,omitempty"`
IsError bool `json:"isError,omitempty"`
}
type Message struct {
Role MessageRole `json:"role,omitempty"`
Content []Part `json:"content" required:"true" nullable:"false"`
FinishReason FinishReason `json:"finishReason,omitempty"`
}
func (PartType) Enum() []any {
return []any{PartTypeText, PartTypeThinking, PartTypeToolCall, PartTypeToolResult, PartTypeGeneric}
}
// normalizeRole keeps an unknown role, lowercased.
func normalizeRole(role string) MessageRole {
if known := knownRole(role); known != "" {
return known
}
return MessageRole(strings.ToLower(strings.TrimSpace(role)))
}
// normalizeFinishReason keeps an unknown reason, lowercased.
func normalizeFinishReason(reason string) FinishReason {
lowered := strings.ToLower(strings.TrimSpace(reason))
switch lowered {
case "stop", "end_turn", "stop_sequence", "completed", "complete", "eos", "finished":
return FinishReasonStop
case "tool_call", "tool_calls", "tool_use", "function_call":
return FinishReasonToolCall
case "length", "max_tokens", "max_output_tokens", "max_completion_tokens", "model_length":
return FinishReasonLength
case "content_filter", "content_filtered", "guardrail_intervened", "safety", "refusal", "recitation", "blocklist", "prohibited_content", "spii":
return FinishReasonContentFilter
case "error", "failed", "incomplete":
return FinishReasonError
}
return FinishReason(lowered)
}
// knownRole maps vendor role names onto MessageRole; anything else is "".
func knownRole(role string) MessageRole {
switch strings.ToLower(strings.TrimSpace(role)) {
case "system", "developer":
return MessageRoleSystem
case "user", "human":
return MessageRoleUser
case "assistant", "ai", "model":
return MessageRoleAssistant
case "tool", "function":
return MessageRoleTool
}
return ""
}

View File

@@ -0,0 +1,971 @@
package aiobservabilitytypes
import (
"encoding/json"
"strings"
)
// Ordered by specificity: earlier converters never match a later format.
var converters = []converter{
convertSemconvMessages,
convertChatMessageList,
convertToolCallList,
convertContentBlockList,
convertChatRequest,
convertChatResponse,
convertResponsesAPIResponse,
convertGeminiResponse,
convertCompletionObject,
convertLangChainGenerations,
convertSingleMessage,
}
var finishReasonKeys = []string{"finish_reason", "finishReason", "stop_reason", "stopReason", "done_reason"}
// NormalizeMessages converts a gen_ai.*.messages value, a JSON string or a
// decoded value; unknown formats become one generic part holding the original.
func NormalizeMessages(raw any) []Message {
var (
value any
original string
)
switch v := raw.(type) {
case nil:
return []Message{}
case string:
original = v
if err := json.Unmarshal([]byte(v), &value); err != nil {
return genericMessages(original)
}
default:
value = v
original = stringOf(v)
}
if list, ok := value.([]any); ok {
if len(list) == 0 {
return []Message{}
}
// [[...]]: some SDKs wrap the conversation in one more list
if _, nested := list[0].([]any); nested {
value = flattenOnce(list)
}
// ["{...}", "{...}"]: an array attribute holding one JSON message per element
if decoded, ok := decodeStringList(list); ok {
value = decoded
}
}
for _, convert := range converters {
if messages, ok := convert(value); ok {
return messages
}
}
return genericMessages(original)
}
type converter func(value any) (messages []Message, ok bool)
func genericMessages(content string) []Message {
return []Message{{Content: []Part{genericPart(content)}}}
}
// convertSemconvMessages handles [{role, parts, finish_reason}] and Gemini contents.
func convertSemconvMessages(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
first, ok := list[0].(map[string]any)
if !ok {
return nil, false
}
if _, ok := first["parts"]; !ok {
return nil, false
}
messages := make([]Message, 0, len(list))
for _, item := range list {
m, ok := item.(map[string]any)
if !ok {
messages = append(messages, genericMessages(stringOf(item))[0])
continue
}
messages = append(messages, partsMessage(m, "")...)
}
return messages, true
}
// partsMessage falls back to chatMessage when m has no parts.
func partsMessage(m map[string]any, defaultRole MessageRole) []Message {
parts, ok := m["parts"].([]any)
if !ok {
return chatMessage(m, defaultRole)
}
role := normalizeRole(stringOf(m["role"]))
if role == "" {
role = defaultRole
}
msg := Message{
Role: role,
Content: []Part{},
FinishReason: normalizeFinishReason(finishReasonOf(m)),
}
for _, p := range parts {
msg.Content = append(msg.Content, semconvPart(p))
}
return []Message{msg}
}
func semconvPart(value any) Part {
p, ok := value.(map[string]any)
if !ok {
if s, ok := value.(string); ok {
return textPart(s)
}
return genericPart(value)
}
switch stringOf(p["type"]) {
case "text":
if boolOf(p["thought"]) {
return Part{Type: PartTypeThinking, Content: stringOf(firstOf(p, "content", "text"))}
}
return Part{Type: PartTypeText, Content: stringOf(firstOf(p, "content", "text"))}
case "reasoning", "thinking":
return Part{Type: PartTypeThinking, Content: stringOf(firstOf(p, "content", "thinking", "text"))}
case "redacted_thinking", "redacted_reasoning":
return Part{Type: PartTypeThinking, Redacted: true}
case "tool_call":
return Part{
Type: PartTypeToolCall,
ID: idOf(p["id"]),
Name: stringOf(p["name"]),
Arguments: parseArguments(firstOf(p, "arguments", "args", "input")),
Server: boolOf(p["server"]),
}
case "tool_call_response":
return Part{
Type: PartTypeToolResult,
ToolCallID: idOf(p["id"]),
Name: stringOf(p["name"]),
Content: stringOf(firstOf(p, "response", "result", "content", "output")),
IsError: boolOf(firstOf(p, "is_error", "isError")),
Server: boolOf(p["server"]),
}
case "":
// Gemini parts carry no type; the field name is the type.
if text, ok := p["text"]; ok {
if boolOf(p["thought"]) {
return Part{Type: PartTypeThinking, Content: stringOf(text)}
}
return Part{Type: PartTypeText, Content: stringOf(text)}
}
if call, ok := firstOf(p, "functionCall", "function_call").(map[string]any); ok {
return Part{Type: PartTypeToolCall, ID: stringOf(call["id"]), Name: stringOf(call["name"]), Arguments: parseArguments(call["args"])}
}
if resp, ok := firstOf(p, "functionResponse", "function_response").(map[string]any); ok {
return Part{Type: PartTypeToolResult, ToolCallID: stringOf(resp["id"]), Name: stringOf(resp["name"]), Content: stringOf(resp["response"])}
}
}
return genericPart(p)
}
// convertChatMessageList handles OpenAI, Anthropic, Bedrock, Vercel and LangChain message lists.
func convertChatMessageList(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
if !isChatMessage(list[0]) {
return nil, false
}
messages := make([]Message, 0, len(list))
for _, item := range list {
messages = append(messages, chatMessage(item, "")...)
}
return messages, true
}
func isChatMessage(value any) bool {
m, ok := value.(map[string]any)
if !ok {
return false
}
if _, ok := m["role"]; ok {
return true
}
if _, ok := m["gen_ai.event.content"]; ok {
return true
}
switch typ := stringOf(m["type"]); typ {
case "message", "reasoning", "human", "ai", "tool", "system":
return true
case "constructor":
_, ok := m["kwargs"]
return ok
default:
return isResponsesItemType(typ)
}
}
// chatMessage returns nil for LangGraph tool definitions.
func chatMessage(value any, defaultRole MessageRole) []Message {
m, ok := value.(map[string]any)
if !ok {
return genericMessages(stringOf(value))
}
if inner, role, ok := unwrapChatEnvelope(m, defaultRole); ok {
return chatMessage(inner, role)
}
if isLangGraphToolDefinition(m) {
return nil
}
typ := stringOf(m["type"])
if isResponsesItemType(typ) {
return responsesItem(m, typ)
}
if typ == "reasoning" {
return []Message{{Role: MessageRoleAssistant, Content: reasoningParts(m)}}
}
role := normalizeRole(stringOf(m["role"]))
if role == "" {
role = knownRole(typ)
}
if role == "" {
role = defaultRole
}
msg := Message{
Role: role,
Content: append(chatContentParts(m, role), chatToolCallParts(m)...),
FinishReason: normalizeFinishReason(finishReasonOf(m)),
}
if refusal := stringOf(m["refusal"]); refusal != "" {
msg.Content = append(msg.Content, textPart(refusal))
}
return []Message{msg}
}
// unwrapChatEnvelope unwraps LangChain serialised messages and Semantic Kernel events.
func unwrapChatEnvelope(m map[string]any, defaultRole MessageRole) (map[string]any, MessageRole, bool) {
if kwargs, ok := m["kwargs"].(map[string]any); ok && stringOf(m["type"]) == "constructor" {
role := langChainRole(m["id"])
if role == "" {
role = knownRole(stringOf(kwargs["type"]))
}
return kwargs, role, true
}
event, ok := m["gen_ai.event.content"].(string)
if !ok {
return nil, "", false
}
var inner map[string]any
if err := json.Unmarshal([]byte(event), &inner); err != nil || inner == nil {
return nil, "", false
}
if message, ok := inner["message"].(map[string]any); ok {
if _, has := message["finish_reason"]; !has {
message["finish_reason"] = inner["finish_reason"]
}
inner = message
}
return inner, defaultRole, true
}
// chatContentParts turns a tool message's text into its result.
func chatContentParts(m map[string]any, role MessageRole) []Part {
parts := []Part{}
switch content := m["content"].(type) {
case nil:
case string:
if role == MessageRoleTool {
parts = append(parts, toolResultOf(m, content))
} else if content != "" {
parts = append(parts, textPart(content))
}
case []any:
for _, item := range content {
part := chatContentPart(item)
if role == MessageRoleTool && part.Type == PartTypeText {
part = toolResultOf(m, part.Content)
}
parts = append(parts, part)
}
case map[string]any:
if contentParts, ok := content["parts"].([]any); ok {
for _, p := range contentParts {
parts = append(parts, semconvPart(p))
}
} else if role == MessageRoleTool {
parts = append(parts, toolResultOf(m, stringOf(content)))
} else {
parts = append(parts, genericPart(content))
}
default:
parts = append(parts, genericPart(content))
}
return parts
}
func chatToolCallParts(m map[string]any) []Part {
calls, _ := firstOf(m, "tool_calls", "toolCalls").([]any)
if kwargs, ok := m["additional_kwargs"].(map[string]any); ok && len(calls) == 0 {
calls, _ = kwargs["tool_calls"].([]any)
}
parts := make([]Part, 0, len(calls)+1)
for _, call := range calls {
parts = append(parts, toolCallPart(call))
}
if call, ok := m["function_call"].(map[string]any); ok {
parts = append(parts, Part{Type: PartTypeToolCall, Name: stringOf(call["name"]), Arguments: parseArguments(call["arguments"])})
}
return parts
}
func toolResultOf(m map[string]any, content string) Part {
return Part{Type: PartTypeToolResult, ToolCallID: stringOf(firstOf(m, "tool_call_id", "toolCallId")), Name: stringOf(m["name"]), Content: content}
}
// isLangGraphToolDefinition matches {role: "tool", content: {type: "function"}} without tool_call_id.
func isLangGraphToolDefinition(m map[string]any) bool {
if normalizeRole(stringOf(m["role"])) != MessageRoleTool {
return false
}
if _, has := m["tool_call_id"]; has {
return false
}
content, ok := m["content"].(map[string]any)
if !ok || stringOf(content["type"]) != "function" {
return false
}
_, ok = content["function"]
return ok
}
func chatContentPart(value any) Part {
p, ok := value.(map[string]any)
if !ok {
if s, ok := value.(string); ok {
return textPart(s)
}
return genericPart(value)
}
typ := stringOf(p["type"])
switch typ {
case "text", "input_text", "output_text", "refusal", "summary_text":
return Part{Type: PartTypeText, Content: stringOf(firstOf(p, "text", "content", "refusal"))}
case "thinking", "reasoning":
return Part{Type: PartTypeThinking, Content: stringOf(firstOf(p, "thinking", "text", "content", "reasoning"))}
case "redacted_thinking":
return Part{Type: PartTypeThinking, Redacted: true}
case "tool-call", "tool_use", "tool_call", "function_call":
return toolCallPart(p)
case "tool-result", "tool_result", "function_call_output":
return toolResultPart(p, firstOf(p, "result", "output", "content"), false)
case "server_tool_use", "mcp_tool_use":
part := toolCallPart(p)
part.Server = true
return part
case "":
if text, ok := p["text"]; ok {
return Part{Type: PartTypeText, Content: stringOf(text)}
}
// Bedrock Converse blocks are typed by field name.
if use, ok := p["toolUse"].(map[string]any); ok {
return Part{Type: PartTypeToolCall, ID: stringOf(use["toolUseId"]), Name: stringOf(use["name"]), Arguments: parseArguments(use["input"])}
}
if result, ok := p["toolResult"].(map[string]any); ok {
part := toolResultPart(result, result["content"], false)
part.ToolCallID = stringOf(result["toolUseId"])
part.IsError = stringOf(result["status"]) == "error"
return part
}
default:
if strings.HasSuffix(typ, "_tool_result") {
return toolResultPart(p, p["content"], true)
}
}
return genericPart(p)
}
// toolCallPart reads the OpenAI, flat, Anthropic and Vercel tool call shapes.
func toolCallPart(value any) Part {
call, ok := value.(map[string]any)
if !ok {
return genericPart(value)
}
part := Part{
Type: PartTypeToolCall,
ID: idOf(firstOf(call, "toolCallId", "call_id", "id")),
Name: stringOf(firstOf(call, "toolName", "name")),
}
if fn, ok := call["function"].(map[string]any); ok {
part.Name = stringOf(fn["name"])
part.Arguments = parseArguments(fn["arguments"])
return part
}
part.Arguments = parseArguments(firstOf(call, "arguments", "args", "input"))
return part
}
// toolResultPart unwraps the Vercel {type, value} result wrapper.
func toolResultPart(p map[string]any, result any, server bool) Part {
if nested, ok := result.(map[string]any); ok && len(nested) <= 2 {
if v, ok := nested["value"]; ok {
result = v
}
}
return Part{
Type: PartTypeToolResult,
ToolCallID: idOf(firstOf(p, "toolCallId", "tool_use_id", "tool_call_id", "call_id", "id")),
Name: stringOf(firstOf(p, "toolName", "name")),
Content: textOf(result),
IsError: boolOf(firstOf(p, "isError", "is_error")),
Server: server,
}
}
// textOf joins a list of text blocks; anything else goes through stringOf.
func textOf(result any) string {
list, ok := result.([]any)
if !ok || len(list) == 0 {
return stringOf(result)
}
texts := make([]string, 0, len(list))
for _, item := range list {
block := asMap(item)
text, ok := block["text"].(string)
if !ok || (len(block) == 2 && stringOf(block["type"]) != "text") || len(block) > 2 {
return stringOf(result)
}
texts = append(texts, text)
}
return strings.Join(texts, "\n")
}
// isResponsesItemType matches role-less Responses API tool and MCP items.
func isResponsesItemType(typ string) bool {
switch typ {
case "":
return false
case "function_call", "function_call_output", "tool_call", "custom_tool_call", "custom_tool_call_output",
"mcp_call", "mcp_list_tools", "mcp_approval_request", "mcp_approval_response":
return true
}
return strings.HasSuffix(typ, "_call") || strings.HasSuffix(typ, "_call_output")
}
// responsesItem maps built-in tools to server tool parts.
func responsesItem(m map[string]any, typ string) []Message {
switch typ {
case "function_call", "tool_call", "custom_tool_call":
return []Message{{Role: MessageRoleAssistant, Content: []Part{{
Type: PartTypeToolCall,
ID: stringOf(firstOf(m, "call_id", "id")),
Name: stringOf(m["name"]),
Arguments: parseArguments(firstOf(m, "arguments", "args", "input")),
}}}}
case "function_call_output", "custom_tool_call_output":
return []Message{{Role: MessageRoleTool, Content: []Part{{
Type: PartTypeToolResult,
ToolCallID: stringOf(firstOf(m, "call_id", "id")),
Content: stringOf(firstOf(m, "output", "result")),
}}}}
}
id := stringOf(firstOf(m, "call_id", "id"))
if strings.HasSuffix(typ, "_output") || typ == "mcp_approval_response" {
return []Message{{Role: MessageRoleTool, Content: []Part{{
Type: PartTypeToolResult,
ToolCallID: id,
Name: strings.TrimSuffix(typ, "_output"),
Content: stringOf(firstOf(m, "output", "result", "results")),
Server: true,
}}}}
}
args := make(map[string]any, len(m))
for k, v := range m {
switch k {
case "type", "id", "call_id", "status", "name", "server_label", "output", "result", "results":
default:
args[k] = v
}
}
name := stringOf(firstOf(m, "name", "server_label"))
if name == "" {
name = typ
}
call := Part{Type: PartTypeToolCall, ID: id, Name: name, Server: true}
if len(args) > 0 {
call.Arguments = args
}
msg := Message{Role: MessageRoleAssistant, Content: []Part{call}}
if result := firstOf(m, "output", "result", "results"); result != nil {
msg.Content = append(msg.Content, Part{Type: PartTypeToolResult, ToolCallID: id, Name: name, Content: stringOf(result), Server: true})
}
return []Message{msg}
}
// reasoningParts marks encrypted reasoning without a summary as redacted.
func reasoningParts(m map[string]any) []Part {
parts := []Part{}
if summary, ok := m["summary"].([]any); ok {
for _, s := range summary {
parts = append(parts, Part{Type: PartTypeThinking, Content: stringOf(firstOf(asMap(s), "text", "content"))})
}
}
if content, ok := m["content"].([]any); ok {
for _, c := range content {
parts = append(parts, Part{Type: PartTypeThinking, Content: stringOf(firstOf(asMap(c), "text", "content"))})
}
}
if len(parts) == 0 {
parts = append(parts, Part{Type: PartTypeThinking, Redacted: true})
}
return parts
}
// convertToolCallList handles a bare tool call list, e.g. Vercel ai.response.toolCalls.
func convertToolCallList(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
msg := Message{Role: MessageRoleAssistant, Content: make([]Part, 0, len(list))}
for _, item := range list {
call, ok := item.(map[string]any)
if !ok || !isToolCall(call) {
return nil, false
}
msg.Content = append(msg.Content, toolCallPart(call))
}
return []Message{msg}, true
}
// isToolCall rejects tool definitions, which carry no arguments.
func isToolCall(call map[string]any) bool {
if _, has := call["toolName"]; has {
return true
}
if fn, ok := call["function"].(map[string]any); ok {
_, has := fn["arguments"]
return has
}
if _, has := call["name"]; !has {
return false
}
_, has := lookup(call, "arguments", "args")
return has
}
// convertContentBlockList handles a bare content block list; the role is unknown.
func convertContentBlockList(value any) ([]Message, bool) {
list, ok := value.([]any)
if !ok || len(list) == 0 {
return nil, false
}
msg := Message{Content: make([]Part, 0, len(list))}
for _, item := range list {
block, ok := item.(map[string]any)
if !ok {
return nil, false
}
if _, has := block["type"].(string); !has {
return nil, false
}
part := chatContentPart(block)
if part.Type == PartTypeGeneric {
return nil, false
}
msg.Content = append(msg.Content, part)
}
return []Message{msg}, true
}
// convertChatRequest handles OpenAI, Anthropic, Vercel, Gemini and LangChain request objects.
func convertChatRequest(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
conversation, ok := lookup(m, "messages", "input", "contents", "prompt")
if !ok || !isConversation(m, conversation) {
return nil, false
}
messages := []Message{}
if system := systemMessage(firstOf(m, "system", "instructions", "system_instruction", "systemInstruction", "system_prompt")); system != nil {
messages = append(messages, *system)
} else if config, ok := m["config"].(map[string]any); ok {
if system := systemMessage(firstOf(config, "system_instruction", "systemInstruction")); system != nil {
messages = append(messages, *system)
}
}
// {messages: "[...]"}: the list arrives JSON-encoded once more from some SDKs
if s, isString := conversation.(string); isString {
var decoded any
if err := json.Unmarshal([]byte(s), &decoded); err == nil {
if _, isList := decoded.([]any); isList {
conversation = decoded
}
}
}
switch c := conversation.(type) {
case string:
messages = append(messages, textMessage(MessageRoleUser, c))
case []any:
for _, item := range flattenOnce(c) {
if s, isString := item.(string); isString {
messages = append(messages, textMessage(MessageRoleUser, s))
continue
}
messages = append(messages, partsMessage(asMap(item), MessageRoleUser)...)
}
case map[string]any:
messages = append(messages, partsMessage(c, MessageRoleUser)...)
default:
return nil, false
}
return messages, true
}
// isConversation rejects embeddings requests.
func isConversation(m map[string]any, conversation any) bool {
if _, isRequestInput := m["input"]; !isRequestInput {
return true
}
if _, hasChatKey := lookup(m, "instructions", "tools", "tool_choice", "parallel_tool_calls", "previous_response_id"); hasChatKey {
return true
}
list, ok := conversation.([]any)
if !ok {
return false
}
for _, item := range list {
if _, isMap := item.(map[string]any); !isMap {
return false
}
}
return true
}
// systemMessage returns nil when value is empty or unknown.
func systemMessage(value any) *Message {
msg := Message{Role: MessageRoleSystem, Content: []Part{}}
switch v := value.(type) {
case string:
if v == "" {
return nil
}
msg.Content = append(msg.Content, textPart(v))
case []any:
for _, item := range v {
msg.Content = append(msg.Content, chatContentPart(item))
}
case map[string]any:
parts, ok := v["parts"].([]any)
if !ok {
return nil
}
for _, p := range parts {
msg.Content = append(msg.Content, semconvPart(p))
}
default:
return nil
}
if len(msg.Content) == 0 {
return nil
}
return &msg
}
// convertChatResponse handles {choices}, {message} and {output: {message}} responses.
func convertChatResponse(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
wrapped, _ := m["message"].(map[string]any)
if output, ok := m["output"].(map[string]any); ok && wrapped == nil {
wrapped, _ = output["message"].(map[string]any)
}
if wrapped != nil {
if !isChatMessage(wrapped) {
return nil, false
}
return withFinishReason(chatMessage(wrapped, MessageRoleAssistant), finishReasonOf(m)), true
}
choices, ok := m["choices"].([]any)
if !ok {
return nil, false
}
messages := make([]Message, 0, len(choices))
for _, c := range choices {
choice := asMap(c)
var converted []Message
if message, ok := firstOf(choice, "message", "delta").(map[string]any); ok {
converted = chatMessage(message, MessageRoleAssistant)
} else if text, ok := choice["text"]; ok {
converted = []Message{textMessage(MessageRoleAssistant, stringOf(text))}
} else {
converted = []Message{{Role: MessageRoleAssistant, Content: []Part{genericPart(choice)}}}
}
messages = append(messages, withFinishReason(converted, stringOf(choice["finish_reason"]))...)
}
return messages, true
}
func convertResponsesAPIResponse(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
output, ok := m["output"].([]any)
if !ok {
return nil, false
}
messages := make([]Message, 0, len(output))
for _, item := range output {
messages = append(messages, chatMessage(item, MessageRoleAssistant)...)
}
if len(messages) == 0 {
return messages, true
}
last := &messages[len(messages)-1]
if last.FinishReason == "" {
if details, ok := m["incomplete_details"].(map[string]any); ok {
last.FinishReason = normalizeFinishReason(stringOf(details["reason"]))
} else if stringOf(m["status"]) == "completed" {
last.FinishReason = FinishReasonStop
}
}
return messages, true
}
// convertGeminiResponse handles Gemini {candidates} and Google ADK {content} responses.
func convertGeminiResponse(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
if content, ok := m["content"].(map[string]any); ok {
if _, hasParts := content["parts"]; hasParts {
return withFinishReason(partsMessage(content, MessageRoleAssistant), finishReasonOf(m)), true
}
}
candidates, ok := m["candidates"].([]any)
if !ok {
return nil, false
}
messages := make([]Message, 0, len(candidates))
for _, c := range candidates {
candidate := asMap(c)
content, ok := candidate["content"].(map[string]any)
if !ok {
messages = append(messages, Message{Role: MessageRoleAssistant, Content: []Part{genericPart(candidate)}})
continue
}
messages = append(messages, withFinishReason(partsMessage(content, MessageRoleAssistant), finishReasonOf(candidate))...)
}
return messages, true
}
func convertCompletionObject(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
completion, ok := m["completion"].(string)
if !ok {
return nil, false
}
msg := Message{Role: MessageRoleAssistant, Content: []Part{}}
if reasoning, ok := m["reasoning"].(string); ok && reasoning != "" {
msg.Content = append(msg.Content, Part{Type: PartTypeThinking, Content: reasoning})
}
msg.Content = append(msg.Content, textPart(completion))
return []Message{msg}, true
}
func convertLangChainGenerations(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
generations, ok := m["generations"].([]any)
if !ok {
return nil, false
}
messages := []Message{}
for _, g := range flattenOnce(generations) {
gen := asMap(g)
var converted []Message
if message, ok := gen["message"].(map[string]any); ok {
converted = chatMessage(message, MessageRoleAssistant)
} else {
converted = []Message{textMessage(MessageRoleAssistant, stringOf(gen["text"]))}
}
messages = append(messages, withFinishReason(converted, stringOf(asMap(gen["generation_info"])["finish_reason"]))...)
}
return messages, true
}
func convertSingleMessage(value any) ([]Message, bool) {
m, ok := value.(map[string]any)
if !ok {
return nil, false
}
if _, ok := m["parts"]; ok {
return partsMessage(m, ""), true
}
if isChatMessage(m) {
return chatMessage(m, ""), true
}
return nil, false
}
// withFinishReason sets reason on the last message that has none.
func withFinishReason(messages []Message, reason string) []Message {
if len(messages) == 0 {
return messages
}
last := &messages[len(messages)-1]
if last.FinishReason == "" {
last.FinishReason = normalizeFinishReason(reason)
}
return messages
}
func langChainRole(id any) MessageRole {
path, ok := id.([]any)
if !ok || len(path) == 0 {
return ""
}
switch class := stringOf(path[len(path)-1]); {
case strings.HasPrefix(class, "System"):
return MessageRoleSystem
case strings.HasPrefix(class, "Human"):
return MessageRoleUser
case strings.HasPrefix(class, "AI"):
return MessageRoleAssistant
case strings.HasPrefix(class, "Tool"), strings.HasPrefix(class, "Function"):
return MessageRoleTool
}
return ""
}
// parseArguments decodes JSON-encoded arguments; anything else is returned as is.
func parseArguments(value any) any {
s, ok := value.(string)
if !ok {
return value
}
var decoded any
if err := json.Unmarshal([]byte(s), &decoded); err != nil {
return s
}
return decoded
}
// idOf picks the call_ entry, else the last, from list ids such as ["run_id", "call_id"].
func idOf(value any) string {
list, ok := value.([]any)
if !ok {
return stringOf(value)
}
if len(list) == 0 {
return ""
}
for _, item := range list {
if s, ok := item.(string); ok && strings.HasPrefix(s, "call_") {
return s
}
}
return stringOf(list[len(list)-1])
}
func decodeStringList(list []any) ([]any, bool) {
out := make([]any, 0, len(list))
for _, item := range list {
s, ok := item.(string)
if !ok {
return nil, false
}
var decoded map[string]any
if err := json.Unmarshal([]byte(s), &decoded); err != nil || decoded == nil {
return nil, false
}
out = append(out, decoded)
}
return out, true
}
func flattenOnce(list []any) []any {
out := make([]any, 0, len(list))
for _, item := range list {
if inner, ok := item.([]any); ok {
out = append(out, inner...)
continue
}
out = append(out, item)
}
return out
}
func lookup(m map[string]any, keys ...string) (any, bool) {
for _, k := range keys {
if v, ok := m[k]; ok && v != nil {
return v, true
}
}
return nil, false
}
func firstOf(m map[string]any, keys ...string) any {
v, _ := lookup(m, keys...)
return v
}
func asMap(value any) map[string]any {
m, _ := value.(map[string]any)
return m
}
func boolOf(value any) bool {
b, _ := value.(bool)
return b
}
// stringOf renders nil as "" and non-strings as compact JSON.
func stringOf(value any) string {
switch v := value.(type) {
case nil:
return ""
case string:
return v
}
data, err := json.Marshal(value)
if err != nil {
return ""
}
return string(data)
}
func finishReasonOf(m map[string]any) string {
return stringOf(firstOf(m, finishReasonKeys...))
}
func textPart(content string) Part {
return Part{Type: PartTypeText, Content: content}
}
func genericPart(value any) Part {
return Part{Type: PartTypeGeneric, Content: stringOf(value)}
}
func textMessage(role MessageRole, content string) Message {
return Message{Role: role, Content: []Part{textPart(content)}}
}

View File

@@ -0,0 +1,387 @@
package aiobservabilitytypes
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestNormalizeMessages(t *testing.T) {
text := func(role MessageRole, content string) Message {
return Message{Role: role, Content: []Part{{Type: PartTypeText, Content: content}}}
}
testCases := []struct {
name string
raw any
want []Message
}{
{
name: "SemconvInput_Litellm",
raw: `[{"role": "system", "parts": [{"type": "text", "content": "You are a concise assistant."}]}, {"role": "user", "parts": [{"type": "text", "content": "Give me a one-line definition of observability."}]}]`,
want: []Message{
text(MessageRoleSystem, "You are a concise assistant."),
text(MessageRoleUser, "Give me a one-line definition of observability."),
},
},
{
name: "SemconvOutput_FinishReason_Bifrost",
raw: `[{"role": "assistant", "parts": [{"content": "Observability is X.", "type": "text"}], "finish_reason": "stop"}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Observability is X."}}, FinishReason: FinishReasonStop}},
},
{
name: "SemconvToolCallAndResponse_Langchain",
raw: `[{"role": "user", "parts": [{"type": "text", "content": "What's the weather in Bengaluru?"}]}, {"role": "assistant", "parts": [{"type": "tool_call", "id": "call_1", "name": "get_weather", "arguments": {"city": "Bengaluru"}}]}, {"role": "tool", "parts": [{"type": "tool_call_response", "id": "call_1", "response": "{\"city\": \"Bengaluru\", \"temp_c\": 18}"}]}]`,
want: []Message{
text(MessageRoleUser, "What's the weather in Bengaluru?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_1", Name: "get_weather", Arguments: map[string]any{"city": "Bengaluru"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_1", Content: `{"city": "Bengaluru", "temp_c": 18}`}}},
},
},
{
name: "SemconvOutput_TwoToolCalls_Openllmetry",
raw: `[{"role": "assistant", "parts": [{"type": "tool_call", "name": "search_web", "id": "call_a", "arguments": {"query": "SigNoz"}}, {"type": "tool_call", "name": "get_weather", "id": "call_b", "arguments": {"city": "Bengaluru"}}], "finish_reason": "tool_call"}]`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonToolCall, Content: []Part{
{Type: PartTypeToolCall, ID: "call_a", Name: "search_web", Arguments: map[string]any{"query": "SigNoz"}},
{Type: PartTypeToolCall, ID: "call_b", Name: "get_weather", Arguments: map[string]any{"city": "Bengaluru"}},
}}},
},
{
name: "OpenAIChatList_FlattenedToolCalls_BifrostGateway",
raw: `[{"role":"user","content":"What's the weather in Bengaluru? Use the tool."},{"role":"assistant","content":"","tool_calls":[{"id":"call_y","type":"function","name":"get_current_weather","args":"{\"city\":\"Bengaluru\"}"}]},{"role":"tool","content":"{\"city\": \"Bengaluru\", \"temp_c\": 28}"}]`,
want: []Message{
text(MessageRoleUser, "What's the weather in Bengaluru? Use the tool."),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_y", Name: "get_current_weather", Arguments: map[string]any{"city": "Bengaluru"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, Content: `{"city": "Bengaluru", "temp_c": 28}`}}},
},
},
{
name: "OpenAIChatRequest_NestedFunctionToolCalls_OpenrouterGateway",
raw: `{"messages":[{"role":"user","content":"Weather?"},{"content":null,"refusal":null,"role":"assistant","tool_calls":[{"id":"call_A","function":{"arguments":"{\"city\":\"Bengaluru\"}","name":"get_current_weather"},"type":"function","index":0}]},{"role":"tool","tool_call_id":"call_A","content":"{\"temp_c\": 28}"}]}`,
want: []Message{
text(MessageRoleUser, "Weather?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_A", Name: "get_current_weather", Arguments: map[string]any{"city": "Bengaluru"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_A", Content: `{"temp_c": 28}`}}},
},
},
{
name: "OpenAIChatRequest_IgnoresModelAndTools_Openinference",
raw: `{"messages": [{"role": "user", "content": "What is the weather in Bengaluru in celsius?"}], "model": "gpt-4o-mini", "tool_choice": "auto", "tools": [{"type": "function", "function": {"name": "get_current_weather"}}]}`,
want: []Message{text(MessageRoleUser, "What is the weather in Bengaluru in celsius?")},
},
{
name: "OpenAIChatResponse_Openinference",
raw: `{"id":"chatcmpl-1","choices":[{"finish_reason":"stop","index":0,"logprobs":null,"message":{"content":"Hello! How are you today?","refusal":null,"role":"assistant","annotations":[]}}],"model":"gpt-4o-mini","object":"chat.completion","usage":{"total_tokens":28}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Hello! How are you today?"}}, FinishReason: FinishReasonStop}},
},
{
name: "OpenAIResponsesRequest_OpenAIAgents",
raw: `{"include": [], "input": [{"content": "What's the weather in Bangalore right now?", "role": "user"}], "instructions": "You are a concise weather assistant.", "model": "gpt-4o-mini", "tools": [{"name": "get_weather", "type": "function"}]}`,
want: []Message{
text(MessageRoleSystem, "You are a concise weather assistant."),
text(MessageRoleUser, "What's the weather in Bangalore right now?"),
},
},
{
name: "OpenAIResponsesResponse_FunctionCall_OpenAIAgents",
raw: `{"id":"resp_1","object":"response","status":"completed","output":[{"arguments":"{\"city\":\"Bangalore\"}","call_id":"call_M","name":"get_weather","type":"function_call","id":"fc_1","status":"completed"}],"usage":{"total_tokens":113}}`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonStop, Content: []Part{{Type: PartTypeToolCall, ID: "call_M", Name: "get_weather", Arguments: map[string]any{"city": "Bangalore"}}}}},
},
{
name: "OpenAIResponsesResponse_MessageAndReasoning",
raw: `{"object":"response","status":"completed","output":[{"type":"reasoning","id":"rs_1","summary":[{"type":"summary_text","text":"Thinking about it"}]},{"type":"message","role":"assistant","content":[{"type":"output_text","text":"Paris."}]}]}`,
want: []Message{
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeThinking, Content: "Thinking about it"}}},
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Paris."}}, FinishReason: FinishReasonStop},
},
},
{
name: "VercelPromptMessages_ToolParts",
raw: `[{"role":"user","content":[{"type":"text","text":"What is the weather in Bengaluru in celsius?"}]},{"role":"assistant","content":[{"type":"tool-call","toolCallId":"call_l","toolName":"get_current_weather","args":{"city":"Bengaluru","unit":"c"}}]},{"role":"tool","content":[{"type":"tool-result","toolCallId":"call_l","toolName":"get_current_weather","result":{"city":"Bengaluru","temperature":27}}]}]`,
want: []Message{
text(MessageRoleUser, "What is the weather in Bengaluru in celsius?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_l", Name: "get_current_weather", Arguments: map[string]any{"city": "Bengaluru", "unit": "c"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_l", Name: "get_current_weather", Content: `{"city":"Bengaluru","temperature":27}`}}},
},
},
{
name: "MastraPromptMessages_MixedContent",
raw: `[{"role":"system","content":"You are a weather assistant."},{"role":"user","content":[{"type":"text","text":"What is the weather in Bengaluru?"}]}]`,
want: []Message{
text(MessageRoleSystem, "You are a weather assistant."),
text(MessageRoleUser, "What is the weather in Bengaluru?"),
},
},
{
name: "AnthropicMessages_ThinkingAndToolUse",
raw: `[{"role":"user","content":"Hi"},{"role":"assistant","content":[{"type":"thinking","thinking":"Let me see"},{"type":"redacted_thinking","data":"x"},{"type":"tool_use","id":"toolu_1","name":"lookup","input":{"q":"a"}}],"stop_reason":"tool_use"},{"role":"user","content":[{"type":"tool_result","tool_use_id":"toolu_1","content":"found","is_error":true}]}]`,
want: []Message{
text(MessageRoleUser, "Hi"),
{Role: MessageRoleAssistant, FinishReason: FinishReasonToolCall, Content: []Part{
{Type: PartTypeThinking, Content: "Let me see"},
{Type: PartTypeThinking, Redacted: true},
{Type: PartTypeToolCall, ID: "toolu_1", Name: "lookup", Arguments: map[string]any{"q": "a"}},
}},
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "toolu_1", Content: "found", IsError: true}}},
},
},
{
name: "GeminiContents_Converted",
raw: `[{"role":"user","parts":[{"text":"Weather in Paris?"}]},{"role":"model","parts":[{"functionCall":{"name":"get_weather","args":{"city":"Paris"}}}]},{"role":"user","parts":[{"functionResponse":{"name":"get_weather","response":{"temp":20}}}]}]`,
want: []Message{
text(MessageRoleUser, "Weather in Paris?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}},
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, Name: "get_weather", Content: `{"temp":20}`}}},
},
},
{
name: "LangChainSerialisedPrompt_Langsmith",
raw: `{"messages":[[{"lc":1,"type":"constructor","id":["langchain","schema","messages","SystemMessage"],"kwargs":{"content":"You are concise.","type":"system"}},{"lc":1,"type":"constructor","id":["langchain","schema","messages","HumanMessage"],"kwargs":{"content":"Define observability.","type":"human"}}]]}`,
want: []Message{
text(MessageRoleSystem, "You are concise."),
text(MessageRoleUser, "Define observability."),
},
},
{
name: "LangChainGenerations_Langsmith",
raw: `{"generations":[[{"text":"Observability is Y.","generation_info":{"finish_reason":"stop","logprobs":null},"type":"ChatGeneration","message":{"lc":1,"type":"constructor","id":["langchain","schema","messages","AIMessage"],"kwargs":{"content":"Observability is Y.","type":"ai","tool_calls":[],"invalid_tool_calls":[]}}}]],"llm_output":{"model_name":"gpt-4o-mini"}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Observability is Y."}}, FinishReason: FinishReasonStop}},
},
{
name: "SemconvTextPartWithTextKey_Litellm",
raw: `[{"role": "user", "parts": [{"type": "text", "text": "What animal is in this image?"}, {"type": "image_url", "image_url": {"url": "data:image/jpeg;base64,AAAA"}}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{
{Type: PartTypeText, Content: "What animal is in this image?"},
{Type: PartTypeGeneric, Content: `{"image_url":{"url":"data:image/jpeg;base64,AAAA"},"type":"image_url"}`},
}}},
},
{
name: "CompletionObject_OpenrouterGateway",
raw: `{"completion":"Paris.","reasoning":"The user asks for a capital.","rawRequest":{"model":"openai/gpt-4o-mini"}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{
{Type: PartTypeThinking, Content: "The user asks for a capital."},
{Type: PartTypeText, Content: "Paris."},
}}},
},
{
name: "VercelPrompt_SystemAndPrompt",
raw: `{"system":"You are concise.","prompt":"Say hello in five words."}`,
want: []Message{text(MessageRoleSystem, "You are concise."), text(MessageRoleUser, "Say hello in five words.")},
},
{
name: "VercelResponseToolCalls_BareList",
raw: `[{"toolCallType":"function","toolCallId":"call_1","toolName":"getWeather","args":"{\"city\":\"Bengaluru\"}"},{"type":"tool-call","toolCallId":"call_2","toolName":"searchWeb","input":{"q":"SigNoz"}}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{
{Type: PartTypeToolCall, ID: "call_1", Name: "getWeather", Arguments: map[string]any{"city": "Bengaluru"}},
{Type: PartTypeToolCall, ID: "call_2", Name: "searchWeb", Arguments: map[string]any{"q": "SigNoz"}},
}}},
},
{
name: "LangChainTypeMessages_AdditionalKwargsToolCalls",
raw: `[{"type":"human","content":"Weather?"},{"type":"ai","content":"","additional_kwargs":{"tool_calls":[{"id":"call_1","type":"function","function":{"name":"get_weather","arguments":"{\"city\":\"Paris\"}"}}]}},{"type":"tool","content":"20C","tool_call_id":"call_1","name":"get_weather"}]`,
want: []Message{
text(MessageRoleUser, "Weather?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_1", Name: "get_weather", Content: "20C"}}},
},
},
{
name: "LangGraphToolDefinitionMessage_Skipped",
raw: `[{"role":"tool","content":{"type":"function","function":{"name":"get_weather","parameters":{}}}},{"role":"user","content":"Hi"}]`,
want: []Message{text(MessageRoleUser, "Hi")},
},
{
name: "GeminiResponse_CandidatesWithFinishReason",
raw: `{"candidates":[{"content":{"parts":[{"text":"Let me check","thought":true},{"function_call":{"name":"get_weather","args":{"city":"Paris"}}}],"role":"model"},"finishReason":"STOP"}],"usageMetadata":{}}`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonStop, Content: []Part{
{Type: PartTypeThinking, Content: "Let me check"},
{Type: PartTypeToolCall, Name: "get_weather", Arguments: map[string]any{"city": "Paris"}},
}}},
},
{
name: "GeminiRequest_ContentsWithSystemInstruction",
raw: `{"model":"gemini-2.0","config":{"system_instruction":"Be brief."},"contents":[{"role":"user","parts":[{"text":"Hi"}]},{"role":"user","parts":[{"function_response":{"name":"get_weather","response":{"temp":20}}}]}]}`,
want: []Message{
text(MessageRoleSystem, "Be brief."),
text(MessageRoleUser, "Hi"),
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, Name: "get_weather", Content: `{"temp":20}`}}},
},
},
{
name: "GeminiRequest_StringContents",
raw: `{"contents":"Hi there","model":"gemini-2.0"}`,
want: []Message{text(MessageRoleUser, "Hi there")},
},
{
name: "MicrosoftAgent_ArrayToolCallIDs",
raw: `[{"role":"assistant","parts":[{"type":"tool_call","id":["run_1","call_9"],"name":"lookup","arguments":{"q":"x"}}]},{"role":"tool","parts":[{"type":"tool_call_response","id":["run_1","call_9"],"response":"found"}]}]`,
want: []Message{
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "call_9", Name: "lookup", Arguments: map[string]any{"q": "x"}}}},
{Role: MessageRoleTool, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_9", Content: "found"}}},
},
},
{
name: "PydanticAI_ToolCallResponseResultKey",
raw: `[{"role":"user","parts":[{"type":"tool_call_response","id":"call_1","name":"lookup","result":{"ok":true}}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "call_1", Name: "lookup", Content: `{"ok":true}`}}}},
},
{
name: "SemanticKernel_EventContentWrapper",
raw: `[{"role":"system","gen_ai.event.content":"{\"role\":\"system\",\"content\":\"Be brief.\",\"tool_calls\":[]}","gen_ai.system":"openai"},{"gen_ai.event.content":"{\"index\":0,\"message\":{\"role\":\"Assistant\",\"content\":\"Paris.\"},\"finish_reason\":\"Stop\"}"}]`,
want: []Message{
text(MessageRoleSystem, "Be brief."),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Paris."}}, FinishReason: FinishReasonStop},
},
},
{
name: "BedrockConverse_ToolUseAndToolResult",
raw: `{"messages":[{"role":"user","content":[{"text":"Weather?"}]},{"role":"assistant","content":[{"toolUse":{"toolUseId":"t1","name":"get_weather","input":{"city":"Paris"}}}]},{"role":"user","content":[{"toolResult":{"toolUseId":"t1","content":[{"text":"20C"}],"status":"error"}}]}],"system":[{"text":"Be brief."}]}`,
want: []Message{
text(MessageRoleSystem, "Be brief."),
text(MessageRoleUser, "Weather?"),
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "t1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}},
{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "t1", Content: "20C", IsError: true}}},
},
},
{
name: "AnthropicRequest_SystemString",
raw: `{"model":"claude","system":"Be brief.","messages":[{"role":"user","content":"Hi"}],"max_tokens":100}`,
want: []Message{text(MessageRoleSystem, "Be brief."), text(MessageRoleUser, "Hi")},
},
{
name: "OpenAIResponses_BuiltInToolCallIsServer",
raw: `{"object":"response","status":"completed","output":[{"type":"web_search_call","id":"ws_1","status":"completed","action":{"type":"search","query":"SigNoz"}},{"type":"custom_tool_call","call_id":"c1","name":"grep","input":"foo"},{"type":"message","role":"assistant","content":[{"type":"output_text","text":"Found it."}]}]}`,
want: []Message{
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "ws_1", Name: "web_search_call", Arguments: map[string]any{"action": map[string]any{"type": "search", "query": "SigNoz"}}, Server: true}}},
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "grep", Arguments: "foo"}}},
{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Found it."}}, FinishReason: FinishReasonStop},
},
},
{
name: "NestedMessageList_Unwrapped",
raw: `[[{"role":"user","content":"Hi"}]]`,
want: []Message{text(MessageRoleUser, "Hi")},
},
{
name: "StringifiedMessages_Decoded",
raw: `{"messages":"[{\"role\":\"user\",\"content\":\"Hi\"}]"}`,
want: []Message{text(MessageRoleUser, "Hi")},
},
{
name: "EmbeddingsRequest_Generic",
raw: `{"input": ["a", "b"], "model": "text-embedding-3-small"}`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `{"input": ["a", "b"], "model": "text-embedding-3-small"}`}}}},
},
{
name: "BedrockConverseResponse_OutputMessageWrapper",
raw: `{"output":{"message":{"role":"assistant","content":[{"text":"20C in Paris."}]}},"stopReason":"end_turn","usage":{"inputTokens":10}}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "20C in Paris."}}, FinishReason: FinishReasonStop}},
},
{
name: "OllamaResponse_MessageWrapper",
raw: `{"model":"llama3","message":{"role":"assistant","content":"Hi!"},"done":true,"done_reason":"stop"}`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeText, Content: "Hi!"}}, FinishReason: FinishReasonStop}},
},
{
name: "CohereV2Response_MessageWrapper",
raw: `{"id":"x","message":{"role":"assistant","tool_calls":[{"id":"c1","type":"function","function":{"name":"get_weather","arguments":"{\"city\":\"Paris\"}"}}]},"finish_reason":"TOOL_CALL"}`,
want: []Message{{Role: MessageRoleAssistant, FinishReason: FinishReasonToolCall, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}}},
},
{
name: "OpenAIToolMessage_TextBlocksBecomeToolResult",
raw: `[{"role":"tool","tool_call_id":"c1","content":[{"type":"text","text":"20C"},{"type":"text","text":"clear"}]}]`,
want: []Message{{Role: MessageRoleTool, Content: []Part{
{Type: PartTypeToolResult, ToolCallID: "c1", Content: "20C"},
{Type: PartTypeToolResult, ToolCallID: "c1", Content: "clear"},
}}},
},
{
name: "AnthropicToolResult_TextBlocksJoined",
raw: `[{"role":"user","content":[{"type":"tool_result","tool_use_id":"t1","content":[{"type":"text","text":"line one"},{"type":"text","text":"line two"}]}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{{Type: PartTypeToolResult, ToolCallID: "t1", Content: "line one\nline two"}}}},
},
{
name: "MessageWithContentParts_GeminiNested",
raw: `[{"role":"model","content":{"parts":[{"text":"Hi"}],"role":"model"}}]`,
want: []Message{text(MessageRoleAssistant, "Hi")},
},
{
name: "ToolCallTypedItem_PlainToolCall",
raw: `[{"type":"tool_call","id":"c1","name":"get_weather","args":{"city":"Paris"}}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}}},
},
{
name: "OpenAIToolCallList_Bare",
raw: `[{"id":"c1","type":"function","function":{"name":"get_weather","arguments":"{\"city\":\"Paris\"}"}}]`,
want: []Message{{Role: MessageRoleAssistant, Content: []Part{{Type: PartTypeToolCall, ID: "c1", Name: "get_weather", Arguments: map[string]any{"city": "Paris"}}}}},
},
{
name: "ToolDefinitionList_Generic",
raw: `[{"type":"function","function":{"name":"get_weather","description":"Weather","parameters":{"type":"object"}}}]`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `[{"type":"function","function":{"name":"get_weather","description":"Weather","parameters":{"type":"object"}}}]`}}}},
},
{
name: "ContentBlockList_RolelessMessage",
raw: `[{"type":"text","text":"Let me check."},{"type":"tool_use","id":"t1","name":"lookup","input":{"q":"x"}}]`,
want: []Message{{Content: []Part{
{Type: PartTypeText, Content: "Let me check."},
{Type: PartTypeToolCall, ID: "t1", Name: "lookup", Arguments: map[string]any{"q": "x"}},
}}},
},
{
name: "ListOfJSONStrings_Decoded",
raw: []any{`{"role":"user","content":"Hi"}`, `{"role":"assistant","content":"Hello"}`},
want: []Message{text(MessageRoleUser, "Hi"), text(MessageRoleAssistant, "Hello")},
},
{
name: "SingleMessageObject_Converted",
raw: `{"role":"assistant","content":"Done."}`,
want: []Message{text(MessageRoleAssistant, "Done.")},
},
{
name: "UnknownRoleAndFinishReason_KeptLowercased",
raw: `[{"role":"Narrator","parts":[{"type":"text","content":"x"}],"finish_reason":"Weird"}]`,
want: []Message{{Role: "narrator", Content: []Part{{Type: PartTypeText, Content: "x"}}, FinishReason: "weird"}},
},
{
name: "UnknownPartType_Generic",
raw: `[{"role":"user","parts":[{"type":"image","url":"http://x/y.png"}]}]`,
want: []Message{{Role: MessageRoleUser, Content: []Part{{Type: PartTypeGeneric, Content: `{"type":"image","url":"http://x/y.png"}`}}}},
},
{
name: "PlainText_Generic",
raw: "Let the cost of the ball be x dollars.",
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: "Let the cost of the ball be x dollars."}}}},
},
{
name: "UnknownJSONShape_GenericWithOriginal",
raw: `{"output": "{\"query\": \"SigNoz\"}", "kwargs": {"name": "search_web"}}`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `{"output": "{\"query\": \"SigNoz\"}", "kwargs": {"name": "search_web"}}`}}}},
},
{
name: "JSONEncodedString_Generic",
raw: `"{\"query\": \"SigNoz\"}"`,
want: []Message{{Content: []Part{{Type: PartTypeGeneric, Content: `"{\"query\": \"SigNoz\"}"`}}}},
},
{
name: "DecodedValue_Converted",
raw: []any{map[string]any{"role": "user", "content": "hi"}},
want: []Message{text(MessageRoleUser, "hi")},
},
{
name: "EmptyList_NoMessages",
raw: `[]`,
want: []Message{},
},
{
name: "Nil_NoMessages",
raw: nil,
want: []Message{},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
assert.Equal(t, testCase.want, NormalizeMessages(testCase.raw))
})
}
}

View File

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

View File

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

View File

@@ -0,0 +1,173 @@
package spantypes
import (
"testing"
"time"
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestNewThreadQuery(t *testing.T) {
cursor := ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}
testCases := []struct {
name string
postable PostableThreadQuery
want *ThreadQuery
wantErr bool
}{
{name: "ZeroLimit_UsesDefault", postable: PostableThreadQuery{}, want: &ThreadQuery{Limit: threadDefaultLimit}},
{name: "PositiveLimit_Kept", postable: PostableThreadQuery{Limit: 25}, want: &ThreadQuery{Limit: 25}},
{name: "MaxLimit_Kept", postable: PostableThreadQuery{Limit: threadMaxLimit}, want: &ThreadQuery{Limit: threadMaxLimit}},
{name: "AboveMaxLimit_Rejected", postable: PostableThreadQuery{Limit: threadMaxLimit + 1}, wantErr: true},
{name: "NegativeLimit_Rejected", postable: PostableThreadQuery{Limit: -1}, wantErr: true},
{name: "Cursor_Decoded", postable: PostableThreadQuery{Limit: 10, Cursor: cursor.Encode()}, want: &ThreadQuery{Limit: 10, Cursor: &cursor}},
{name: "InvalidCursor_Rejected", postable: PostableThreadQuery{Cursor: "not base64!"}, wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
got, err := NewThreadQuery(&testCase.postable)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, got)
})
}
}
func TestDecodeThreadCursor(t *testing.T) {
testCases := []struct {
name string
cursor string
want *ThreadCursor
wantErr bool
}{
{name: "EncodedCursor_RoundTrips", cursor: ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}.Encode(), want: &ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}},
{name: "NotBase64_Rejected", cursor: "not base64!", wantErr: true},
{name: "NotJSON_Rejected", cursor: "bm90IGpzb24", wantErr: true},
{name: "MissingSpanID_Rejected", cursor: "eyJ0IjogMX0", wantErr: true},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
got, err := DecodeThreadCursor(testCase.cursor)
if testCase.wantErr {
assert.Error(t, err)
return
}
require.NoError(t, err)
assert.Equal(t, testCase.want, got)
})
}
}
func TestNewGettableTraceThread(t *testing.T) {
spans := []StorableSpan{
{SpanID: "a", StartTime: time.Unix(1, 500_000_000)},
{SpanID: "b", StartTime: time.Unix(2, 0)},
{SpanID: "c", StartTime: time.Unix(3, 0)},
}
testCases := []struct {
name string
spans []StorableSpan
limit int
wantSpanIDs []string
wantTimeUnix []uint64
wantNextCursor string
}{
{name: "MoreThanLimit_TrimsAndSetsCursor", spans: spans, limit: 2, wantSpanIDs: []string{"a", "b"}, wantTimeUnix: []uint64{1500, 2000}, wantNextCursor: ThreadCursor{TimeUnixNano: 2_000_000_000, SpanID: "b"}.Encode()},
{name: "WithinLimit_NoCursor", spans: spans, limit: 3, wantSpanIDs: []string{"a", "b", "c"}, wantTimeUnix: []uint64{1500, 2000, 3000}},
{name: "NoSpans_EmptyList", limit: 3, wantSpanIDs: []string{}, wantTimeUnix: []uint64{}},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
thread := NewGettableTraceThread("trace-1", testCase.spans, testCase.limit)
require.NotNil(t, thread.Spans)
spanIDs := make([]string, len(thread.Spans))
timeUnix := make([]uint64, len(thread.Spans))
for i, span := range thread.Spans {
spanIDs[i] = span.SpanID
timeUnix[i] = span.TimeUnix
assert.Equal(t, "trace-1", span.TraceID)
}
assert.Equal(t, testCase.wantSpanIDs, spanIDs)
assert.Equal(t, testCase.wantTimeUnix, timeUnix)
assert.Equal(t, testCase.wantNextCursor, thread.NextCursor)
})
}
}
func TestNewThreadSpan(t *testing.T) {
userHi := []aiobservabilitytypes.Message{{
Role: aiobservabilitytypes.MessageRoleUser,
Content: []aiobservabilitytypes.Part{{Type: aiobservabilitytypes.PartTypeText, Content: "hi"}},
}}
assistantHello := []aiobservabilitytypes.Message{{
Role: aiobservabilitytypes.MessageRoleAssistant,
Content: []aiobservabilitytypes.Part{{Type: aiobservabilitytypes.PartTypeText, Content: "hello"}},
FinishReason: aiobservabilitytypes.FinishReasonStop,
}}
testCases := []struct {
name string
span StorableSpan
wantInput []aiobservabilitytypes.Message
wantOutput []aiobservabilitytypes.Message
wantAttrs map[string]any
}{
{
name: "MessagesInLegacyMap",
span: StorableSpan{AttributesString: map[string]string{
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
}},
wantInput: userHi,
wantOutput: assistantHello,
wantAttrs: map[string]any{
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
},
},
{
name: "MessagesInJSONColumn_FlattenedToDottedKeys",
span: StorableSpan{AttributesJSON: telemetrystoretypes.JSONValue{
"gen_ai": map[string]any{
"input": map[string]any{"messages": `[{"role":"user","content":"hi"}]`},
"request": map[string]any{"model": "gpt-4o"},
},
}},
wantInput: userHi,
wantAttrs: map[string]any{
"gen_ai.input.messages": `[{"role":"user","content":"hi"}]`,
"gen_ai.request.model": "gpt-4o",
},
},
{
name: "LegacyMapWinsOverJSONColumn",
span: StorableSpan{
AttributesJSON: telemetrystoretypes.JSONValue{"gen_ai": map[string]any{"request": map[string]any{"model": "json"}}},
AttributesString: map[string]string{"gen_ai.request.model": "map"},
},
wantAttrs: map[string]any{"gen_ai.request.model": "map"},
},
{
name: "NoMessages_FieldsUnset",
span: StorableSpan{AttributesString: map[string]string{"http.method": "GET"}},
wantAttrs: map[string]any{"http.method": "GET"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
span := newThreadSpan("trace-1", &testCase.span)
assert.Equal(t, testCase.wantInput, span.FormattedInput)
assert.Equal(t, testCase.wantOutput, span.FormattedOutput)
assert.Equal(t, testCase.wantAttrs, span.Attributes)
})
}
}

View File

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

View File

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