Compare commits

...

8 Commits

Author SHA1 Message Date
nityanandagohain
f5f019f61b feat: trace detail thread endpoint 2026-09-28 17:32:55 +05:30
Aditya Singh
ed1bf7ab89 fix(bottom-strip): size pages from the layout instead of the viewport (#12940)
Some checks failed
build-staging / prepare (push) Has been cancelled
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
#### Description

- pages that hardcoded `100vh` minus a guess at what sits above them
came out taller
than the pane they live in, which showed up as scroll that should not be
there. they
  now take what the layout gives them.
- most of the `100vh` in the app turned out to be harmless.. either
flex-shrink absorbs
it or the pane scrolls anyway. those are left alone, only the ones with
a real symptom
  are changed here.
- alert rules and triggered alerts also needed the AlertList tabs chain
to hand height
  down, that page uses antd `Tabs` directly instead of `RouteTab`.
- licenses, status and support pages had `max-height: 100vh` with
`overflow: hidden` and
no inner scroller, so anything past a viewport was clipped with no way
to reach it.
  removed the cap on all three.

#### Issues closed by this PR

Part of https://github.com/SigNoz/engineering-pod/issues/6074


#### Screenshots/ Screen recording

Home page


https://github.com/user-attachments/assets/8ba3e3c9-1959-4393-b503-ddf579e9d139

Without bottom strip


https://github.com/user-attachments/assets/2ee08181-710a-4c34-a4e7-893c4a320453



Status page

<img width="1728" height="1000" alt="status"
src="https://github.com/user-attachments/assets/f44ec7d7-4ee0-4866-9a90-fe25622fe25b"
/>

Without bottom strip

<img width="1728" height="997" alt="status2"
src="https://github.com/user-attachments/assets/be57dce6-65ca-4450-8f0d-0337625a1d64"
/>

Alert rules


https://github.com/user-attachments/assets/7cc25f56-54e4-4489-bdb4-453409151bac

Triggered alerts


https://github.com/user-attachments/assets/2d8e1885-7cad-4659-90f6-59960f4afce8

Without bottom strip



https://github.com/user-attachments/assets/5247f6c0-3675-4910-b1e4-4076bf93c16e



Support
<img width="1728" height="997" alt="support"
src="https://github.com/user-attachments/assets/c2b3c0b2-2abd-4538-bb52-2660e75817b3"
/>

Trace funnel



https://github.com/user-attachments/assets/c4872135-3cb5-4348-b6c6-2d3e4dabdc1b

Without bottom strip


https://github.com/user-attachments/assets/aef94d5d-02e5-4154-995b-76819d55c2d8


Trace details



https://github.com/user-attachments/assets/24d701b0-296f-4fb4-ba9d-615a47b33285

Without bottom strip


https://github.com/user-attachments/assets/23083f20-67bf-4b06-aaec-aba390a6f594



#### Additional Information

- every page here was checked on screen before and after. the ones left
untouched
(infra hosts/k8s, traces + llm explorer list views, llm settings tables,
the k8s logs
  drawer) were checked too and are fine.
2026-09-25 15:14:44 +00:00
Aditya Singh
6e979c8318 feat(bottom-strip): add the layout shell behind a feature flag (#12936)
#### Description

- adds the bottom strip to the app layout behind a localStorage flag.
shows the build
version on the left for now.. right side actions and the per page count
come in the
  next tickets.
- `.app-content` is a column flex now and `LayoutContent` takes the
height left over
instead of `height: 100%`, so the strip has a stable box to sit under.
this is the
  only bit not behind the flag.
- fixed bottom elements read `--bottom-strip-height`. the var only
exists while the
strip is mounted, so with the flag off everything falls back to where it
is today.
- hides nothing. each later ticket hides the piece it replaces.

#### Issues closed by this PR

Part of https://github.com/SigNoz/engineering-pod/issues/6074

<img width="3084" height="1566" alt="image"
src="https://github.com/user-attachments/assets/b1821fda-5c33-40e7-926a-5d91fedb797e"
/>


#### Additional Information

- pages that still hardcode `100vh` (infra hosts/k8s, trace details,
traces and llm
list views) push the strip off screen. that is the next PR on this
ticket.
- pylon chat window offset is not here.. needs a pylon enabled tenant to
verify so it
  goes with the right side actions ticket.
2026-09-25 12:04:21 +00:00
praneeth-signoz
3b6becff7a chore(tests): split alerts test suite as per domain boundaries (#12983)
Some checks failed
build-staging / prepare (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description
Segregated alert-related test suite in to alert manager and ruler as per
their domain boundaries.

---------

Co-authored-by: Praneeth Lingam <praneethlingam@Ollys-MacBook-Pro.local>
2026-09-25 07:30:53 +00:00
Aditya Singh
ab715533b9 feat(saved-views): read saved views from the v2 api on home, noz and column sync (#12970)
#### Description
- moved the home saved views widget, noz open saved view and the saved
view column/format sync (`usePreferenceSync`) from
`/api/v1/explorer/views` to `/api/v2/saved_views`.. generated client and
DTOs used as is, no adapter. labels read `spec.displayName`, columns
`spec.selectedFields`, formatting `spec.display`.
- small `container/SavedViews/utils.ts` for the two things every v2
consumer needs.. shaping the v2 spec for the existing v5 reverse mapper,
and the `DataSource` → api source map. rest of the saved views hooks
come with the sidebar work.
- explorer bottom bar, the `/saved-views` pages and `ExplorerCard` stay
on v1 on purpose.. they get deleted with the bottom strip work, no point
migrating something with a death date. v1 and v2 run in parallel till
then.
- home widget drops the tags badges (nothing ever wrote tags) and the
extra lookup on click. functionalities kept same.

#### Issues closed by this PR

Closes https://github.com/SigNoz/engineering-pod/issues/6095
Part of https://github.com/SigNoz/engineering-pod/issues/5918

#### Additional Information
- traces view with no saved columns now falls back to the typed
`defaultTraceSelectedColumns` (what the loader uses) instead of the
string list from `ListView/configs`.. old one was strings in a
`TelemetryFieldKey[]` hidden by `JSON.parse`.
- `viewName` is still written to the url on open so the old bar shows
the view as selected.. goes away when the bar does.
- noz open view could not be tested locally, covered by unit tests only.
- home storybook mocks regenerated for the v2 endpoint.
2026-09-25 04:37:54 +00:00
Vinicius Lourenço
8e2da68fc6 test(api-monitoring): mock /fields/keys for quick filters settings stories (#12980)
Some checks failed
build-staging / prepare (push) Has been cancelled
build-staging / js-build (push) Has been cancelled
build-staging / go-build (push) Has been cancelled
build-staging / staging (push) Has been cancelled
cacheci / tests (push) Has been cancelled
Release Drafter / update_release_draft (push) Has been cancelled
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description

Add missing mocks for stories on api monitoring after
https://github.com/SigNoz/signoz/pull/12968
2026-09-24 14:24:56 +00:00
Nityananda Gohain
8371a70801 perf(querybuilder): compare materialized exists columns explicitly (#12978)
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description

Materialized existence checks now render as an explicit comparison
instead of a bare bool column. Results are unchanged; only skip-index
usage improves.

  ```sql
  -- before
  WHERE `attribute_string_gen_ai$$request$$model_exists`
     OR `attribute_string_gen_ai$$provider$$name` = 'anthropic'

  -- after
  WHERE `attribute_string_gen_ai$$request$$model_exists` = true
     OR `attribute_string_gen_ai$$provider$$name` = 'anthropic'
  ```

  <details>
<summary>EXPLAIN indexes = 1 (trace-matching phase, 123M
spans)</summary>

  Before: bare `col_exists`
  ```
  Name: idx_gen_ai_span_exists
  Granules: 15193/15193
  Name: <Combined skip indexes>
  Granules: 15193/15193
  ```

  After: `col_exists = true`
  ```
  Name: idx_gen_ai_span_exists
  Granules: 15193/15193
  Name: <Combined skip indexes>
  Granules: 488/15193
  ```
  </details>

----
- ClickHouse can use a different skip index for each side of an OR and
union the results, but it can't when one side is a bare bool column.
Comparing with `= true` fixes that.
- This shape comes from the AI explorer trace list with a span filter: a
trace qualifies when it has a gen_ai span *and* a span matching the
filter (possibly different spans), so the WHERE is `(gen_ai gate) OR
<filter>` followed by a HAVING.
- Needs the gen_ai materialized columns and `idx_gen_ai_span_exists`
from SigNoz/signoz-otel-collector#929; without them there's no index to
combine.

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

<!--Anything reviewers should keep in mind while reviewing -->
#### Additional Information
- Benchmarked the AI trace list filtered on `gen_ai.provider.name`
against a 123M-span table (direct I/O, caches off): from ~30M spans in
the window, latency drops 16–17% and CPU 35–38%, with ~25x fewer rows
read (123M spans: 510 → 427 ms, 1.5 → 0.9 sCPU). The saved time and CPU
keep growing with span count, so larger windows save more.
- Single-condition filters (`gen_ai.request.model EXISTS` in dashboard
panels, the AND-ed gate in AI aggregations) already pruned with the bare
form; no change there.
2026-09-24 13:15:32 +00:00
Naman Verma
9d9b0e194a chore: add ability to mark API stability as beta/alpha (#12957)
<!--A few plain bullets saying what changed and why, for a reviewer
skimming it - not a wall of text, not a restatement of the diff, not
generated boilerplate.-->
#### Description

If an API that is already deployed is currently being tested via UI
integration or any other means, we should mark such APIs as under
development so that other external clients know that these APIs aren't
fully stable. This is especially required if we are working on v2
versions of APIs for any entity.

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
#### Issues closed by this PR

Part of https://github.com/SigNoz/pulse-pod/issues/369

<!--Anything reviewers should keep in mind while reviewing -->
#### Additional Information

This PR adds the development flag on the v2 notification channel APIs

<!--Please delete paragraphs that you did not use before submitting.-->
2026-09-24 12:19:02 +00:00
138 changed files with 4061 additions and 405 deletions

View File

@@ -38,7 +38,6 @@ jobs:
fail-fast: false
matrix:
suite:
- alerts
- alertmanager
- alertmanagerrotation
- basepath
@@ -64,10 +63,12 @@ jobs:
- querierauthz
- role
- rootuser
- ruler
- savedview
- semconvfamilies
- serviceaccount
- spanmapper
- tracedetail
- querier_json_body
- querier_skip_resource_fingerprint
- ttl

File diff suppressed because it is too large Load Diff

View File

@@ -179,6 +179,7 @@ The `handler.New` function ties the HTTP handler to OpenAPI metadata via `OpenAP
- **SuccessStatusCode**: The HTTP status for successful responses (for example, `http.StatusOK`, `http.StatusCreated`, `http.StatusNoContent`).
- **ErrorStatusCodes**: Additional error status codes beyond the standard ones automatically added by `handler.New`.
- **SecuritySchemes**: Auth mechanisms and scopes required by the operation.
- **Stability**: Maturity marker (`handler.StabilityDevelopment`, `handler.StabilityAlpha`, `handler.StabilityBeta`, `handler.StabilityStable`, the OpenTelemetry Collector levels) emitted as the `x-signoz-stability` extension on every operation. Unset is emitted as `alpha`.
The generic handler:

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

@@ -47,4 +47,5 @@ export enum LOCALSTORAGE {
DASHBOARDS_LIST_VIEWS = 'DASHBOARDS_LIST_VIEWS',
DASHBOARD_V2_PANEL_COLUMN_WIDTHS = 'DASHBOARD_V2_PANEL_COLUMN_WIDTHS',
LLM_ATTRIBUTE_MAPPING_TEST_SPAN = 'LLM_ATTRIBUTE_MAPPING_TEST_SPAN',
SAVED_VIEW_ENABLED = 'SAVED_VIEW_ENABLED',
}

View File

@@ -3,15 +3,22 @@ import {
MessageActionKindDTO,
SavedViewEntityDTO,
} from 'api/ai-assistant/sigNozAIAssistantAPI.schemas';
import { getAllViews } from 'api/saveView/getAllViews';
import { getViewById } from 'api/saveView/getViewById';
import {
getSavedView,
listSavedViews,
} from 'api/generated/services/saved-view';
import {
GetSavedView200,
ListSavedViews200,
SavedviewtypesPanelTypeDTO,
SavedviewtypesSavedViewDTO,
SavedviewtypesSchemaVersionDTO,
SavedviewtypesSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
import ROUTES from 'constants/routes';
import { QueryParams } from 'constants/query';
import { PANEL_TYPES } from 'constants/queryBuilder';
import { ICompositeMetricQuery } from 'types/api/alerts/compositeQuery';
import { AllViewsProps, ViewProps } from 'types/api/saveViews/types';
import { DataSource } from 'types/common/queryBuilder';
import { AxiosResponse } from 'axios';
import type { History } from 'history';
import {
@@ -31,8 +38,7 @@ import {
} from '../resolveOpenResource';
import { resourceRoute, ResourceType } from '../resourceRoute';
jest.mock('api/saveView/getAllViews');
jest.mock('api/saveView/getViewById');
jest.mock('api/generated/services/saved-view');
jest.mock(
'lib/newQueryBuilder/queryBuilderMappers/mapQueryDataFromApi',
@@ -48,43 +54,45 @@ jest.mock(
}),
);
const mockedGetAllViews = getAllViews as jest.MockedFunction<
typeof getAllViews
const mockedListSavedViews = listSavedViews as jest.MockedFunction<
typeof listSavedViews
>;
const mockedGetViewById = getViewById as jest.MockedFunction<
typeof getViewById
const mockedGetSavedView = getSavedView as jest.MockedFunction<
typeof getSavedView
>;
function makeView(id: string, sourcePage: DataSource): ViewProps {
function makeView(
id: string,
source: SavedviewtypesSourceDTO,
): SavedviewtypesSavedViewDTO {
return {
id,
name: `View ${id}`,
category: 'test',
name: `view-${id}`,
source,
schemaVersion: SavedviewtypesSchemaVersionDTO.v2,
createdAt: '2021-07-07T06:31:00.000Z',
createdBy: 'user',
updatedAt: '2021-07-07T06:33:00.000Z',
updatedBy: 'user',
sourcePage,
tags: [],
extraData: '',
compositeQuery: {
panelType: PANEL_TYPES.LIST,
} as ICompositeMetricQuery,
};
spec: {
displayName: `View ${id}`,
panelType: SavedviewtypesPanelTypeDTO.list,
requestType: 'raw',
queries: [{ type: 'builder_query', spec: { name: 'A', signal: source } }],
},
} as unknown as SavedviewtypesSavedViewDTO;
}
function mockViewsResponse(views: ViewProps[]): AxiosResponse<AllViewsProps> {
return {
data: { status: 'success', data: views },
} as AxiosResponse<AllViewsProps>;
function mockViewsResponse(
views: SavedviewtypesSavedViewDTO[],
): ListSavedViews200 {
return { status: 'success', data: views };
}
function mockViewByIdResponse(
view: ViewProps,
): AxiosResponse<{ status: string; data: ViewProps }> {
return {
data: { status: 'success', data: view },
} as AxiosResponse<{ status: string; data: ViewProps }>;
view: SavedviewtypesSavedViewDTO,
): GetSavedView200 {
return { status: 'success', data: view };
}
describe('resourceRoute', () => {
@@ -190,18 +198,33 @@ describe('resolveOpenResource', () => {
describe('findSavedViewInLists', () => {
beforeEach(() => {
mockedGetAllViews.mockReset();
mockedListSavedViews.mockReset();
});
it('loads only the hinted source when entity is provided', async () => {
const tracesView = makeView('view-traces', DataSource.TRACES);
mockedGetAllViews.mockResolvedValueOnce(mockViewsResponse([tracesView]));
const tracesView = makeView('view-traces', SavedviewtypesSourceDTO.traces);
mockedListSavedViews.mockResolvedValueOnce(mockViewsResponse([tracesView]));
const result = await findSavedViewInLists('view-traces', DataSource.TRACES);
expect(result).toStrictEqual(tracesView);
expect(mockedGetAllViews).toHaveBeenCalledTimes(1);
expect(mockedGetAllViews).toHaveBeenCalledWith(DataSource.TRACES);
expect(mockedListSavedViews).toHaveBeenCalledTimes(1);
expect(mockedListSavedViews).toHaveBeenCalledWith({
source: SavedviewtypesSourceDTO.traces,
});
});
it('treats a null list as empty and probes the next source', async () => {
const metricsView = makeView('view-metrics', SavedviewtypesSourceDTO.metrics);
mockedListSavedViews
.mockResolvedValueOnce({ status: 'success', data: null })
.mockResolvedValueOnce(mockViewsResponse([]))
.mockResolvedValueOnce(mockViewsResponse([metricsView]));
const result = await findSavedViewInLists('view-metrics');
expect(result).toStrictEqual(metricsView);
expect(mockedListSavedViews).toHaveBeenCalledTimes(3);
});
});
@@ -227,52 +250,75 @@ describe('openSavedView', () => {
it('navigates with history.push and view query params', () => {
const push = jest.fn();
const history = { push } as unknown as History;
const view = makeView('view-logs', DataSource.LOGS);
const view = makeView('view-logs', SavedviewtypesSourceDTO.logs);
openSavedView(view, history);
expect(push).toHaveBeenCalledTimes(1);
const pushedUrl = push.mock.calls[0][0] as string;
expect(pushedUrl).toContain(ROUTES.LOGS_EXPLORER);
expect(pushedUrl).toContain(QueryParams.viewKey);
const params = new URLSearchParams(pushedUrl.split('?')[1]);
expect(params.get(QueryParams.viewKey)).toBe('"view-logs"');
expect(params.get(QueryParams.viewName)).toBe('"View view-logs"');
expect(params.get(QueryParams.panelTypes)).toBe('"list"');
});
it('throws when the view has no source', () => {
const view = makeView('view-logs', SavedviewtypesSourceDTO.logs);
delete view.source;
expect(() =>
openSavedView(view, { push: jest.fn() } as unknown as History),
).toThrow('Unsupported saved view source');
});
it('throws when the view has no queries', () => {
const view = makeView('view-logs', SavedviewtypesSourceDTO.logs);
view.spec.queries = [];
expect(() =>
openSavedView(view, { push: jest.fn() } as unknown as History),
).toThrow('Saved view is missing query data');
});
});
describe('openSavedViewByKey', () => {
beforeEach(() => {
mockedGetAllViews.mockReset();
mockedGetViewById.mockReset();
mockedListSavedViews.mockReset();
mockedGetSavedView.mockReset();
});
it('prefers the direct view lookup endpoint', async () => {
const view = makeView('view-logs', DataSource.LOGS);
mockedGetViewById.mockResolvedValueOnce(mockViewByIdResponse(view));
const view = makeView('view-logs', SavedviewtypesSourceDTO.logs);
mockedGetSavedView.mockResolvedValueOnce(mockViewByIdResponse(view));
const push = jest.fn();
const history = { push } as unknown as History;
await openSavedViewByKey('view-logs', DataSource.LOGS, history);
expect(mockedGetViewById).toHaveBeenCalledWith('view-logs');
expect(mockedGetAllViews).not.toHaveBeenCalled();
expect(mockedGetSavedView).toHaveBeenCalledWith({ id: 'view-logs' });
expect(mockedListSavedViews).not.toHaveBeenCalled();
expect(push).toHaveBeenCalled();
});
it('falls back to list probing when direct lookup fails', async () => {
const view = makeView('view-traces', DataSource.TRACES);
mockedGetViewById.mockRejectedValueOnce(new Error('not found'));
mockedGetAllViews.mockResolvedValueOnce(mockViewsResponse([view]));
const view = makeView('view-traces', SavedviewtypesSourceDTO.traces);
mockedGetSavedView.mockRejectedValueOnce(new Error('not found'));
mockedListSavedViews.mockResolvedValueOnce(mockViewsResponse([view]));
const push = jest.fn();
const history = { push } as unknown as History;
await openSavedViewByKey('view-traces', DataSource.TRACES, history);
expect(mockedGetAllViews).toHaveBeenCalledWith(DataSource.TRACES);
expect(mockedListSavedViews).toHaveBeenCalledWith({
source: SavedviewtypesSourceDTO.traces,
});
expect(push).toHaveBeenCalled();
});
it('throws when the saved view does not exist', async () => {
mockedGetViewById.mockRejectedValueOnce(new Error('not found'));
mockedGetAllViews.mockResolvedValue(mockViewsResponse([]));
mockedGetSavedView.mockRejectedValueOnce(new Error('not found'));
mockedListSavedViews.mockResolvedValue(mockViewsResponse([]));
await expect(
openSavedViewByKey('missing', DataSource.LOGS, {

View File

@@ -1,15 +1,22 @@
import { getAllViews } from 'api/saveView/getAllViews';
import { getViewById } from 'api/saveView/getViewById';
import {
getSavedView,
listSavedViews,
} from 'api/generated/services/saved-view';
import { SavedviewtypesSavedViewDTO } from 'api/generated/services/sigNoz.schemas';
import { QueryParams } from 'constants/query';
import { PANEL_TYPES } from 'constants/queryBuilder';
import { mapQueryDataFromApi } from 'lib/newQueryBuilder/queryBuilderMappers/mapQueryDataFromApi';
import {
findSavedView,
getSavedViewQuery,
SavedViewSourcePage,
toSavedViewSource,
} from 'container/SavedViews/utils';
import { SOURCEPAGE_VS_ROUTES } from 'pages/SaveView/constants';
import { ViewProps } from 'types/api/saveViews/types';
import { DataSource } from 'types/common/queryBuilder';
import { Query } from 'types/api/queryBuilder/queryBuilderData';
import { History } from 'history';
type SavedViewSourceHint = DataSource | 'meter';
type SavedViewSourceHint = SavedViewSourcePage;
const DEFAULT_PROBE_SOURCES: SavedViewSourceHint[] = [
DataSource.LOGS,
@@ -20,13 +27,15 @@ const DEFAULT_PROBE_SOURCES: SavedViewSourceHint[] = [
export async function findSavedViewInLists(
viewKey: string,
sourceHint?: SavedViewSourceHint | null,
): Promise<ViewProps | null> {
): Promise<SavedviewtypesSavedViewDTO | null> {
const sources = sourceHint ? [sourceHint] : DEFAULT_PROBE_SOURCES;
for (const source of sources) {
try {
const response = await getAllViews(source);
const match = response.data.data.find((view) => view.id === viewKey);
const response = await listSavedViews({
source: toSavedViewSource(source),
});
const match = findSavedView(response.data, viewKey);
if (match) {
return match;
}
@@ -41,11 +50,11 @@ export async function findSavedViewInLists(
async function loadSavedView(
viewKey: string,
sourceHint?: SavedViewSourceHint | null,
): Promise<ViewProps> {
): Promise<SavedviewtypesSavedViewDTO> {
try {
const response = await getViewById(viewKey);
if (response.data?.data) {
return response.data.data;
const response = await getSavedView({ id: viewKey });
if (response.data) {
return response.data;
}
} catch {
// Fall back to list probing when the direct lookup fails.
@@ -85,20 +94,23 @@ export function buildExplorerNavigationUrl(
return `${route}?${params.toString()}`;
}
export function openSavedView(view: ViewProps, history: History): void {
const route = explorerRouteForSourcePage(view.sourcePage);
export function openSavedView(
view: SavedviewtypesSavedViewDTO,
history: History,
): void {
const route = view.source ? explorerRouteForSourcePage(view.source) : null;
if (!route) {
throw new Error('Unsupported saved view source');
}
if (!view.compositeQuery) {
if (!view.spec.queries?.length) {
throw new Error('Saved view is missing query data');
}
const query = mapQueryDataFromApi(view.compositeQuery);
const query = getSavedViewQuery(view);
const url = buildExplorerNavigationUrl(route, query, {
[QueryParams.panelTypes]: view.compositeQuery.panelType as PANEL_TYPES,
[QueryParams.viewName]: view.name,
[QueryParams.panelTypes]: view.spec.panelType as unknown as PANEL_TYPES,
[QueryParams.viewName]: view.spec.displayName,
[QueryParams.viewKey]: view.id,
});
history.push(url);
@@ -112,6 +124,3 @@ export async function openSavedViewByKey(
const view = await loadSavedView(viewKey, sourceHint);
openSavedView(view, history);
}
/** @deprecated Use findSavedViewInLists — kept for tests. */
export const findSavedView = findSavedViewInLists;

View File

@@ -53,6 +53,10 @@
z-index: 0;
background: var(--l1-background);
// Column so the bottom strip sits under the scrolling content, not inside it.
display: flex;
flex-direction: column;
&.full-screen-content {
width: 100%;
}
@@ -70,7 +74,9 @@
.chat-support-gateway {
position: fixed;
bottom: 20px;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: calc(20px + var(--bottom-strip-height, 0px));
right: 20px;
z-index: 1000;

View File

@@ -43,6 +43,7 @@ import { USER_PREFERENCES } from 'constants/userPreferences';
import AIAssistantModal from 'container/AIAssistant/AIAssistantModal';
import AIAssistantPanel from 'container/AIAssistant/AIAssistantPanel';
import { useAIAssistantStore } from 'container/AIAssistant/store/useAIAssistantStore';
import BottomStrip from 'container/BottomStrip';
import SideNav from 'container/SideNav';
import TopNav from 'container/TopNav';
import dayjs from 'dayjs';
@@ -51,6 +52,7 @@ import { useIsDarkMode } from 'hooks/useDarkMode';
import { useGetTenantLicense } from 'hooks/useGetTenantLicense';
import { useIsAIAssistantEnabled } from 'hooks/useIsAIAssistantEnabled';
import { useNotifications } from 'hooks/useNotifications';
import { useSavedViewEnabled } from 'hooks/useSavedViewEnabled';
import useTabVisibility from 'hooks/useTabFocus';
import history from 'lib/history';
import { isNull } from 'lodash-es';
@@ -402,6 +404,7 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
}, [pathname]);
const isToDisplayLayout = isLoggedIn;
const isSavedViewEnabled = useSavedViewEnabled();
const routeKey = useMemo(() => getRouteKey(pathname), [pathname]);
const pageTitle = t(routeKey);
@@ -868,6 +871,10 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
</OverlayScrollbar>
</LayoutContent>
</Sentry.ErrorBoundary>
{isSavedViewEnabled && isToDisplayLayout && !renderFullScreen && (
<BottomStrip />
)}
</div>
{isLoggedIn && isAIAssistantEnabled && (

View File

@@ -12,8 +12,12 @@ export const Layout = styled(LayoutComponent)`
}
`;
// Takes the height left in `.app-content` after the bottom strip.
// `min-height: 0` is not needed right now, overlayscrollbars already sets
// `overflow: auto` here. Kept so this does not break if that goes away.
export const LayoutContent = styled(LayoutComponent.Content)`
height: 100%;
flex: 1;
min-height: 0;
&::-webkit-scrollbar {
width: 0.1rem;
}

View File

@@ -0,0 +1,36 @@
.strip {
display: flex;
align-items: center;
justify-content: space-between;
gap: var(--spacing-6);
flex-shrink: 0;
height: var(--bottom-strip-height);
padding: 0 var(--spacing-6);
background: var(--l2-background);
border-top: 1px solid var(--l2-border);
font-family: var(--font-family-sf-mono, monospace);
// Above page content, below the body-portalled overlays that are meant to
// cover the strip.
position: relative;
z-index: 1;
}
.left,
.right {
display: flex;
align-items: center;
gap: var(--spacing-6);
min-width: 0;
}
// Temporary placeholder for the left slot. Replaced later.
.version {
color: var(--l2-foreground);
white-space: nowrap;
overflow: hidden;
text-overflow: ellipsis;
}

View File

@@ -0,0 +1,49 @@
import { render } from 'tests/test-utils';
import BottomStrip, {
BOTTOM_STRIP_HEIGHT,
BOTTOM_STRIP_HEIGHT_VAR,
BOTTOM_STRIP_ON_CLASS,
} from '..';
describe('BottomStrip', () => {
it('publishes the body class and height property while mounted', () => {
const { unmount } = render(<BottomStrip />);
expect(document.body.classList.contains(BOTTOM_STRIP_ON_CLASS)).toBe(true);
expect(document.body.style.getPropertyValue(BOTTOM_STRIP_HEIGHT_VAR)).toBe(
`${BOTTOM_STRIP_HEIGHT}px`,
);
unmount();
expect(document.body.classList.contains(BOTTOM_STRIP_ON_CLASS)).toBe(false);
expect(document.body.style.getPropertyValue(BOTTOM_STRIP_HEIGHT_VAR)).toBe(
'',
);
});
// The string is whatever the Go build injected, so it is rendered untouched —
// same as SideNav. Release tags carry the "v", local builds do not.
it.each([['v0.134.67'], ['main-64f1c2a']])(
'renders the build version %p exactly as given',
(version) => {
const { getByTestId } = render(<BottomStrip />, undefined, {
appContextOverrides: {
versionData: { version, ee: 'Y', setupCompleted: true },
},
});
expect(getByTestId('bottom-strip-version')).toHaveTextContent(version);
},
);
it('renders the strip without a version when none is available', () => {
const { getByTestId, queryByTestId } = render(<BottomStrip />, undefined, {
appContextOverrides: { versionData: null },
});
expect(getByTestId('bottom-strip')).toBeInTheDocument();
expect(queryByTestId('bottom-strip-version')).not.toBeInTheDocument();
});
});

View File

@@ -0,0 +1,42 @@
import { useLayoutEffect } from 'react';
import { useAppContext } from 'providers/App/App';
import styles from './BottomStrip.module.scss';
export const BOTTOM_STRIP_HEIGHT = 24;
export const BOTTOM_STRIP_ON_CLASS = 'bottom-strip-on';
export const BOTTOM_STRIP_HEIGHT_VAR = '--bottom-strip-height';
function BottomStrip(): JSX.Element {
const { versionData } = useAppContext();
const version = versionData?.version?.trim();
useLayoutEffect(() => {
document.body.classList.add(BOTTOM_STRIP_ON_CLASS);
document.body.style.setProperty(
BOTTOM_STRIP_HEIGHT_VAR,
`${BOTTOM_STRIP_HEIGHT}px`,
);
return (): void => {
document.body.classList.remove(BOTTOM_STRIP_ON_CLASS);
document.body.style.removeProperty(BOTTOM_STRIP_HEIGHT_VAR);
};
}, []);
return (
<div className={styles.strip} data-testid="bottom-strip">
<div className={styles.left}>
{version && (
<span className={styles.version} data-testid="bottom-strip-version">
{version}
</span>
)}
</div>
<div className={styles.right} />
</div>
);
}
export default BottomStrip;

View File

@@ -1,6 +1,8 @@
.create-alert-v2-footer {
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
left: 63px;
right: 0;
background-color: var(--l1-background);

View File

@@ -1,6 +1,8 @@
.explorer-options-container {
position: fixed;
bottom: 0px;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
left: calc(50% + 240px);
transform: translate(calc(-50% - 120px), 0);
transition: left 0.2s linear;

View File

@@ -1,6 +1,8 @@
.explorer-option-droppable-container {
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
width: -webkit-fill-available;
height: 24px;
display: flex;

View File

@@ -1,7 +1,6 @@
.home-container {
display: flex;
flex-direction: column;
min-height: 100vh;
overflow-y: auto;
height: 100%;
width: 100%;

View File

@@ -1,17 +1,18 @@
import { useEffect, useMemo, useState } from 'react';
import { Link } from 'react-router-dom';
import { Button, Skeleton } from 'antd';
import { Badge } from '@signozhq/ui/badge';
import logEvent from 'api/common/logEvent';
import { getViewDetailsUsingViewKey } from 'components/ExplorerCard/utils';
import { useListSavedViews } from 'api/generated/services/saved-view';
import {
SavedviewtypesSavedViewDTO,
SavedviewtypesSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
import ROUTES from 'constants/routes';
import { useGetAllViews } from 'hooks/saveViews/useGetAllViews';
import { getSavedViewQuery } from 'container/SavedViews/utils';
import { useHandleExplorerTabChange } from 'hooks/useHandleExplorerTabChange';
import { SOURCEPAGE_VS_ROUTES } from 'pages/SaveView/constants';
import Card from 'periscope/components/Card/Card';
import { useAppContext } from 'providers/App/App';
import { ViewProps } from 'types/api/saveViews/types';
import { DataSource } from 'types/common/queryBuilder';
import { USER_ROLES } from 'types/roles';
import floppyDiscUrl from '@/assets/Icons/floppy-disc.svg';
@@ -35,38 +36,40 @@ export default function SavedViews({
}): JSX.Element {
const { user } = useAppContext();
const [selectedEntity, setSelectedEntity] = useState<string>('logs');
const [selectedEntityViews, setSelectedEntityViews] = useState<any[]>([]);
const [selectedEntityViews, setSelectedEntityViews] = useState<
SavedviewtypesSavedViewDTO[]
>([]);
const {
data: logsViewsData,
isLoading: logsViewsLoading,
isError: logsViewsError,
} = useGetAllViews(DataSource.LOGS);
} = useListSavedViews({ source: SavedviewtypesSourceDTO.logs });
const {
data: tracesViewsData,
isLoading: tracesViewsLoading,
isError: tracesViewsError,
} = useGetAllViews(DataSource.TRACES);
} = useListSavedViews({ source: SavedviewtypesSourceDTO.traces });
const {
data: metricsViewsData,
isLoading: metricsViewsLoading,
isError: metricsViewsError,
} = useGetAllViews(DataSource.METRICS);
} = useListSavedViews({ source: SavedviewtypesSourceDTO.metrics });
const logsViews = useMemo(
() => [...(logsViewsData?.data.data || [])],
() => [...(logsViewsData?.data || [])],
[logsViewsData],
);
const tracesViews = useMemo(
() => [...(tracesViewsData?.data.data || [])],
() => [...(tracesViewsData?.data || [])],
[tracesViewsData],
);
const metricsViews = useMemo(
() => [...(metricsViewsData?.data.data || [])],
() => [...(metricsViewsData?.data || [])],
[metricsViewsData],
);
@@ -88,39 +91,22 @@ export default function SavedViews({
const { handleExplorerTabChange } = useHandleExplorerTabChange();
const handleRedirectQuery = (view: ViewProps): void => {
const handleRedirectQuery = (view: SavedviewtypesSavedViewDTO): void => {
logEvent('Homepage: Saved view clicked', {
viewId: view.id,
viewName: view.name,
viewName: view.spec.displayName,
entity: selectedEntity,
});
let currentViews: ViewProps[] = [];
if (selectedEntity === 'logs') {
currentViews = logsViews;
} else if (selectedEntity === 'traces') {
currentViews = tracesViews;
} else if (selectedEntity === 'metrics') {
currentViews = metricsViews;
}
const currentViewDetails = getViewDetailsUsingViewKey(view.id, currentViews);
if (!currentViewDetails) {
return;
}
const { query, name, id, panelType: currentPanelType } = currentViewDetails;
if (selectedEntity) {
handleExplorerTabChange(
currentPanelType,
{
query,
viewName: name,
viewKey: id,
},
SOURCEPAGE_VS_ROUTES[selectedEntity],
);
}
handleExplorerTabChange(
view.spec.panelType,
{
query: getSavedViewQuery(view),
viewName: view.spec.displayName,
viewKey: view.id,
},
SOURCEPAGE_VS_ROUTES[selectedEntity],
);
};
useEffect(() => {
@@ -239,24 +225,10 @@ export default function SavedViews({
/>
<div className="saved-view-item-name home-data-item-name">
{view.name}
{view.spec.displayName}
</div>
</div>
<div className="saved-view-item-description home-data-item-tag">
{view.tags?.map((tag: string) => {
if (tag === '') {
return null;
}
return (
<Badge color="sienna" key={tag}>
{tag}
</Badge>
);
})}
</div>
<Button
type="link"
size="small"
@@ -307,7 +279,7 @@ export default function SavedViews({
logEvent('Homepage: Saved views switched', {
tab,
});
let currentViews: ViewProps[] = [];
let currentViews: SavedviewtypesSavedViewDTO[] = [];
if (tab === 'logs') {
currentViews = logsViews;
} else if (tab === 'traces') {

View File

@@ -1,7 +1,4 @@
.licenses-page {
max-height: 100vh;
overflow: hidden;
.licenses-page-header {
border-bottom: 1px solid var(--l1-border);
background: var(--l1-background);
@@ -32,7 +29,6 @@
.licenses-page-content {
flex: 1;
height: calc(100vh - 48px);
background: var(--l1-background);
padding: 10px 8px;
overflow-y: auto;

View File

@@ -2,7 +2,7 @@
display: flex;
flex-direction: column;
gap: 1rem;
height: calc(100vh - 62px);
flex: 1;
min-height: 400px;
}

View File

@@ -181,7 +181,9 @@
.ant-pagination {
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new
// fixed-bottom UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
width: calc(100% - 54px);
background: var(--l1-background);
padding: 16px;

View File

@@ -0,0 +1,126 @@
import {
SavedviewtypesPanelTypeDTO,
SavedviewtypesSavedViewDTO,
SavedviewtypesSchemaVersionDTO,
SavedviewtypesSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
import { PANEL_TYPES } from 'constants/queryBuilder';
import { EQueryType } from 'types/common/dashboard';
import { DataSource } from 'types/common/queryBuilder';
import { findSavedView, getSavedViewQuery, toSavedViewSource } from '../utils';
jest.mock('uuid', () => ({
v4: (): string => 'test-id',
}));
function makeView(): SavedviewtypesSavedViewDTO {
return {
id: 'view-1',
name: 'errors-by-service-abc123',
source: SavedviewtypesSourceDTO.traces,
schemaVersion: SavedviewtypesSchemaVersionDTO.v2,
createdBy: 'a@b.c',
updatedBy: 'a@b.c',
spec: {
displayName: 'Errors by service',
panelType: SavedviewtypesPanelTypeDTO.list,
requestType: 'raw',
queries: [
{
type: 'builder_query',
spec: {
name: 'A',
signal: 'traces',
stepInterval: 60,
filter: { expression: 'has_error = true' },
// v2 reads back fully defaulted envelopes; nulls must not break the mapper
groupBy: null,
order: null,
selectFields: null,
functions: null,
legend: '',
disabled: false,
},
},
],
selectedFields: [{ name: 'service.name' }],
display: { color: 'red' },
},
} as SavedviewtypesSavedViewDTO;
}
describe('getSavedViewQuery', () => {
it('maps the v2 spec through the v5 branch of mapQueryDataFromApi', () => {
const query = getSavedViewQuery(makeView());
expect(query.queryType).toBe(EQueryType.QUERY_BUILDER);
expect(query.promql).toStrictEqual([]);
expect(query.clickhouse_sql).toStrictEqual([]);
expect(query.builder.queryData).toHaveLength(1);
const [queryData] = query.builder.queryData;
expect(queryData.queryName).toBe('A');
expect(queryData.dataSource).toBe(DataSource.TRACES);
expect(queryData.filter).toStrictEqual({ expression: 'has_error = true' });
expect(queryData.groupBy).toStrictEqual([]);
expect(queryData.orderBy).toStrictEqual([]);
});
it('keeps formulas alongside builder queries', () => {
const view = makeView();
view.spec.queries.push({
type: 'builder_formula',
spec: { name: 'F1', expression: 'A / 2' },
} as SavedviewtypesSavedViewDTO['spec']['queries'][number]);
const query = getSavedViewQuery(view);
expect(query.builder.queryData).toHaveLength(1);
expect(query.builder.queryFormulas).toHaveLength(1);
expect(query.builder.queryFormulas[0].queryName).toBe('F1');
});
it('does not read the panel type into the query', () => {
const view = makeView();
view.spec.panelType = SavedviewtypesPanelTypeDTO.graph;
const query = getSavedViewQuery(view);
// panelType travels separately (url param), the Query itself has no such field
expect(query).not.toHaveProperty('panelType', PANEL_TYPES.TIME_SERIES);
});
});
describe('toSavedViewSource', () => {
it('maps every explorer source page to the v2 source', () => {
expect(toSavedViewSource(DataSource.LOGS)).toBe(SavedviewtypesSourceDTO.logs);
expect(toSavedViewSource(DataSource.TRACES)).toBe(
SavedviewtypesSourceDTO.traces,
);
expect(toSavedViewSource(DataSource.METRICS)).toBe(
SavedviewtypesSourceDTO.metrics,
);
expect(toSavedViewSource('meter')).toBe(SavedviewtypesSourceDTO.meter);
});
});
describe('findSavedView', () => {
const views = [
{ ...makeView(), id: 'a' },
{ ...makeView(), id: 'b' },
];
it('returns the view with the matching id', () => {
expect(findSavedView(views, 'b')?.id).toBe('b');
});
it('returns undefined when the id is not in the list', () => {
expect(findSavedView(views, 'c')).toBeUndefined();
});
it('returns undefined for a null or not yet loaded list', () => {
expect(findSavedView(null, 'a')).toBeUndefined();
expect(findSavedView(undefined, 'a')).toBeUndefined();
});
});

View File

@@ -0,0 +1,49 @@
import {
SavedviewtypesSavedViewDTO,
SavedviewtypesSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
import { PANEL_TYPES } from 'constants/queryBuilder';
import { mapQueryDataFromApi } from 'lib/newQueryBuilder/queryBuilderMappers/mapQueryDataFromApi';
import { Query } from 'types/api/queryBuilder/queryBuilderData';
import { QueryEnvelope } from 'types/api/v5/queryRange';
import { EQueryType } from 'types/common/dashboard';
import { DataSource } from 'types/common/queryBuilder';
export type SavedViewSourcePage = DataSource | 'meter';
// Explorers and the preferences module are keyed by DataSource (the signal),
// the api keys views by source page. Same values today, so this is the one
// place they meet. AI observability views will come with their own source and
// DataSource cannot tell them apart from traces, so preferences should move to
// source page at that point and this map goes with it.
const SAVED_VIEW_SOURCE: Record<SavedViewSourcePage, SavedviewtypesSourceDTO> =
{
[DataSource.LOGS]: SavedviewtypesSourceDTO.logs,
[DataSource.TRACES]: SavedviewtypesSourceDTO.traces,
[DataSource.METRICS]: SavedviewtypesSourceDTO.metrics,
meter: SavedviewtypesSourceDTO.meter,
};
export function toSavedViewSource(
sourcePage: SavedViewSourcePage,
): SavedviewtypesSourceDTO {
return SAVED_VIEW_SOURCE[sourcePage];
}
// Explorers only save builder queries; v2 carries no queryType, so it is fixed here.
export function getSavedViewQuery(view: SavedviewtypesSavedViewDTO): Query {
const { queries, panelType } = view.spec;
return mapQueryDataFromApi({
queries: queries as QueryEnvelope[],
panelType: panelType as unknown as PANEL_TYPES,
queryType: EQueryType.QUERY_BUILDER,
unit: undefined,
});
}
export function findSavedView(
views: SavedviewtypesSavedViewDTO[] | null | undefined,
id: string,
): SavedviewtypesSavedViewDTO | undefined {
return views?.find((view) => view.id === id);
}

View File

@@ -2,7 +2,7 @@
display: flex;
flex-direction: column;
gap: 1rem;
height: calc(100vh - 62px);
flex: 1;
min-height: 400px;
padding-top: var(--spacing-8);
}

View File

@@ -1,7 +1,4 @@
.version-container {
max-height: 100vh;
overflow: hidden;
.version-page-header {
border-bottom: 1px solid var(--l1-border);
background: var(--l1-background);

View File

@@ -1,11 +1,18 @@
import { useMutation, UseMutationResult } from 'react-query';
import { useMutation, UseMutationResult, useQueryClient } from 'react-query';
import { invalidateListSavedViews } from 'api/generated/services/saved-view';
import { deleteView } from 'api/saveView/deleteView';
import { DeleteViewPayloadProps } from 'types/api/saveViews/types';
export const useDeleteView = (
uuid: string,
): UseMutationResult<DeleteViewPayloadProps, Error, string> =>
useMutation({
): UseMutationResult<DeleteViewPayloadProps, Error, string> => {
const queryClient = useQueryClient();
return useMutation({
mutationKey: [uuid],
mutationFn: () => deleteView(uuid),
// v1 and v2 share storage; consumers already on v2 must see this write.
// Temporary till the v1 client is deleted with the explorer bar.
onSuccess: () => invalidateListSavedViews(queryClient),
});
};

View File

@@ -1,4 +1,5 @@
import { useMutation, UseMutationResult } from 'react-query';
import { useMutation, UseMutationResult, useQueryClient } from 'react-query';
import { invalidateListSavedViews } from 'api/generated/services/saved-view';
import { saveView } from 'api/saveView/saveView';
import { AxiosResponse } from 'axios';
import { SaveViewPayloadProps, SaveViewProps } from 'types/api/saveViews/types';
@@ -13,8 +14,14 @@ export const useSaveView = ({
Error,
SaveViewProps,
SaveViewPayloadProps
> =>
useMutation({
> => {
const queryClient = useQueryClient();
return useMutation({
mutationKey: [viewName, sourcePage, compositeQuery, extraData],
mutationFn: saveView,
// v1 and v2 share storage; consumers already on v2 must see this write.
// Temporary till the v1 client is deleted with the explorer bar.
onSuccess: () => invalidateListSavedViews(queryClient),
});
};

View File

@@ -1,4 +1,5 @@
import { useMutation, UseMutationResult } from 'react-query';
import { useMutation, UseMutationResult, useQueryClient } from 'react-query';
import { invalidateListSavedViews } from 'api/generated/services/saved-view';
import { updateView } from 'api/saveView/updateView';
import {
UpdateViewPayloadProps,
@@ -16,8 +17,10 @@ export const useUpdateView = ({
Error,
UpdateViewProps,
UpdateViewPayloadProps
> =>
useMutation({
> => {
const queryClient = useQueryClient();
return useMutation({
mutationKey: [viewName, sourcePage, compositeQuery, extraData],
mutationFn: () =>
updateView({
@@ -27,4 +30,8 @@ export const useUpdateView = ({
sourcePage,
viewKey,
}),
// v1 and v2 share storage; consumers already on v2 must see this write.
// Temporary till the v1 client is deleted with the explorer bar.
onSuccess: () => invalidateListSavedViews(queryClient),
});
};

View File

@@ -0,0 +1,11 @@
import getLocalStorageKey from 'api/browser/localstorage/get';
import { LOCALSTORAGE } from 'constants/localStorage';
import { useState } from 'react';
export function useSavedViewEnabled(): boolean {
const [isEnabled] = useState(
() => getLocalStorageKey(LOCALSTORAGE.SAVED_VIEW_ENABLED) === 'true',
);
return isEnabled;
}

View File

@@ -1,4 +1,29 @@
.alerts-container {
// Hands the page height down to the active tab so its content can bound itself
// instead of guessing with 100vh. Child combinators only, nested Tabs
// (Configuration) must not be caught.
flex: 1;
min-height: 0;
> .ant-tabs-content-holder {
display: flex;
flex-direction: column;
> .ant-tabs-content {
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
> .ant-tabs-tabpane-active {
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}
}
}
.top-level-tab.periscope-tab {
padding: 2px 0;
}
@@ -40,5 +65,9 @@
.alert-rules-container {
margin-top: 10px;
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}
}

View File

@@ -8,6 +8,7 @@ import {
QuickfiltertypesSourceDTO,
TelemetrytypesFieldContextDTO,
TelemetrytypesFieldDataTypeDTO,
TelemetrytypesSignalDTO,
} from 'api/generated/services/sigNoz.schemas';
import ROUTES from 'constants/routes';
import { VIEWS } from 'container/ApiMonitoring/Explorer/Domains/DomainDetails/constants';
@@ -24,7 +25,10 @@ import {
toggleControl,
} from '@/storybook/controls/controls';
import { defineStoryMocks } from '@/storybook/controls/defineStoryMocks';
import { fieldValuesResponse } from '@/storybook/msw/__story_mockdata__/fields';
import {
fieldKeysResponse,
fieldValuesResponse,
} from '@/storybook/msw/__story_mockdata__/fields';
import { quickFiltersResponse } from '@/storybook/msw/__story_mockdata__/quickFilters';
import {
@@ -317,6 +321,21 @@ export const apiMonitoringMocks = defineStoryMocks({
})),
),
rest.get(
'http://localhost/api/v1/fields/keys',
response.json((req) =>
fieldKeysResponse(
groupByAttributeKeys(req.url.searchParams.get('searchText') ?? '').map(
({ key }) => key,
),
{
signal: TelemetrytypesSignalDTO.traces,
fieldContext: TelemetrytypesFieldContextDTO.attribute,
},
),
),
),
rest.get(
'http://localhost/api/v1/fields/values',
response.json((req) =>

View File

@@ -2,7 +2,9 @@
display: flex;
flex-direction: column;
position: fixed;
bottom: 0;
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
// UI belongs in the bounded layout, not in another offset here.
bottom: var(--bottom-strip-height, 0px);
left: 0;
width: 100%;
z-index: 100;

View File

@@ -164,10 +164,10 @@ export const homeMocks = defineStoryMocks({
),
rest.get(
'http://localhost/api/v1/explorer/views',
'http://localhost/api/v2/saved_views',
response.json((req) => {
const sourcePage = req.url.searchParams.get('sourcePage') ?? 'logs';
const signal = isSavedViewSignal(sourcePage) ? sourcePage : 'logs';
const source = req.url.searchParams.get('source') ?? 'logs';
const signal = isSavedViewSignal(source) ? source : 'logs';
return savedViewsResponse(
values.savedViewSignals.includes(signal) ? values.savedViews : 0,

View File

@@ -6,10 +6,21 @@
import { FeatureKeys } from 'constants/features';
import { ORG_PREFERENCES } from 'constants/orgPreferences';
import { checkListStepToPreferenceKeyMap } from 'container/Home/constants';
import type { RuletypesRuleDTO } from 'api/generated/services/sigNoz.schemas';
import {
type ListSavedViews200,
Querybuildertypesv5QueryBuilderQueryGithubComSigNozSignozPkgTypesQuerybuildertypesQuerybuildertypesv5LogAggregationDTOSignal as LogsSignal,
Querybuildertypesv5QueryBuilderQueryGithubComSigNozSignozPkgTypesQuerybuildertypesQuerybuildertypesv5MetricAggregationDTOSignal as MetricsSignal,
Querybuildertypesv5QueryBuilderQueryGithubComSigNozSignozPkgTypesQuerybuildertypesQuerybuildertypesv5TraceAggregationDTOSignal as TracesSignal,
Querybuildertypesv5QueryEnvelopeBuilderDTOType,
type Querybuildertypesv5QueryEnvelopeDTO,
Querybuildertypesv5RequestTypeDTO,
type RuletypesRuleDTO,
SavedviewtypesPanelTypeDTO,
SavedviewtypesSchemaVersionDTO,
SavedviewtypesSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
import type { ServiceDataProps } from 'api/metrics/getTopLevelOperations';
import { alertRulesFixture } from 'mocks-server/__mockdata__/alert_rules';
import { explorerView } from 'mocks-server/__mockdata__/explorer_views';
import { defaultFeatureFlags } from 'tests/fixtures/appContextMock';
import type { FeatureFlagProps } from 'types/api/features/getFeaturesFlags';
import type { MetricRangePayloadV3 } from 'types/api/metrics/getQueryRange';
@@ -165,20 +176,53 @@ const VIEW_NAMES: Record<SavedViewSignal, string[]> = {
export const isSavedViewSignal = (value: string): value is SavedViewSignal =>
SAVED_VIEW_SIGNALS.includes(value as SavedViewSignal);
const SAVED_VIEW_SOURCE: Record<SavedViewSignal, SavedviewtypesSourceDTO> = {
logs: SavedviewtypesSourceDTO.logs,
traces: SavedviewtypesSourceDTO.traces,
metrics: SavedviewtypesSourceDTO.metrics,
};
const SAVED_VIEW_QUERY: Record<
SavedViewSignal,
Querybuildertypesv5QueryEnvelopeDTO
> = {
logs: {
type: Querybuildertypesv5QueryEnvelopeBuilderDTOType.builder_query,
spec: { name: 'A', signal: LogsSignal.logs },
},
traces: {
type: Querybuildertypesv5QueryEnvelopeBuilderDTOType.builder_query,
spec: { name: 'A', signal: TracesSignal.traces },
},
metrics: {
type: Querybuildertypesv5QueryEnvelopeBuilderDTOType.builder_query,
spec: { name: 'A', signal: MetricsSignal.metrics },
},
};
export const savedViewsResponse = (
count: number,
sourcePage: SavedViewSignal,
): Record<string, unknown> => {
const names = VIEW_NAMES[sourcePage];
signal: SavedViewSignal,
): ListSavedViews200 => {
const names = VIEW_NAMES[signal];
return {
status: 'success',
data: Array.from({ length: Math.min(count, names.length) }, (_, index) => ({
...explorerView.data[0],
id: `storybook-${sourcePage}-view-${index + 1}`,
name: names[index],
sourcePage,
tags: [sourcePage],
id: `storybook-${signal}-view-${index + 1}`,
name: `storybook-${signal}-view-${index + 1}`,
source: SAVED_VIEW_SOURCE[signal],
schemaVersion: SavedviewtypesSchemaVersionDTO.v2,
createdAt: '2026-08-20T09:00:00Z',
createdBy: 'storybook@signoz.io',
updatedAt: '2026-08-20T09:00:00Z',
updatedBy: 'storybook@signoz.io',
spec: {
displayName: names[index],
panelType: SavedviewtypesPanelTypeDTO.list,
requestType: Querybuildertypesv5RequestTypeDTO.raw,
queries: [SAVED_VIEW_QUERY[signal]],
},
})),
};
};

View File

@@ -1,7 +1,4 @@
.support-page-container {
max-height: 100vh;
overflow: hidden;
.support-page-header {
border-bottom: 1px solid var(--l1-border);
background: var(--l1-background);

View File

@@ -1,5 +1,6 @@
.root {
height: calc(100vh);
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}

View File

@@ -1,13 +1,24 @@
.traces-funnel-details {
display: flex;
// 45px -> height of the tab bar
height: calc(100vh - 45px);
height: 100%;
&__steps-config {
flex-shrink: 0;
width: 600px;
border-right: 1px solid var(--l1-border);
// Positioning context for the absolute .steps-footer.
position: relative;
display: flex;
flex-direction: column;
// Scoped here so the modal usage of FunnelConfiguration on trace details
// stays in normal flow.
.funnel-configuration {
flex: 1;
min-height: 0;
display: flex;
flex-direction: column;
}
}
&__steps-results {
width: 100%;

View File

@@ -4,14 +4,17 @@
flex-direction: column;
justify-content: flex-start;
&.funnel-details-page {
height: calc(
100vh - 170px
); // 64px bottom bar + 61px configuration header + 45px page navbar
flex: 1;
min-height: 0;
// .steps-footer is absolute against the config column, so its 64px is
// reserved rather than laid out.
margin-bottom: 64px;
overflow: auto;
}
}
&__header {
flex-shrink: 0;
display: flex;
align-items: center;
justify-content: space-between;

View File

@@ -0,0 +1,202 @@
import { renderHook } from '@testing-library/react';
import { useListSavedViews } from 'api/generated/services/saved-view';
import {
SavedviewtypesSavedViewDTO,
SavedviewtypesSourceDTO,
} from 'api/generated/services/sigNoz.schemas';
import {
defaultLogsSelectedColumns,
defaultTraceSelectedColumns,
ensureLogsRequiredColumns,
} from 'container/OptionsMenu/constants';
import { DataSource } from 'types/common/queryBuilder';
import { usePreferenceSync } from '../sync/usePreferenceSync';
import { PreferenceMode } from '../types';
jest.mock('api/generated/services/saved-view');
const loaderPreferences = { columns: [{ name: 'from-loader' }] };
jest.mock('../loader/usePreferenceLoader', () => ({
usePreferenceLoader: jest.fn(() => ({
preferences: loaderPreferences,
loading: false,
error: null,
})),
}));
jest.mock('../updater/usePreferenceUpdater', () => ({
usePreferenceUpdater: jest.fn(() => ({
updateColumns: jest.fn(),
updateFormatting: jest.fn(),
})),
}));
const mockedUseListSavedViews = useListSavedViews as jest.MockedFunction<
typeof useListSavedViews
>;
function makeView(
id: string,
source: SavedviewtypesSourceDTO,
spec: Partial<SavedviewtypesSavedViewDTO['spec']>,
): SavedviewtypesSavedViewDTO {
return {
id,
source,
schemaVersion: 'v2',
spec: {
displayName: id,
panelType: 'list',
requestType: 'raw',
queries: [],
...spec,
},
} as unknown as SavedviewtypesSavedViewDTO;
}
function mockViews(views: SavedviewtypesSavedViewDTO[]): void {
mockedUseListSavedViews.mockReturnValue({
data: { status: 'success', data: views },
} as unknown as ReturnType<typeof useListSavedViews>);
}
describe('usePreferenceSync in saved view mode', () => {
beforeEach(() => {
mockedUseListSavedViews.mockReset();
});
it('fetches the list for the data source only in saved view mode', () => {
mockViews([]);
renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.DIRECT,
dataSource: DataSource.LOGS,
savedViewId: undefined,
}),
);
expect(mockedUseListSavedViews).toHaveBeenCalledWith(
{ source: 'logs' },
{ query: { enabled: false } },
);
});
it('returns loader preferences outside saved view mode', () => {
mockViews([]);
const { result } = renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.DIRECT,
dataSource: DataSource.LOGS,
savedViewId: undefined,
}),
);
expect(result.current.preferences).toBe(loaderPreferences);
});
it('applies selectedFields and display of the active logs view', () => {
mockViews([
makeView('view-1', SavedviewtypesSourceDTO.logs, {
selectedFields: [{ name: 'service.name' }, { name: 'body' }],
display: { maxLines: 3, format: 'raw', fontSize: 'large', color: 'red' },
}),
]);
const { result } = renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.SAVED_VIEW,
dataSource: DataSource.LOGS,
savedViewId: 'view-1',
}),
);
expect(result.current.preferences?.columns).toStrictEqual(
ensureLogsRequiredColumns([{ name: 'service.name' }, { name: 'body' }]),
);
expect(result.current.preferences?.formatting).toStrictEqual({
maxLines: 3,
format: 'raw',
fontSize: 'large',
version: 1,
});
});
it('falls back to defaults when the view has zero-valued display and no fields', () => {
mockViews([
makeView('view-1', SavedviewtypesSourceDTO.logs, {
selectedFields: undefined,
display: { maxLines: 0, format: '', fontSize: '', color: '' },
}),
]);
const { result } = renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.SAVED_VIEW,
dataSource: DataSource.LOGS,
savedViewId: 'view-1',
}),
);
expect(result.current.preferences?.columns).toStrictEqual(
ensureLogsRequiredColumns(defaultLogsSelectedColumns),
);
expect(result.current.preferences?.formatting).toStrictEqual({
maxLines: 1,
format: 'table',
fontSize: 'small',
version: 1,
});
});
it('passes trace selectedFields through and defaults when absent', () => {
mockViews([
makeView('with-fields', SavedviewtypesSourceDTO.traces, {
selectedFields: [{ name: 'name' }, { name: 'durationNano' }],
}),
makeView('without-fields', SavedviewtypesSourceDTO.traces, {}),
]);
const withFields = renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.SAVED_VIEW,
dataSource: DataSource.TRACES,
savedViewId: 'with-fields',
}),
);
const withoutFields = renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.SAVED_VIEW,
dataSource: DataSource.TRACES,
savedViewId: 'without-fields',
}),
);
expect(withFields.result.current.preferences?.columns).toStrictEqual([
{ name: 'name' },
{ name: 'durationNano' },
]);
expect(withFields.result.current.preferences?.formatting).toBeUndefined();
expect(withoutFields.result.current.preferences?.columns).toBe(
defaultTraceSelectedColumns,
);
});
it('uses defaults when the saved view id is not in the list', () => {
mockViews([makeView('other', SavedviewtypesSourceDTO.logs, {})]);
const { result } = renderHook(() =>
usePreferenceSync({
mode: PreferenceMode.SAVED_VIEW,
dataSource: DataSource.LOGS,
savedViewId: 'missing',
}),
);
expect(result.current.preferences?.columns).toStrictEqual(
ensureLogsRequiredColumns(defaultLogsSelectedColumns),
);
});
});

View File

@@ -1,12 +1,14 @@
/* eslint-disable sonarjs/cognitive-complexity */
import { useEffect, useState } from 'react';
import { useListSavedViews } from 'api/generated/services/saved-view';
import { TelemetryFieldKey } from 'api/v5/v5';
import {
defaultLogsSelectedColumns,
defaultTraceSelectedColumns,
ensureLogsRequiredColumns,
} from 'container/OptionsMenu/constants';
import { defaultSelectedColumns as defaultTracesSelectedColumns } from 'container/TracesExplorer/ListView/configs';
import { useGetAllViews } from 'hooks/saveViews/useGetAllViews';
import { FontSize, LogViewMode } from 'container/OptionsMenu/types';
import { findSavedView, toSavedViewSource } from 'container/SavedViews/utils';
import { DataSource } from 'types/common/queryBuilder';
import { usePreferenceLoader } from '../loader/usePreferenceLoader';
@@ -28,16 +30,16 @@ export function usePreferenceSync({
updateColumns: (newColumns: TelemetryFieldKey[]) => void;
updateFormatting: (newFormatting: FormattingOptions) => void;
} {
const { data: viewsData } = useGetAllViews(
dataSource,
mode === PreferenceMode.SAVED_VIEW,
const { data: viewsData } = useListSavedViews(
{ source: toSavedViewSource(dataSource) },
{ query: { enabled: mode === PreferenceMode.SAVED_VIEW } },
);
const [savedViewPreferences, setSavedViewPreferences] =
useState<Preferences | null>(null);
const updateExtraDataSelectColumns = (
columns: TelemetryFieldKey[],
const withColumnNames = (
columns: TelemetryFieldKey[] | undefined,
): TelemetryFieldKey[] | null => {
if (!columns) {
return null;
@@ -49,27 +51,28 @@ export function usePreferenceSync({
};
useEffect(() => {
const extraData = viewsData?.data?.data?.find(
(view) => view.id === savedViewId,
)?.extraData;
const spec = savedViewId
? findSavedView(viewsData?.data, savedViewId)?.spec
: undefined;
const selectedFields = spec?.selectedFields as
| TelemetryFieldKey[]
| undefined;
const parsedExtraData = JSON.parse(extraData || '{}');
let columns: TelemetryFieldKey[] = [];
let formatting: FormattingOptions | undefined;
if (dataSource === DataSource.LOGS) {
columns = ensureLogsRequiredColumns(
updateExtraDataSelectColumns(parsedExtraData?.selectColumns) ||
defaultLogsSelectedColumns,
withColumnNames(selectedFields) || defaultLogsSelectedColumns,
);
formatting = {
maxLines: parsedExtraData?.maxLines ?? 1,
format: parsedExtraData?.format ?? 'table',
fontSize: parsedExtraData?.fontSize ?? 'small',
version: parsedExtraData?.version ?? 1,
maxLines: spec?.display?.maxLines || 1,
format: (spec?.display?.format as LogViewMode) || 'table',
fontSize: (spec?.display?.fontSize as FontSize) || FontSize.SMALL,
version: 1,
};
}
if (dataSource === DataSource.TRACES) {
columns = parsedExtraData?.selectColumns || defaultTracesSelectedColumns;
columns = selectedFields || defaultTraceSelectedColumns;
}
setSavedViewPreferences({ columns, formatting });
}, [viewsData, dataSource, savedViewId, mode]);

View File

@@ -145,6 +145,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusCreated,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusConflict},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbCreate)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
@@ -173,6 +174,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbList)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
@@ -199,6 +201,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbRead)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
@@ -226,6 +229,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbUpdate)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
@@ -253,6 +257,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbDelete)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
@@ -281,6 +286,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusOK,
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbUpdate)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{
@@ -308,6 +314,7 @@ func (provider *provider) addAlertmanagerRoutes(router *mux.Router) error {
SuccessStatusCode: http.StatusNoContent,
ErrorStatusCodes: []int{http.StatusBadRequest},
Deprecated: false,
Stability: handler.StabilityDevelopment,
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceNotificationChannel.Scope(coretypes.VerbCreate)}),
},
handler.WithResourceDefs(handler.BasicResourceDef{

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

@@ -0,0 +1,75 @@
package handler
import (
"net/http"
"testing"
"github.com/gorilla/mux"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/swaggest/openapi-go"
"github.com/swaggest/openapi-go/openapi3"
)
type bespokeOpenAPIHandler struct{}
func (bespokeOpenAPIHandler) ServeHTTP(http.ResponseWriter, *http.Request) {}
func (bespokeOpenAPIHandler) ServeOpenAPI(opCtx openapi.OperationContext) {
opCtx.SetID("Bespoke")
opCtx.AddRespStructure(nil, openapi.WithHTTPStatus(http.StatusOK))
}
func (bespokeOpenAPIHandler) ResourceDefs() []ResourceDef { return nil }
func TestAttachStabilities(t *testing.T) {
router := mux.NewRouter()
router.Handle("/development", New(func(http.ResponseWriter, *http.Request) {}, OpenAPIDef{ID: "Development", SuccessStatusCode: http.StatusOK, Stability: StabilityDevelopment})).Methods(http.MethodGet)
router.Handle("/beta/{id}", New(func(http.ResponseWriter, *http.Request) {}, OpenAPIDef{ID: "Beta", SuccessStatusCode: http.StatusOK, Stability: StabilityBeta})).Methods(http.MethodPut)
router.Handle("/unset", New(func(http.ResponseWriter, *http.Request) {}, OpenAPIDef{ID: "Unset", SuccessStatusCode: http.StatusOK})).Methods(http.MethodGet)
router.Handle("/bespoke", bespokeOpenAPIHandler{}).Methods(http.MethodGet)
reflector := openapi3.NewReflector()
collector := NewOpenAPICollector(reflector)
require.NoError(t, router.Walk(collector.Walker))
collector.AttachStabilities(reflector.Spec)
testCases := []struct {
subtestName string
path string
method string
expectedExtensionValue any
}{
{
subtestName: "development handler",
path: "/development",
method: "get",
expectedExtensionValue: "development",
},
{
subtestName: "beta handler with path parameter",
path: "/beta/{id}",
method: "put",
expectedExtensionValue: "beta",
},
{
subtestName: "unset handler defaults to alpha",
path: "/unset",
method: "get",
expectedExtensionValue: "alpha",
},
{
subtestName: "handler built outside New defaults to alpha",
path: "/bespoke",
method: "get",
expectedExtensionValue: "alpha",
},
}
for _, testCase := range testCases {
t.Run(testCase.subtestName, func(t *testing.T) {
operation := reflector.Spec.Paths.MapOfPathItemValues[testCase.path].MapOfOperationValues[testCase.method]
assert.Equal(t, testCase.expectedExtensionValue, operation.MapOfAnything["x-signoz-stability"])
})
}
}

View File

@@ -1,14 +1,37 @@
package handler
import (
"net/http"
"reflect"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/gorilla/mux"
"github.com/swaggest/jsonschema-go"
openapigo "github.com/swaggest/openapi-go"
"github.com/swaggest/openapi-go/openapi3"
"github.com/swaggest/rest/openapi"
)
const signozStabilityKey string = "x-signoz-stability"
var (
StabilityDevelopment = Stability{valuer.NewString("development")}
StabilityAlpha = Stability{valuer.NewString("alpha")}
StabilityBeta = Stability{valuer.NewString("beta")}
StabilityStable = Stability{valuer.NewString("stable")}
)
// Stability is emitted as the x-signoz-stability extension on every operation; unset means alpha.
type Stability struct{ valuer.String }
func (stability Stability) StringValue() string {
if stability.IsZero() {
return StabilityAlpha.String.StringValue()
}
return stability.String.StringValue()
}
// OpenAPIExample is a named example for an OpenAPI operation.
type OpenAPIExample struct {
Name string
@@ -32,6 +55,7 @@ type OpenAPIDef struct {
SuccessStatusCode int
ErrorStatusCodes []int
Deprecated bool
Stability Stability
SecuritySchemes []OpenAPISecurityScheme
}
@@ -42,14 +66,16 @@ type OpenAPISecurityScheme struct {
// OpenAPICollector is a collector for OpenAPI operations.
type OpenAPICollector struct {
collector *openapi.Collector
collector *openapi.Collector
stabilities map[operationKey]Stability
}
func NewOpenAPICollector(reflector openapigo.Reflector) *OpenAPICollector {
c := openapi.NewCollector(reflector)
return &OpenAPICollector{
collector: c,
collector: c,
stabilities: make(map[operationKey]Stability),
}
}
@@ -77,6 +103,9 @@ func (c *OpenAPICollector) Walker(route *mux.Route, _ *mux.Router, _ []*mux.Rout
if err := c.collector.CollectOperation(method, path, c.collect(method, path, handler.ServeOpenAPI)); err != nil {
return err
}
if err := c.recordStability(method, path, httpHandler); err != nil {
return err
}
}
return nil
}
@@ -84,6 +113,17 @@ func (c *OpenAPICollector) Walker(route *mux.Route, _ *mux.Router, _ []*mux.Rout
return nil
}
// AttachStabilities stamps every operation in spec, so handlers built outside New
// carry the unset stability rather than none.
func (c *OpenAPICollector) AttachStabilities(spec *openapi3.Spec) {
for path, pathItem := range spec.Paths.MapOfPathItemValues {
for method, operation := range pathItem.MapOfOperationValues {
operation.WithMapOfAnythingItem(signozStabilityKey, c.stabilities[operationKey{method: method, path: path}].StringValue())
pathItem.MapOfOperationValues[method] = operation
}
}
}
func (c *OpenAPICollector) collect(method string, path string, serveOpenAPIFunc ServeOpenAPIFunc) func(oc openapigo.OperationContext) error {
return func(oc openapigo.OperationContext) error {
// Serve the OpenAPI documentation for the handler
@@ -117,3 +157,23 @@ func (c *OpenAPICollector) collect(method string, path string, serveOpenAPIFunc
return nil
}
}
func (c *OpenAPICollector) recordStability(method string, path string, httpHandler http.Handler) error {
generic, ok := httpHandler.(*handler)
if !ok {
return nil
}
cleanMethod, cleanPath, _, err := openapigo.SanitizeMethodPath(method, path)
if err != nil {
return err
}
c.stabilities[operationKey{method: cleanMethod, path: cleanPath}] = generic.openAPIDef.Stability
return nil
}
type operationKey struct {
method string
path string
}

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

@@ -94,10 +94,10 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumBool, schema.ColumnTypeEnumFloat64:
leftOperand := fmt.Sprintf("mapContains(%s, %s)", column.Name, clickhousesql.StringLiteral(key.Name))
if key.Materialized {
leftOperand = telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key)
return telemetrytypes.FieldKeyToMaterializedExistsCondition(key, exists), nil
}
leftOperand := fmt.Sprintf("mapContains(%s, %s)", column.Name, clickhousesql.StringLiteral(key.Name))
if exists {
return leftOperand, nil
}

View File

@@ -174,6 +174,7 @@ func (openapi *OpenAPI) CreateAndWrite(path string) error {
}
attachDiscriminators(openapi.reflector.Spec)
openapi.collector.AttachStabilities(openapi.reflector.Spec)
// The library's MarshalYAML does a JSON round-trip that converts all numbers
// to float64, causing large integers (e.g. epoch millisecond timestamps) to

View File

@@ -237,13 +237,13 @@ func TestBuild_FullSQL_TraceList_MaterializedColumns(t *testing.T) {
assertSQLEqual(t, `
WITH matched AS (
SELECT trace_id,
maxIf(timestamp, (attribute_string_gen_ai$$request$$model_exists OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name'))) AS last_activity_time
maxIf(timestamp, (attribute_string_gen_ai$$request$$model_exists = true OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name'))) AS last_activity_time
FROM signoz_traces.distributed_signoz_index_v3
WHERE timestamp >= '1747947419000000000'
AND timestamp < '1747983448000000000'
AND ts_bucket_start >= 1747945619
AND ts_bucket_start <= 1747983448
AND ((attribute_string_gen_ai$$request$$model_exists OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name')))
AND ((attribute_string_gen_ai$$request$$model_exists = true OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name')))
GROUP BY trace_id
ORDER BY last_activity_time DESC, trace_id DESC
LIMIT 20
@@ -268,16 +268,16 @@ SELECT trace_id,
count() AS span_count,
anyIf(name, parent_span_id = '') AS root_span_name,
any(multiIf(resource.service.name IS NOT NULL, resource.service.name::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS service.name,
countIf(attribute_string_gen_ai$$request$$model_exists) AS llm_call_count,
countIf(attribute_string_gen_ai$$request$$model_exists = true) AS llm_call_count,
countIf(mapContains(attributes_string, 'gen_ai.tool.name')) AS tool_call_count,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name'), attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name')) AS distinct_tool_count,
sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)) AS input_tokens,
sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists = true, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens'), toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
coalesce(sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens'), toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
coalesce(sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists = true, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens'), toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, 'signoz.gen_ai.usage.tokens.cost'), toFloat64(attributes_number['signoz.gen_ai.usage.tokens.cost']), NULL)) AS estimated_total_cost,
maxIf(duration_nano, attribute_string_gen_ai$$request$$model_exists) AS max_llm_duration_nano,
maxIf(duration_nano, attribute_string_gen_ai$$request$$model_exists = true) AS max_llm_duration_nano,
countIf(has_error = true) AS error_count,
maxIf(timestamp, (attribute_string_gen_ai$$request$$model_exists OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name'))) AS last_activity_time,
maxIf(timestamp, (attribute_string_gen_ai$$request$$model_exists = true OR mapContains(attributes_string, 'gen_ai.tool.name') OR mapContains(attributes_string, 'gen_ai.agent.name'))) AS last_activity_time,
argMinIf(multiIf(mapContains(attributes_string, 'gen_ai.input.messages'), attributes_string['gen_ai.input.messages'], NULL), timestamp, mapContains(attributes_string, 'gen_ai.input.messages')) AS input,
argMaxIf(multiIf(mapContains(attributes_string, 'gen_ai.output.messages'), attributes_string['gen_ai.output.messages'], NULL), timestamp, mapContains(attributes_string, 'gen_ai.output.messages')) AS output
FROM signoz_traces.distributed_signoz_index_v3

View File

@@ -92,7 +92,7 @@ func TestStatementBuilder(t *testing.T) {
Limit: 100,
},
expected: qbtypes.Statement{
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$principal$$id` = ? AND `attribute_string_signoz$$audit$$principal$$id_exists`) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$principal$$id` = ? AND `attribute_string_signoz$$audit$$principal$$id_exists` = true) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"019a-1234-abcd-5678", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 100},
},
},
@@ -109,7 +109,7 @@ func TestStatementBuilder(t *testing.T) {
Limit: 100,
},
expected: qbtypes.Statement{
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists`) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists` = true) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"failure", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 100},
},
},
@@ -143,7 +143,7 @@ func TestStatementBuilder(t *testing.T) {
Limit: 100,
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_audit.distributed_logs_resource WHERE (simpleJSONExtractString(labels, 'signoz.audit.resource.kind') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND (`attribute_string_signoz$$audit$$action` = ? AND `attribute_string_signoz$$audit$$action_exists`) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_audit.distributed_logs_resource WHERE (simpleJSONExtractString(labels, 'signoz.audit.resource.kind') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND (`attribute_string_signoz$$audit$$action` = ? AND `attribute_string_signoz$$audit$$action_exists` = true) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"dashboard", "%signoz.audit.resource.kind%", "%signoz.audit.resource.kind\":\"dashboard%", uint64(1747945619), uint64(1747983448), "delete", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 100},
},
},
@@ -160,7 +160,7 @@ func TestStatementBuilder(t *testing.T) {
Limit: 100,
},
expected: qbtypes.Statement{
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$principal$$type` = ? AND `attribute_string_signoz$$audit$$principal$$type_exists`) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, event_name, attributes_string, attributes_number, attributes_bool, resource, scope_string FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$principal$$type` = ? AND `attribute_string_signoz$$audit$$principal$$type_exists` = true) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"service_account", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 100},
},
},
@@ -180,7 +180,7 @@ func TestStatementBuilder(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "SELECT count() AS __result_0 FROM signoz_audit.distributed_logs WHERE ((`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists`) AND (`attribute_string_signoz$$audit$$action` = ? AND `attribute_string_signoz$$audit$$action_exists`)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC",
Query: "SELECT count() AS __result_0 FROM signoz_audit.distributed_logs WHERE ((`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists` = true) AND (`attribute_string_signoz$$audit$$action` = ? AND `attribute_string_signoz$$audit$$action_exists` = true)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC",
Args: []any{"failure", "update", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448)},
},
},
@@ -204,7 +204,7 @@ func TestStatementBuilder(t *testing.T) {
Limit: 5,
},
expected: qbtypes.Statement{
Query: "WITH __limit_cte AS (SELECT toString(multiIf(`attribute_string_signoz$$audit$$principal$$email_exists`, `attribute_string_signoz$$audit$$principal$$email`, NULL)) AS `signoz.audit.principal.email`, count() AS __result_0 FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists`) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `signoz.audit.principal.email` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 60 SECOND) AS ts, toString(multiIf(`attribute_string_signoz$$audit$$principal$$email_exists`, `attribute_string_signoz$$audit$$principal$$email`, NULL)) AS `signoz.audit.principal.email`, count() AS __result_0 FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists`) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? AND (`signoz.audit.principal.email`) GLOBAL IN (SELECT `signoz.audit.principal.email` FROM __limit_cte) GROUP BY ts, `signoz.audit.principal.email`",
Query: "WITH __limit_cte AS (SELECT toString(multiIf(`attribute_string_signoz$$audit$$principal$$email_exists` = true, `attribute_string_signoz$$audit$$principal$$email`, NULL)) AS `signoz.audit.principal.email`, count() AS __result_0 FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists` = true) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `signoz.audit.principal.email` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 60 SECOND) AS ts, toString(multiIf(`attribute_string_signoz$$audit$$principal$$email_exists` = true, `attribute_string_signoz$$audit$$principal$$email`, NULL)) AS `signoz.audit.principal.email`, count() AS __result_0 FROM signoz_audit.distributed_logs WHERE (`attribute_string_signoz$$audit$$outcome` = ? AND `attribute_string_signoz$$audit$$outcome_exists` = true) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? AND (`signoz.audit.principal.email`) GLOBAL IN (SELECT `signoz.audit.principal.email` FROM __limit_cte) GROUP BY ts, `signoz.audit.principal.email`",
Args: []any{"failure", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 5, "failure", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448)},
},
},

View File

@@ -180,7 +180,7 @@ func TestStatementBuilderTimeSeries(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint), __limit_cte AS (SELECT toString(multiIf(`attribute_string_materialized$$key$$name_exists`, `attribute_string_materialized$$key$$name`, NULL)) AS `__GROUP_BY_KEY_0_materialized.key.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_materialized.key.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 30 SECOND) AS ts, toString(multiIf(`attribute_string_materialized$$key$$name_exists`, `attribute_string_materialized$$key$$name`, NULL)) AS `__GROUP_BY_KEY_0_materialized.key.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_materialized.key.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_materialized.key.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_materialized.key.name`",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint), __limit_cte AS (SELECT toString(multiIf(`attribute_string_materialized$$key$$name_exists` = true, `attribute_string_materialized$$key$$name`, NULL)) AS `__GROUP_BY_KEY_0_materialized.key.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_materialized.key.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 30 SECOND) AS ts, toString(multiIf(`attribute_string_materialized$$key$$name_exists` = true, `attribute_string_materialized$$key$$name`, NULL)) AS `__GROUP_BY_KEY_0_materialized.key.name`, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_materialized.key.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_materialized.key.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_materialized.key.name`",
Args: []any{"cartservice", "%service.name%", "%service.name\":\"cartservice%", uint64(1705397400), uint64(1705485600), "1705399200000000000", uint64(1705397400), "1705485600000000000", uint64(1705485600), 10, "1705399200000000000", uint64(1705397400), "1705485600000000000", uint64(1705485600)},
},
},
@@ -203,7 +203,7 @@ func TestStatementBuilderTimeSeries(t *testing.T) {
Limit: 10,
},
expected: qbtypes.Statement{
Query: "SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 30 SECOND) AS ts, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists`) OR (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists`)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY ts",
Query: "SELECT toStartOfInterval(fromUnixTimestamp64Nano(timestamp), INTERVAL 30 SECOND) AS ts, count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists` = true) OR (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists` = true)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? GROUP BY ts",
Args: []any{"redis.*", "memcached", "1705399200000000000", uint64(1705397400), "1705485600000000000", uint64(1705485600)},
},
expectedErr: nil,
@@ -300,7 +300,7 @@ func TestStatementBuilderListQuery(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists`, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists` = true, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Args: []any{"cartservice", "%service.name%", "%service.name\":\"cartservice%", uint64(1747945619), uint64(1747983448), "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 10},
},
expectedErr: nil,
@@ -328,7 +328,7 @@ func TestStatementBuilderListQuery(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists`) OR (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists`)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists`, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Query: "SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists` = true) OR (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists` = true)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists` = true, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Args: []any{"redis.*", "memcached", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 10},
},
expectedErr: nil,
@@ -442,7 +442,7 @@ func TestStatementBuilderListQueryResourceTests(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND match(LOWER(body), LOWER(?)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists`, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND match(LOWER(body), LOWER(?)) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists` = true, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Args: []any{"cartservice", "%service.name%", "%service.name\":\"cartservice%", uint64(1747945619), uint64(1747983448), "hello", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 10},
},
expectedErr: nil,
@@ -666,7 +666,7 @@ func TestStatementBuilderListQueryServiceCollision(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND LOWER(body) LIKE LOWER(?) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists`, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_logs.distributed_logs_v2_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp, id, trace_id, span_id, trace_flags, severity_text, severity_number, scope_name, scope_version, body, attributes_string, attributes_number, attributes_bool, resources_string, scope_string FROM signoz_logs.distributed_logs_v2 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND LOWER(body) LIKE LOWER(?) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY multiIf(`attribute_string_materialized$$key$$name_exists` = true, `attribute_string_materialized$$key$$name`, NULL) desc LIMIT ?",
Args: []any{"cartservice", "%service.name%", "%service.name\":\"cartservice%", uint64(1747945619), uint64(1747983448), "%error%", "1747947419000000000", uint64(1747945619), "1747983448000000000", uint64(1747983448), 10},
},
expectedErr: nil,

View File

@@ -129,7 +129,7 @@ func TestStatementBuilder(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __limit_cte AS (SELECT toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists`) OR multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) = ?) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_service.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists`) OR multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) = ?) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_service.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_service.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_service.name`",
Query: "WITH __limit_cte AS (SELECT toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists` = true) OR multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) = ?) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_service.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, count() AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists` = true) OR multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) = ?) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_service.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_service.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_service.name`",
Args: []any{"redis-manual", "redis-manual", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10, "redis-manual", "redis-manual", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
expectedErr: nil,
@@ -268,7 +268,7 @@ func TestStatementBuilder(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint), __limit_cte AS (SELECT toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists`, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_service.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists`, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_service.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_service.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_service.name`",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint), __limit_cte AS (SELECT toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists` = true, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_service.name` ORDER BY __result_0 DESC LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists` = true, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_service.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_service.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_service.name`",
Args: []any{"redis-manual", "%service.name%", "%service.name\":\"redis-manual%", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10, "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
expectedErr: nil,
@@ -307,7 +307,7 @@ func TestStatementBuilder(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint), __limit_cte AS (SELECT toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists`, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_service.name` ORDER BY `__GROUP_BY_KEY_0_service.name` desc LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists`, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_service.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_service.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_service.name` ORDER BY `__GROUP_BY_KEY_0_service.name` desc, ts desc",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint), __limit_cte AS (SELECT toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists` = true, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? GROUP BY `__GROUP_BY_KEY_0_service.name` ORDER BY `__GROUP_BY_KEY_0_service.name` desc LIMIT ?) SELECT toStartOfInterval(timestamp, INTERVAL 30 SECOND) AS ts, toString(multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL)) AS `__GROUP_BY_KEY_0_service.name`, sum(multiIf(`attribute_number_cart$$items_count_exists` = true, toFloat64(`attribute_number_cart$$items_count`), NULL)) AS __result_0 FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? AND (`__GROUP_BY_KEY_0_service.name`) GLOBAL IN (SELECT `__GROUP_BY_KEY_0_service.name` FROM __limit_cte) GROUP BY ts, `__GROUP_BY_KEY_0_service.name` ORDER BY `__GROUP_BY_KEY_0_service.name` desc, ts desc",
Args: []any{"redis-manual", "%service.name%", "%service.name\":\"redis-manual%", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10, "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448)},
},
expectedErr: nil,
@@ -552,7 +552,7 @@ func TestStatementBuilderListQuery(t *testing.T) {
},
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp AS `__SELECT_KEY_0_timestamp`, trace_id AS `__SELECT_KEY_1_trace_id`, span_id AS `__SELECT_KEY_2_span_id`, name AS `__SELECT_KEY_3_name`, multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) AS `__SELECT_KEY_4_service.name`, duration_nano AS `__SELECT_KEY_5_duration_nano`, multiIf(`attribute_number_cart$$items_count_exists`, `attribute_number_cart$$items_count`, NULL) AS `__SELECT_KEY_6_cart.items_count` FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp AS `__SELECT_KEY_0_timestamp`, trace_id AS `__SELECT_KEY_1_trace_id`, span_id AS `__SELECT_KEY_2_span_id`, name AS `__SELECT_KEY_3_name`, multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) AS `__SELECT_KEY_4_service.name`, duration_nano AS `__SELECT_KEY_5_duration_nano`, multiIf(`attribute_number_cart$$items_count_exists` = true, `attribute_number_cart$$items_count`, NULL) AS `__SELECT_KEY_6_cart.items_count` FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"redis-manual", "%service.name%", "%service.name\":\"redis-manual%", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
expectedErr: nil,
@@ -669,7 +669,7 @@ func TestStatementBuilderListQuery(t *testing.T) {
Limit: 10,
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp AS `__SELECT_KEY_0_timestamp`, trace_id AS `__SELECT_KEY_1_trace_id`, span_id AS `__SELECT_KEY_2_span_id`, name AS `__SELECT_KEY_3_name`, resource_string_service$$name AS `__SELECT_KEY_4_serviceName`, duration_nano AS `__SELECT_KEY_5_durationNano`, http_method AS `__SELECT_KEY_6_httpMethod`, multiIf(`attribute_string_mixed$$materialization$$key_exists`, `attribute_string_mixed$$materialization$$key`, multiIf(resource.`mixed.materialization.key` IS NOT NULL, resource.`mixed.materialization.key`::String, mapContains(resources_string, 'mixed.materialization.key'), resources_string['mixed.materialization.key'], NULL) IS NOT NULL, multiIf(resource.`mixed.materialization.key` IS NOT NULL, resource.`mixed.materialization.key`::String, mapContains(resources_string, 'mixed.materialization.key'), resources_string['mixed.materialization.key'], NULL), NULL) AS `__SELECT_KEY_7_mixed.materialization.key` FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp AS `__SELECT_KEY_0_timestamp`, trace_id AS `__SELECT_KEY_1_trace_id`, span_id AS `__SELECT_KEY_2_span_id`, name AS `__SELECT_KEY_3_name`, resource_string_service$$name AS `__SELECT_KEY_4_serviceName`, duration_nano AS `__SELECT_KEY_5_durationNano`, http_method AS `__SELECT_KEY_6_httpMethod`, multiIf(`attribute_string_mixed$$materialization$$key_exists` = true, `attribute_string_mixed$$materialization$$key`, multiIf(resource.`mixed.materialization.key` IS NOT NULL, resource.`mixed.materialization.key`::String, mapContains(resources_string, 'mixed.materialization.key'), resources_string['mixed.materialization.key'], NULL) IS NOT NULL, multiIf(resource.`mixed.materialization.key` IS NOT NULL, resource.`mixed.materialization.key`::String, mapContains(resources_string, 'mixed.materialization.key'), resources_string['mixed.materialization.key'], NULL), NULL) AS `__SELECT_KEY_7_mixed.materialization.key` FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"redis-manual", "%service.name%", "%service.name\":\"redis-manual%", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
expectedErr: nil,
@@ -714,7 +714,7 @@ func TestStatementBuilderListQuery(t *testing.T) {
Limit: 10,
},
expected: qbtypes.Statement{
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp AS `__SELECT_KEY_0_timestamp`, trace_id AS `__SELECT_KEY_1_trace_id`, span_id AS `__SELECT_KEY_2_span_id`, name AS `__SELECT_KEY_3_name`, resource_string_service$$name AS `__SELECT_KEY_4_serviceName`, duration_nano AS `__SELECT_KEY_5_durationNano`, http_method AS `__SELECT_KEY_6_httpMethod`, multiIf(`attribute_string_mixed$$materialization$$key_exists`, `attribute_string_mixed$$materialization$$key`, NULL) AS `__SELECT_KEY_7_mixed.materialization.key` FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Query: "WITH __resource_filter AS (SELECT fingerprint FROM signoz_traces.distributed_traces_v3_resource WHERE (simpleJSONExtractString(labels, 'service.name') = ? AND labels LIKE ? AND labels LIKE ?) AND seen_at_ts_bucket_start >= ? AND seen_at_ts_bucket_start <= ? GROUP BY fingerprint) SELECT timestamp AS `__SELECT_KEY_0_timestamp`, trace_id AS `__SELECT_KEY_1_trace_id`, span_id AS `__SELECT_KEY_2_span_id`, name AS `__SELECT_KEY_3_name`, resource_string_service$$name AS `__SELECT_KEY_4_serviceName`, duration_nano AS `__SELECT_KEY_5_durationNano`, http_method AS `__SELECT_KEY_6_httpMethod`, multiIf(`attribute_string_mixed$$materialization$$key_exists` = true, `attribute_string_mixed$$materialization$$key`, NULL) AS `__SELECT_KEY_7_mixed.materialization.key` FROM signoz_traces.distributed_signoz_index_v3 WHERE resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? LIMIT ?",
Args: []any{"redis-manual", "%service.name%", "%service.name\":\"redis-manual%", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
expectedErr: nil,
@@ -1178,7 +1178,7 @@ func TestStatementBuilderTraceQuery(t *testing.T) {
Limit: 10,
},
expected: qbtypes.Statement{
Query: "WITH __toe AS (SELECT trace_id FROM signoz_traces.distributed_signoz_index_v3 WHERE (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists`) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ?), __toe_duration_sorted AS (SELECT trace_id, duration_nano, resource_string_service$$name as `service.name`, name FROM signoz_traces.distributed_signoz_index_v3 WHERE parent_span_id = '' AND trace_id GLOBAL IN __toe AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? ORDER BY duration_nano DESC LIMIT 1 BY trace_id) SELECT __toe_duration_sorted.`service.name` AS `service.name`, __toe_duration_sorted.name AS `name`, count() AS span_count, __toe_duration_sorted.duration_nano AS `duration_nano`, __toe_duration_sorted.trace_id AS `trace_id` FROM __toe INNER JOIN __toe_duration_sorted ON __toe.trace_id = __toe_duration_sorted.trace_id GROUP BY trace_id, duration_nano, name, `service.name` ORDER BY duration_nano DESC LIMIT 1 BY trace_id LIMIT ? SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000",
Query: "WITH __toe AS (SELECT trace_id FROM signoz_traces.distributed_signoz_index_v3 WHERE (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists` = true) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ?), __toe_duration_sorted AS (SELECT trace_id, duration_nano, resource_string_service$$name as `service.name`, name FROM signoz_traces.distributed_signoz_index_v3 WHERE parent_span_id = '' AND trace_id GLOBAL IN __toe AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? ORDER BY duration_nano DESC LIMIT 1 BY trace_id) SELECT __toe_duration_sorted.`service.name` AS `service.name`, __toe_duration_sorted.name AS `name`, count() AS span_count, __toe_duration_sorted.duration_nano AS `duration_nano`, __toe_duration_sorted.trace_id AS `trace_id` FROM __toe INNER JOIN __toe_duration_sorted ON __toe.trace_id = __toe_duration_sorted.trace_id GROUP BY trace_id, duration_nano, name, `service.name` ORDER BY duration_nano DESC LIMIT 1 BY trace_id LIMIT ? SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000",
Args: []any{"redis-manual", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
expectedErr: nil,
@@ -1194,7 +1194,7 @@ func TestStatementBuilderTraceQuery(t *testing.T) {
Limit: 10,
},
expected: qbtypes.Statement{
Query: "WITH __toe AS (SELECT trace_id FROM signoz_traces.distributed_signoz_index_v3 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists`) OR multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) = ?) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ?), __toe_duration_sorted AS (SELECT trace_id, duration_nano, resource_string_service$$name as `service.name`, name FROM signoz_traces.distributed_signoz_index_v3 WHERE parent_span_id = '' AND trace_id GLOBAL IN __toe AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? ORDER BY duration_nano DESC LIMIT 1 BY trace_id) SELECT __toe_duration_sorted.`service.name` AS `service.name`, __toe_duration_sorted.name AS `name`, count() AS span_count, __toe_duration_sorted.duration_nano AS `duration_nano`, __toe_duration_sorted.trace_id AS `trace_id` FROM __toe INNER JOIN __toe_duration_sorted ON __toe.trace_id = __toe_duration_sorted.trace_id GROUP BY trace_id, duration_nano, name, `service.name` ORDER BY duration_nano DESC LIMIT 1 BY trace_id LIMIT ? SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000",
Query: "WITH __toe AS (SELECT trace_id FROM signoz_traces.distributed_signoz_index_v3 WHERE ((match(`attribute_string_materialized$$key$$name`, ?) AND `attribute_string_materialized$$key$$name_exists` = true) OR multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, mapContains(resources_string, 'service.name'), resources_string['service.name'], NULL) = ?) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ?), __toe_duration_sorted AS (SELECT trace_id, duration_nano, resource_string_service$$name as `service.name`, name FROM signoz_traces.distributed_signoz_index_v3 WHERE parent_span_id = '' AND trace_id GLOBAL IN __toe AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? ORDER BY duration_nano DESC LIMIT 1 BY trace_id) SELECT __toe_duration_sorted.`service.name` AS `service.name`, __toe_duration_sorted.name AS `name`, count() AS span_count, __toe_duration_sorted.duration_nano AS `duration_nano`, __toe_duration_sorted.trace_id AS `trace_id` FROM __toe INNER JOIN __toe_duration_sorted ON __toe.trace_id = __toe_duration_sorted.trace_id GROUP BY trace_id, duration_nano, name, `service.name` ORDER BY duration_nano DESC LIMIT 1 BY trace_id LIMIT ? SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000",
Args: []any{"redis-manual", "redis-manual", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
expectedErr: nil,
@@ -1240,7 +1240,7 @@ func TestStatementBuilderTraceQuery(t *testing.T) {
Limit: 10,
},
expected: qbtypes.Statement{
Query: "WITH __toe AS (SELECT trace_id FROM signoz_traces.distributed_signoz_index_v3 WHERE (((name, resource_string_service$$name) GLOBAL IN (SELECT DISTINCT name, serviceName from signoz_traces.distributed_top_level_operations WHERE time >= toDateTime(1747947419))) AND parent_span_id != '' OR (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists`)) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ?), __toe_duration_sorted AS (SELECT trace_id, duration_nano, resource_string_service$$name as `service.name`, name FROM signoz_traces.distributed_signoz_index_v3 WHERE parent_span_id = '' AND trace_id GLOBAL IN __toe AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? ORDER BY duration_nano DESC LIMIT 1 BY trace_id) SELECT __toe_duration_sorted.`service.name` AS `service.name`, __toe_duration_sorted.name AS `name`, count() AS span_count, __toe_duration_sorted.duration_nano AS `duration_nano`, __toe_duration_sorted.trace_id AS `trace_id` FROM __toe INNER JOIN __toe_duration_sorted ON __toe.trace_id = __toe_duration_sorted.trace_id GROUP BY trace_id, duration_nano, name, `service.name` ORDER BY duration_nano DESC LIMIT 1 BY trace_id LIMIT ? SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000",
Query: "WITH __toe AS (SELECT trace_id FROM signoz_traces.distributed_signoz_index_v3 WHERE (((name, resource_string_service$$name) GLOBAL IN (SELECT DISTINCT name, serviceName from signoz_traces.distributed_top_level_operations WHERE time >= toDateTime(1747947419))) AND parent_span_id != '' OR (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists` = true)) AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ?), __toe_duration_sorted AS (SELECT trace_id, duration_nano, resource_string_service$$name as `service.name`, name FROM signoz_traces.distributed_signoz_index_v3 WHERE parent_span_id = '' AND trace_id GLOBAL IN __toe AND timestamp >= ? AND timestamp < ? AND ts_bucket_start >= ? AND ts_bucket_start <= ? ORDER BY duration_nano DESC LIMIT 1 BY trace_id) SELECT __toe_duration_sorted.`service.name` AS `service.name`, __toe_duration_sorted.name AS `name`, count() AS span_count, __toe_duration_sorted.duration_nano AS `duration_nano`, __toe_duration_sorted.trace_id AS `trace_id` FROM __toe INNER JOIN __toe_duration_sorted ON __toe.trace_id = __toe_duration_sorted.trace_id GROUP BY trace_id, duration_nano, name, `service.name` ORDER BY duration_nano DESC LIMIT 1 BY trace_id LIMIT ? SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000",
Args: []any{"redis-manual", "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), "1747947419000000000", "1747983448000000000", uint64(1747945619), uint64(1747983448), 10},
},
expectedErr: nil,

View File

@@ -461,7 +461,7 @@ func TestConditionFor(t *testing.T) {
evolutions: mockEvolution,
operator: qbtypes.FilterOperatorRegexp,
value: "frontend-.*",
expectedSQL: "WHERE (match(`resource_string_service$$name`, ?) AND `resource_string_service$$name_exists`)",
expectedSQL: "WHERE (match(`resource_string_service$$name`, ?) AND `resource_string_service$$name_exists` = true)",
expectedArgs: []any{"frontend-.*"},
expectedError: nil,
},

View File

@@ -1596,7 +1596,7 @@ func TestFilterExprLogs(t *testing.T) {
category: "Materialized key",
query: "materialized.key.name=\"test\"",
shouldPass: true,
expectedQuery: "WHERE (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists`)",
expectedQuery: "WHERE (`attribute_string_materialized$$key$$name` = ? AND `attribute_string_materialized$$key$$name_exists` = true)",
expectedArgs: []any{"test"},
expectedErrorContains: "",
},

View File

@@ -182,7 +182,7 @@ func (m *storage) read(_ context.Context, q qbtypes.QueryInfo, key *telemetrytyp
// a key could have been materialized, if so return the materialized column name
if key.Materialized {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExpr = append(existExpr, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
existExpr = append(existExpr, telemetrytypes.FieldKeyToMaterializedExistsCondition(key, true))
} else {
exprs = append(exprs, fmt.Sprintf("%s[%s]", columnName, clickhousesql.StringLiteral(key.Name)))
existExpr = append(existExpr, fmt.Sprintf("mapContains(%s, %s)", columnName, clickhousesql.StringLiteral(key.Name)))

View File

@@ -580,7 +580,7 @@ func TestFieldForWithMaterialized(t *testing.T) {
name: "Multi evolution - both columns (JSON + materialized)",
start: time.Date(2024, 2, 1, 0, 0, 0, 0, time.UTC),
end: time.Date(2024, 4, 2, 0, 0, 0, 0, time.UTC),
expectedResult: "multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, `resource_string_service$$name_exists`, `resource_string_service$$name`, NULL)",
expectedResult: "multiIf(resource.`service.name` IS NOT NULL, resource.`service.name`::String, `resource_string_service$$name_exists` = true, `resource_string_service$$name`, NULL)",
},
}

View File

@@ -306,7 +306,7 @@ func (m *storage) resolveColumnExprs(
// a key could have been materialized, if so return the materialized column name
if key.Materialized {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedExistsCondition(key, true))
} else {
exprs = append(exprs, fmt.Sprintf("%s[%s]", columnName, clickhousesql.StringLiteral(key.Name)))
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, %s)", columnName, clickhousesql.StringLiteral(key.Name)))

View File

@@ -80,7 +80,7 @@ func TestGetFieldKeyName(t *testing.T) {
Materialized: true,
Evolutions: mockEvolution,
},
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists` = true, `resource_string_deployment$$environment`, NULL)",
expectedError: nil,
},
{
@@ -228,7 +228,7 @@ func TestFieldForResourceWithEvolution(t *testing.T) {
},
tsStart: uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
tsEnd: uint64(time.Date(2025, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists` = true, `resource_string_deployment$$environment`, NULL)",
},
}

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

@@ -218,6 +218,12 @@ func FieldKeyToMaterializedColumnNameForExists(key *TelemetryFieldKey) string {
))
}
// FieldKeyToMaterializedExistsCondition compares the exists column explicitly: a bare bool
// column defeats skip-index pruning across OR.
func FieldKeyToMaterializedExistsCondition(key *TelemetryFieldKey, exists bool) string {
return fmt.Sprintf("%s = %t", FieldKeyToMaterializedColumnNameForExists(key), exists)
}
type TelemetryFieldValues struct {
StringValues []string `json:"stringValues,omitempty"`
BoolValues []bool `json:"boolValues,omitempty"`

View File

@@ -108,14 +108,23 @@ def delete_all_rules(signoz: types.SigNoz, token: str) -> None:
def seed_alert_rules(
signoz: types.SigNoz,
get_token: Callable[[str, str], str],
notification_channel: types.TestContainerDocker,
create_notification_channel: Callable[[dict], str],
create_alert_rule: Callable[[dict], str],
) -> Callable[[dict, list[dict]], None]:
) -> Callable[[str, list[dict]], None]:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
def _seed_alert_rules(channel_config: dict, rules: list[dict]) -> None:
# create_notification_channel rather than create_webhook_notification_channel:
# only the former deletes on teardown, and callers reuse one channel name
# across tests, so a leaked channel fails the next create as a duplicate.
def _seed_alert_rules(channel_name: str, rules: list[dict]) -> None:
delete_all_rules(signoz, admin_token)
create_notification_channel(channel_config)
create_notification_channel(
{
"name": channel_name,
"webhook_configs": [{"url": notification_channel.container_configs["8080"].get(f"/alert/{channel_name}"), "send_resolved": False}],
}
)
for rule in rules:
create_alert_rule(rule)

Some files were not shown because too many files have changed in this diff Show More