mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-02 01:20:41 +01:00
Compare commits
23 Commits
feat/text-
...
issue_5947
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f537154df9 | ||
|
|
7c50fe3763 | ||
|
|
160a1b018c | ||
|
|
1c6d966afa | ||
|
|
f22a18d0a7 | ||
|
|
752899a7a4 | ||
|
|
eb4c570d53 | ||
|
|
7ce405aff7 | ||
|
|
72618d1d83 | ||
|
|
80b7edc22a | ||
|
|
45a8bb424c | ||
|
|
17afa7a3cf | ||
|
|
6dc9bf7b16 | ||
|
|
d15a1452f8 | ||
|
|
2fe033ec78 | ||
|
|
02a2800f89 | ||
|
|
915aa2eb70 | ||
|
|
518caff0f2 | ||
|
|
0d3ac28286 | ||
|
|
47beef07de | ||
|
|
5f2891bd6c | ||
|
|
97f0e832ab | ||
|
|
770a8f7b0a |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -50,6 +50,7 @@ jobs:
|
||||
- logspipelines
|
||||
- passwordauthn
|
||||
- preference
|
||||
- quickfilter
|
||||
- querierlogs
|
||||
- queriertraces
|
||||
- queriermetrics
|
||||
|
||||
@@ -109,6 +109,57 @@ components:
|
||||
webhook_url:
|
||||
$ref: '#/components/schemas/ConfigSecretURL'
|
||||
type: object
|
||||
AlertmanagertypesJSMOpsReceiverConfig:
|
||||
properties:
|
||||
api_key:
|
||||
type: string
|
||||
description:
|
||||
type: string
|
||||
http_config:
|
||||
$ref: '#/components/schemas/ConfigHTTPClientConfig'
|
||||
message:
|
||||
type: string
|
||||
priority:
|
||||
type: string
|
||||
send_resolved:
|
||||
type: boolean
|
||||
tags:
|
||||
type: string
|
||||
type: object
|
||||
AlertmanagertypesJiraReceiverConfig:
|
||||
properties:
|
||||
custom_fields:
|
||||
additionalProperties: {}
|
||||
type: object
|
||||
description:
|
||||
type: string
|
||||
http_config:
|
||||
$ref: '#/components/schemas/ConfigHTTPClientConfig'
|
||||
issue_type:
|
||||
type: string
|
||||
labels:
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
priority:
|
||||
type: string
|
||||
project:
|
||||
type: string
|
||||
reopen_duration:
|
||||
$ref: '#/components/schemas/ModelDuration'
|
||||
reopen_transition:
|
||||
type: string
|
||||
resolve_transition:
|
||||
type: string
|
||||
send_resolved:
|
||||
type: boolean
|
||||
site:
|
||||
type: string
|
||||
summary:
|
||||
type: string
|
||||
wont_fix_resolution:
|
||||
type: string
|
||||
type: object
|
||||
AlertmanagertypesMaintenanceKind:
|
||||
enum:
|
||||
- fixed
|
||||
@@ -162,6 +213,10 @@ components:
|
||||
oneOf:
|
||||
- required:
|
||||
- googlechat_configs
|
||||
- required:
|
||||
- jira_configs
|
||||
- required:
|
||||
- jsmops_configs
|
||||
- required:
|
||||
- discord_configs
|
||||
- required:
|
||||
@@ -192,8 +247,6 @@ components:
|
||||
- msteams_configs
|
||||
- required:
|
||||
- msteamsv2_configs
|
||||
- required:
|
||||
- jira_configs
|
||||
- required:
|
||||
- rocketchat_configs
|
||||
- required:
|
||||
@@ -217,7 +270,11 @@ components:
|
||||
type: array
|
||||
jira_configs:
|
||||
items:
|
||||
$ref: '#/components/schemas/ConfigJiraConfig'
|
||||
$ref: '#/components/schemas/AlertmanagertypesJiraReceiverConfig'
|
||||
type: array
|
||||
jsmops_configs:
|
||||
items:
|
||||
$ref: '#/components/schemas/AlertmanagertypesJSMOpsReceiverConfig'
|
||||
type: array
|
||||
mattermost_configs:
|
||||
items:
|
||||
@@ -344,7 +401,11 @@ components:
|
||||
type: array
|
||||
jira_configs:
|
||||
items:
|
||||
$ref: '#/components/schemas/ConfigJiraConfig'
|
||||
$ref: '#/components/schemas/AlertmanagertypesJiraReceiverConfig'
|
||||
type: array
|
||||
jsmops_configs:
|
||||
items:
|
||||
$ref: '#/components/schemas/AlertmanagertypesJSMOpsReceiverConfig'
|
||||
type: array
|
||||
mattermost_configs:
|
||||
items:
|
||||
@@ -7565,6 +7626,37 @@ components:
|
||||
- custom
|
||||
- text
|
||||
type: string
|
||||
QuickfiltertypesSourceFilters:
|
||||
properties:
|
||||
createdAt:
|
||||
format: date-time
|
||||
type: string
|
||||
filters:
|
||||
items:
|
||||
$ref: '#/components/schemas/TelemetrytypesTelemetryFieldKey'
|
||||
type: array
|
||||
id:
|
||||
type: string
|
||||
orgId:
|
||||
type: string
|
||||
source:
|
||||
type: string
|
||||
updatedAt:
|
||||
format: date-time
|
||||
type: string
|
||||
required:
|
||||
- id
|
||||
- filters
|
||||
type: object
|
||||
QuickfiltertypesUpdatableQuickFilters:
|
||||
properties:
|
||||
filters:
|
||||
items:
|
||||
$ref: '#/components/schemas/TelemetrytypesTelemetryFieldKey'
|
||||
type: array
|
||||
required:
|
||||
- filters
|
||||
type: object
|
||||
RenderErrorResponse:
|
||||
properties:
|
||||
error:
|
||||
@@ -18568,6 +18660,171 @@ paths:
|
||||
summary: Get query range result (v2)
|
||||
tags:
|
||||
- dashboard
|
||||
/api/v2/quick_filters:
|
||||
get:
|
||||
deprecated: false
|
||||
description: Returns the org's quick filters for every source, each filter as
|
||||
a telemetry field key.
|
||||
operationId: ListQuickFilters
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
items:
|
||||
$ref: '#/components/schemas/QuickfiltertypesSourceFilters'
|
||||
nullable: true
|
||||
type: array
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- quick-filter:list
|
||||
- tokenizer:
|
||||
- quick-filter:list
|
||||
summary: List quick filters
|
||||
tags:
|
||||
- quick_filter
|
||||
/api/v2/quick_filters/{source}:
|
||||
get:
|
||||
deprecated: false
|
||||
description: Returns the org's quick filters for one source, each filter as
|
||||
a telemetry field key.
|
||||
operationId: GetQuickFilters
|
||||
parameters:
|
||||
- in: path
|
||||
name: source
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/QuickfiltertypesSourceFilters'
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- quick-filter:read
|
||||
- tokenizer:
|
||||
- quick-filter:read
|
||||
summary: Get a source's quick filters
|
||||
tags:
|
||||
- quick_filter
|
||||
put:
|
||||
deprecated: false
|
||||
description: Replaces the org's quick filters for the source named in the path.
|
||||
operationId: UpdateQuickFilters
|
||||
parameters:
|
||||
- in: path
|
||||
name: source
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
requestBody:
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/QuickfiltertypesUpdatableQuickFilters'
|
||||
responses:
|
||||
"204":
|
||||
description: No Content
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- quick-filter:update
|
||||
- tokenizer:
|
||||
- quick-filter:update
|
||||
summary: Update quick filters
|
||||
tags:
|
||||
- quick_filter
|
||||
/api/v2/readyz:
|
||||
get:
|
||||
operationId: Readyz
|
||||
|
||||
316
frontend/src/api/generated/services/quick-filter/index.ts
Normal file
316
frontend/src/api/generated/services/quick-filter/index.ts
Normal file
@@ -0,0 +1,316 @@
|
||||
/**
|
||||
* ! Do not edit manually
|
||||
* * The file has been auto-generated using Orval for SigNoz
|
||||
* * regenerate with 'pnpm generate:api'
|
||||
* SigNoz
|
||||
*/
|
||||
import { useMutation, useQuery } from 'react-query';
|
||||
import type {
|
||||
InvalidateOptions,
|
||||
MutationFunction,
|
||||
QueryClient,
|
||||
QueryFunction,
|
||||
QueryKey,
|
||||
UseMutationOptions,
|
||||
UseMutationResult,
|
||||
UseQueryOptions,
|
||||
UseQueryResult,
|
||||
} from 'react-query';
|
||||
|
||||
import type {
|
||||
GetQuickFilters200,
|
||||
GetQuickFiltersPathParameters,
|
||||
ListQuickFilters200,
|
||||
QuickfiltertypesUpdatableQuickFiltersDTO,
|
||||
RenderErrorResponseDTO,
|
||||
UpdateQuickFiltersPathParameters,
|
||||
} from '../sigNoz.schemas';
|
||||
|
||||
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
|
||||
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
|
||||
|
||||
/**
|
||||
* Returns the org's quick filters for every source, each filter as a telemetry field key.
|
||||
* @summary List quick filters
|
||||
*/
|
||||
export const listQuickFilters = (signal?: AbortSignal) => {
|
||||
return GeneratedAPIInstance<ListQuickFilters200>({
|
||||
url: `/api/v2/quick_filters`,
|
||||
method: 'GET',
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getListQuickFiltersQueryKey = () => {
|
||||
return [`/api/v2/quick_filters`] as const;
|
||||
};
|
||||
|
||||
export const getListQuickFiltersQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof listQuickFilters>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listQuickFilters>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
}) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey = queryOptions?.queryKey ?? getListQuickFiltersQueryKey();
|
||||
|
||||
const queryFn: QueryFunction<Awaited<ReturnType<typeof listQuickFilters>>> = ({
|
||||
signal,
|
||||
}) => listQuickFilters(signal);
|
||||
|
||||
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listQuickFilters>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type ListQuickFiltersQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof listQuickFilters>>
|
||||
>;
|
||||
export type ListQuickFiltersQueryError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary List quick filters
|
||||
*/
|
||||
|
||||
export function useListQuickFilters<
|
||||
TData = Awaited<ReturnType<typeof listQuickFilters>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listQuickFilters>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
}): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getListQuickFiltersQueryOptions(options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return { ...query, queryKey: queryOptions.queryKey };
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary List quick filters
|
||||
*/
|
||||
export const invalidateListQuickFilters = async (
|
||||
queryClient: QueryClient,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getListQuickFiltersQueryKey() },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* Returns the org's quick filters for one source, each filter as a telemetry field key.
|
||||
* @summary Get a source's quick filters
|
||||
*/
|
||||
export const getQuickFilters = (
|
||||
{ source }: GetQuickFiltersPathParameters,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<GetQuickFilters200>({
|
||||
url: `/api/v2/quick_filters/${source}`,
|
||||
method: 'GET',
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getGetQuickFiltersQueryKey = ({
|
||||
source,
|
||||
}: GetQuickFiltersPathParameters) => {
|
||||
return [`/api/v2/quick_filters/${source}`] as const;
|
||||
};
|
||||
|
||||
export const getGetQuickFiltersQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof getQuickFilters>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
{ source }: GetQuickFiltersPathParameters,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getQuickFilters>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey =
|
||||
queryOptions?.queryKey ?? getGetQuickFiltersQueryKey({ source });
|
||||
|
||||
const queryFn: QueryFunction<Awaited<ReturnType<typeof getQuickFilters>>> = ({
|
||||
signal,
|
||||
}) => getQuickFilters({ source }, signal);
|
||||
|
||||
return {
|
||||
queryKey,
|
||||
queryFn,
|
||||
enabled: !!source,
|
||||
...queryOptions,
|
||||
} as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getQuickFilters>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type GetQuickFiltersQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof getQuickFilters>>
|
||||
>;
|
||||
export type GetQuickFiltersQueryError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Get a source's quick filters
|
||||
*/
|
||||
|
||||
export function useGetQuickFilters<
|
||||
TData = Awaited<ReturnType<typeof getQuickFilters>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
{ source }: GetQuickFiltersPathParameters,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getQuickFilters>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getGetQuickFiltersQueryOptions({ source }, options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return { ...query, queryKey: queryOptions.queryKey };
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary Get a source's quick filters
|
||||
*/
|
||||
export const invalidateGetQuickFilters = async (
|
||||
queryClient: QueryClient,
|
||||
{ source }: GetQuickFiltersPathParameters,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getGetQuickFiltersQueryKey({ source }) },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* Replaces the org's quick filters for the source named in the path.
|
||||
* @summary Update quick filters
|
||||
*/
|
||||
export const updateQuickFilters = (
|
||||
{ source }: UpdateQuickFiltersPathParameters,
|
||||
quickfiltertypesUpdatableQuickFiltersDTO?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<void>({
|
||||
url: `/api/v2/quick_filters/${source}`,
|
||||
method: 'PUT',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
data: quickfiltertypesUpdatableQuickFiltersDTO,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getUpdateQuickFiltersMutationOptions = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof updateQuickFilters>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateQuickFiltersPathParameters;
|
||||
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
|
||||
},
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof updateQuickFilters>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateQuickFiltersPathParameters;
|
||||
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
|
||||
},
|
||||
TContext
|
||||
> => {
|
||||
const mutationKey = ['updateQuickFilters'];
|
||||
const { mutation: mutationOptions } = options
|
||||
? options.mutation &&
|
||||
'mutationKey' in options.mutation &&
|
||||
options.mutation.mutationKey
|
||||
? options
|
||||
: { ...options, mutation: { ...options.mutation, mutationKey } }
|
||||
: { mutation: { mutationKey } };
|
||||
|
||||
const mutationFn: MutationFunction<
|
||||
Awaited<ReturnType<typeof updateQuickFilters>>,
|
||||
{
|
||||
pathParams: UpdateQuickFiltersPathParameters;
|
||||
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
|
||||
}
|
||||
> = (props) => {
|
||||
const { pathParams, data } = props ?? {};
|
||||
|
||||
return updateQuickFilters(pathParams, data);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type UpdateQuickFiltersMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof updateQuickFilters>>
|
||||
>;
|
||||
export type UpdateQuickFiltersMutationBody =
|
||||
| BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>
|
||||
| undefined;
|
||||
export type UpdateQuickFiltersMutationError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Update quick filters
|
||||
*/
|
||||
export const useUpdateQuickFilters = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof updateQuickFilters>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateQuickFiltersPathParameters;
|
||||
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
|
||||
},
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof updateQuickFilters>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateQuickFiltersPathParameters;
|
||||
data?: BodyType<QuickfiltertypesUpdatableQuickFiltersDTO>;
|
||||
},
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getUpdateQuickFiltersMutationOptions(options));
|
||||
};
|
||||
@@ -385,6 +385,93 @@ export interface AlertmanagertypesGoogleChatReceiverConfigDTO {
|
||||
webhook_url?: ConfigSecretURLDTO;
|
||||
}
|
||||
|
||||
export interface AlertmanagertypesJSMOpsReceiverConfigDTO {
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
api_key?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
description?: string;
|
||||
http_config?: ConfigHTTPClientConfigDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
message?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
priority?: string;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
send_resolved?: boolean;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
tags?: string;
|
||||
}
|
||||
|
||||
export type AlertmanagertypesJiraReceiverConfigDTOCustomFields = {
|
||||
[key: string]: unknown;
|
||||
};
|
||||
|
||||
export type ModelDurationDTO = number;
|
||||
|
||||
export interface AlertmanagertypesJiraReceiverConfigDTO {
|
||||
/**
|
||||
* @type object
|
||||
*/
|
||||
custom_fields?: AlertmanagertypesJiraReceiverConfigDTOCustomFields;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
description?: string;
|
||||
http_config?: ConfigHTTPClientConfigDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
issue_type?: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
labels?: string[];
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
priority?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
project?: string;
|
||||
reopen_duration?: ModelDurationDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
reopen_transition?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
resolve_transition?: string;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
send_resolved?: boolean;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
site?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
summary?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
wont_fix_resolution?: string;
|
||||
}
|
||||
|
||||
export enum AlertmanagertypesMaintenanceKindDTO {
|
||||
fixed = 'fixed',
|
||||
recurring = 'recurring',
|
||||
@@ -631,69 +718,6 @@ export interface ConfigIncidentioConfigDTO {
|
||||
url_file?: string;
|
||||
}
|
||||
|
||||
export interface ConfigJiraFieldConfigDTO {
|
||||
/**
|
||||
* @type boolean,null
|
||||
*/
|
||||
enable_update?: boolean | null;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
template?: string;
|
||||
}
|
||||
|
||||
export type ModelDurationDTO = number;
|
||||
|
||||
export type ConfigJiraConfigDTOCustomFields = { [key: string]: unknown };
|
||||
|
||||
export interface ConfigJiraConfigDTO {
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
api_type?: string;
|
||||
api_url?: ConfigURLType2DTO;
|
||||
/**
|
||||
* @type object
|
||||
*/
|
||||
custom_fields?: ConfigJiraConfigDTOCustomFields;
|
||||
description?: ConfigJiraFieldConfigDTO;
|
||||
http_config?: ConfigHTTPClientConfigDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
issue_type?: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
labels?: string[];
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
priority?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
project?: string;
|
||||
reopen_duration?: ModelDurationDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
reopen_transition?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
resolve_transition?: string;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
send_resolved?: boolean;
|
||||
summary?: ConfigJiraFieldConfigDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
wont_fix_resolution?: string;
|
||||
}
|
||||
|
||||
export interface ConfigMattermostFieldDTO {
|
||||
/**
|
||||
* @type boolean,null
|
||||
@@ -1652,7 +1676,11 @@ export type AlertmanagertypesPostableChannelDTO = unknown & {
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
jira_configs?: ConfigJiraConfigDTO[];
|
||||
jira_configs?: AlertmanagertypesJiraReceiverConfigDTO[];
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
jsmops_configs?: AlertmanagertypesJSMOpsReceiverConfigDTO[];
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
@@ -1779,7 +1807,11 @@ export interface AlertmanagertypesReceiverDTO {
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
jira_configs?: ConfigJiraConfigDTO[];
|
||||
jira_configs?: AlertmanagertypesJiraReceiverConfigDTO[];
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
jsmops_configs?: AlertmanagertypesJSMOpsReceiverConfigDTO[];
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
@@ -3268,6 +3300,67 @@ export interface CommonJSONRefDTO {
|
||||
$ref?: string;
|
||||
}
|
||||
|
||||
export type ConfigJiraConfigDTOCustomFields = { [key: string]: unknown };
|
||||
|
||||
export interface ConfigJiraFieldConfigDTO {
|
||||
/**
|
||||
* @type boolean,null
|
||||
*/
|
||||
enable_update?: boolean | null;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
template?: string;
|
||||
}
|
||||
|
||||
export interface ConfigJiraConfigDTO {
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
api_type?: string;
|
||||
api_url?: ConfigURLType2DTO;
|
||||
/**
|
||||
* @type object
|
||||
*/
|
||||
custom_fields?: ConfigJiraConfigDTOCustomFields;
|
||||
description?: ConfigJiraFieldConfigDTO;
|
||||
http_config?: ConfigHTTPClientConfigDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
issue_type?: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
labels?: string[];
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
priority?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
project?: string;
|
||||
reopen_duration?: ModelDurationDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
reopen_transition?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
resolve_transition?: string;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
send_resolved?: boolean;
|
||||
summary?: ConfigJiraFieldConfigDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
wont_fix_resolution?: string;
|
||||
}
|
||||
|
||||
export interface DashboardGridItemDTO {
|
||||
content?: CommonJSONRefDTO;
|
||||
/**
|
||||
@@ -8642,6 +8735,42 @@ export enum Querybuildertypesv5QueryTypeDTO {
|
||||
clickhouse_sql = 'clickhouse_sql',
|
||||
promql = 'promql',
|
||||
}
|
||||
export interface QuickfiltertypesSourceFiltersDTO {
|
||||
/**
|
||||
* @type string
|
||||
* @format date-time
|
||||
*/
|
||||
createdAt?: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
filters: TelemetrytypesTelemetryFieldKeyDTO[];
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
id: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
orgId?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
source?: string;
|
||||
/**
|
||||
* @type string
|
||||
* @format date-time
|
||||
*/
|
||||
updatedAt?: string;
|
||||
}
|
||||
|
||||
export interface QuickfiltertypesUpdatableQuickFiltersDTO {
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
filters: TelemetrytypesTelemetryFieldKeyDTO[];
|
||||
}
|
||||
|
||||
export interface RenderErrorResponseDTO {
|
||||
error: ErrorsJSONDTO;
|
||||
/**
|
||||
@@ -12043,6 +12172,31 @@ export type GetPublicDashboardPanelQueryRangeV2200 = {
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListQuickFilters200 = {
|
||||
/**
|
||||
* @type array,null
|
||||
*/
|
||||
data: QuickfiltertypesSourceFiltersDTO[] | null;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type GetQuickFiltersPathParameters = {
|
||||
source: string;
|
||||
};
|
||||
export type GetQuickFilters200 = {
|
||||
data: QuickfiltertypesSourceFiltersDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type UpdateQuickFiltersPathParameters = {
|
||||
source: string;
|
||||
};
|
||||
export type Readyz200 = {
|
||||
data: FactoryResponseDTO;
|
||||
/**
|
||||
|
||||
557
pkg/alertmanager/alertmanagernotify/jira/jira.go
Normal file
557
pkg/alertmanager/alertmanagernotify/jira/jira.go
Normal file
@@ -0,0 +1,557 @@
|
||||
// Copyright (c) 2026 SigNoz, Inc.
|
||||
// Copyright 2023 Prometheus Team
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package jira
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode/utf16"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagertemplate"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/templating/markdownrenderer/adf"
|
||||
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/ruletypes"
|
||||
"github.com/prometheus/alertmanager/notify"
|
||||
"github.com/prometheus/alertmanager/template"
|
||||
"github.com/prometheus/alertmanager/types"
|
||||
)
|
||||
|
||||
const Integration = "jira"
|
||||
|
||||
const (
|
||||
maxSummaryLenRunes = 255
|
||||
maxDescriptionLenRunes = 32767
|
||||
)
|
||||
|
||||
// Notifier implements notify.Notifier for Jira.
|
||||
type Notifier struct {
|
||||
conf *alertmanagertypes.JiraReceiverConfig
|
||||
logger *slog.Logger
|
||||
client *http.Client
|
||||
retrier *notify.Retrier
|
||||
templater alertmanagertypes.Templater
|
||||
}
|
||||
|
||||
func New(conf *alertmanagertypes.JiraReceiverConfig, _ *template.Template, l *slog.Logger, templater alertmanagertypes.Templater) (*Notifier, error) {
|
||||
if conf.HTTPConfig == nil {
|
||||
return nil, errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "jira http_config is nil")
|
||||
}
|
||||
client, err := notify.NewClientWithTracing(*conf.HTTPConfig, Integration)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Notifier{
|
||||
conf: conf,
|
||||
logger: l,
|
||||
client: client,
|
||||
retrier: ¬ify.Retrier{RetryCodes: []int{http.StatusTooManyRequests}},
|
||||
templater: templater,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (n *Notifier) Notify(ctx context.Context, as ...*types.Alert) (bool, error) {
|
||||
key, err := notify.ExtractGroupKey(ctx)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
groupID := key.Hash()
|
||||
firing := types.Alerts(as...).HasFiring()
|
||||
n.logger.DebugContext(ctx, "sending jira notification", slog.String("group_key", key.String()), slog.Bool("firing", firing))
|
||||
|
||||
customTitle, customBody := alertmanagertemplate.ExtractTemplatesFromAnnotations(as)
|
||||
result, err := n.templater.Expand(ctx, alertmanagertypes.ExpandRequest{
|
||||
TitleTemplate: customTitle,
|
||||
BodyTemplate: customBody,
|
||||
DefaultTitleTemplate: n.conf.Summary,
|
||||
DefaultBodyTemplate: n.conf.Description,
|
||||
}, as)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
summary := truncateRunes(result.Title, maxSummaryLenRunes)
|
||||
|
||||
var parts []string
|
||||
for _, body := range result.Body {
|
||||
if body != "" {
|
||||
parts = append(parts, body)
|
||||
}
|
||||
}
|
||||
// custom body templates render per alert; join them under ADF rule dividers.
|
||||
// The default body is a single combined part, so the join is a no-op there.
|
||||
descText := truncateRunes(strings.Join(parts, "\n\n---\n\n"), maxDescriptionLenRunes)
|
||||
|
||||
baseURL, retry, err := n.resolveAPIBaseURL(ctx)
|
||||
if err != nil {
|
||||
return retry, err
|
||||
}
|
||||
|
||||
existing, retry, err := n.searchIssue(ctx, baseURL, groupID, firing)
|
||||
if err != nil {
|
||||
return retry, err
|
||||
}
|
||||
|
||||
fields := n.buildFields(groupID, summary, descText, as, firing)
|
||||
|
||||
// No existing issue: create for firing groups; never create for resolved-only.
|
||||
if existing == nil {
|
||||
if !firing {
|
||||
return false, nil
|
||||
}
|
||||
return n.createIssue(ctx, baseURL, fields)
|
||||
}
|
||||
|
||||
// Existing issue: refresh it, then transition + comment based on the new state.
|
||||
if retry, err := n.updateIssue(ctx, baseURL, existing, fields); err != nil {
|
||||
return retry, err
|
||||
}
|
||||
|
||||
// Each state-change comment carries the same rich snapshot as the description
|
||||
// (panel + details + deep-links), so the comment timeline mirrors the card
|
||||
// Google Chat re-posts on every notification.
|
||||
switch {
|
||||
case firing && existing.isDone(): // re-fired after resolution → reopen
|
||||
if retry, err := n.applyTransition(ctx, baseURL, existing.Key, false, n.conf.ReopenTransition); err != nil {
|
||||
return retry, err
|
||||
}
|
||||
case !firing: // resolved (search returns only open issues, so this one is open)
|
||||
if retry, err := n.applyTransition(ctx, baseURL, existing.Key, true, n.conf.ResolveTransition); err != nil {
|
||||
return retry, err
|
||||
}
|
||||
}
|
||||
// firing && !isDone (still firing) needs no transition.
|
||||
return n.addComment(ctx, baseURL, existing.Key, fields.Description)
|
||||
}
|
||||
|
||||
func (n *Notifier) buildFields(groupID, summary, descText string, alerts []*types.Alert, firing bool) *issueFields {
|
||||
f := &issueFields{
|
||||
Project: &idKey{Key: n.conf.Project},
|
||||
Issuetype: &idName{Name: n.conf.IssueType},
|
||||
Summary: summary,
|
||||
Labels: n.labels(groupID),
|
||||
Description: n.buildBoundedDescription(descText, alerts, firing),
|
||||
}
|
||||
if n.conf.Priority != "" {
|
||||
f.Priority = &idName{Name: n.conf.Priority}
|
||||
}
|
||||
return f
|
||||
}
|
||||
|
||||
// buildBoundedDescription builds the ADF issue body and keeps it within Jira's
|
||||
// description limit, which counts text characters plus per-node overhead — so a
|
||||
// text-only markdown cap is not enough. Over-limit bodies are shrunk at the
|
||||
// markdown level and rebuilt; the panel and deep-links are part of the measured
|
||||
// document, so the result always fits.
|
||||
func (n *Notifier) buildBoundedDescription(descText string, alerts []*types.Alert, firing bool) map[string]any {
|
||||
doc := n.buildDescription(descText, alerts, firing)
|
||||
for range 4 {
|
||||
size := adfDocLen(doc)
|
||||
if size <= maxDescriptionLenRunes {
|
||||
return doc
|
||||
}
|
||||
runes := []rune(descText)
|
||||
keep := len(runes) * maxDescriptionLenRunes / size * 9 / 10
|
||||
if keep >= len(runes) {
|
||||
keep = len(runes) - 1
|
||||
}
|
||||
if keep <= 0 {
|
||||
break
|
||||
}
|
||||
descText = string(runes[:keep]) + "…"
|
||||
doc = n.buildDescription(descText, alerts, firing)
|
||||
}
|
||||
if adfDocLen(doc) <= maxDescriptionLenRunes {
|
||||
return doc
|
||||
}
|
||||
// still over after shrinking: keep just the panel and deep-links
|
||||
return n.buildDescription("", alerts, firing)
|
||||
}
|
||||
|
||||
// adfDocLen approximates how Jira measures an ADF document against the 32767
|
||||
// limit: text length in UTF-16 code units, plus per-node overhead (block
|
||||
// boundaries count like newlines), plus link targets. Deliberately counts on
|
||||
// the high side so a passing measurement never 400s.
|
||||
func adfDocLen(node any) int {
|
||||
m, ok := node.(map[string]any)
|
||||
if !ok {
|
||||
return 0
|
||||
}
|
||||
size := 2
|
||||
if text, ok := m["text"].(string); ok {
|
||||
for _, r := range text {
|
||||
size += utf16.RuneLen(r)
|
||||
}
|
||||
}
|
||||
if marks, ok := m["marks"].([]any); ok {
|
||||
for _, mark := range marks {
|
||||
if mm, ok := mark.(map[string]any); ok {
|
||||
if attrs, ok := mm["attrs"].(map[string]any); ok {
|
||||
if href, ok := attrs["href"].(string); ok {
|
||||
size += len(href)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if content, ok := m["content"].([]any); ok {
|
||||
for _, child := range content {
|
||||
size += adfDocLen(child)
|
||||
}
|
||||
}
|
||||
return size
|
||||
}
|
||||
|
||||
// buildDescription assembles the ADF issue body: a firing/resolved status panel,
|
||||
// the rendered markdown body, and SigNoz deep-links.
|
||||
func (n *Notifier) buildDescription(descText string, alerts []*types.Alert, firing bool) map[string]any {
|
||||
content := []any{statusPanel(firing)}
|
||||
content = append(content, adf.Render(descText)...)
|
||||
if links := deepLinks(alerts); links != nil {
|
||||
content = append(content, links)
|
||||
}
|
||||
return map[string]any{"type": "doc", "version": 1, "content": content}
|
||||
}
|
||||
|
||||
func statusPanel(firing bool) map[string]any {
|
||||
panelType, label := "success", "🟢 RESOLVED"
|
||||
if firing {
|
||||
panelType, label = "error", "🔴 FIRING"
|
||||
}
|
||||
return map[string]any{
|
||||
"type": "panel",
|
||||
"attrs": map[string]any{"panelType": panelType},
|
||||
"content": []any{map[string]any{
|
||||
"type": "paragraph",
|
||||
"content": []any{map[string]any{"type": "text", "text": label, "marks": []any{map[string]any{"type": "strong"}}}},
|
||||
}},
|
||||
}
|
||||
}
|
||||
|
||||
// deepLinks builds a paragraph of SigNoz links from the per-rule ruleSource label
|
||||
// and the related-logs/traces annotations. Returns nil when none are present.
|
||||
func deepLinks(alerts []*types.Alert) map[string]any {
|
||||
if len(alerts) == 0 {
|
||||
return nil
|
||||
}
|
||||
a := alerts[0]
|
||||
var parts []any
|
||||
add := func(label, url string) {
|
||||
if url == "" {
|
||||
return
|
||||
}
|
||||
if len(parts) > 0 {
|
||||
parts = append(parts, map[string]any{"type": "text", "text": " · "})
|
||||
}
|
||||
parts = append(parts, map[string]any{
|
||||
"type": "text",
|
||||
"text": label,
|
||||
"marks": []any{map[string]any{"type": "link", "attrs": map[string]any{"href": url}}},
|
||||
})
|
||||
}
|
||||
add("Open in SigNoz", string(a.Labels[ruletypes.LabelRuleSource]))
|
||||
add("View Related Logs", string(a.Annotations[ruletypes.AnnotationRelatedLogs]))
|
||||
add("View Related Traces", string(a.Annotations[ruletypes.AnnotationRelatedTraces]))
|
||||
if len(parts) == 0 {
|
||||
return nil
|
||||
}
|
||||
return map[string]any{"type": "paragraph", "content": parts}
|
||||
}
|
||||
|
||||
func (n *Notifier) labels(groupID string) []string {
|
||||
out := append([]string{}, n.conf.Labels...)
|
||||
out = append(out, "signoz-alert", fmt.Sprintf("ALERT{%s}", groupID))
|
||||
sort.Strings(out)
|
||||
return out
|
||||
}
|
||||
|
||||
func (n *Notifier) searchIssue(ctx context.Context, baseURL, groupID string, firing bool) (*issue, bool, error) {
|
||||
var jql strings.Builder
|
||||
if n.conf.WontFixResolution != "" {
|
||||
// != alone also drops unresolved (EMPTY) issues, so keep those explicitly.
|
||||
fmt.Fprintf(&jql, `(resolution is EMPTY or resolution != %q) and `, n.conf.WontFixResolution)
|
||||
}
|
||||
if reopenMin := int64(time.Duration(n.conf.ReopenDuration).Minutes()); firing && reopenMin > 0 {
|
||||
fmt.Fprintf(&jql, `(resolutiondate is EMPTY OR resolutiondate >= -%dm) and `, reopenMin)
|
||||
} else {
|
||||
jql.WriteString(`statusCategory != Done and `)
|
||||
}
|
||||
fmt.Fprintf(&jql, `project=%q and labels=%q order by status ASC, resolutiondate DESC`, n.conf.Project, fmt.Sprintf("ALERT{%s}", groupID))
|
||||
|
||||
body, retry, err := n.callAPI(ctx, http.MethodPost, baseURL+"/search/jql", searchRequest{
|
||||
JQL: jql.String(), MaxResults: 2, Fields: []string{"status", "labels"},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, retry, err
|
||||
}
|
||||
var res searchResult
|
||||
if err := json.Unmarshal(body, &res); err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
if len(res.Issues) == 0 {
|
||||
return nil, false, nil
|
||||
}
|
||||
// the JQL order is not category-aware, so prefer an open issue over a done
|
||||
// one; all done falls back to the most recently resolved (resolutiondate DESC)
|
||||
for i := range res.Issues {
|
||||
if !res.Issues[i].isDone() {
|
||||
return &res.Issues[i], false, nil
|
||||
}
|
||||
}
|
||||
return &res.Issues[0], false, nil
|
||||
}
|
||||
|
||||
func (n *Notifier) createIssue(ctx context.Context, baseURL string, fields *issueFields) (bool, error) {
|
||||
_, retry, err := n.callAPI(ctx, http.MethodPost, baseURL+"/issue", issue{Fields: fields})
|
||||
return retry, err
|
||||
}
|
||||
|
||||
func (n *Notifier) updateIssue(ctx context.Context, baseURL string, existing *issue, fields *issueFields) (bool, error) {
|
||||
// project and issue type are set at creation and cannot be edited.
|
||||
upd := *fields
|
||||
upd.Project = nil
|
||||
upd.Issuetype = nil
|
||||
// Jira replaces the labels array wholesale, so union in the labels already
|
||||
// on the issue to keep user-added ones.
|
||||
if existing.Fields != nil {
|
||||
upd.Labels = mergeLabels(existing.Fields.Labels, fields.Labels)
|
||||
}
|
||||
_, retry, err := n.callAPI(ctx, http.MethodPut, n.issueURL(baseURL, existing.Key, ""), issue{Fields: &upd})
|
||||
return retry, err
|
||||
}
|
||||
|
||||
func mergeLabels(existing, ours []string) []string {
|
||||
seen := make(map[string]bool, len(existing)+len(ours))
|
||||
var merged []string
|
||||
for _, label := range append(append([]string{}, existing...), ours...) {
|
||||
if !seen[label] {
|
||||
seen[label] = true
|
||||
merged = append(merged, label)
|
||||
}
|
||||
}
|
||||
sort.Strings(merged)
|
||||
return merged
|
||||
}
|
||||
|
||||
// applyTransition moves the issue into (toDone) or out of (!toDone) the "done"
|
||||
// status category, preferring the named override, else the first matching
|
||||
// transition, else skipping without error when none is available.
|
||||
func (n *Notifier) applyTransition(ctx context.Context, baseURL, key string, toDone bool, override string) (bool, error) {
|
||||
transitions, retry, err := n.getTransitions(ctx, baseURL, key)
|
||||
if err != nil {
|
||||
return retry, err
|
||||
}
|
||||
id := selectTransition(transitions, toDone, override)
|
||||
if id == "" {
|
||||
n.logger.WarnContext(ctx, "jira: no matching transition, leaving issue as-is", slog.String("issue", key), slog.Bool("to_done", toDone))
|
||||
return false, nil
|
||||
}
|
||||
_, retry, err = n.callAPI(ctx, http.MethodPost, n.issueURL(baseURL, key, "transitions"), issue{Transition: &idName{ID: id}})
|
||||
return retry, err
|
||||
}
|
||||
|
||||
func (n *Notifier) getTransitions(ctx context.Context, baseURL, key string) ([]jiraTransition, bool, error) {
|
||||
body, retry, err := n.callAPI(ctx, http.MethodGet, n.issueURL(baseURL, key, "transitions"), nil)
|
||||
if err != nil {
|
||||
return nil, retry, err
|
||||
}
|
||||
var tr transitionsResponse
|
||||
if err := json.Unmarshal(body, &tr); err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
return tr.Transitions, false, nil
|
||||
}
|
||||
|
||||
func (n *Notifier) addComment(ctx context.Context, baseURL, key string, body any) (bool, error) {
|
||||
_, retry, err := n.callAPI(ctx, http.MethodPost, n.issueURL(baseURL, key, "comment"), comment{Body: body})
|
||||
return retry, err
|
||||
}
|
||||
|
||||
func (n *Notifier) issueURL(baseURL, key, sub string) string {
|
||||
u := baseURL + "/issue/" + key
|
||||
if sub != "" {
|
||||
u += "/" + sub
|
||||
}
|
||||
return u
|
||||
}
|
||||
|
||||
func (n *Notifier) callAPI(ctx context.Context, method, url string, reqBody any) ([]byte, bool, error) {
|
||||
var body io.Reader
|
||||
if reqBody != nil {
|
||||
var buf bytes.Buffer
|
||||
if err := json.NewEncoder(&buf).Encode(reqBody); err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
body = &buf
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, method, url, body)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
resp, err := n.client.Do(req) //nolint:bodyclose // notify.Drain closes the body
|
||||
if err != nil {
|
||||
return nil, true, notify.RedactURL(err)
|
||||
}
|
||||
defer notify.Drain(resp)
|
||||
|
||||
respBody, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
shouldRetry, err := n.retrier.Check(resp.StatusCode, bytes.NewReader(respBody))
|
||||
if err != nil {
|
||||
return respBody, shouldRetry, notify.NewErrorWithReason(notify.GetFailureReasonFromStatusCode(resp.StatusCode), err)
|
||||
}
|
||||
return respBody, false, nil
|
||||
}
|
||||
|
||||
// resolveAPIBaseURL resolves the service-account cloud id per notification (it
|
||||
// is never persisted); personal API tokens use the site host directly.
|
||||
func (n *Notifier) resolveAPIBaseURL(ctx context.Context) (string, bool, error) {
|
||||
if !n.conf.IsServiceAccount() {
|
||||
return n.conf.APIBaseURL(""), false, nil
|
||||
}
|
||||
cloudID, retry, err := n.resolveCloudID(ctx)
|
||||
if err != nil {
|
||||
return "", retry, err
|
||||
}
|
||||
return n.conf.APIBaseURL(cloudID), false, nil
|
||||
}
|
||||
|
||||
// resolveCloudID fetches the site's cloud id from its unauthenticated
|
||||
// tenant_info endpoint; transport failures are retryable, bad responses are not.
|
||||
func (n *Notifier) resolveCloudID(ctx context.Context) (string, bool, error) {
|
||||
url := strings.TrimRight(n.conf.Site, "/") + "/_edge/tenant_info"
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
resp, err := n.client.Do(req)
|
||||
if err != nil {
|
||||
return "", true, errors.WrapInternalf(err, errors.CodeInternal, "failed to fetch jira cloud id")
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", true, err
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return "", false, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "failed to resolve jira cloud id from %s: status %d", url, resp.StatusCode)
|
||||
}
|
||||
|
||||
var out struct {
|
||||
CloudID string `json:"cloudId"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &out); err != nil {
|
||||
return "", false, errors.WrapInternalf(err, errors.CodeInternal, "failed to parse jira tenant_info response")
|
||||
}
|
||||
if out.CloudID == "" {
|
||||
return "", false, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "jira tenant_info returned an empty cloud id for %s", n.conf.Site)
|
||||
}
|
||||
return out.CloudID, false, nil
|
||||
}
|
||||
|
||||
// selectTransition returns the id of the transition whose target status category
|
||||
// matches toDone, preferring one named override when present.
|
||||
func selectTransition(transitions []jiraTransition, toDone bool, override string) string {
|
||||
if override != "" {
|
||||
for _, t := range transitions {
|
||||
if t.Name == override {
|
||||
return t.ID
|
||||
}
|
||||
}
|
||||
}
|
||||
for _, t := range transitions {
|
||||
if (t.To.StatusCategory.Key == "done") == toDone {
|
||||
return t.ID
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Jira API types.
|
||||
type issue struct {
|
||||
Key string `json:"key,omitempty"`
|
||||
Fields *issueFields `json:"fields,omitempty"`
|
||||
Transition *idName `json:"transition,omitempty"`
|
||||
}
|
||||
|
||||
type issueFields struct {
|
||||
Project *idKey `json:"project,omitempty"`
|
||||
Issuetype *idName `json:"issuetype,omitempty"`
|
||||
Summary string `json:"summary,omitempty"`
|
||||
Labels []string `json:"labels,omitempty"`
|
||||
Priority *idName `json:"priority,omitempty"`
|
||||
Description any `json:"description,omitempty"`
|
||||
Status *issueStatus `json:"status,omitempty"`
|
||||
}
|
||||
|
||||
type idKey struct {
|
||||
Key string `json:"key"`
|
||||
}
|
||||
|
||||
type idName struct {
|
||||
ID string `json:"id,omitempty"`
|
||||
Name string `json:"name,omitempty"`
|
||||
}
|
||||
|
||||
type issueStatus struct {
|
||||
StatusCategory struct {
|
||||
Key string `json:"key"`
|
||||
} `json:"statusCategory"`
|
||||
}
|
||||
|
||||
func (i *issue) isDone() bool {
|
||||
return i.Fields != nil && i.Fields.Status != nil && i.Fields.Status.StatusCategory.Key == "done"
|
||||
}
|
||||
|
||||
type searchRequest struct {
|
||||
JQL string `json:"jql"`
|
||||
MaxResults int `json:"maxResults"`
|
||||
Fields []string `json:"fields"`
|
||||
}
|
||||
|
||||
type searchResult struct {
|
||||
Issues []issue `json:"issues"`
|
||||
}
|
||||
|
||||
type transitionsResponse struct {
|
||||
Transitions []jiraTransition `json:"transitions"`
|
||||
}
|
||||
|
||||
type jiraTransition struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
To struct {
|
||||
StatusCategory struct {
|
||||
Key string `json:"key"`
|
||||
} `json:"statusCategory"`
|
||||
} `json:"to"`
|
||||
}
|
||||
|
||||
type comment struct {
|
||||
Body any `json:"body"`
|
||||
}
|
||||
|
||||
func truncateRunes(s string, max int) string {
|
||||
r := []rune(s)
|
||||
if len(r) <= max {
|
||||
return s
|
||||
}
|
||||
return string(r[:max])
|
||||
}
|
||||
489
pkg/alertmanager/alertmanagernotify/jira/jira_test.go
Normal file
489
pkg/alertmanager/alertmanagernotify/jira/jira_test.go
Normal file
@@ -0,0 +1,489 @@
|
||||
package jira
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagertemplate"
|
||||
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/ruletypes"
|
||||
"github.com/prometheus/alertmanager/notify"
|
||||
"github.com/prometheus/alertmanager/notify/test"
|
||||
"github.com/prometheus/alertmanager/types"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
"github.com/prometheus/common/model"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
type mockReq struct {
|
||||
method string
|
||||
path string
|
||||
body map[string]any
|
||||
}
|
||||
|
||||
type mockJira struct {
|
||||
srv *httptest.Server
|
||||
mu sync.Mutex
|
||||
reqs []mockReq
|
||||
searchIssues []issue
|
||||
transitions []jiraTransition
|
||||
createStatus int
|
||||
}
|
||||
|
||||
func newMockJira(t *testing.T) *mockJira {
|
||||
t.Helper()
|
||||
m := &mockJira{}
|
||||
m.srv = httptest.NewServer(http.HandlerFunc(m.handle))
|
||||
t.Cleanup(m.srv.Close)
|
||||
return m
|
||||
}
|
||||
|
||||
func (m *mockJira) handle(w http.ResponseWriter, r *http.Request) {
|
||||
var body map[string]any
|
||||
_ = json.NewDecoder(r.Body).Decode(&body)
|
||||
m.mu.Lock()
|
||||
m.reqs = append(m.reqs, mockReq{r.Method, r.URL.Path, body})
|
||||
m.mu.Unlock()
|
||||
|
||||
p := r.URL.Path
|
||||
switch {
|
||||
case strings.HasSuffix(p, "/search/jql"):
|
||||
_ = json.NewEncoder(w).Encode(searchResult{Issues: m.searchIssues})
|
||||
case strings.HasSuffix(p, "/transitions") && r.Method == http.MethodGet:
|
||||
_ = json.NewEncoder(w).Encode(transitionsResponse{Transitions: m.transitions})
|
||||
case strings.HasSuffix(p, "/transitions"):
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
case strings.HasSuffix(p, "/comment"):
|
||||
w.WriteHeader(http.StatusCreated)
|
||||
_, _ = w.Write([]byte(`{"id":"1"}`))
|
||||
case strings.HasSuffix(p, "/issue") && r.Method == http.MethodPost:
|
||||
st := m.createStatus
|
||||
if st == 0 {
|
||||
st = http.StatusCreated
|
||||
}
|
||||
w.WriteHeader(st)
|
||||
_, _ = w.Write([]byte(`{"key":"KAN-1"}`))
|
||||
case r.Method == http.MethodPut:
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
default:
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
}
|
||||
}
|
||||
|
||||
func (m *mockJira) countPost(suffix string) int {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
c := 0
|
||||
for _, r := range m.reqs {
|
||||
if r.method == http.MethodPost && strings.HasSuffix(r.path, suffix) {
|
||||
c++
|
||||
}
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
func (m *mockJira) countPuts() int {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
c := 0
|
||||
for _, r := range m.reqs {
|
||||
if r.method == http.MethodPut {
|
||||
c++
|
||||
}
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
func newNotifier(t *testing.T, m *mockJira) *Notifier {
|
||||
t.Helper()
|
||||
tmpl := test.CreateTmpl(t)
|
||||
n, err := New(&alertmanagertypes.JiraReceiverConfig{
|
||||
Site: m.srv.URL,
|
||||
Project: "KAN",
|
||||
IssueType: "Task",
|
||||
Summary: alertmanagertypes.DefaultJiraSummaryTemplate,
|
||||
Description: alertmanagertypes.DefaultJiraDescriptionTemplate,
|
||||
HTTPConfig: &commoncfg.HTTPClientConfig{},
|
||||
ReopenDuration: model.Duration(3 * 24 * time.Hour),
|
||||
}, tmpl, slog.New(slog.DiscardHandler), alertmanagertemplate.New(tmpl, slog.New(slog.DiscardHandler)))
|
||||
require.NoError(t, err)
|
||||
return n
|
||||
}
|
||||
|
||||
func alert(firing bool) *types.Alert {
|
||||
a := &types.Alert{Alert: model.Alert{
|
||||
Labels: model.LabelSet{"alertname": "HighCPU", "severity": "critical"},
|
||||
Annotations: model.LabelSet{"summary": "cpu high"},
|
||||
StartsAt: time.Now().Add(-time.Minute),
|
||||
}}
|
||||
if firing {
|
||||
a.EndsAt = time.Now().Add(time.Hour)
|
||||
} else {
|
||||
a.EndsAt = time.Now().Add(-time.Minute)
|
||||
}
|
||||
return a
|
||||
}
|
||||
|
||||
func ctx() context.Context {
|
||||
return notify.WithGroupKey(context.Background(), "test-jira")
|
||||
}
|
||||
|
||||
func doneIssue() issue {
|
||||
i := issue{Key: "KAN-1", Fields: &issueFields{Status: &issueStatus{}}}
|
||||
i.Fields.Status.StatusCategory.Key = "done"
|
||||
return i
|
||||
}
|
||||
|
||||
func openIssue() issue {
|
||||
i := issue{Key: "KAN-1", Fields: &issueFields{Status: &issueStatus{}}}
|
||||
i.Fields.Status.StatusCategory.Key = "new"
|
||||
return i
|
||||
}
|
||||
|
||||
func transition(id, name, category string) jiraTransition {
|
||||
tr := jiraTransition{ID: id, Name: name}
|
||||
tr.To.StatusCategory.Key = category
|
||||
return tr
|
||||
}
|
||||
|
||||
func TestNotifyCreatesWhenNoExistingIssue(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 1, m.countPost("/issue"))
|
||||
assert.Equal(t, 0, m.countPost("/comment")) // no comment on create
|
||||
assert.Equal(t, 0, m.countPuts()) // no update
|
||||
}
|
||||
|
||||
func TestNotifyResolvedOnlyWithNoIssueIsNoop(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(false))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 1, m.countPost("/search/jql"))
|
||||
assert.Equal(t, 0, m.countPost("/issue"))
|
||||
}
|
||||
|
||||
func TestNotifyStillFiringUpdatesAndComments(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
m.searchIssues = []issue{openIssue()}
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 0, m.countPost("/issue")) // no create
|
||||
assert.Equal(t, 1, m.countPuts()) // update
|
||||
assert.Equal(t, 1, m.countPost("/comment"))
|
||||
assert.Equal(t, 0, m.countPost("/transitions")) // still open, no transition
|
||||
|
||||
// comment carries the full rich snapshot (panel + labeled body), not a one-liner.
|
||||
cjs, err := json.Marshal(m.lastBody(t, "/comment"))
|
||||
require.NoError(t, err)
|
||||
assert.Contains(t, string(cjs), `"panel"`)
|
||||
assert.Contains(t, string(cjs), "Summary:")
|
||||
}
|
||||
|
||||
func TestNotifyResolveTransitionsToDoneAndComments(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
m.searchIssues = []issue{openIssue()}
|
||||
m.transitions = []jiraTransition{transition("11", "To Do", "new"), transition("41", "Done", "done")}
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(false))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 1, m.countPuts()) // update
|
||||
assert.Equal(t, 1, m.countPost("/transitions")) // resolve transition
|
||||
assert.Equal(t, 1, m.countPost("/comment"))
|
||||
}
|
||||
|
||||
func TestNotifyReopensDoneIssue(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
m.searchIssues = []issue{doneIssue()}
|
||||
m.transitions = []jiraTransition{transition("11", "To Do", "new"), transition("41", "Done", "done")}
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 1, m.countPost("/transitions")) // reopen transition
|
||||
assert.Equal(t, 1, m.countPost("/comment"))
|
||||
}
|
||||
|
||||
func TestNotifySafeSkipsWhenNoMatchingTransition(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
m.searchIssues = []issue{openIssue()}
|
||||
m.transitions = []jiraTransition{transition("11", "To Do", "new")} // no done-category transition
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(false))
|
||||
require.NoError(t, err) // must not error
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 0, m.countPost("/transitions")) // skipped
|
||||
assert.Equal(t, 1, m.countPost("/comment")) // comment still posted
|
||||
}
|
||||
|
||||
func TestNotifyPrefersOpenIssueOverRecentlyDone(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
open := openIssue()
|
||||
open.Key = "KAN-2"
|
||||
// the JQL order can put a recently-done issue first; the open one must win
|
||||
m.searchIssues = []issue{doneIssue(), open}
|
||||
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, retry)
|
||||
assert.Equal(t, 0, m.countPost("/issue")) // no duplicate create
|
||||
assert.Equal(t, 0, m.countPost("/transitions")) // open issue → no reopen
|
||||
assert.Equal(t, 1, m.countPuts())
|
||||
assert.Equal(t, 1, m.countPost("/comment"))
|
||||
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
for _, r := range m.reqs {
|
||||
if r.method == http.MethodPut || strings.HasSuffix(r.path, "/comment") {
|
||||
assert.Contains(t, r.path, "KAN-2")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestNotifyRetriesOn429(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
m.createStatus = http.StatusTooManyRequests
|
||||
retry, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.Error(t, err)
|
||||
assert.True(t, retry)
|
||||
}
|
||||
|
||||
func (m *mockJira) lastBody(t *testing.T, suffix string) map[string]any {
|
||||
t.Helper()
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
for i := len(m.reqs) - 1; i >= 0; i-- {
|
||||
if m.reqs[i].method == http.MethodPost && strings.HasSuffix(m.reqs[i].path, suffix) {
|
||||
return m.reqs[i].body
|
||||
}
|
||||
}
|
||||
t.Fatalf("no POST request to %s", suffix)
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestNotifyRichDescriptionPanelAndLinks(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
a := alert(true)
|
||||
a.Labels[ruletypes.LabelRuleSource] = model.LabelValue("https://app.signoz.io/alerts?ruleId=1")
|
||||
a.Annotations[ruletypes.AnnotationRelatedLogs] = model.LabelValue("https://app.signoz.io/logs")
|
||||
|
||||
_, err := newNotifier(t, m).Notify(ctx(), a)
|
||||
require.NoError(t, err)
|
||||
|
||||
body := m.lastBody(t, "/issue")
|
||||
js, err := json.Marshal(body)
|
||||
require.NoError(t, err)
|
||||
s := string(js)
|
||||
assert.Contains(t, s, `"panel"`) // status panel present
|
||||
assert.Contains(t, s, `"error"`) // firing → error panel
|
||||
assert.Contains(t, s, "Open in SigNoz") // rule deep-link
|
||||
assert.Contains(t, s, "https://app.signoz.io/alerts?ruleId=1") // rule url
|
||||
assert.Contains(t, s, "View Related Logs") // related-logs deep-link
|
||||
assert.Contains(t, s, "Summary:") // labeled body section
|
||||
assert.Contains(t, s, "cpu high") // rendered annotation
|
||||
}
|
||||
|
||||
func TestNotifyCustomTemplateAnnotationsOverrideDefaults(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
a1 := alert(true)
|
||||
a1.Labels["service"] = "payment"
|
||||
a1.Labels["namespace"] = "ns-one"
|
||||
a1.Annotations[ruletypes.AnnotationTitleTemplate] = "High throughput for $service"
|
||||
a1.Annotations[ruletypes.AnnotationBodyTemplate] = "Firing in NS: $labels.namespace"
|
||||
a2 := alert(true)
|
||||
a2.Labels["service"] = "payment"
|
||||
a2.Labels["namespace"] = "ns-two"
|
||||
a2.Annotations[ruletypes.AnnotationTitleTemplate] = "High throughput for $service"
|
||||
a2.Annotations[ruletypes.AnnotationBodyTemplate] = "Firing in NS: $labels.namespace"
|
||||
|
||||
_, err := newNotifier(t, m).Notify(ctx(), a1, a2)
|
||||
require.NoError(t, err)
|
||||
|
||||
body := m.lastBody(t, "/issue")
|
||||
fields, ok := body["fields"].(map[string]any)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, "High throughput for payment", fields["summary"])
|
||||
|
||||
js, err := json.Marshal(fields["description"])
|
||||
require.NoError(t, err)
|
||||
s := string(js)
|
||||
assert.Contains(t, s, "Firing in NS: ns-one")
|
||||
assert.Contains(t, s, "Firing in NS: ns-two")
|
||||
// per-alert custom bodies are separated by an ADF rule divider
|
||||
assert.Contains(t, s, `"rule"`)
|
||||
assert.NotContains(t, s, "Summary:") // default body template not used
|
||||
}
|
||||
|
||||
// Jira replaces labels wholesale on PUT, so the update must union in the
|
||||
// labels already on the issue or user-added ones get wiped.
|
||||
func TestNotifyUpdatePreservesUserAddedLabels(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
existing := openIssue()
|
||||
existing.Fields.Labels = []string{"user-added-label", "signoz-alert"}
|
||||
m.searchIssues = []issue{existing}
|
||||
|
||||
_, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.NoError(t, err)
|
||||
|
||||
search := m.lastBody(t, "/search/jql")
|
||||
assert.Contains(t, search["fields"], "labels")
|
||||
|
||||
m.mu.Lock()
|
||||
var putLabels []any
|
||||
for _, r := range m.reqs {
|
||||
if r.method == http.MethodPut {
|
||||
putLabels, _ = r.body["fields"].(map[string]any)["labels"].([]any)
|
||||
}
|
||||
}
|
||||
m.mu.Unlock()
|
||||
assert.Contains(t, putLabels, "user-added-label")
|
||||
assert.Contains(t, putLabels, "signoz-alert")
|
||||
assert.Equal(t, 1, strings.Count(fmt.Sprint(putLabels), "signoz-alert")) // no duplicates
|
||||
// the dedup label is re-asserted
|
||||
found := false
|
||||
for _, l := range putLabels {
|
||||
if s, ok := l.(string); ok && strings.HasPrefix(s, "ALERT{") {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
assert.True(t, found)
|
||||
}
|
||||
|
||||
func TestADFDocLen(t *testing.T) {
|
||||
text := func(s string) map[string]any { return map[string]any{"type": "text", "text": s} }
|
||||
para := func(children ...any) map[string]any {
|
||||
return map[string]any{"type": "paragraph", "content": children}
|
||||
}
|
||||
cases := []struct {
|
||||
name string
|
||||
node any
|
||||
want int
|
||||
}{
|
||||
{"text node", text("hello"), 7}, // 5 utf16 + 2 overhead
|
||||
{"emoji counts utf16", text("🔴"), 4}, // 2 utf16 units + 2 overhead
|
||||
{"paragraph wraps text", para(text("hi")), 6}, // 2 + (2+2)
|
||||
{"link href counted", map[string]any{"type": "text", "text": "a", "marks": []any{map[string]any{"type": "link", "attrs": map[string]any{"href": "https://x"}}}}, 12}, // 1 + 9 href + 2
|
||||
{"non-map is zero", "junk", 0},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
assert.Equal(t, c.want, adfDocLen(c.node))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// 30 fat custom bodies overflow Jira's description accounting (text + per-node
|
||||
// overhead); the built doc must be shrunk under the limit, never rejected.
|
||||
func TestNotifyDescriptionShrunkUnderJiraLimit(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
filler := strings.Repeat("This is a long runbook detail line used to inflate the alert body. ", 25)
|
||||
alerts := make([]*types.Alert, 0, 30)
|
||||
for i := range 30 {
|
||||
a := alert(true)
|
||||
a.Labels["service"] = model.LabelValue(strings.Repeat("s", 3) + string(rune('a'+i%26)))
|
||||
a.Annotations[ruletypes.AnnotationTitleTemplate] = "overflow probe"
|
||||
a.Annotations[ruletypes.AnnotationBodyTemplate] = model.LabelValue("**Alert in service** $labels.service\n\n" + filler)
|
||||
alerts = append(alerts, a)
|
||||
}
|
||||
|
||||
_, err := newNotifier(t, m).Notify(ctx(), alerts...)
|
||||
require.NoError(t, err)
|
||||
|
||||
body := m.lastBody(t, "/issue")
|
||||
fields, ok := body["fields"].(map[string]any)
|
||||
require.True(t, ok)
|
||||
desc := fields["description"]
|
||||
assert.LessOrEqual(t, adfDocLen(desc), maxDescriptionLenRunes)
|
||||
|
||||
js, err := json.Marshal(desc)
|
||||
require.NoError(t, err)
|
||||
assert.Contains(t, string(js), "FIRING") // status panel survives the shrink
|
||||
assert.Contains(t, string(js), "…") // body ends with the shrink marker
|
||||
}
|
||||
|
||||
func TestFiringSearchJQLHasReopenWindow(t *testing.T) {
|
||||
m := newMockJira(t)
|
||||
_, err := newNotifier(t, m).Notify(ctx(), alert(true))
|
||||
require.NoError(t, err)
|
||||
|
||||
body := m.lastBody(t, "/search/jql")
|
||||
jql, ok := body["jql"].(string)
|
||||
require.True(t, ok)
|
||||
// newNotifier uses a 3d window → 4320 minutes.
|
||||
assert.Contains(t, jql, "resolutiondate >= -4320m")
|
||||
}
|
||||
|
||||
func TestSelectTransition(t *testing.T) {
|
||||
ts := []jiraTransition{
|
||||
transition("41", "Done", "done"),
|
||||
transition("51", "Won't Do", "done"),
|
||||
transition("11", "To Do", "new"),
|
||||
}
|
||||
assert.Equal(t, "41", selectTransition(ts, true, "")) // first done-category
|
||||
assert.Equal(t, "51", selectTransition(ts, true, "Won't Do")) // named override
|
||||
assert.Equal(t, "11", selectTransition(ts, false, "")) // first non-done
|
||||
assert.Equal(t, "41", selectTransition(ts, true, "Nonexistent")) // bad override → fallback
|
||||
assert.Equal(t, "", selectTransition([]jiraTransition{transition("11", "To Do", "new")}, true, "")) // none → skip
|
||||
}
|
||||
|
||||
func TestResolveCloudID(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
handler http.HandlerFunc
|
||||
want string
|
||||
wantErr bool
|
||||
wantRetry bool
|
||||
}{
|
||||
{
|
||||
name: "success",
|
||||
handler: func(w http.ResponseWriter, r *http.Request) {
|
||||
assert.Equal(t, "/_edge/tenant_info", r.URL.Path)
|
||||
_, _ = w.Write([]byte(`{"cloudId":"abc-123"}`))
|
||||
},
|
||||
want: "abc-123",
|
||||
},
|
||||
{
|
||||
name: "non-200 is not retryable",
|
||||
handler: func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusNotFound) },
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "empty cloud id",
|
||||
handler: func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write([]byte(`{"cloudId":""}`)) },
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "bad json",
|
||||
handler: func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write([]byte(`not json`)) },
|
||||
wantErr: true,
|
||||
},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
srv := httptest.NewServer(c.handler)
|
||||
defer srv.Close()
|
||||
|
||||
n, err := New(&alertmanagertypes.JiraReceiverConfig{Site: srv.URL, HTTPConfig: &commoncfg.HTTPClientConfig{}}, nil, slog.New(slog.DiscardHandler), nil)
|
||||
require.NoError(t, err)
|
||||
|
||||
got, retry, err := n.resolveCloudID(context.Background())
|
||||
if c.wantErr {
|
||||
assert.Error(t, err)
|
||||
assert.Equal(t, c.wantRetry, retry)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, c.want, got)
|
||||
})
|
||||
}
|
||||
}
|
||||
61
pkg/alertmanager/alertmanagernotify/jsmops/jsmops.go
Normal file
61
pkg/alertmanager/alertmanagernotify/jsmops/jsmops.go
Normal file
@@ -0,0 +1,61 @@
|
||||
// Copyright (c) 2026 SigNoz, Inc.
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package jsmops delivers Jira Service Management Ops alerts by reusing the
|
||||
// Opsgenie notifier: JSM Ops is the ex-Opsgenie alert API, so we map the JSM
|
||||
// config onto config.OpsGenieConfig with APIURL pinned to the JSM native
|
||||
// integration-events gateway.
|
||||
package jsmops
|
||||
|
||||
import (
|
||||
"log/slog"
|
||||
"net/url"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/opsgenie"
|
||||
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
"github.com/prometheus/alertmanager/template"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
)
|
||||
|
||||
const (
|
||||
Integration = "jsmops"
|
||||
source = "SigNoz"
|
||||
)
|
||||
|
||||
// New builds an Opsgenie notifier pointed at the JSM native endpoint.
|
||||
// advancedFeatures enables the rich treatment: HTML body and a note timeline
|
||||
// (per fire and on resolve).
|
||||
func New(c *alertmanagertypes.JSMOpsReceiverConfig, t *template.Template, l *slog.Logger, templater alertmanagertypes.Templater, advancedFeatures bool) (*opsgenie.Notifier, error) {
|
||||
conf, err := toOpsGenieConfig(c)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return opsgenie.New(conf, t, l, templater, advancedFeatures)
|
||||
}
|
||||
|
||||
// toOpsGenieConfig maps the JSM config onto config.OpsGenieConfig with APIURL
|
||||
// pinned to the JSM native gateway.
|
||||
func toOpsGenieConfig(c *alertmanagertypes.JSMOpsReceiverConfig) (*config.OpsGenieConfig, error) {
|
||||
apiURL, err := url.Parse(alertmanagertypes.JSMOpsAPIBaseURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
httpConfig := c.HTTPConfig
|
||||
if httpConfig == nil {
|
||||
httpConfig = &commoncfg.HTTPClientConfig{}
|
||||
}
|
||||
|
||||
return &config.OpsGenieConfig{
|
||||
NotifierConfig: c.NotifierConfig,
|
||||
HTTPConfig: httpConfig,
|
||||
APIKey: c.APIKey,
|
||||
APIURL: &config.URL{URL: apiURL},
|
||||
Message: c.Message,
|
||||
Description: c.Description,
|
||||
Priority: c.Priority,
|
||||
Tags: c.Tags,
|
||||
Source: source,
|
||||
}, nil
|
||||
}
|
||||
41
pkg/alertmanager/alertmanagernotify/jsmops/jsmops_test.go
Normal file
41
pkg/alertmanager/alertmanagernotify/jsmops/jsmops_test.go
Normal file
@@ -0,0 +1,41 @@
|
||||
package jsmops
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestToOpsGenieConfig(t *testing.T) {
|
||||
c := &alertmanagertypes.JSMOpsReceiverConfig{
|
||||
APIKey: "key-123",
|
||||
Message: "msg",
|
||||
Description: "desc",
|
||||
Priority: "P1",
|
||||
Tags: "signoz",
|
||||
HTTPConfig: &commoncfg.HTTPClientConfig{},
|
||||
}
|
||||
|
||||
og, err := toOpsGenieConfig(c)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Trailing slash is required: the Opsgenie notifier appends "v2/alerts..."
|
||||
// with no separator, yielding /jsm/ops/integration/v2/alerts.
|
||||
assert.Equal(t, "https://api.atlassian.com/jsm/ops/integration/", og.APIURL.String())
|
||||
assert.Equal(t, "key-123", string(og.APIKey))
|
||||
assert.Equal(t, "msg", og.Message)
|
||||
assert.Equal(t, "desc", og.Description)
|
||||
assert.Equal(t, "P1", og.Priority)
|
||||
assert.Equal(t, "signoz", og.Tags)
|
||||
assert.Equal(t, source, og.Source)
|
||||
assert.Same(t, c.HTTPConfig, og.HTTPConfig)
|
||||
}
|
||||
|
||||
func TestToOpsGenieConfigNilHTTPConfig(t *testing.T) {
|
||||
og, err := toOpsGenieConfig(&alertmanagertypes.JSMOpsReceiverConfig{APIKey: "k"})
|
||||
require.NoError(t, err)
|
||||
assert.NotNil(t, og.HTTPConfig)
|
||||
}
|
||||
@@ -14,6 +14,7 @@ import (
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagertemplate"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
@@ -32,8 +33,13 @@ const (
|
||||
Integration = "opsgenie"
|
||||
)
|
||||
|
||||
// https://docs.opsgenie.com/docs/alert-api - 130 characters meaning runes.
|
||||
const maxMessageLenRunes = 130
|
||||
// https://support.atlassian.com/opsgenie/docs/alert-fields/ - message 130,
|
||||
// description 15000, note 25000 runes.
|
||||
const (
|
||||
maxMessageLenRunes = 130
|
||||
maxDescriptionLenRunes = 15000
|
||||
maxNoteLenRunes = 25000
|
||||
)
|
||||
|
||||
// Notifier implements a Notifier for OpsGenie notifications.
|
||||
type Notifier struct {
|
||||
@@ -43,21 +49,29 @@ type Notifier struct {
|
||||
client *http.Client
|
||||
retrier *notify.Retrier
|
||||
templater alertmanagertypes.Templater
|
||||
// advancedFeatures bundles the JSM Ops enrichments: render the default body as
|
||||
// HTML (markdown -> HTML), and post a note per fire and on resolve to build an
|
||||
// immutable timeline. Off for plain OpsGenie. The alert-refresh-on-refire part
|
||||
// rides on the upstream UpdateAlerts config flag, set alongside this.
|
||||
advancedFeatures bool
|
||||
}
|
||||
|
||||
// New returns a new OpsGenie notifier.
|
||||
func New(c *config.OpsGenieConfig, t *template.Template, l *slog.Logger, templater alertmanagertypes.Templater, httpOpts ...commoncfg.HTTPClientOption) (*Notifier, error) {
|
||||
// New returns a new OpsGenie notifier. advancedFeatures enables the JSM Ops
|
||||
// enrichments (HTML default body + a note timeline per fire and on resolve);
|
||||
// pass false for plain OpsGenie.
|
||||
func New(c *config.OpsGenieConfig, t *template.Template, l *slog.Logger, templater alertmanagertypes.Templater, advancedFeatures bool, httpOpts ...commoncfg.HTTPClientOption) (*Notifier, error) {
|
||||
client, err := notify.NewClientWithTracing(*c.HTTPConfig, Integration, httpOpts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Notifier{
|
||||
conf: c,
|
||||
tmpl: t,
|
||||
logger: l,
|
||||
client: client,
|
||||
retrier: ¬ify.Retrier{RetryCodes: []int{http.StatusTooManyRequests}},
|
||||
templater: templater,
|
||||
conf: c,
|
||||
tmpl: t,
|
||||
logger: l,
|
||||
client: client,
|
||||
retrier: ¬ify.Retrier{RetryCodes: []int{http.StatusTooManyRequests}},
|
||||
templater: templater,
|
||||
advancedFeatures: advancedFeatures,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -94,6 +108,30 @@ type opsGenieUpdateDescriptionMessage struct {
|
||||
Description string `json:"description,omitempty"`
|
||||
}
|
||||
|
||||
type opsGenieAddNoteMessage struct {
|
||||
Note string `json:"note"`
|
||||
Source string `json:"source"`
|
||||
}
|
||||
|
||||
// noteRequest builds a POST to the alert's notes endpoint (append-only timeline).
|
||||
func (n *Notifier) noteRequest(ctx context.Context, alias, note, source string) (*http.Request, error) {
|
||||
noteEndpointURL := n.conf.APIURL.Copy()
|
||||
noteEndpointURL.Path += fmt.Sprintf("v2/alerts/%s/notes", alias)
|
||||
q := noteEndpointURL.Query()
|
||||
q.Set("identifierType", "alias")
|
||||
noteEndpointURL.RawQuery = q.Encode()
|
||||
|
||||
var buf bytes.Buffer
|
||||
if err := json.NewEncoder(&buf).Encode(&opsGenieAddNoteMessage{Note: note, Source: source}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
req, err := http.NewRequest("POST", noteEndpointURL.String(), &buf)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return req.WithContext(ctx), nil
|
||||
}
|
||||
|
||||
// Notify implements the Notifier interface.
|
||||
func (n *Notifier) Notify(ctx context.Context, as ...*types.Alert) (bool, error) {
|
||||
requests, retry, err := n.createRequests(ctx, as...)
|
||||
@@ -110,12 +148,24 @@ func (n *Notifier) Notify(ctx context.Context, as ...*types.Alert) (bool, error)
|
||||
shouldRetry, err := n.retrier.Check(resp.StatusCode, resp.Body)
|
||||
notify.Drain(resp)
|
||||
if err != nil {
|
||||
// notes are enrichment; a permanently-failed note (e.g. the first-fire
|
||||
// note racing JSM's async alert create) must not fail the notification
|
||||
if !shouldRetry && isNoteRequest(req) {
|
||||
n.logger.WarnContext(ctx, "dropping failed note", slog.Int("status_code", resp.StatusCode), errors.Attr(err))
|
||||
continue
|
||||
}
|
||||
return shouldRetry, notify.NewErrorWithReason(notify.GetFailureReasonFromStatusCode(resp.StatusCode), err)
|
||||
}
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// isNoteRequest reports whether req targets the notes endpoint, the only one
|
||||
// built by noteRequest.
|
||||
func isNoteRequest(req *http.Request) bool {
|
||||
return strings.HasSuffix(req.URL.Path, "/notes")
|
||||
}
|
||||
|
||||
// Like Split but filter out empty strings.
|
||||
func safeSplit(s, sep string) []string {
|
||||
a := strings.Split(strings.TrimSpace(s), sep)
|
||||
@@ -145,28 +195,13 @@ func (n *Notifier) prepareContent(ctx context.Context, alerts []*types.Alert) (s
|
||||
}
|
||||
|
||||
var description string
|
||||
if result.IsDefaultBody {
|
||||
if result.IsDefaultBody && !n.advancedFeatures {
|
||||
description = strings.Join(result.Body, "\n")
|
||||
} else {
|
||||
var b strings.Builder
|
||||
first := true
|
||||
for _, part := range result.Body {
|
||||
if part == "" {
|
||||
continue
|
||||
}
|
||||
rendered, renderErr := markdownrenderer.RenderHTML(part)
|
||||
if renderErr != nil {
|
||||
return "", "", renderErr
|
||||
}
|
||||
if !first {
|
||||
b.WriteString("<hr>")
|
||||
}
|
||||
b.WriteString("<div>")
|
||||
b.WriteString(rendered)
|
||||
b.WriteString("</div>")
|
||||
first = false
|
||||
description, err = buildHTMLDescription(result.Body, maxDescriptionLenRunes)
|
||||
if err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
description = b.String()
|
||||
}
|
||||
|
||||
title, truncated := notify.TruncateInRunes(result.Title, maxMessageLenRunes)
|
||||
@@ -174,9 +209,141 @@ func (n *Notifier) prepareContent(ctx context.Context, alerts []*types.Alert) (s
|
||||
n.logger.WarnContext(ctx, "Truncated message", slog.Int("max_runes", maxMessageLenRunes))
|
||||
}
|
||||
|
||||
// The API silently truncates over-limit descriptions, which would drop the
|
||||
// trailing SigNoz link; cap here with an ellipsis instead. The HTML path is
|
||||
// pre-fitted above, so this only ever cuts the plain-text default body.
|
||||
description, descTruncated := notify.TruncateInRunes(description, maxDescriptionLenRunes)
|
||||
if descTruncated {
|
||||
n.logger.WarnContext(ctx, "Truncated description", slog.Int("max_runes", maxDescriptionLenRunes))
|
||||
}
|
||||
|
||||
return title, description, nil
|
||||
}
|
||||
|
||||
const (
|
||||
// room reserved for the "+N more" trailer appended when parts are dropped.
|
||||
descriptionTrailerReserveRunes = 80
|
||||
// below this rendering budget a shrunk part carries no signal; drop it instead.
|
||||
minShrinkBudgetRunes = 64
|
||||
)
|
||||
|
||||
// buildHTMLDescription renders each markdown part to HTML (<div>-wrapped,
|
||||
// <hr>-joined) while keeping the total within budget runes. An over-budget part
|
||||
// is shrunk at the markdown level and re-rendered so the HTML stays well-formed;
|
||||
// fully dropped parts are summarized by a "+N more" trailer.
|
||||
func buildHTMLDescription(parts []string, budget int) (string, error) {
|
||||
rendering := make([]string, 0, len(parts))
|
||||
for _, part := range parts {
|
||||
if part != "" {
|
||||
rendering = append(rendering, part)
|
||||
}
|
||||
}
|
||||
|
||||
budget -= descriptionTrailerReserveRunes
|
||||
var b strings.Builder
|
||||
used, included := 0, 0
|
||||
for _, part := range rendering {
|
||||
rendered, err := markdownrenderer.RenderHTML(part)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
overhead := len("<div></div>")
|
||||
if included > 0 {
|
||||
overhead += len("<hr>")
|
||||
}
|
||||
if used+overhead+utf8.RuneCountInString(rendered) > budget {
|
||||
rendered, err = shrinkMarkdownToFit(part, budget-used-overhead)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if rendered == "" {
|
||||
break
|
||||
}
|
||||
}
|
||||
if included > 0 {
|
||||
b.WriteString("<hr>")
|
||||
}
|
||||
b.WriteString("<div>")
|
||||
b.WriteString(rendered)
|
||||
b.WriteString("</div>")
|
||||
used += overhead + utf8.RuneCountInString(rendered)
|
||||
included++
|
||||
}
|
||||
if dropped := len(rendering) - included; dropped > 0 {
|
||||
fmt.Fprintf(&b, "<hr><div><i>…and %d more alerts. Open in SigNoz for the full list.</i></div>", dropped)
|
||||
}
|
||||
return b.String(), nil
|
||||
}
|
||||
|
||||
// shrinkMarkdownToFit cuts markdown until its rendered HTML fits within budget
|
||||
// runes, returning "" when the budget is too small to carry anything useful.
|
||||
// Only the markdown is ever cut, never the rendered HTML, so goldmark always
|
||||
// emits balanced markup.
|
||||
func shrinkMarkdownToFit(md string, budget int) (string, error) {
|
||||
if budget < minShrinkBudgetRunes {
|
||||
return "", nil
|
||||
}
|
||||
for range 4 {
|
||||
rendered, err := markdownrenderer.RenderHTML(md)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
renderedLen := utf8.RuneCountInString(rendered)
|
||||
if renderedLen <= budget {
|
||||
return rendered, nil
|
||||
}
|
||||
runes := []rune(md)
|
||||
keep := len(runes) * budget / renderedLen * 9 / 10
|
||||
if keep >= len(runes) {
|
||||
keep = len(runes) - 1
|
||||
}
|
||||
if keep < minShrinkBudgetRunes {
|
||||
return "", nil
|
||||
}
|
||||
md = string(runes[:keep]) + "…"
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
|
||||
// prepareNote renders the same body template as plain text for a timeline note.
|
||||
// JSM Ops notes render neither HTML nor markdown, so links flatten to
|
||||
// "text (url)" and all markers are stripped.
|
||||
func (n *Notifier) prepareNote(ctx context.Context, alerts []*types.Alert) (string, error) {
|
||||
customTitle, customBody := alertmanagertemplate.ExtractTemplatesFromAnnotations(alerts)
|
||||
result, err := n.templater.Expand(ctx, alertmanagertypes.ExpandRequest{
|
||||
TitleTemplate: customTitle,
|
||||
BodyTemplate: customBody,
|
||||
DefaultTitleTemplate: n.conf.Message,
|
||||
DefaultBodyTemplate: n.conf.Description,
|
||||
}, alerts)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
var b strings.Builder
|
||||
first := true
|
||||
for _, part := range result.Body {
|
||||
text, renderErr := markdownrenderer.RenderPlainText(part)
|
||||
if renderErr != nil {
|
||||
return "", renderErr
|
||||
}
|
||||
if text = strings.TrimSpace(text); text == "" {
|
||||
continue
|
||||
}
|
||||
if !first {
|
||||
b.WriteString("\n\n")
|
||||
}
|
||||
b.WriteString(text)
|
||||
first = false
|
||||
}
|
||||
|
||||
note, truncated := notify.TruncateInRunes(b.String(), maxNoteLenRunes)
|
||||
if truncated {
|
||||
n.logger.WarnContext(ctx, "Truncated note", slog.Int("max_runes", maxNoteLenRunes))
|
||||
}
|
||||
return note, nil
|
||||
}
|
||||
|
||||
// Create requests for a list of alerts.
|
||||
func (n *Notifier) createRequests(ctx context.Context, as ...*types.Alert) ([]*http.Request, bool, error) {
|
||||
key, err := notify.ExtractGroupKey(ctx)
|
||||
@@ -206,6 +373,21 @@ func (n *Notifier) createRequests(ctx context.Context, as ...*types.Alert) ([]*h
|
||||
)
|
||||
switch alerts.Status() {
|
||||
case model.AlertResolved:
|
||||
// Post the resolved snapshot to the timeline before closing (closed alerts
|
||||
// reject notes), so the note lands first.
|
||||
if n.advancedFeatures {
|
||||
note, err := n.prepareNote(ctx, as)
|
||||
if err != nil {
|
||||
n.logger.ErrorContext(ctx, "failed to prepare notification content", errors.Attr(err))
|
||||
return nil, false, err
|
||||
}
|
||||
noteReq, err := n.noteRequest(ctx, alias, note, tmpl(n.conf.Source))
|
||||
if err != nil {
|
||||
return nil, true, err
|
||||
}
|
||||
requests = append(requests, noteReq)
|
||||
}
|
||||
|
||||
resolvedEndpointURL := n.conf.APIURL.Copy()
|
||||
resolvedEndpointURL.Path += fmt.Sprintf("v2/alerts/%s/close", alias)
|
||||
q := resolvedEndpointURL.Query()
|
||||
@@ -322,6 +504,21 @@ func (n *Notifier) createRequests(ctx context.Context, as ...*types.Alert) ([]*h
|
||||
}
|
||||
requests = append(requests, req.WithContext(ctx))
|
||||
}
|
||||
|
||||
// Append this fire's snapshot to the timeline (every fire, including the
|
||||
// first, so no datapoint is lost when the description is overwritten).
|
||||
// Notes are plain text, so this uses the plain-text render, not the HTML body.
|
||||
if n.advancedFeatures {
|
||||
note, err := n.prepareNote(ctx, as)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
noteReq, err := n.noteRequest(ctx, alias, note, tmpl(n.conf.Source))
|
||||
if err != nil {
|
||||
return nil, true, err
|
||||
}
|
||||
requests = append(requests, noteReq)
|
||||
}
|
||||
}
|
||||
|
||||
var apiKey string
|
||||
|
||||
@@ -6,14 +6,18 @@ package opsgenie
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagertemplate"
|
||||
"github.com/SigNoz/signoz/pkg/types/alertmanagertypes"
|
||||
@@ -44,6 +48,7 @@ func TestOpsGenieRetry(t *testing.T) {
|
||||
tmpl,
|
||||
promslog.NewNopLogger(),
|
||||
newTestTemplater(tmpl),
|
||||
false,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -69,6 +74,7 @@ func TestOpsGenieRedactedURL(t *testing.T) {
|
||||
tmpl,
|
||||
promslog.NewNopLogger(),
|
||||
newTestTemplater(tmpl),
|
||||
false,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -96,6 +102,7 @@ func TestGettingOpsGegineApikeyFromFile(t *testing.T) {
|
||||
tmpl,
|
||||
promslog.NewNopLogger(),
|
||||
newTestTemplater(tmpl),
|
||||
false,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -216,7 +223,7 @@ func TestOpsGenie(t *testing.T) {
|
||||
},
|
||||
} {
|
||||
t.Run(tc.title, func(t *testing.T) {
|
||||
notifier, err := New(tc.cfg, tmpl, logger, newTestTemplater(tmpl))
|
||||
notifier, err := New(tc.cfg, tmpl, logger, newTestTemplater(tmpl), false)
|
||||
require.NoError(t, err)
|
||||
|
||||
ctx := context.Background()
|
||||
@@ -292,7 +299,7 @@ func TestOpsGenieWithUpdate(t *testing.T) {
|
||||
APIURL: &config.URL{URL: u},
|
||||
HTTPConfig: &commoncfg.HTTPClientConfig{},
|
||||
}
|
||||
notifierWithUpdate, err := New(&opsGenieConfigWithUpdate, tmpl, promslog.NewNopLogger(), newTestTemplater(tmpl))
|
||||
notifierWithUpdate, err := New(&opsGenieConfigWithUpdate, tmpl, promslog.NewNopLogger(), newTestTemplater(tmpl), false)
|
||||
alert := &types.Alert{
|
||||
Alert: model.Alert{
|
||||
StartsAt: time.Now(),
|
||||
@@ -324,6 +331,111 @@ func TestOpsGenieWithUpdate(t *testing.T) {
|
||||
assert.JSONEq(t, `{"description":"new description"}`, body2)
|
||||
}
|
||||
|
||||
func TestOpsGenieAdvancedFeatures(t *testing.T) {
|
||||
u, err := url.Parse("https://test-opsgenie-url")
|
||||
require.NoError(t, err)
|
||||
tmpl := test.CreateTmpl(t)
|
||||
ctx := notify.WithGroupKey(context.Background(), "1")
|
||||
key, _ := notify.ExtractGroupKey(ctx)
|
||||
alias := key.Hash()
|
||||
|
||||
cfg := &config.OpsGenieConfig{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: true},
|
||||
Message: `{{ .CommonLabels.Message }}`,
|
||||
Description: `{{ .CommonLabels.Description }}`,
|
||||
UpdateAlerts: true,
|
||||
APIKey: "k",
|
||||
APIURL: &config.URL{URL: u},
|
||||
HTTPConfig: &commoncfg.HTTPClientConfig{},
|
||||
}
|
||||
notifier, err := New(cfg, tmpl, promslog.NewNopLogger(), newTestTemplater(tmpl), true)
|
||||
require.NoError(t, err)
|
||||
|
||||
firing := &types.Alert{Alert: model.Alert{
|
||||
StartsAt: time.Now(),
|
||||
EndsAt: time.Now().Add(time.Hour),
|
||||
Labels: model.LabelSet{"Message": "m", "Description": "**Alert:** d [View](https://s.io/a)"},
|
||||
}}
|
||||
|
||||
// Fire: create + update message + update description + a timeline note.
|
||||
reqs, _, err := notifier.createRequests(ctx, firing)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, reqs, 4)
|
||||
assert.Equal(t, "https://test-opsgenie-url/v2/alerts", reqs[0].URL.String())
|
||||
assert.Equal(t, fmt.Sprintf("https://test-opsgenie-url/v2/alerts/%s/notes?identifierType=alias", alias), reqs[3].URL.String())
|
||||
assert.Equal(t, http.MethodPost, reqs[3].Method)
|
||||
|
||||
// the note body is the plain-text render: markers stripped, link flattened
|
||||
var noteMsg opsGenieAddNoteMessage
|
||||
require.NoError(t, json.Unmarshal([]byte(readBody(t, reqs[3])), ¬eMsg))
|
||||
assert.Equal(t, "Alert: d View (https://s.io/a)", noteMsg.Note)
|
||||
|
||||
// Resolve: note posted before the close.
|
||||
resolved := &types.Alert{Alert: model.Alert{
|
||||
StartsAt: time.Now().Add(-time.Hour),
|
||||
EndsAt: time.Now().Add(-time.Minute),
|
||||
Labels: model.LabelSet{"Message": "m", "Description": "d"},
|
||||
}}
|
||||
reqs, _, err = notifier.createRequests(ctx, resolved)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, reqs, 2)
|
||||
assert.Equal(t, fmt.Sprintf("https://test-opsgenie-url/v2/alerts/%s/notes?identifierType=alias", alias), reqs[0].URL.String())
|
||||
assert.Equal(t, fmt.Sprintf("https://test-opsgenie-url/v2/alerts/%s/close?identifierType=alias", alias), reqs[1].URL.String())
|
||||
}
|
||||
|
||||
func TestOpsGenieNotifyBestEffortNote(t *testing.T) {
|
||||
tmpl := test.CreateTmpl(t)
|
||||
ctx := notify.WithGroupKey(context.Background(), "1")
|
||||
|
||||
firing := &types.Alert{Alert: model.Alert{
|
||||
StartsAt: time.Now(),
|
||||
EndsAt: time.Now().Add(time.Hour),
|
||||
Labels: model.LabelSet{"Message": "m", "Description": "d"},
|
||||
}}
|
||||
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
createStatus int
|
||||
noteStatus int
|
||||
wantErr bool
|
||||
wantRetry bool
|
||||
}{
|
||||
{name: "note_404_is_dropped", createStatus: http.StatusAccepted, noteStatus: http.StatusNotFound, wantErr: false, wantRetry: true},
|
||||
{name: "note_429_still_retries", createStatus: http.StatusAccepted, noteStatus: http.StatusTooManyRequests, wantErr: true, wantRetry: true},
|
||||
{name: "create_404_still_fails", createStatus: http.StatusNotFound, noteStatus: http.StatusAccepted, wantErr: true, wantRetry: false},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if strings.HasSuffix(r.URL.Path, "/notes") {
|
||||
w.WriteHeader(tc.noteStatus)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(tc.createStatus)
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
u, err := url.Parse(srv.URL)
|
||||
require.NoError(t, err)
|
||||
notifier, err := New(&config.OpsGenieConfig{
|
||||
Message: `{{ .CommonLabels.Message }}`,
|
||||
Description: `{{ .CommonLabels.Description }}`,
|
||||
APIKey: "k",
|
||||
APIURL: &config.URL{URL: u},
|
||||
HTTPConfig: &commoncfg.HTTPClientConfig{},
|
||||
}, tmpl, promslog.NewNopLogger(), newTestTemplater(tmpl), true)
|
||||
require.NoError(t, err)
|
||||
|
||||
retry, err := notifier.Notify(ctx, firing)
|
||||
if tc.wantErr {
|
||||
require.Error(t, err)
|
||||
} else {
|
||||
require.NoError(t, err)
|
||||
}
|
||||
assert.Equal(t, tc.wantRetry, retry)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpsGenieApiKeyFile(t *testing.T) {
|
||||
u, err := url.Parse("https://test-opsgenie-url")
|
||||
require.NoError(t, err)
|
||||
@@ -335,7 +447,7 @@ func TestOpsGenieApiKeyFile(t *testing.T) {
|
||||
APIURL: &config.URL{URL: u},
|
||||
HTTPConfig: &commoncfg.HTTPClientConfig{},
|
||||
}
|
||||
notifierWithUpdate, err := New(&opsGenieConfigWithUpdate, tmpl, promslog.NewNopLogger(), newTestTemplater(tmpl))
|
||||
notifierWithUpdate, err := New(&opsGenieConfigWithUpdate, tmpl, promslog.NewNopLogger(), newTestTemplater(tmpl), false)
|
||||
|
||||
require.NoError(t, err)
|
||||
requests, _, err := notifierWithUpdate.createRequests(ctx)
|
||||
@@ -437,6 +549,109 @@ func TestPrepareContent(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
func TestShrinkMarkdownToFit(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
md string
|
||||
budget int
|
||||
wantEmpty bool
|
||||
}{
|
||||
{"fits untouched", "**bold** text", 1000, false},
|
||||
{"shrinks to fit", strings.Repeat("lorem ipsum ", 500), 1000, false},
|
||||
{"budget too small", strings.Repeat("lorem ipsum ", 500), 10, true},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
got, err := shrinkMarkdownToFit(c.md, c.budget)
|
||||
require.NoError(t, err)
|
||||
if c.wantEmpty {
|
||||
assert.Empty(t, got)
|
||||
return
|
||||
}
|
||||
assert.NotEmpty(t, got)
|
||||
assert.LessOrEqual(t, utf8.RuneCountInString(got), c.budget)
|
||||
assert.Equal(t, strings.Count(got, "<p>"), strings.Count(got, "</p>"))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildHTMLDescriptionOverflow(t *testing.T) {
|
||||
bigPart := strings.Repeat("alpha beta gamma ", 100)
|
||||
cases := []struct {
|
||||
name string
|
||||
parts []string
|
||||
budget int
|
||||
wantTrailer bool
|
||||
}{
|
||||
{"all parts fit", []string{"**a**", "**b**"}, maxDescriptionLenRunes, false},
|
||||
{"empty parts skipped", []string{"", "hello", ""}, maxDescriptionLenRunes, false},
|
||||
{"overflow drops parts with trailer", repeatParts(bigPart, 12), maxDescriptionLenRunes, true},
|
||||
{"single huge part shrunk without trailer", []string{strings.Repeat(bigPart, 20)}, maxDescriptionLenRunes, false},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
got, err := buildHTMLDescription(c.parts, c.budget)
|
||||
require.NoError(t, err)
|
||||
assert.LessOrEqual(t, utf8.RuneCountInString(got), c.budget)
|
||||
assert.Equal(t, strings.Count(got, "<div>"), strings.Count(got, "</div>"))
|
||||
assert.True(t, strings.HasSuffix(got, "</div>"))
|
||||
if c.wantTrailer {
|
||||
assert.Regexp(t, `…and \d+ more alerts\. Open in SigNoz for the full list\.`, got)
|
||||
} else {
|
||||
assert.NotContains(t, got, "more alerts")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// prepareContent end-to-end: 40 custom-template alerts overflow the description
|
||||
// budget yet the posted HTML stays within limits and well-formed.
|
||||
func TestPrepareContentDescriptionOverflow(t *testing.T) {
|
||||
tmpl := test.CreateTmpl(t)
|
||||
notifier := &Notifier{
|
||||
conf: &config.OpsGenieConfig{
|
||||
Message: `{{ .CommonLabels.alertname }}`,
|
||||
Description: `{{ .CommonLabels.alertname }}`,
|
||||
},
|
||||
tmpl: tmpl,
|
||||
logger: promslog.NewNopLogger(),
|
||||
templater: newTestTemplater(tmpl),
|
||||
advancedFeatures: true,
|
||||
}
|
||||
|
||||
bodyTemplate := "**Alert in** $labels.namespace\n\n" + strings.Repeat("detail line for the runbook ", 30)
|
||||
alerts := make([]*types.Alert, 0, 40)
|
||||
for i := range 40 {
|
||||
alerts = append(alerts, &types.Alert{
|
||||
Alert: model.Alert{
|
||||
Labels: model.LabelSet{
|
||||
"alertname": "overflow",
|
||||
"namespace": model.LabelValue(fmt.Sprintf("ns-%d", i)),
|
||||
},
|
||||
Annotations: model.LabelSet{
|
||||
ruletypes.AnnotationBodyTemplate: model.LabelValue(bodyTemplate),
|
||||
},
|
||||
StartsAt: time.Now(),
|
||||
EndsAt: time.Now().Add(time.Hour),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
_, desc, err := notifier.prepareContent(notify.WithGroupKey(context.Background(), "1"), alerts)
|
||||
require.NoError(t, err)
|
||||
assert.LessOrEqual(t, utf8.RuneCountInString(desc), maxDescriptionLenRunes)
|
||||
assert.Equal(t, strings.Count(desc, "<div>"), strings.Count(desc, "</div>"))
|
||||
assert.Regexp(t, `…and \d+ more alerts\. Open in SigNoz for the full list\.`, desc)
|
||||
}
|
||||
|
||||
func repeatParts(part string, n int) []string {
|
||||
parts := make([]string, n)
|
||||
for i := range parts {
|
||||
parts[i] = part
|
||||
}
|
||||
return parts
|
||||
}
|
||||
|
||||
func readBody(t *testing.T, r *http.Request) string {
|
||||
t.Helper()
|
||||
body, err := io.ReadAll(r.Body)
|
||||
|
||||
@@ -6,6 +6,8 @@ import (
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/email"
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/googlechat"
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/jira"
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/jsmops"
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/msteamsv2"
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/opsgenie"
|
||||
"github.com/SigNoz/signoz/pkg/alertmanager/alertmanagernotify/pagerduty"
|
||||
@@ -26,6 +28,8 @@ var customNotifierIntegrations = []string{
|
||||
slack.Integration,
|
||||
msteamsv2.Integration,
|
||||
googlechat.Integration,
|
||||
jira.Integration,
|
||||
jsmops.Integration,
|
||||
}
|
||||
|
||||
func NewReceiverIntegrations(nc *alertmanagertypes.Receiver, tmpl *template.Template, logger *slog.Logger, templater alertmanagertypes.Templater) ([]notify.Integration, error) {
|
||||
@@ -66,7 +70,7 @@ func NewReceiverIntegrations(nc *alertmanagertypes.Receiver, tmpl *template.Temp
|
||||
add(pagerduty.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) { return pagerduty.New(c, tmpl, l, templater) })
|
||||
}
|
||||
for i, c := range nc.OpsGenieConfigs {
|
||||
add(opsgenie.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) { return opsgenie.New(c, tmpl, l, templater) })
|
||||
add(opsgenie.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) { return opsgenie.New(c, tmpl, l, templater, false) })
|
||||
}
|
||||
for i, c := range nc.SlackConfigs {
|
||||
add(slack.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) { return slack.New(c, tmpl, l, templater) })
|
||||
@@ -81,6 +85,16 @@ func NewReceiverIntegrations(nc *alertmanagertypes.Receiver, tmpl *template.Temp
|
||||
return googlechat.New(c, tmpl, l, templater)
|
||||
})
|
||||
}
|
||||
for i, c := range nc.JiraConfigs {
|
||||
add(jira.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) {
|
||||
return jira.New(c, tmpl, l, templater)
|
||||
})
|
||||
}
|
||||
for i, c := range nc.JSMOpsConfigs {
|
||||
add(jsmops.Integration, i, c, func(l *slog.Logger) (notify.Notifier, error) {
|
||||
return jsmops.New(c, tmpl, l, templater, true)
|
||||
})
|
||||
}
|
||||
|
||||
if errs.Len() > 0 {
|
||||
return nil, &errs
|
||||
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/modules/organization"
|
||||
"github.com/SigNoz/signoz/pkg/modules/preference"
|
||||
"github.com/SigNoz/signoz/pkg/modules/promote"
|
||||
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
|
||||
"github.com/SigNoz/signoz/pkg/modules/rawdataexport"
|
||||
"github.com/SigNoz/signoz/pkg/modules/rulestatehistory"
|
||||
"github.com/SigNoz/signoz/pkg/modules/savedview"
|
||||
@@ -82,6 +83,8 @@ type provider struct {
|
||||
llmPricingRuleHandler llmpricingrule.Handler
|
||||
statsHandler statsreporter.Handler
|
||||
savedViewHandler savedview.Handler
|
||||
quickFilterModule quickfilter.Module
|
||||
quickFilterHandler quickfilter.Handler
|
||||
}
|
||||
|
||||
func NewFactory(
|
||||
@@ -121,6 +124,8 @@ func NewFactory(
|
||||
rulerHandler ruler.Handler,
|
||||
statsHandler statsreporter.Handler,
|
||||
savedViewHandler savedview.Handler,
|
||||
quickFilterModule quickfilter.Module,
|
||||
quickFilterHandler quickfilter.Handler,
|
||||
) factory.ProviderFactory[apiserver.APIServer, apiserver.Config] {
|
||||
return factory.NewProviderFactory(factory.MustNewName("signoz"), func(ctx context.Context, providerSettings factory.ProviderSettings, config apiserver.Config) (apiserver.APIServer, error) {
|
||||
return newProvider(
|
||||
@@ -163,6 +168,8 @@ func NewFactory(
|
||||
rulerHandler,
|
||||
statsHandler,
|
||||
savedViewHandler,
|
||||
quickFilterModule,
|
||||
quickFilterHandler,
|
||||
)
|
||||
})
|
||||
}
|
||||
@@ -207,6 +214,8 @@ func newProvider(
|
||||
rulerHandler ruler.Handler,
|
||||
statsHandler statsreporter.Handler,
|
||||
savedViewHandler savedview.Handler,
|
||||
quickFilterModule quickfilter.Module,
|
||||
quickFilterHandler quickfilter.Handler,
|
||||
) (apiserver.APIServer, error) {
|
||||
settings := factory.NewScopedProviderSettings(providerSettings, "github.com/SigNoz/signoz/pkg/apiserver/signozapiserver")
|
||||
router := mux.NewRouter().UseEncodedPath()
|
||||
@@ -250,6 +259,8 @@ func newProvider(
|
||||
llmPricingRuleHandler: llmPricingRuleHandler,
|
||||
statsHandler: statsHandler,
|
||||
savedViewHandler: savedViewHandler,
|
||||
quickFilterModule: quickFilterModule,
|
||||
quickFilterHandler: quickFilterHandler,
|
||||
}
|
||||
|
||||
provider.authzMiddleware = middleware.NewAuthZ(settings.Logger(), orgGetter, authzService)
|
||||
@@ -394,6 +405,10 @@ func (provider *provider) AddToRouter(router *mux.Router) error {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := provider.addQuickFilterRoutes(router); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
120
pkg/apiserver/signozapiserver/quickfilter.go
Normal file
120
pkg/apiserver/signozapiserver/quickfilter.go
Normal file
@@ -0,0 +1,120 @@
|
||||
package signozapiserver
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/http/handler"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/coretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
|
||||
func (provider *provider) addQuickFilterRoutes(router *mux.Router) error {
|
||||
if err := router.Handle("/api/v2/quick_filters", handler.New(
|
||||
provider.authzMiddleware.CheckResources(provider.quickFilterHandler.ListQuickFiltersV2, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName),
|
||||
handler.OpenAPIDef{
|
||||
ID: "ListQuickFilters",
|
||||
Tags: []string{"quick_filter"},
|
||||
Summary: "List quick filters",
|
||||
Description: "Returns the org's quick filters for every source, each filter as a telemetry field key.",
|
||||
Request: nil,
|
||||
RequestContentType: "",
|
||||
Response: new([]*quickfiltertypes.SourceFilters),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusOK,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest},
|
||||
Deprecated: false,
|
||||
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceQuickFilter.Scope(coretypes.VerbList)}),
|
||||
},
|
||||
handler.WithResourceDefs(handler.BasicResourceDef{
|
||||
Resource: coretypes.ResourceMetaResourceQuickFilter,
|
||||
Verb: coretypes.VerbList,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
Selector: coretypes.WildcardSelector,
|
||||
}),
|
||||
)).Methods(http.MethodGet).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v2/quick_filters/{source}", handler.New(
|
||||
provider.authzMiddleware.CheckResources(provider.quickFilterHandler.GetQuickFiltersV2, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName),
|
||||
handler.OpenAPIDef{
|
||||
ID: "GetQuickFilters",
|
||||
Tags: []string{"quick_filter"},
|
||||
Summary: "Get a source's quick filters",
|
||||
Description: "Returns the org's quick filters for one source, each filter as a telemetry field key.",
|
||||
Request: nil,
|
||||
RequestContentType: "",
|
||||
Response: new(quickfiltertypes.SourceFilters),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusOK,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest},
|
||||
Deprecated: false,
|
||||
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceQuickFilter.Scope(coretypes.VerbRead)}),
|
||||
},
|
||||
handler.WithResourceDefs(handler.BasicResourceDef{
|
||||
Resource: coretypes.ResourceMetaResourceQuickFilter,
|
||||
Verb: coretypes.VerbRead,
|
||||
Category: coretypes.ActionCategoryDataAccess,
|
||||
ID: coretypes.PathParam("source"),
|
||||
Selector: provider.quickFilterSelector,
|
||||
}),
|
||||
)).Methods(http.MethodGet).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v2/quick_filters/{source}", handler.New(
|
||||
provider.authzMiddleware.CheckResources(provider.quickFilterHandler.UpdateQuickFiltersV2, authtypes.SigNozAdminRoleName),
|
||||
handler.OpenAPIDef{
|
||||
ID: "UpdateQuickFilters",
|
||||
Tags: []string{"quick_filter"},
|
||||
Summary: "Update quick filters",
|
||||
Description: "Replaces the org's quick filters for the source named in the path.",
|
||||
Request: new(quickfiltertypes.UpdatableQuickFilters),
|
||||
RequestContentType: "application/json",
|
||||
Response: nil,
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusNoContent,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest},
|
||||
Deprecated: false,
|
||||
SecuritySchemes: newScopedSecuritySchemes([]string{coretypes.ResourceMetaResourceQuickFilter.Scope(coretypes.VerbUpdate)}),
|
||||
},
|
||||
handler.WithResourceDefs(handler.BasicResourceDef{
|
||||
Resource: coretypes.ResourceMetaResourceQuickFilter,
|
||||
Verb: coretypes.VerbUpdate,
|
||||
Category: coretypes.ActionCategoryConfigurationChange,
|
||||
ID: coretypes.PathParam("source"),
|
||||
Selector: provider.quickFilterSelector,
|
||||
}),
|
||||
)).Methods(http.MethodPut).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (provider *provider) quickFilterSelector(ctx context.Context, resource coretypes.Resource, source string, orgID valuer.UUID) ([]coretypes.Selector, error) {
|
||||
validatedSource, err := quickfiltertypes.NewSource(source)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// A source can have no stored row yet: GET serves it as empty and PUT
|
||||
// creates it, so only the wildcard grant applies until the row exists.
|
||||
quickFilter, err := provider.quickFilterModule.Get(ctx, orgID, validatedSource)
|
||||
if err != nil {
|
||||
if errors.Ast(err, errors.TypeNotFound) {
|
||||
return []coretypes.Selector{resource.Type().MustSelector(coretypes.WildCardSelectorString)}, nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return []coretypes.Selector{
|
||||
resource.Type().MustSelector(quickFilter.ID.StringValue()),
|
||||
resource.Type().MustSelector(coretypes.WildCardSelectorString),
|
||||
}, nil
|
||||
}
|
||||
@@ -4,10 +4,13 @@ import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/http/render"
|
||||
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
@@ -20,6 +23,13 @@ func NewHandler(module quickfilter.Module) quickfilter.Handler {
|
||||
return &handler{module: module}
|
||||
}
|
||||
|
||||
// legacySourceFilters is the v1 API shape: filters as v3 attribute keys,
|
||||
// with the source still spelled "signal" on the wire.
|
||||
type legacySourceFilters struct {
|
||||
Source quickfiltertypes.Source `json:"signal"`
|
||||
Filters []v3.AttributeKey `json:"filters"`
|
||||
}
|
||||
|
||||
func (handler *handler) GetQuickFilters(rw http.ResponseWriter, r *http.Request) {
|
||||
claims, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
@@ -27,13 +37,41 @@ func (handler *handler) GetQuickFilters(rw http.ResponseWriter, r *http.Request)
|
||||
return
|
||||
}
|
||||
|
||||
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID))
|
||||
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), quickfiltertypes.Source{})
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, filters)
|
||||
legacyFilters := make([]*legacySourceFilters, 0, len(filters))
|
||||
for _, sourceFilters := range filters {
|
||||
legacyFilters = append(legacyFilters, newLegacySourceFilters(sourceFilters))
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, legacyFilters)
|
||||
}
|
||||
|
||||
func (handler *handler) GetSourceFilters(rw http.ResponseWriter, r *http.Request) {
|
||||
claims, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
source := mux.Vars(r)["signal"]
|
||||
validatedSource, err := quickfiltertypes.NewSource(source)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSource)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, newLegacySourceFilters(handler.sourceFiltersOrEmpty(filters, validatedSource)))
|
||||
}
|
||||
|
||||
func (handler *handler) UpdateQuickFilters(rw http.ResponseWriter, r *http.Request) {
|
||||
@@ -43,14 +81,19 @@ func (handler *handler) UpdateQuickFilters(rw http.ResponseWriter, r *http.Reque
|
||||
return
|
||||
}
|
||||
|
||||
var req quickfiltertypes.UpdatableQuickFilters
|
||||
decodeErr := json.NewDecoder(r.Body).Decode(&req)
|
||||
if decodeErr != nil {
|
||||
render.Error(rw, decodeErr)
|
||||
var req legacySourceFilters
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
err = handler.module.UpdateQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), req.Signal, req.Filters)
|
||||
fieldKeys, err := newTelemetryFieldKeysFromLegacy(req.Source, req.Filters)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
err = handler.module.UpsertQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), req.Source, fieldKeys)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
@@ -59,21 +102,14 @@ func (handler *handler) UpdateQuickFilters(rw http.ResponseWriter, r *http.Reque
|
||||
render.Success(rw, http.StatusNoContent, nil)
|
||||
}
|
||||
|
||||
func (handler *handler) GetSignalFilters(rw http.ResponseWriter, r *http.Request) {
|
||||
func (handler *handler) ListQuickFiltersV2(rw http.ResponseWriter, r *http.Request) {
|
||||
claims, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
signal := mux.Vars(r)["signal"]
|
||||
validatedSignal, err := quickfiltertypes.NewSignal(signal)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
filters, err := handler.module.GetSignalFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSignal)
|
||||
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), quickfiltertypes.Source{})
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
@@ -81,3 +117,141 @@ func (handler *handler) GetSignalFilters(rw http.ResponseWriter, r *http.Request
|
||||
|
||||
render.Success(rw, http.StatusOK, filters)
|
||||
}
|
||||
|
||||
func (handler *handler) UpdateQuickFiltersV2(rw http.ResponseWriter, r *http.Request) {
|
||||
claims, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
source := mux.Vars(r)["source"]
|
||||
validatedSource, err := quickfiltertypes.NewSource(source)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
var req quickfiltertypes.UpdatableQuickFilters
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
err = handler.module.UpsertQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSource, req.Filters)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusNoContent, nil)
|
||||
}
|
||||
|
||||
func (handler *handler) GetQuickFiltersV2(rw http.ResponseWriter, r *http.Request) {
|
||||
claims, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
source := mux.Vars(r)["source"]
|
||||
validatedSource, err := quickfiltertypes.NewSource(source)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
filters, err := handler.module.GetQuickFilters(r.Context(), valuer.MustNewUUID(claims.OrgID), validatedSource)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, handler.sourceFiltersOrEmpty(filters, validatedSource))
|
||||
}
|
||||
|
||||
// sourceFiltersOrEmpty keeps the single-source response contract: a source
|
||||
// with no stored filters is served as an empty filter list, not an error.
|
||||
func (handler *handler) sourceFiltersOrEmpty(filters []*quickfiltertypes.SourceFilters, source quickfiltertypes.Source) *quickfiltertypes.SourceFilters {
|
||||
if len(filters) == 0 {
|
||||
return quickfiltertypes.NewSourceFiltersFromSource(source)
|
||||
}
|
||||
return filters[0]
|
||||
}
|
||||
|
||||
// newTelemetryFieldKeysFromLegacy converts a v1 write payload with the same
|
||||
// normalizations as the storage migration: alias contexts, numerics to number.
|
||||
// The v1 shape carries no per filter signal, so meter keys get it restored.
|
||||
func newTelemetryFieldKeysFromLegacy(source quickfiltertypes.Source, filters []v3.AttributeKey) ([]telemetrytypes.TelemetryFieldKey, error) {
|
||||
var fieldSignal telemetrytypes.Signal
|
||||
if source == quickfiltertypes.SourceMeter {
|
||||
fieldSignal = telemetrytypes.SignalMetrics
|
||||
}
|
||||
|
||||
fieldKeys := make([]telemetrytypes.TelemetryFieldKey, 0, len(filters))
|
||||
for _, filter := range filters {
|
||||
if err := filter.Validate(); err != nil {
|
||||
return nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter: %v", err)
|
||||
}
|
||||
|
||||
fieldContext, ok := telemetrytypes.FieldContextFromText(string(filter.Type))
|
||||
if !ok {
|
||||
fieldContext = telemetrytypes.FieldContextUnspecified
|
||||
}
|
||||
|
||||
var fieldDataType telemetrytypes.FieldDataType
|
||||
if err := fieldDataType.Scan(string(filter.DataType)); err != nil {
|
||||
fieldDataType = telemetrytypes.FieldDataTypeUnspecified
|
||||
}
|
||||
if fieldDataType == telemetrytypes.FieldDataTypeInt64 {
|
||||
fieldDataType = telemetrytypes.FieldDataTypeNumber
|
||||
}
|
||||
|
||||
fieldKeys = append(fieldKeys, telemetrytypes.TelemetryFieldKey{
|
||||
Name: filter.Key,
|
||||
Signal: fieldSignal,
|
||||
FieldContext: fieldContext,
|
||||
FieldDataType: fieldDataType,
|
||||
})
|
||||
}
|
||||
|
||||
return fieldKeys, nil
|
||||
}
|
||||
|
||||
// newLegacySourceFilters renders stored telemetry field keys
|
||||
// back into the v1 shape, restoring the legacy spellings v1 clients expect.
|
||||
func newLegacySourceFilters(sourceFilters *quickfiltertypes.SourceFilters) *legacySourceFilters {
|
||||
filters := make([]v3.AttributeKey, 0, len(sourceFilters.Filters))
|
||||
for _, fieldKey := range sourceFilters.Filters {
|
||||
// Only tag and resource exist in the v3 enum; other contexts render as
|
||||
// unspecified so v1 clients never see spellings their queries can't use.
|
||||
var attributeType v3.AttributeKeyType
|
||||
switch fieldKey.FieldContext {
|
||||
case telemetrytypes.FieldContextAttribute:
|
||||
attributeType = v3.AttributeKeyTypeTag
|
||||
case telemetrytypes.FieldContextResource:
|
||||
attributeType = v3.AttributeKeyTypeResource
|
||||
default:
|
||||
attributeType = v3.AttributeKeyTypeUnspecified
|
||||
}
|
||||
|
||||
var dataType v3.AttributeKeyDataType
|
||||
switch fieldKey.FieldDataType {
|
||||
case telemetrytypes.FieldDataTypeNumber:
|
||||
dataType = v3.AttributeKeyDataTypeFloat64
|
||||
default:
|
||||
dataType = v3.AttributeKeyDataType(fieldKey.FieldDataType.StringValue())
|
||||
}
|
||||
|
||||
filters = append(filters, v3.AttributeKey{
|
||||
Key: fieldKey.Name,
|
||||
Type: attributeType,
|
||||
DataType: dataType,
|
||||
})
|
||||
}
|
||||
|
||||
return &legacySourceFilters{
|
||||
Source: sourceFilters.Source,
|
||||
Filters: filters,
|
||||
}
|
||||
}
|
||||
|
||||
62
pkg/modules/quickfilter/implquickfilter/handler_test.go
Normal file
62
pkg/modules/quickfilter/implquickfilter/handler_test.go
Normal file
@@ -0,0 +1,62 @@
|
||||
package implquickfilter
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestNewTelemetryFieldKeysFromLegacy(t *testing.T) {
|
||||
fieldKeys, err := newTelemetryFieldKeysFromLegacy(quickfiltertypes.SourceTraces, []v3.AttributeKey{
|
||||
{Key: "service.name", Type: v3.AttributeKeyTypeResource, DataType: v3.AttributeKeyDataTypeString},
|
||||
{Key: "http.method", Type: v3.AttributeKeyTypeTag, DataType: v3.AttributeKeyDataTypeString},
|
||||
{Key: "duration_nano", Type: v3.AttributeKeyTypeTag, DataType: v3.AttributeKeyDataTypeFloat64},
|
||||
{Key: "code_line", Type: v3.AttributeKeyTypeTag, DataType: v3.AttributeKeyDataTypeInt64},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, fieldKeys, 4)
|
||||
|
||||
assert.Equal(t, telemetrytypes.TelemetryFieldKey{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString}, fieldKeys[0])
|
||||
assert.Equal(t, telemetrytypes.FieldContextAttribute, fieldKeys[1].FieldContext)
|
||||
assert.Equal(t, telemetrytypes.FieldDataTypeNumber, fieldKeys[2].FieldDataType)
|
||||
assert.Equal(t, telemetrytypes.FieldDataTypeNumber, fieldKeys[3].FieldDataType)
|
||||
|
||||
t.Run("meter writes restore the per-filter telemetry signal", func(t *testing.T) {
|
||||
fieldKeys, err := newTelemetryFieldKeysFromLegacy(quickfiltertypes.SourceMeter, []v3.AttributeKey{
|
||||
{Key: "host.name", DataType: v3.AttributeKeyDataTypeString},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, fieldKeys, 1)
|
||||
assert.Equal(t, telemetrytypes.SignalMetrics, fieldKeys[0].Signal)
|
||||
})
|
||||
|
||||
t.Run("rejects a filter without a key", func(t *testing.T) {
|
||||
_, err := newTelemetryFieldKeysFromLegacy(quickfiltertypes.SourceTraces, []v3.AttributeKey{{DataType: v3.AttributeKeyDataTypeString}})
|
||||
require.Error(t, err)
|
||||
})
|
||||
}
|
||||
|
||||
func TestNewLegacySourceFilters(t *testing.T) {
|
||||
legacy := newLegacySourceFilters(&quickfiltertypes.SourceFilters{
|
||||
Source: quickfiltertypes.SourceLogs,
|
||||
Filters: []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "http.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "duration_nano", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeNumber},
|
||||
{Name: "severity_text", FieldContext: telemetrytypes.FieldContextLog, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "host.name", Signal: telemetrytypes.SignalMetrics},
|
||||
},
|
||||
})
|
||||
|
||||
assert.Equal(t, quickfiltertypes.SourceLogs, legacy.Source)
|
||||
require.Len(t, legacy.Filters, 5)
|
||||
assert.Equal(t, v3.AttributeKey{Key: "service.name", Type: v3.AttributeKeyTypeResource, DataType: v3.AttributeKeyDataTypeString}, legacy.Filters[0])
|
||||
assert.Equal(t, v3.AttributeKeyTypeTag, legacy.Filters[1].Type)
|
||||
assert.Equal(t, v3.AttributeKeyDataTypeFloat64, legacy.Filters[2].DataType)
|
||||
assert.Equal(t, v3.AttributeKeyTypeUnspecified, legacy.Filters[3].Type, "contexts outside the v3 enum must render as unspecified")
|
||||
assert.Equal(t, v3.AttributeKey{Key: "host.name"}, legacy.Filters[4])
|
||||
}
|
||||
@@ -2,12 +2,11 @@ package implquickfilter
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
|
||||
@@ -19,91 +18,54 @@ func NewModule(store quickfiltertypes.QuickFilterStore) quickfilter.Module {
|
||||
return &module{store: store}
|
||||
}
|
||||
|
||||
// GetQuickFilters returns all quick filters for an organization.
|
||||
func (module *module) GetQuickFilters(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes.SignalFilters, error) {
|
||||
storedFilters, err := module.store.Get(ctx, orgID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error fetching organization filters")
|
||||
}
|
||||
|
||||
result := make([]*quickfiltertypes.SignalFilters, 0, len(storedFilters))
|
||||
for _, storedFilter := range storedFilters {
|
||||
signalFilter, err := quickfiltertypes.NewSignalFilterFromStorableQuickFilter(storedFilter)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for signal: %s", storedFilter.Signal)
|
||||
}
|
||||
result = append(result, signalFilter)
|
||||
}
|
||||
|
||||
return result, nil
|
||||
func (module *module) Get(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) (*quickfiltertypes.StorableQuickFilter, error) {
|
||||
return module.store.GetBySource(ctx, orgID, source.StringValue())
|
||||
}
|
||||
|
||||
// GetSignalFilters returns quick filters for a specific signal in an organization.
|
||||
func (m *module) GetSignalFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal) (*quickfiltertypes.SignalFilters, error) {
|
||||
storedFilter, err := m.store.GetBySignal(ctx, orgID, signal.StringValue())
|
||||
// GetQuickFilters returns quick filters for a source, or for every source when source is zero.
|
||||
func (module *module) GetQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) ([]*quickfiltertypes.SourceFilters, error) {
|
||||
if source.IsZero() {
|
||||
storedFilters, err := module.store.Get(ctx, orgID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error fetching organization filters")
|
||||
}
|
||||
|
||||
result := make([]*quickfiltertypes.SourceFilters, 0, len(storedFilters))
|
||||
for _, storedFilter := range storedFilters {
|
||||
sourceFilter, err := quickfiltertypes.NewSourceFilterFromStorableQuickFilter(storedFilter)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for source: %s", storedFilter.Source)
|
||||
}
|
||||
result = append(result, sourceFilter)
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
storedFilter, err := module.store.GetBySource(ctx, orgID, source.StringValue())
|
||||
if err != nil {
|
||||
if errors.Ast(err, errors.TypeNotFound) {
|
||||
return []*quickfiltertypes.SourceFilters{}, nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// If no filter exists for this signal, return empty filters with the requested signal
|
||||
if storedFilter == nil {
|
||||
return &quickfiltertypes.SignalFilters{
|
||||
Signal: signal,
|
||||
Filters: []v3.AttributeKey{},
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Convert stored filter to signal filter
|
||||
signalFilter, err := quickfiltertypes.NewSignalFilterFromStorableQuickFilter(storedFilter)
|
||||
sourceFilter, err := quickfiltertypes.NewSourceFilterFromStorableQuickFilter(storedFilter)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for signal: %s", storedFilter.Signal)
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error processing filter for source: %s", storedFilter.Source)
|
||||
}
|
||||
|
||||
return signalFilter, nil
|
||||
return []*quickfiltertypes.SourceFilters{sourceFilter}, nil
|
||||
}
|
||||
|
||||
// UpdateQuickFilters updates quick filters for a specific signal in an organization.
|
||||
func (module *module) UpdateQuickFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal, filters []v3.AttributeKey) error {
|
||||
// Validate each filter
|
||||
for _, filter := range filters {
|
||||
if err := filter.Validate(); err != nil {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Marshal filters to JSON
|
||||
filterJSON, err := json.Marshal(filters)
|
||||
// UpsertQuickFilters replaces quick filters for a specific source in an organization, creating them if absent.
|
||||
func (module *module) UpsertQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source, filters []telemetrytypes.TelemetryFieldKey) error {
|
||||
filter, err := quickfiltertypes.NewStorableQuickFilter(orgID, source, filters)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error marshalling filters")
|
||||
}
|
||||
|
||||
// Check if filter exists
|
||||
existingFilter, err := module.store.GetBySignal(ctx, orgID, signal.StringValue())
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error checking existing filters")
|
||||
}
|
||||
|
||||
var filter *quickfiltertypes.StorableQuickFilter
|
||||
if existingFilter != nil {
|
||||
// Update in place
|
||||
if err := existingFilter.Update(filterJSON); err != nil {
|
||||
return errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "error updating existing filter")
|
||||
}
|
||||
filter = existingFilter
|
||||
} else {
|
||||
// Create new
|
||||
filter, err = quickfiltertypes.NewStorableQuickFilter(orgID, signal, filterJSON)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "error creating new filter")
|
||||
}
|
||||
}
|
||||
|
||||
// Persist filter
|
||||
if err := module.store.Upsert(ctx, filter); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
return module.store.Upsert(ctx, filter)
|
||||
}
|
||||
|
||||
func (module *module) SetDefaultConfig(ctx context.Context, orgID valuer.UUID) error {
|
||||
|
||||
@@ -26,7 +26,7 @@ func (s *store) Get(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes
|
||||
NewSelect().
|
||||
Model(&filters).
|
||||
Where("org_id = ?", orgID).
|
||||
Order("signal ASC").
|
||||
Order("source ASC").
|
||||
Scan(ctx)
|
||||
|
||||
if err != nil {
|
||||
@@ -36,7 +36,7 @@ func (s *store) Get(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes
|
||||
return filters, nil
|
||||
}
|
||||
|
||||
func (s *store) GetBySignal(ctx context.Context, orgID valuer.UUID, signal string) (*quickfiltertypes.StorableQuickFilter, error) {
|
||||
func (s *store) GetBySource(ctx context.Context, orgID valuer.UUID, source string) (*quickfiltertypes.StorableQuickFilter, error) {
|
||||
filter := new(quickfiltertypes.StorableQuickFilter)
|
||||
|
||||
err := s.store.
|
||||
@@ -44,12 +44,12 @@ func (s *store) GetBySignal(ctx context.Context, orgID valuer.UUID, signal strin
|
||||
NewSelect().
|
||||
Model(filter).
|
||||
Where("org_id = ?", orgID).
|
||||
Where("signal = ?", signal).
|
||||
Where("source = ?", source).
|
||||
Scan(ctx)
|
||||
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, s.store.WrapNotFoundErrf(err, errors.CodeNotFound, "No rows found for org_id: "+orgID.StringValue()+" signal: "+signal)
|
||||
return nil, s.store.WrapNotFoundErrf(err, errors.CodeNotFound, "No rows found for org_id: "+orgID.StringValue()+" source: "+source)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
@@ -62,7 +62,7 @@ func (s *store) Upsert(ctx context.Context, filter *quickfiltertypes.StorableQui
|
||||
BunDB().
|
||||
NewInsert().
|
||||
Model(filter).
|
||||
On("CONFLICT (id) DO UPDATE").
|
||||
On("CONFLICT (org_id, source) DO UPDATE").
|
||||
Set("filter = EXCLUDED.filter").
|
||||
Set("updated_at = EXCLUDED.updated_at").
|
||||
Exec(ctx)
|
||||
@@ -78,7 +78,7 @@ func (s *store) Create(ctx context.Context, filters []*quickfiltertypes.Storable
|
||||
BunDBCtx(ctx).
|
||||
NewInsert().
|
||||
Model(&filters).
|
||||
On("CONFLICT (org_id, signal) DO NOTHING").
|
||||
On("CONFLICT (org_id, source) DO NOTHING").
|
||||
Exec(ctx)
|
||||
|
||||
if err != nil {
|
||||
|
||||
@@ -4,20 +4,27 @@ import (
|
||||
"context"
|
||||
"net/http"
|
||||
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
|
||||
type Module interface {
|
||||
GetQuickFilters(ctx context.Context, orgID valuer.UUID) ([]*quickfiltertypes.SignalFilters, error)
|
||||
UpdateQuickFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal, filters []v3.AttributeKey) error
|
||||
GetSignalFilters(ctx context.Context, orgID valuer.UUID, signal quickfiltertypes.Signal) (*quickfiltertypes.SignalFilters, error)
|
||||
// Get returns the stored quick filter row for a source.
|
||||
Get(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) (*quickfiltertypes.StorableQuickFilter, error)
|
||||
// GetQuickFilters returns quick filters for a source, or for every source when source is zero.
|
||||
GetQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source) ([]*quickfiltertypes.SourceFilters, error)
|
||||
UpsertQuickFilters(ctx context.Context, orgID valuer.UUID, source quickfiltertypes.Source, filters []telemetrytypes.TelemetryFieldKey) error
|
||||
SetDefaultConfig(ctx context.Context, orgID valuer.UUID) error
|
||||
}
|
||||
|
||||
type Handler interface {
|
||||
// Legacy v1 endpoints, served by converting to and from the v3 attribute key shape.
|
||||
GetQuickFilters(http.ResponseWriter, *http.Request)
|
||||
UpdateQuickFilters(http.ResponseWriter, *http.Request)
|
||||
GetSignalFilters(http.ResponseWriter, *http.Request)
|
||||
GetSourceFilters(http.ResponseWriter, *http.Request)
|
||||
|
||||
ListQuickFiltersV2(http.ResponseWriter, *http.Request)
|
||||
GetQuickFiltersV2(http.ResponseWriter, *http.Request)
|
||||
UpdateQuickFiltersV2(http.ResponseWriter, *http.Request)
|
||||
}
|
||||
|
||||
@@ -451,9 +451,9 @@ func (aH *APIHandler) RegisterRoutes(router *mux.Router, am *middleware.AuthZ) {
|
||||
|
||||
router.HandleFunc("/api/v1/disks", am.ViewAccess(aH.getDisks)).Methods(http.MethodGet)
|
||||
|
||||
// Quick Filters
|
||||
// Quick Filters (v1 routes serve the legacy v3 shape; v2 lives in signozapiserver)
|
||||
router.HandleFunc("/api/v1/orgs/me/filters", am.ViewAccess(aH.Signoz.Handlers.QuickFilter.GetQuickFilters)).Methods(http.MethodGet)
|
||||
router.HandleFunc("/api/v1/orgs/me/filters/{signal}", am.ViewAccess(aH.Signoz.Handlers.QuickFilter.GetSignalFilters)).Methods(http.MethodGet)
|
||||
router.HandleFunc("/api/v1/orgs/me/filters/{signal}", am.ViewAccess(aH.Signoz.Handlers.QuickFilter.GetSourceFilters)).Methods(http.MethodGet)
|
||||
router.HandleFunc("/api/v1/orgs/me/filters", am.AdminAccess(aH.Signoz.Handlers.QuickFilter.UpdateQuickFilters)).Methods(http.MethodPut)
|
||||
|
||||
router.HandleFunc("/api/v1/register", am.OpenAccess(aH.registerUser)).Methods(http.MethodPost)
|
||||
|
||||
@@ -29,6 +29,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/modules/organization"
|
||||
"github.com/SigNoz/signoz/pkg/modules/preference"
|
||||
"github.com/SigNoz/signoz/pkg/modules/promote"
|
||||
"github.com/SigNoz/signoz/pkg/modules/quickfilter"
|
||||
"github.com/SigNoz/signoz/pkg/modules/rawdataexport"
|
||||
"github.com/SigNoz/signoz/pkg/modules/rulestatehistory"
|
||||
"github.com/SigNoz/signoz/pkg/modules/savedview"
|
||||
@@ -95,6 +96,8 @@ func NewOpenAPI(ctx context.Context, instrumentation instrumentation.Instrumenta
|
||||
struct{ ruler.Handler }{},
|
||||
struct{ statsreporter.Handler }{},
|
||||
struct{ savedview.Handler }{},
|
||||
struct{ quickfilter.Module }{},
|
||||
struct{ quickfilter.Handler }{},
|
||||
).New(ctx, instrumentation.ToProviderSettings(), apiserver.Config{})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@@ -246,6 +246,8 @@ func NewSQLMigrationProviderFactories(
|
||||
sqlmigration.NewAddAuthDomainTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddDeploymentHostTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddSystemDashboardFactory(sqlstore, sqlschema),
|
||||
sqlmigration.NewMigrateQuickFiltersFactory(sqlstore),
|
||||
sqlmigration.NewAddQuickFilterTuplesFactory(sqlstore),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -350,6 +352,8 @@ func NewAPIServerProviderFactories(orgGetter organization.Getter, authz authz.Au
|
||||
handlers.RulerHandler,
|
||||
handlers.StatsHandler,
|
||||
handlers.SavedView,
|
||||
modules.QuickFilter,
|
||||
handlers.QuickFilter,
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -3,12 +3,13 @@ package sqlmigration
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
@@ -39,6 +40,52 @@ func (m *createQuickFilters) Register(migrations *migrate.Migrations) error {
|
||||
}
|
||||
|
||||
func (m *createQuickFilters) Up(ctx context.Context, db *bun.DB) error {
|
||||
// Frozen copy of the defaults as this migration shipped (hence the old
|
||||
// camelCase keys); migrations must not read live types. 031 replaces these rows.
|
||||
defaultFilters := []struct {
|
||||
signal string
|
||||
filters []map[string]any
|
||||
}{
|
||||
{"traces", []map[string]any{
|
||||
{"key": "duration_nano", "dataType": "float64", "type": "tag"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "hasError", "dataType": "bool", "type": "tag"},
|
||||
{"key": "serviceName", "dataType": "string", "type": "tag"},
|
||||
{"key": "name", "dataType": "string", "type": "resource"},
|
||||
{"key": "rpcMethod", "dataType": "string", "type": "tag"},
|
||||
{"key": "responseStatusCode", "dataType": "string", "type": "resource"},
|
||||
{"key": "httpHost", "dataType": "string", "type": "tag"},
|
||||
{"key": "httpMethod", "dataType": "string", "type": "tag"},
|
||||
{"key": "httpRoute", "dataType": "string", "type": "tag"},
|
||||
{"key": "httpUrl", "dataType": "string", "type": "tag"},
|
||||
{"key": "traceID", "dataType": "string", "type": "tag"},
|
||||
}},
|
||||
{"logs", []map[string]any{
|
||||
{"key": "severity_text", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "serviceName", "dataType": "string", "type": "tag"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
|
||||
}},
|
||||
{"api_monitoring", []map[string]any{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "serviceName", "dataType": "string", "type": "tag"},
|
||||
{"key": "rpcMethod", "dataType": "string", "type": "tag"},
|
||||
}},
|
||||
{"exceptions", []map[string]any{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "serviceName", "dataType": "string", "type": "tag"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "tag"},
|
||||
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.namespace.name", "dataType": "string", "type": "tag"},
|
||||
{"key": "k8s.pod.name", "dataType": "string", "type": "tag"},
|
||||
}},
|
||||
}
|
||||
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -72,15 +119,31 @@ func (m *createQuickFilters) Up(ctx context.Context, db *bun.DB) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// Get the default quick filters
|
||||
storableQuickFilters, err := quickfiltertypes.NewDefaultQuickFilter(defaultOrg)
|
||||
if err != nil {
|
||||
return err
|
||||
now := time.Now()
|
||||
quickFilters := make([]*quickFilter, 0, len(defaultFilters))
|
||||
for _, defaultFilter := range defaultFilters {
|
||||
filterJSON, err := json.Marshal(defaultFilter.filters)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
quickFilters = append(quickFilters, &quickFilter{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: defaultOrg.StringValue(),
|
||||
Filter: string(filterJSON),
|
||||
Signal: defaultFilter.signal,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: now,
|
||||
UpdatedAt: now,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// Insert all filters at once
|
||||
_, err = tx.NewInsert().
|
||||
Model(&storableQuickFilters).
|
||||
Model(&quickFilters).
|
||||
Exec(ctx)
|
||||
|
||||
if err != nil {
|
||||
|
||||
@@ -3,11 +3,13 @@ package sqlmigration
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
@@ -38,6 +40,61 @@ func (migration *updateQuickFilters) Register(migrations *migrate.Migrations) er
|
||||
}
|
||||
|
||||
func (migration *updateQuickFilters) Up(ctx context.Context, db *bun.DB) error {
|
||||
// Frozen copy of the defaults as this migration shipped; migrations must not
|
||||
// read live types. api_monitoring's service.name is "tag" here — 035 fixes it.
|
||||
defaultFilters := []struct {
|
||||
signal string
|
||||
filters []map[string]any
|
||||
}{
|
||||
{"traces", []map[string]any{
|
||||
{"key": "duration_nano", "dataType": "float64", "type": "tag"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "hasError", "dataType": "bool", "type": "tag"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "name", "dataType": "string", "type": "tag"},
|
||||
{"key": "rpc.method", "dataType": "string", "type": "tag"},
|
||||
{"key": "response_status_code", "dataType": "string", "type": "tag"},
|
||||
{"key": "http_host", "dataType": "string", "type": "tag"},
|
||||
{"key": "http.method", "dataType": "string", "type": "tag"},
|
||||
{"key": "http.route", "dataType": "string", "type": "tag"},
|
||||
{"key": "http_url", "dataType": "string", "type": "tag"},
|
||||
{"key": "trace_id", "dataType": "string", "type": "tag"},
|
||||
}},
|
||||
{"logs", []map[string]any{
|
||||
{"key": "severity_text", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
|
||||
}},
|
||||
{"api_monitoring", []map[string]any{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "tag"},
|
||||
{"key": "rpc.method", "dataType": "string", "type": "tag"},
|
||||
}},
|
||||
{"exceptions", []map[string]any{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
|
||||
}},
|
||||
}
|
||||
|
||||
signalFilters := make([]struct{ signal, filter string }, 0, len(defaultFilters))
|
||||
for _, defaultFilter := range defaultFilters {
|
||||
filterJSON, err := json.Marshal(defaultFilter.filters)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
signalFilters = append(signalFilters, struct{ signal, filter string }{defaultFilter.signal, string(filterJSON)})
|
||||
}
|
||||
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -73,17 +130,28 @@ func (migration *updateQuickFilters) Up(ctx context.Context, db *bun.DB) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// For each organization, create new quick filters with the updated NewDefaultQuickFilter function
|
||||
// For each organization, create new quick filters with the updated defaults
|
||||
for _, orgID := range orgIDs {
|
||||
// Get the updated default quick filters
|
||||
storableQuickFilters, err := quickfiltertypes.NewDefaultQuickFilter(valuer.MustNewUUID(orgID))
|
||||
if err != nil {
|
||||
return err
|
||||
now := time.Now()
|
||||
quickFilters := make([]*quickFilter, 0, len(signalFilters))
|
||||
for _, signalFilter := range signalFilters {
|
||||
quickFilters = append(quickFilters, &quickFilter{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: signalFilter.filter,
|
||||
Signal: signalFilter.signal,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: now,
|
||||
UpdatedAt: now,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// Insert all filters for this organization
|
||||
_, err = tx.NewInsert().
|
||||
Model(&storableQuickFilters).
|
||||
Model(&quickFilters).
|
||||
Exec(ctx)
|
||||
|
||||
if err != nil {
|
||||
|
||||
@@ -2,20 +2,16 @@ package sqlmigration
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
"github.com/SigNoz/signoz/pkg/types/quickfiltertypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
)
|
||||
|
||||
type updateApiMonitoringFilters struct {
|
||||
store sqlstore.SQLStore
|
||||
}
|
||||
type updateApiMonitoringFilters struct{}
|
||||
|
||||
func NewUpdateApiMonitoringFiltersFactory(store sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
|
||||
return factory.NewProviderFactory(factory.MustNewName("update_api_monitoring_filters"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
|
||||
@@ -23,10 +19,8 @@ func NewUpdateApiMonitoringFiltersFactory(store sqlstore.SQLStore) factory.Provi
|
||||
})
|
||||
}
|
||||
|
||||
func newUpdateApiMonitoringFilters(_ context.Context, _ factory.ProviderSettings, _ Config, store sqlstore.SQLStore) (SQLMigration, error) {
|
||||
return &updateApiMonitoringFilters{
|
||||
store: store,
|
||||
}, nil
|
||||
func newUpdateApiMonitoringFilters(_ context.Context, _ factory.ProviderSettings, _ Config, _ sqlstore.SQLStore) (SQLMigration, error) {
|
||||
return &updateApiMonitoringFilters{}, nil
|
||||
}
|
||||
|
||||
func (migration *updateApiMonitoringFilters) Register(migrations *migrate.Migrations) error {
|
||||
@@ -38,63 +32,29 @@ func (migration *updateApiMonitoringFilters) Register(migrations *migrate.Migrat
|
||||
}
|
||||
|
||||
func (migration *updateApiMonitoringFilters) Up(ctx context.Context, db *bun.DB) error {
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
// Frozen copy of the api_monitoring defaults as this migration shipped; the
|
||||
// change over 031 is service.name moving from "tag" to "resource".
|
||||
apiMonitoringFilters := []map[string]any{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "rpc.method", "dataType": "string", "type": "tag"},
|
||||
}
|
||||
|
||||
apiMonitoringFilterJSON, err := json.Marshal(apiMonitoringFilters)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
defer func() {
|
||||
_ = tx.Rollback()
|
||||
}()
|
||||
|
||||
// Get all organization IDs as strings
|
||||
var orgIDs []string
|
||||
err = tx.NewSelect().
|
||||
Table("organizations").
|
||||
Column("id").
|
||||
Scan(ctx, &orgIDs)
|
||||
// The filter JSON is org-independent, so one update covers every org's row.
|
||||
_, err = db.NewUpdate().
|
||||
Table("quick_filter").
|
||||
Set("filter = ?, updated_at = ?", string(apiMonitoringFilterJSON), time.Now()).
|
||||
Where("signal = ?", "api_monitoring").
|
||||
Exec(ctx)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
if err := tx.Commit(); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
for _, orgID := range orgIDs {
|
||||
// Get the updated default quick filters which includes the new API monitoring filters
|
||||
storableQuickFilters, err := quickfiltertypes.NewDefaultQuickFilter(valuer.MustNewUUID(orgID))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Find the API monitoring filter from the storable quick filters
|
||||
var apiMonitoringFilterJSON string
|
||||
for _, filter := range storableQuickFilters {
|
||||
if filter.Signal == quickfiltertypes.SignalApiMonitoring {
|
||||
apiMonitoringFilterJSON = filter.Filter
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
if apiMonitoringFilterJSON != "" {
|
||||
_, err = tx.NewUpdate().
|
||||
Table("quick_filter").
|
||||
Set("filter = ?, updated_at = ?", apiMonitoringFilterJSON, time.Now()).
|
||||
Where("signal = ? AND org_id = ?", quickfiltertypes.SignalApiMonitoring, orgID).
|
||||
Exec(ctx)
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if err := tx.Commit(); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
180
pkg/sqlmigration/120_migrate_quick_filters.go
Normal file
180
pkg/sqlmigration/120_migrate_quick_filters.go
Normal file
@@ -0,0 +1,180 @@
|
||||
package sqlmigration
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"strings"
|
||||
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
)
|
||||
|
||||
type storableQuickFilterRow struct {
|
||||
bun.BaseModel `bun:"table:quick_filter"`
|
||||
|
||||
ID string `bun:"id,pk"`
|
||||
Filter string `bun:"filter"`
|
||||
}
|
||||
|
||||
// legacyQuickFilterEntry carries both shapes a stored entry can be in: the
|
||||
// legacy key/type/dataType shape and the current name-carrying shape.
|
||||
type legacyQuickFilterEntry struct {
|
||||
Name string `json:"name"`
|
||||
Key string `json:"key"`
|
||||
Type string `json:"type"`
|
||||
DataType string `json:"dataType"`
|
||||
Signal string `json:"signal"`
|
||||
}
|
||||
|
||||
// quickFilterLegacyTypeToFieldContext maps the v3 attribute key types the v1
|
||||
// write path could store. Materialized top-level fields carried no type, and
|
||||
// anything unknown (e.g. "Sum" in the old meter defaults) normalizes to
|
||||
// unspecified, matching what the v1 write path does at runtime.
|
||||
var quickFilterLegacyTypeToFieldContext = map[string]string{
|
||||
"tag": "attribute",
|
||||
"resource": "resource",
|
||||
"scope": "scope",
|
||||
}
|
||||
|
||||
// quickFilterLegacyDataTypeToFieldDataType maps the v3 attribute key data
|
||||
// types the v1 write path could store, with numerics collapsed to number,
|
||||
// matching the fields API and the v1 write path.
|
||||
var quickFilterLegacyDataTypeToFieldDataType = map[string]string{
|
||||
"string": "string",
|
||||
"bool": "bool",
|
||||
"int64": "number",
|
||||
"float64": "number",
|
||||
}
|
||||
|
||||
type migrateQuickFilters struct {
|
||||
sqlstore sqlstore.SQLStore
|
||||
settings factory.ProviderSettings
|
||||
}
|
||||
|
||||
func NewMigrateQuickFiltersFactory(sqlstore sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
|
||||
return factory.NewProviderFactory(factory.MustNewName("migrate_quick_filters"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
|
||||
return &migrateQuickFilters{sqlstore: sqlstore, settings: ps}, nil
|
||||
})
|
||||
}
|
||||
|
||||
func (migration *migrateQuickFilters) Register(migrations *migrate.Migrations) error {
|
||||
return migrations.Register(migration.Up, migration.Down)
|
||||
}
|
||||
|
||||
func (migration *migrateQuickFilters) Up(ctx context.Context, db *bun.DB) error {
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
var rows []*storableQuickFilterRow
|
||||
if err := tx.NewSelect().Model(&rows).Scan(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var migrated, skipped int
|
||||
for _, row := range rows {
|
||||
migratedFilter, changed, ok := migrateQuickFilterEntries(row.Filter)
|
||||
if !ok {
|
||||
migration.settings.Logger.WarnContext(ctx, "quick filter could not be parsed, leaving it untouched", slog.String("quick_filter_id", row.ID), slog.String("raw_filter", row.Filter))
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
if !changed {
|
||||
continue
|
||||
}
|
||||
|
||||
migrated++
|
||||
if _, err := tx.NewUpdate().Model((*storableQuickFilterRow)(nil)).Set("filter = ?", migratedFilter).Where("id = ?", row.ID).Exec(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
migration.settings.Logger.InfoContext(ctx, "migrated quick filters to telemetry field keys", slog.Int("total", len(rows)), slog.Int("migrated", migrated), slog.Int("skipped", skipped))
|
||||
|
||||
if _, err := migration.sqlstore.Dialect().RenameColumn(ctx, tx, "quick_filter", "signal", "source"); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, column := range []string{"created_by", "updated_by"} {
|
||||
if err := migration.sqlstore.Dialect().DropColumn(ctx, tx, "quick_filter", column); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
func (migration *migrateQuickFilters) Down(context.Context, *bun.DB) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// migrateQuickFilterEntries rewrites a stored filter list from the legacy
|
||||
// key/dataType/type shape to telemetry field keys; ok=false means unparseable.
|
||||
func migrateQuickFilterEntries(filter string) (migrated string, changed bool, ok bool) {
|
||||
var entriesRaw []json.RawMessage
|
||||
if err := json.Unmarshal([]byte(filter), &entriesRaw); err != nil {
|
||||
return "", false, false
|
||||
}
|
||||
|
||||
migratedEntries := make([]json.RawMessage, 0, len(entriesRaw))
|
||||
for _, rawEntry := range entriesRaw {
|
||||
var entry legacyQuickFilterEntry
|
||||
if err := json.Unmarshal(rawEntry, &entry); err != nil {
|
||||
// Some stored entries are plain strings rather than objects; treat
|
||||
// the string as the filter key name, dropping empty ones.
|
||||
var name string
|
||||
if err := json.Unmarshal(rawEntry, &name); err != nil {
|
||||
return "", false, false
|
||||
}
|
||||
entry = legacyQuickFilterEntry{Key: name}
|
||||
}
|
||||
|
||||
switch {
|
||||
case entry.Name != "":
|
||||
migratedEntries = append(migratedEntries, rawEntry)
|
||||
case entry.Key != "":
|
||||
migratedJSON, err := marshalUnescaped(telemetryFieldKeyOutput{
|
||||
Name: entry.Key,
|
||||
Signal: entry.Signal,
|
||||
FieldContext: quickFilterFieldContext(entry.Type),
|
||||
FieldDataType: quickFilterFieldDataType(entry.DataType),
|
||||
})
|
||||
if err != nil {
|
||||
return "", false, false
|
||||
}
|
||||
migratedEntries = append(migratedEntries, migratedJSON)
|
||||
changed = true
|
||||
default:
|
||||
changed = true
|
||||
}
|
||||
}
|
||||
|
||||
if !changed {
|
||||
return "", false, true
|
||||
}
|
||||
|
||||
migratedJSON, err := marshalUnescaped(migratedEntries)
|
||||
if err != nil {
|
||||
return "", false, false
|
||||
}
|
||||
|
||||
return string(migratedJSON), true, true
|
||||
}
|
||||
|
||||
// quickFilterFieldDataType resolves legacy datatype spellings, with unknowns
|
||||
// normalized to unspecified.
|
||||
func quickFilterFieldDataType(legacyDataType string) string {
|
||||
return quickFilterLegacyDataTypeToFieldDataType[strings.ToLower(strings.TrimSpace(legacyDataType))]
|
||||
}
|
||||
|
||||
// quickFilterFieldContext resolves legacy type spellings, with unknowns
|
||||
// normalized to unspecified.
|
||||
func quickFilterFieldContext(legacyType string) string {
|
||||
return quickFilterLegacyTypeToFieldContext[strings.ToLower(strings.TrimSpace(legacyType))]
|
||||
}
|
||||
139
pkg/sqlmigration/121_add_quick_filter_tuples.go
Normal file
139
pkg/sqlmigration/121_add_quick_filter_tuples.go
Normal file
@@ -0,0 +1,139 @@
|
||||
package sqlmigration
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
"github.com/SigNoz/signoz/pkg/types/authtypes"
|
||||
"github.com/oklog/ulid/v2"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/dialect"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
)
|
||||
|
||||
type addQuickFilterTuples struct {
|
||||
sqlstore sqlstore.SQLStore
|
||||
}
|
||||
|
||||
func NewAddQuickFilterTuplesFactory(sqlstore sqlstore.SQLStore) factory.ProviderFactory[SQLMigration, Config] {
|
||||
return factory.NewProviderFactory(factory.MustNewName("add_quick_filter_tuples"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
|
||||
return &addQuickFilterTuples{sqlstore: sqlstore}, nil
|
||||
})
|
||||
}
|
||||
|
||||
func (migration *addQuickFilterTuples) Register(migrations *migrate.Migrations) error {
|
||||
return migrations.Register(migration.Up, migration.Down)
|
||||
}
|
||||
|
||||
func (migration *addQuickFilterTuples) Up(ctx context.Context, db *bun.DB) error {
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
var storeID string
|
||||
err = tx.QueryRowContext(ctx, `SELECT id FROM store WHERE name = ? LIMIT 1`, "signoz").Scan(&storeID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var orgIDs []string
|
||||
err = tx.NewSelect().
|
||||
Table("organizations").
|
||||
Column("id").
|
||||
Scan(ctx, &orgIDs)
|
||||
if err != nil && err != sql.ErrNoRows {
|
||||
return err
|
||||
}
|
||||
|
||||
isPG := migration.sqlstore.BunDB().Dialect().Name() == dialect.PG
|
||||
|
||||
// quick-filter moved from the legacy ViewAccess/AdminAccess role gate to
|
||||
// CheckResources, which on enterprise requires real tuples -- existing orgs
|
||||
// never had these written, only new orgs get them from the registry at bootstrap.
|
||||
tuples := []migrationTuple{
|
||||
{authtypes.SigNozAdminRoleName, "metaresource", "quick-filter", "read"},
|
||||
{authtypes.SigNozAdminRoleName, "metaresource", "quick-filter", "update"},
|
||||
{authtypes.SigNozAdminRoleName, "metaresource", "quick-filter", "list"},
|
||||
{authtypes.SigNozEditorRoleName, "metaresource", "quick-filter", "read"},
|
||||
{authtypes.SigNozEditorRoleName, "metaresource", "quick-filter", "list"},
|
||||
{authtypes.SigNozViewerRoleName, "metaresource", "quick-filter", "read"},
|
||||
{authtypes.SigNozViewerRoleName, "metaresource", "quick-filter", "list"},
|
||||
}
|
||||
|
||||
for _, orgID := range orgIDs {
|
||||
for _, tuple := range tuples {
|
||||
entropy := ulid.DefaultEntropy()
|
||||
now := time.Now().UTC()
|
||||
tupleID := ulid.MustNew(ulid.Timestamp(now), entropy).String()
|
||||
|
||||
objectID := "organization/" + orgID + "/" + tuple.objectName + "/*"
|
||||
roleSubject := "organization/" + orgID + "/role/" + tuple.roleName
|
||||
|
||||
if isPG {
|
||||
user := "role:" + roleSubject + "#assignee"
|
||||
result, err := tx.ExecContext(ctx, `
|
||||
INSERT INTO tuple (store, object_type, object_id, relation, _user, user_type, ulid, inserted_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT (store, object_type, object_id, relation, _user) DO NOTHING`,
|
||||
storeID, tuple.objectType, objectID, tuple.relation, user, "userset", tupleID, now,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rowsAffected, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if rowsAffected == 0 {
|
||||
continue
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, `
|
||||
INSERT INTO changelog (store, object_type, object_id, relation, _user, operation, ulid, inserted_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT (store, ulid, object_type) DO NOTHING`,
|
||||
storeID, tuple.objectType, objectID, tuple.relation, user, 0, tupleID, now,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
result, err := tx.ExecContext(ctx, `
|
||||
INSERT INTO tuple (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation, user_type, ulid, inserted_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation) DO NOTHING`,
|
||||
storeID, tuple.objectType, objectID, tuple.relation, "role", roleSubject, "assignee", "userset", tupleID, now,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rowsAffected, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if rowsAffected == 0 {
|
||||
continue
|
||||
}
|
||||
_, err = tx.ExecContext(ctx, `
|
||||
INSERT INTO changelog (store, object_type, object_id, relation, user_object_type, user_object_id, user_relation, operation, ulid, inserted_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT (store, ulid, object_type) DO NOTHING`,
|
||||
storeID, tuple.objectType, objectID, tuple.relation, "role", roleSubject, "assignee", 0, tupleID, now,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
func (migration *addQuickFilterTuples) Down(context.Context, *bun.DB) error {
|
||||
return nil
|
||||
}
|
||||
173
pkg/templating/markdownrenderer/adf/adf.go
Normal file
173
pkg/templating/markdownrenderer/adf/adf.go
Normal file
@@ -0,0 +1,173 @@
|
||||
// Package adf converts Markdown into Atlassian Document Format (ADF) nodes,
|
||||
// the JSON rich-text format used by Jira Cloud's v3 API.
|
||||
package adf
|
||||
|
||||
import (
|
||||
"strings"
|
||||
|
||||
"github.com/yuin/goldmark"
|
||||
"github.com/yuin/goldmark/ast"
|
||||
"github.com/yuin/goldmark/extension"
|
||||
extast "github.com/yuin/goldmark/extension/ast"
|
||||
"github.com/yuin/goldmark/text"
|
||||
)
|
||||
|
||||
// parser is stateless across Parse calls and safe for concurrent use; only
|
||||
// goldmark's renderers hold per-document state (which we don't use). Strikethrough
|
||||
// is included for the strike mark; linkify is deliberately omitted since it
|
||||
// fragments plain text into word tokens while scanning for bare URLs.
|
||||
var parser = goldmark.New(goldmark.WithExtensions(extension.Strikethrough)).Parser()
|
||||
|
||||
// Render returns the ADF block nodes for markdown (without the doc wrapper),
|
||||
// so callers can embed them alongside their own nodes (panels, links, …).
|
||||
func Render(markdown string) []any {
|
||||
src := []byte(markdown)
|
||||
return blockChildren(parser.Parse(text.NewReader(src)), src)
|
||||
}
|
||||
|
||||
func blockChildren(parent ast.Node, src []byte) []any {
|
||||
var out []any
|
||||
for c := parent.FirstChild(); c != nil; c = c.NextSibling() {
|
||||
if b := block(c, src); b != nil {
|
||||
out = append(out, b)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func block(n ast.Node, src []byte) any {
|
||||
switch node := n.(type) {
|
||||
case *ast.Heading:
|
||||
return map[string]any{"type": "heading", "attrs": map[string]any{"level": node.Level}, "content": inlineChildren(node, src, nil)}
|
||||
case *ast.Paragraph:
|
||||
return paragraph(inlineChildren(node, src, nil))
|
||||
case *ast.TextBlock:
|
||||
return paragraph(inlineChildren(node, src, nil))
|
||||
case *ast.List:
|
||||
typ := "bulletList"
|
||||
if node.IsOrdered() {
|
||||
typ = "orderedList"
|
||||
}
|
||||
return map[string]any{"type": typ, "content": blockChildren(node, src)}
|
||||
case *ast.ListItem:
|
||||
return map[string]any{"type": "listItem", "content": blockChildren(node, src)}
|
||||
case *ast.Blockquote:
|
||||
return map[string]any{"type": "blockquote", "content": blockChildren(node, src)}
|
||||
case *ast.FencedCodeBlock:
|
||||
return codeBlock(codeText(node, src), string(node.Language(src)))
|
||||
case *ast.CodeBlock:
|
||||
return codeBlock(codeText(node, src), "")
|
||||
case *ast.ThematicBreak:
|
||||
return map[string]any{"type": "rule"}
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func paragraph(content []any) map[string]any {
|
||||
p := map[string]any{"type": "paragraph"}
|
||||
if len(content) > 0 {
|
||||
p["content"] = content
|
||||
}
|
||||
return p
|
||||
}
|
||||
|
||||
func codeBlock(code, lang string) map[string]any {
|
||||
cb := map[string]any{"type": "codeBlock"}
|
||||
if lang != "" {
|
||||
cb["attrs"] = map[string]any{"language": lang}
|
||||
}
|
||||
if code = strings.TrimRight(code, "\n"); code != "" {
|
||||
cb["content"] = []any{map[string]any{"type": "text", "text": code}}
|
||||
}
|
||||
return cb
|
||||
}
|
||||
|
||||
// inlineChildren flattens an inline subtree into ADF text nodes, carrying the
|
||||
// active marks (strong/em/code/strike/link) down the tree.
|
||||
func inlineChildren(parent ast.Node, src []byte, marks []any) []any {
|
||||
var out []any
|
||||
for c := parent.FirstChild(); c != nil; c = c.NextSibling() {
|
||||
switch node := c.(type) {
|
||||
case *ast.Text:
|
||||
if t := string(node.Segment.Value(src)); t != "" {
|
||||
out = append(out, textNode(t, marks))
|
||||
}
|
||||
if node.HardLineBreak() {
|
||||
out = append(out, map[string]any{"type": "hardBreak"})
|
||||
} else if node.SoftLineBreak() {
|
||||
out = append(out, textNode(" ", marks))
|
||||
}
|
||||
case *ast.String:
|
||||
if len(node.Value) > 0 {
|
||||
out = append(out, textNode(string(node.Value), marks))
|
||||
}
|
||||
case *ast.CodeSpan:
|
||||
if t := rawText(node, src); t != "" {
|
||||
out = append(out, textNode(t, withMark(marks, mark("code"))))
|
||||
}
|
||||
case *ast.Emphasis:
|
||||
m := "em"
|
||||
if node.Level == 2 {
|
||||
m = "strong"
|
||||
}
|
||||
out = append(out, inlineChildren(node, src, withMark(marks, mark(m)))...)
|
||||
case *extast.Strikethrough:
|
||||
out = append(out, inlineChildren(node, src, withMark(marks, mark("strike")))...)
|
||||
case *ast.Link:
|
||||
out = append(out, inlineChildren(node, src, withMark(marks, linkMark(string(node.Destination))))...)
|
||||
case *ast.AutoLink:
|
||||
if u := string(node.URL(src)); u != "" {
|
||||
out = append(out, textNode(u, withMark(marks, linkMark(u))))
|
||||
}
|
||||
default:
|
||||
out = append(out, inlineChildren(c, src, marks)...)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func textNode(s string, marks []any) map[string]any {
|
||||
tn := map[string]any{"type": "text", "text": s}
|
||||
if len(marks) > 0 {
|
||||
tn["marks"] = marks
|
||||
}
|
||||
return tn
|
||||
}
|
||||
|
||||
func mark(typ string) any { return map[string]any{"type": typ} }
|
||||
|
||||
func linkMark(href string) any {
|
||||
return map[string]any{"type": "link", "attrs": map[string]any{"href": href}}
|
||||
}
|
||||
|
||||
func withMark(marks []any, m any) []any {
|
||||
out := make([]any, 0, len(marks)+1)
|
||||
out = append(out, marks...)
|
||||
return append(out, m)
|
||||
}
|
||||
|
||||
func rawText(n ast.Node, src []byte) string {
|
||||
var b strings.Builder
|
||||
for c := n.FirstChild(); c != nil; c = c.NextSibling() {
|
||||
switch t := c.(type) {
|
||||
case *ast.Text:
|
||||
b.Write(t.Segment.Value(src))
|
||||
case *ast.String:
|
||||
b.Write(t.Value)
|
||||
default:
|
||||
b.WriteString(rawText(c, src))
|
||||
}
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
|
||||
func codeText(n ast.Node, src []byte) string {
|
||||
var b strings.Builder
|
||||
lines := n.Lines()
|
||||
for i := 0; i < lines.Len(); i++ {
|
||||
seg := lines.At(i)
|
||||
b.Write(seg.Value(src))
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
85
pkg/templating/markdownrenderer/adf/adf_test.go
Normal file
85
pkg/templating/markdownrenderer/adf/adf_test.go
Normal file
@@ -0,0 +1,85 @@
|
||||
package adf
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func toJSON(t *testing.T, v any) string {
|
||||
t.Helper()
|
||||
b, err := json.Marshal(v)
|
||||
require.NoError(t, err)
|
||||
return string(b)
|
||||
}
|
||||
|
||||
func TestRenderInlineMarks(t *testing.T) {
|
||||
js := toJSON(t, Render("**bold** and *em* and `code` and [txt](https://x.io)"))
|
||||
assert.Contains(t, js, `"type":"strong"`)
|
||||
assert.Contains(t, js, `"type":"em"`)
|
||||
assert.Contains(t, js, `"type":"code"`)
|
||||
assert.Contains(t, js, `"type":"link"`)
|
||||
assert.Contains(t, js, `"href":"https://x.io"`)
|
||||
assert.Contains(t, js, `"text":"bold"`)
|
||||
}
|
||||
|
||||
func TestRenderHeadingAndList(t *testing.T) {
|
||||
js := toJSON(t, Render("# Title\n\n- a\n- b"))
|
||||
assert.Contains(t, js, `"type":"heading"`)
|
||||
assert.Contains(t, js, `"level":1`)
|
||||
assert.Contains(t, js, `"type":"bulletList"`)
|
||||
assert.Contains(t, js, `"type":"listItem"`)
|
||||
}
|
||||
|
||||
func TestRenderOrderedList(t *testing.T) {
|
||||
js := toJSON(t, Render("1. one\n2. two"))
|
||||
assert.Contains(t, js, `"type":"orderedList"`)
|
||||
}
|
||||
|
||||
func TestRenderCodeBlock(t *testing.T) {
|
||||
js := toJSON(t, Render("```go\nx := 1\n```"))
|
||||
assert.Contains(t, js, `"type":"codeBlock"`)
|
||||
assert.Contains(t, js, `"language":"go"`)
|
||||
assert.Contains(t, js, `x := 1`)
|
||||
}
|
||||
|
||||
func TestRenderStrikethrough(t *testing.T) {
|
||||
js := toJSON(t, Render("~~gone~~"))
|
||||
assert.Contains(t, js, `"type":"strike"`)
|
||||
assert.Contains(t, js, `"text":"gone"`)
|
||||
}
|
||||
|
||||
func TestRenderBlockquote(t *testing.T) {
|
||||
js := toJSON(t, Render("> quoted"))
|
||||
assert.Contains(t, js, `"type":"blockquote"`)
|
||||
assert.Contains(t, js, `"text":"quoted"`)
|
||||
}
|
||||
|
||||
func TestRenderAutoLink(t *testing.T) {
|
||||
js := toJSON(t, Render("see <https://signoz.io>"))
|
||||
assert.Contains(t, js, `"type":"link"`)
|
||||
assert.Contains(t, js, `"href":"https://signoz.io"`)
|
||||
assert.Contains(t, js, `"text":"https://signoz.io"`)
|
||||
}
|
||||
|
||||
func TestRenderLineBreaks(t *testing.T) {
|
||||
js := toJSON(t, Render("one \ntwo"))
|
||||
assert.Contains(t, js, `"type":"hardBreak"`)
|
||||
|
||||
// a soft break renders as a space, keeping the paragraph intact
|
||||
js = toJSON(t, Render("one\ntwo"))
|
||||
assert.NotContains(t, js, `"type":"hardBreak"`)
|
||||
assert.Contains(t, js, `"text":" "`)
|
||||
}
|
||||
|
||||
func TestRenderPlainText(t *testing.T) {
|
||||
js := toJSON(t, Render("just text"))
|
||||
assert.Contains(t, js, `"type":"paragraph"`)
|
||||
assert.Contains(t, js, `"text":"just text"`)
|
||||
}
|
||||
|
||||
func TestRenderEmptyIsEmpty(t *testing.T) {
|
||||
assert.Empty(t, Render(""))
|
||||
}
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/templating/markdownrenderer/blockkit"
|
||||
"github.com/SigNoz/signoz/pkg/templating/markdownrenderer/mrkdwn"
|
||||
"github.com/SigNoz/signoz/pkg/templating/markdownrenderer/plaintext"
|
||||
"github.com/yuin/goldmark"
|
||||
"github.com/yuin/goldmark/extension"
|
||||
)
|
||||
@@ -32,6 +33,11 @@ var (
|
||||
return goldmark.New(goldmark.WithExtensions(mrkdwn.Extender))
|
||||
},
|
||||
}
|
||||
plaintextPool = sync.Pool{
|
||||
New: func() any {
|
||||
return goldmark.New(goldmark.WithExtensions(plaintext.Extender))
|
||||
},
|
||||
}
|
||||
)
|
||||
|
||||
// RenderHTML converts markdown to HTML.
|
||||
@@ -53,6 +59,14 @@ func RenderSlackMrkdwn(markdown string) (string, error) {
|
||||
return render(md, markdown, "Slack mrkdwn")
|
||||
}
|
||||
|
||||
// RenderPlainText converts markdown to plain text: no markers, links flattened
|
||||
// to "text (url)".
|
||||
func RenderPlainText(markdown string) (string, error) {
|
||||
md := plaintextPool.Get().(goldmark.Markdown)
|
||||
defer plaintextPool.Put(md)
|
||||
return render(md, markdown, "plain text")
|
||||
}
|
||||
|
||||
func render(md goldmark.Markdown, markdown string, format string) (string, error) {
|
||||
var buf bytes.Buffer
|
||||
if err := md.Convert([]byte(markdown), &buf); err != nil {
|
||||
|
||||
301
pkg/templating/markdownrenderer/plaintext/plaintext.go
Normal file
301
pkg/templating/markdownrenderer/plaintext/plaintext.go
Normal file
@@ -0,0 +1,301 @@
|
||||
// Package plaintext provides a goldmark node renderer that emits plain text:
|
||||
// no markdown or HTML markers, and links flattened to "text (url)". It is used
|
||||
// for JSM Ops timeline notes, which render neither HTML nor markdown.
|
||||
package plaintext
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/yuin/goldmark"
|
||||
"github.com/yuin/goldmark/ast"
|
||||
"github.com/yuin/goldmark/extension"
|
||||
extensionast "github.com/yuin/goldmark/extension/ast"
|
||||
"github.com/yuin/goldmark/renderer"
|
||||
"github.com/yuin/goldmark/util"
|
||||
)
|
||||
|
||||
// Extender registers the plain-text node renderer plus the GFM extensions it
|
||||
// handles (tables, strikethrough).
|
||||
var Extender goldmark.Extender = &extender{}
|
||||
|
||||
type extender struct{}
|
||||
|
||||
func (e *extender) Extend(m goldmark.Markdown) {
|
||||
extension.Table.Extend(m)
|
||||
extension.Strikethrough.Extend(m)
|
||||
m.Renderer().AddOptions(
|
||||
renderer.WithNodeRenderers(util.Prioritized(newRenderer(), 1)),
|
||||
)
|
||||
}
|
||||
|
||||
// nodeRenderer holds per-document nesting prefixes, so it is not safe for
|
||||
// concurrent Convert calls; callers pool one instance per goroutine.
|
||||
type nodeRenderer struct {
|
||||
prefixes []string
|
||||
}
|
||||
|
||||
func newRenderer() renderer.NodeRenderer {
|
||||
return &nodeRenderer{}
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) RegisterFuncs(reg renderer.NodeRendererFuncRegisterer) {
|
||||
// Blocks
|
||||
reg.Register(ast.KindDocument, r.renderDocument)
|
||||
reg.Register(ast.KindHeading, r.renderBlock)
|
||||
reg.Register(ast.KindBlockquote, r.renderBlock)
|
||||
reg.Register(ast.KindCodeBlock, r.renderCodeBlock)
|
||||
reg.Register(ast.KindFencedCodeBlock, r.renderCodeBlock)
|
||||
reg.Register(ast.KindHTMLBlock, r.renderHTMLBlock)
|
||||
reg.Register(ast.KindList, r.renderList)
|
||||
reg.Register(ast.KindListItem, r.renderListItem)
|
||||
reg.Register(ast.KindParagraph, r.renderBlock)
|
||||
reg.Register(ast.KindTextBlock, r.renderTextBlock)
|
||||
reg.Register(ast.KindThematicBreak, r.renderThematicBreak)
|
||||
|
||||
// Inlines
|
||||
reg.Register(ast.KindAutoLink, r.renderAutoLink)
|
||||
reg.Register(ast.KindCodeSpan, r.renderCodeSpan)
|
||||
reg.Register(ast.KindEmphasis, r.renderPassthrough)
|
||||
reg.Register(ast.KindImage, r.renderLink)
|
||||
reg.Register(ast.KindLink, r.renderLink)
|
||||
reg.Register(ast.KindText, r.renderText)
|
||||
reg.Register(ast.KindString, r.renderString)
|
||||
reg.Register(ast.KindRawHTML, r.renderRawHTML)
|
||||
|
||||
// Extensions
|
||||
reg.Register(extensionast.KindStrikethrough, r.renderPassthrough)
|
||||
reg.Register(extensionast.KindTable, r.renderTable)
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) writePrefix(w util.BufWriter) {
|
||||
for _, p := range r.prefixes {
|
||||
_, _ = w.WriteString(p)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) writeLineSeparator(w util.BufWriter) {
|
||||
_ = w.WriteByte('\n')
|
||||
r.writePrefix(w)
|
||||
}
|
||||
|
||||
// writeBlockSeparator writes a blank line between block-level elements.
|
||||
func (r *nodeRenderer) writeBlockSeparator(w util.BufWriter) {
|
||||
r.writeLineSeparator(w)
|
||||
r.writeLineSeparator(w)
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) separateFromPrevious(w util.BufWriter, n ast.Node) {
|
||||
if n.PreviousSibling() != nil {
|
||||
r.writeBlockSeparator(w)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderDocument(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
// The renderer is pooled; wipe any prefix stack left over from a prior
|
||||
// document (e.g. one that errored mid-walk) before starting fresh.
|
||||
r.prefixes = r.prefixes[:0]
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
// renderBlock separates block-level nodes (paragraph, heading, blockquote) from
|
||||
// their previous sibling with a blank line, emitting no markers of their own.
|
||||
func (r *nodeRenderer) renderBlock(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
r.separateFromPrevious(w, node)
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderCodeBlock(w util.BufWriter, source []byte, n ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
r.separateFromPrevious(w, n)
|
||||
l := n.Lines().Len()
|
||||
for i := 0; i < l; i++ {
|
||||
line := n.Lines().At(i)
|
||||
_, _ = w.Write(line.Value(source))
|
||||
}
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderList(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering && node.PreviousSibling() != nil {
|
||||
r.writeLineSeparator(w)
|
||||
if node.Parent() == nil || node.Parent().Kind() != ast.KindListItem {
|
||||
r.writeLineSeparator(w)
|
||||
}
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderListItem(w util.BufWriter, source []byte, n ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
if n.PreviousSibling() != nil {
|
||||
r.writeLineSeparator(w)
|
||||
}
|
||||
parent := n.Parent().(*ast.List)
|
||||
var prefixStr string
|
||||
if parent.IsOrdered() {
|
||||
index := parent.Start
|
||||
for c := parent.FirstChild(); c != nil && c != n; c = c.NextSibling() {
|
||||
index++
|
||||
}
|
||||
prefixStr = fmt.Sprintf("%d. ", index)
|
||||
} else {
|
||||
prefixStr = "- "
|
||||
}
|
||||
_, _ = w.WriteString(prefixStr)
|
||||
r.prefixes = append(r.prefixes, " ") // indent wrapped/nested lines
|
||||
} else {
|
||||
r.prefixes = r.prefixes[:len(r.prefixes)-1]
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderTextBlock(w util.BufWriter, source []byte, n ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering && n.PreviousSibling() != nil {
|
||||
r.writeLineSeparator(w)
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderThematicBreak(w util.BufWriter, source []byte, n ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
r.separateFromPrevious(w, n)
|
||||
_, _ = w.WriteString("---")
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderAutoLink(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if !entering {
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
n := node.(*ast.AutoLink)
|
||||
url := string(n.URL(source))
|
||||
if n.AutoLinkType == ast.AutoLinkEmail && !strings.HasPrefix(strings.ToLower(url), "mailto:") {
|
||||
url = "mailto:" + url
|
||||
}
|
||||
_, _ = w.WriteString(url)
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderCodeSpan(w util.BufWriter, source []byte, n ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
for c := n.FirstChild(); c != nil; c = c.NextSibling() {
|
||||
segment := c.(*ast.Text).Segment
|
||||
value := segment.Value(source)
|
||||
if bytes.HasSuffix(value, []byte("\n")) {
|
||||
_, _ = w.Write(value[:len(value)-1])
|
||||
_ = w.WriteByte(' ')
|
||||
} else {
|
||||
_, _ = w.Write(value)
|
||||
}
|
||||
}
|
||||
return ast.WalkSkipChildren, nil
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
// renderPassthrough emits no markers; the node's children render as plain text
|
||||
// (used for emphasis/strong and strikethrough).
|
||||
func (r *nodeRenderer) renderPassthrough(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
// renderLink flattens links and images to "text (url)": children render the
|
||||
// label, then the destination is appended in parentheses on exit.
|
||||
func (r *nodeRenderer) renderLink(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
var dest []byte
|
||||
switch n := node.(type) {
|
||||
case *ast.Link:
|
||||
dest = n.Destination
|
||||
case *ast.Image:
|
||||
dest = n.Destination
|
||||
}
|
||||
if !entering && len(dest) > 0 {
|
||||
_, _ = fmt.Fprintf(w, " (%s)", dest)
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderText(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if !entering {
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
n := node.(*ast.Text)
|
||||
_, _ = w.Write(n.Segment.Value(source))
|
||||
if n.HardLineBreak() || n.SoftLineBreak() {
|
||||
r.writeLineSeparator(w)
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderString(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if entering {
|
||||
_, _ = w.Write(node.(*ast.String).Value)
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderRawHTML(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
// Drop inline raw HTML tags; a plain-text note should never carry markup.
|
||||
return ast.WalkSkipChildren, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderHTMLBlock(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
// Drop block-level raw HTML for the same reason as inline raw HTML.
|
||||
return ast.WalkSkipChildren, nil
|
||||
}
|
||||
|
||||
func (r *nodeRenderer) renderTable(w util.BufWriter, source []byte, node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if !entering {
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
r.separateFromPrevious(w, node)
|
||||
|
||||
first := true
|
||||
for c := node.FirstChild(); c != nil; c = c.NextSibling() {
|
||||
if c.Kind() != extensionast.KindTableHeader && c.Kind() != extensionast.KindTableRow {
|
||||
continue
|
||||
}
|
||||
if !first {
|
||||
r.writeLineSeparator(w)
|
||||
}
|
||||
first = false
|
||||
cellFirst := true
|
||||
for cc := c.FirstChild(); cc != nil; cc = cc.NextSibling() {
|
||||
if cc.Kind() != extensionast.KindTableCell {
|
||||
continue
|
||||
}
|
||||
if !cellFirst {
|
||||
_, _ = w.WriteString(" | ")
|
||||
}
|
||||
cellFirst = false
|
||||
_, _ = w.WriteString(extractPlainText(cc, source))
|
||||
}
|
||||
}
|
||||
return ast.WalkSkipChildren, nil
|
||||
}
|
||||
|
||||
// extractPlainText collects the text content of a node.
|
||||
func extractPlainText(n ast.Node, source []byte) string {
|
||||
var buf bytes.Buffer
|
||||
_ = ast.Walk(n, func(node ast.Node, entering bool) (ast.WalkStatus, error) {
|
||||
if !entering {
|
||||
return ast.WalkContinue, nil
|
||||
}
|
||||
switch t := node.(type) {
|
||||
case *ast.Text:
|
||||
buf.Write(t.Segment.Value(source))
|
||||
case *ast.String:
|
||||
buf.Write(t.Value)
|
||||
}
|
||||
return ast.WalkContinue, nil
|
||||
})
|
||||
return strings.TrimSpace(buf.String())
|
||||
}
|
||||
55
pkg/templating/markdownrenderer/plaintext/plaintext_test.go
Normal file
55
pkg/templating/markdownrenderer/plaintext/plaintext_test.go
Normal file
@@ -0,0 +1,55 @@
|
||||
package plaintext
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"github.com/yuin/goldmark"
|
||||
)
|
||||
|
||||
func render(t *testing.T, md string) string {
|
||||
t.Helper()
|
||||
var b []byte
|
||||
buf := bytesBuffer{&b}
|
||||
g := goldmark.New(goldmark.WithExtensions(Extender))
|
||||
require.NoError(t, g.Convert([]byte(md), &buf))
|
||||
return string(b)
|
||||
}
|
||||
|
||||
// bytesBuffer is a tiny io.Writer so the test needs no extra imports.
|
||||
type bytesBuffer struct{ b *[]byte }
|
||||
|
||||
func (w bytesBuffer) Write(p []byte) (int, error) {
|
||||
*w.b = append(*w.b, p...)
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
func TestPlainText(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
in string
|
||||
want string
|
||||
}{
|
||||
{"strips bold and italic", "**bold** and *italic*", "bold and italic"},
|
||||
{"link becomes text (url)", "[View in SigNoz](https://signoz.io/alert)", "View in SigNoz (https://signoz.io/alert)"},
|
||||
{"bold label kept, marker dropped", "**Alert:** name (critical)", "Alert: name (critical)"},
|
||||
{"strikethrough stripped", "~~gone~~", "gone"},
|
||||
{"inline code unwrapped", "run `foo bar`", "run foo bar"},
|
||||
{"paragraphs separated by blank line", "one\n\ntwo", "one\n\ntwo"},
|
||||
{"unordered list", "- a\n- b", "- a\n- b"},
|
||||
{"ordered list keeps numbering", "1. a\n2. b", "1. a\n2. b"},
|
||||
{"nested list indents under parent", "- a\n - b", "- a\n - b"},
|
||||
{"fenced code block unwrapped", "```go\nx := 1\n```", "x := 1\n"},
|
||||
{"table flattens to pipe-separated rows", "| h1 | h2 |\n|---|---|\n| a | b |\n| c | d |", "h1 | h2\na | b\nc | d"},
|
||||
{"autolink kept as bare url", "see <https://signoz.io>", "see https://signoz.io"},
|
||||
{"inline raw html dropped", "a <b>bold</b> word", "a bold word"},
|
||||
{"html block dropped", "before\n\n<div>markup</div>\n\nafter", "before\n\nafter"},
|
||||
{"hard break becomes newline", "one \ntwo", "one\ntwo"},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
assert.Equal(t, c.want, render(t, c.in))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -216,13 +216,20 @@ func (PostableChannel) JSONSchema() (jsonschema.Schema, error) {
|
||||
schema.WithRequired("name")
|
||||
|
||||
var oneOf []jsonschema.SchemaOrBool
|
||||
// Walk both halves: native fields on Receiver, upstream on the embed.
|
||||
seen := map[string]struct{}{}
|
||||
// Walk both halves: native fields on Receiver, upstream on the embed. A native
|
||||
// field can shadow an upstream one with the same tag (e.g. jira_configs), so
|
||||
// dedupe to avoid emitting two identical oneOf branches.
|
||||
collect := func(t reflect.Type) {
|
||||
for i := 0; i < t.NumField(); i++ {
|
||||
jsonTag := strings.Split(t.Field(i).Tag.Get("json"), ",")[0]
|
||||
if !strings.HasSuffix(jsonTag, "_configs") {
|
||||
continue
|
||||
}
|
||||
if _, ok := seen[jsonTag]; ok {
|
||||
continue
|
||||
}
|
||||
seen[jsonTag] = struct{}{}
|
||||
branch := (&jsonschema.Schema{}).WithRequired(jsonTag)
|
||||
oneOf = append(oneOf, branch.ToSchemaOrBool())
|
||||
}
|
||||
|
||||
@@ -70,15 +70,19 @@ type Config struct {
|
||||
// on Receiver, and extensions to customConfigsOf + isEmpty.
|
||||
type customReceiverConfigs struct {
|
||||
GoogleChat []*GoogleChatReceiverConfig
|
||||
Jira []*JiraReceiverConfig
|
||||
JSMOps []*JSMOpsReceiverConfig
|
||||
}
|
||||
|
||||
func (c customReceiverConfigs) isEmpty() bool {
|
||||
return len(c.GoogleChat) == 0
|
||||
return len(c.GoogleChat) == 0 && len(c.Jira) == 0 && len(c.JSMOps) == 0
|
||||
}
|
||||
|
||||
func customConfigsOf(receiver *Receiver) customReceiverConfigs {
|
||||
return customReceiverConfigs{
|
||||
GoogleChat: receiver.GoogleChatConfigs,
|
||||
Jira: receiver.JiraConfigs,
|
||||
JSMOps: receiver.JSMOpsConfigs,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -187,6 +191,8 @@ func extendedReceivers(c *config.Config, customConfigs map[string]customReceiver
|
||||
receivers[i] = &Receiver{
|
||||
Receiver: &base,
|
||||
GoogleChatConfigs: custom.GoogleChat,
|
||||
JiraConfigs: custom.Jira,
|
||||
JSMOpsConfigs: custom.JSMOps,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -362,6 +368,8 @@ func (c *Config) GetReceiver(name string) (*Receiver, error) {
|
||||
return &Receiver{
|
||||
Receiver: &base,
|
||||
GoogleChatConfigs: custom.GoogleChat,
|
||||
JiraConfigs: custom.Jira,
|
||||
JSMOpsConfigs: custom.JSMOps,
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
@@ -441,6 +449,16 @@ func (c *Config) applyNativeDefaults() {
|
||||
gc.HTTPConfig = httpDefault
|
||||
}
|
||||
}
|
||||
for _, jc := range custom.Jira {
|
||||
if jc.HTTPConfig == nil {
|
||||
jc.HTTPConfig = httpDefault
|
||||
}
|
||||
}
|
||||
for _, jc := range custom.JSMOps {
|
||||
if jc.HTTPConfig == nil {
|
||||
jc.HTTPConfig = httpDefault
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
119
pkg/types/alertmanagertypes/jira.go
Normal file
119
pkg/types/alertmanagertypes/jira.go
Normal file
@@ -0,0 +1,119 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/url"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
"github.com/prometheus/common/model"
|
||||
)
|
||||
|
||||
const defaultJiraReopenDuration = model.Duration(3 * 24 * time.Hour)
|
||||
|
||||
// Service accounts authenticate against the api.atlassian.com gateway (keyed by
|
||||
// cloud id) instead of the site host; they are identified by their email domain.
|
||||
const (
|
||||
jiraCloudHostSuffix = ".atlassian.net"
|
||||
jiraServiceAccountEmailDomain = "@serviceaccount.atlassian.com"
|
||||
jiraGatewayBaseURL = "https://api.atlassian.com/ex/jira/"
|
||||
)
|
||||
|
||||
// Default templates for the issue title and body. The body is rendered to
|
||||
// markdown and then wrapped in the ADF status panel + deep-links by the notifier.
|
||||
const (
|
||||
DefaultJiraSummaryTemplate = `[{{ .Status | toUpper }}{{ if eq .Status "firing" }}:{{ .Alerts.Firing | len }}{{ end }}] {{ .CommonLabels.alertname }}`
|
||||
|
||||
DefaultJiraDescriptionTemplate = `{{ range .Alerts -}}
|
||||
**Alert:** {{ .Labels.alertname }}{{ if .Labels.severity }} ({{ .Labels.severity }}){{ end }}
|
||||
{{ if .Annotations.summary }}
|
||||
**Summary:** {{ .Annotations.summary }}
|
||||
{{ end }}{{ if .Annotations.description }}
|
||||
**Description:** {{ .Annotations.description }}
|
||||
{{ end }}
|
||||
{{ end }}`
|
||||
)
|
||||
|
||||
// JiraReceiverConfig is the SigNoz Jira receiver. Fields are declared explicitly
|
||||
// instead of embedding upstream config.JiraConfig because that type's own
|
||||
// UnmarshalYAML would reset our defaults and drop sibling fields on the yaml
|
||||
// round-trip. Only Jira Cloud (v3/ADF) is supported, so api_url is derived from Site.
|
||||
type JiraReceiverConfig struct {
|
||||
config.NotifierConfig `yaml:",inline"`
|
||||
|
||||
Site string `json:"site,omitempty" yaml:"site,omitempty"`
|
||||
Project string `json:"project,omitempty" yaml:"project,omitempty"`
|
||||
IssueType string `json:"issue_type,omitempty" yaml:"issue_type,omitempty"`
|
||||
Summary string `json:"summary,omitempty" yaml:"summary,omitempty"`
|
||||
Description string `json:"description,omitempty" yaml:"description,omitempty"`
|
||||
Priority string `json:"priority,omitempty" yaml:"priority,omitempty"`
|
||||
Labels []string `json:"labels,omitempty" yaml:"labels,omitempty"`
|
||||
ResolveTransition string `json:"resolve_transition,omitempty" yaml:"resolve_transition,omitempty"`
|
||||
ReopenTransition string `json:"reopen_transition,omitempty" yaml:"reopen_transition,omitempty"`
|
||||
ReopenDuration model.Duration `json:"reopen_duration" yaml:"reopen_duration"`
|
||||
WontFixResolution string `json:"wont_fix_resolution,omitempty" yaml:"wont_fix_resolution,omitempty"`
|
||||
CustomFields map[string]any `json:"custom_fields,omitempty" yaml:"custom_fields,omitempty"`
|
||||
HTTPConfig *commoncfg.HTTPClientConfig `json:"http_config,omitempty" yaml:"http_config,omitempty"`
|
||||
}
|
||||
|
||||
func (c *JiraReceiverConfig) UnmarshalYAML(unmarshal func(any) error) error {
|
||||
type plain JiraReceiverConfig
|
||||
if err := unmarshal((*plain)(c)); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if c.ReopenDuration <= 0 {
|
||||
c.ReopenDuration = defaultJiraReopenDuration
|
||||
}
|
||||
// sub-minute windows truncate to 0 in the reopen JQL and silently disable
|
||||
// reopening, so reject them.
|
||||
if c.ReopenDuration < model.Duration(time.Minute) {
|
||||
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "jira reopen_duration must be at least 1m")
|
||||
}
|
||||
if c.Summary == "" {
|
||||
c.Summary = DefaultJiraSummaryTemplate
|
||||
}
|
||||
if c.Description == "" {
|
||||
c.Description = DefaultJiraDescriptionTemplate
|
||||
}
|
||||
|
||||
site := strings.TrimRight(strings.TrimSpace(c.Site), "/")
|
||||
u, err := url.Parse(site)
|
||||
if site == "" || err != nil || u.Scheme != "https" || !strings.HasSuffix(strings.ToLower(u.Hostname()), jiraCloudHostSuffix) {
|
||||
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, fmt.Sprintf("jira site must be a Jira Cloud URL (https://<site>%s)", jiraCloudHostSuffix))
|
||||
}
|
||||
c.Site = site
|
||||
|
||||
if c.Project == "" {
|
||||
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "jira project is required")
|
||||
}
|
||||
if c.IssueType == "" {
|
||||
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "jira issue_type is required")
|
||||
}
|
||||
if c.HTTPConfig == nil || c.HTTPConfig.BasicAuth == nil {
|
||||
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "jira requires basic auth (email + API token)")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// IsServiceAccount reports whether the basic-auth user is an Atlassian service
|
||||
// account, identified by its email domain. Service accounts must go through the
|
||||
// api.atlassian.com gateway; personal API tokens use the site host directly.
|
||||
func (c *JiraReceiverConfig) IsServiceAccount() bool {
|
||||
if c.HTTPConfig == nil || c.HTTPConfig.BasicAuth == nil {
|
||||
return false
|
||||
}
|
||||
return strings.HasSuffix(strings.ToLower(c.HTTPConfig.BasicAuth.Username), jiraServiceAccountEmailDomain)
|
||||
}
|
||||
|
||||
// APIBaseURL returns the Jira Cloud REST v3 base URL: the api.atlassian.com
|
||||
// gateway when a cloud id is given (service accounts), else the site host.
|
||||
func (c *JiraReceiverConfig) APIBaseURL(cloudID string) string {
|
||||
if cloudID != "" {
|
||||
return fmt.Sprintf("%s%s/rest/api/3", jiraGatewayBaseURL, cloudID)
|
||||
}
|
||||
return fmt.Sprintf("%s/rest/api/3", strings.TrimRight(c.Site, "/"))
|
||||
}
|
||||
119
pkg/types/alertmanagertypes/jira_test.go
Normal file
119
pkg/types/alertmanagertypes/jira_test.go
Normal file
@@ -0,0 +1,119 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
"github.com/prometheus/common/model"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func jiraReceiverJSON(site, project, issueType string, withAuth bool) string {
|
||||
auth := ""
|
||||
if withAuth {
|
||||
auth = `,"http_config":{"basic_auth":{"username":"me@acme.com","password":"token"}}`
|
||||
}
|
||||
return fmt.Sprintf(
|
||||
`{"name":"jira","jira_configs":[{"site":%q,"project":%q,"issue_type":%q%s}]}`,
|
||||
site, project, issueType, auth,
|
||||
)
|
||||
}
|
||||
|
||||
func TestJiraReceiverConfigDefaults(t *testing.T) {
|
||||
r, err := NewReceiver(jiraReceiverJSON("https://acme.atlassian.net", "KAN", "Task", true))
|
||||
require.NoError(t, err)
|
||||
require.Len(t, r.JiraConfigs, 1)
|
||||
|
||||
jc := r.JiraConfigs[0]
|
||||
assert.Equal(t, "https://acme.atlassian.net", jc.Site)
|
||||
assert.Equal(t, "https://acme.atlassian.net/rest/api/3", jc.APIBaseURL(""))
|
||||
assert.False(t, jc.SendResolved()) // default off when omitted, like other channels
|
||||
assert.Equal(t, defaultJiraReopenDuration, jc.ReopenDuration)
|
||||
assert.Equal(t, DefaultJiraSummaryTemplate, jc.Summary)
|
||||
assert.Equal(t, DefaultJiraDescriptionTemplate, jc.Description)
|
||||
|
||||
ch, err := NewChannelFromReceiver(r, "org-1")
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "jira", ch.Type)
|
||||
}
|
||||
|
||||
func TestJiraReceiverConfigSendResolved(t *testing.T) {
|
||||
withSendResolved := func(v bool) string {
|
||||
return fmt.Sprintf(
|
||||
`{"name":"j","jira_configs":[{"site":"https://acme.atlassian.net","project":"KAN","issue_type":"Task","send_resolved":%t,"http_config":{"basic_auth":{"username":"e","password":"t"}}}]}`,
|
||||
v,
|
||||
)
|
||||
}
|
||||
on, err := NewReceiver(withSendResolved(true))
|
||||
require.NoError(t, err)
|
||||
assert.True(t, on.JiraConfigs[0].SendResolved())
|
||||
|
||||
off, err := NewReceiver(withSendResolved(false))
|
||||
require.NoError(t, err)
|
||||
assert.False(t, off.JiraConfigs[0].SendResolved())
|
||||
}
|
||||
|
||||
func TestJiraReceiverConfigReopenDurationMinimum(t *testing.T) {
|
||||
withReopen := func(v string) string {
|
||||
return fmt.Sprintf(
|
||||
`{"name":"j","jira_configs":[{"site":"https://acme.atlassian.net","project":"KAN","issue_type":"Task","reopen_duration":%q,"http_config":{"basic_auth":{"username":"e","password":"t"}}}]}`,
|
||||
v,
|
||||
)
|
||||
}
|
||||
r, err := NewReceiver(withReopen("1m"))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, model.Duration(time.Minute), r.JiraConfigs[0].ReopenDuration)
|
||||
|
||||
_, err = NewReceiver(withReopen("30s"))
|
||||
assert.Error(t, err)
|
||||
}
|
||||
|
||||
func TestJiraAPIBaseURL(t *testing.T) {
|
||||
c := &JiraReceiverConfig{Site: "https://acme.atlassian.net"}
|
||||
assert.Equal(t, "https://acme.atlassian.net/rest/api/3", c.APIBaseURL(""))
|
||||
assert.Equal(t, "https://api.atlassian.com/ex/jira/09851b38-1a40-4c01-a36a-0a9336293200/rest/api/3", c.APIBaseURL("09851b38-1a40-4c01-a36a-0a9336293200"))
|
||||
}
|
||||
|
||||
func TestJiraIsServiceAccount(t *testing.T) {
|
||||
withUser := func(username string) *JiraReceiverConfig {
|
||||
return &JiraReceiverConfig{HTTPConfig: &commoncfg.HTTPClientConfig{BasicAuth: &commoncfg.BasicAuth{Username: username}}}
|
||||
}
|
||||
assert.True(t, withUser("bot@serviceaccount.atlassian.com").IsServiceAccount())
|
||||
assert.True(t, withUser("Bot@ServiceAccount.Atlassian.Com").IsServiceAccount())
|
||||
assert.False(t, withUser("temp@signoz.io").IsServiceAccount())
|
||||
assert.False(t, (&JiraReceiverConfig{}).IsServiceAccount())
|
||||
}
|
||||
|
||||
func TestJiraReceiverConfigTrailingSlashSite(t *testing.T) {
|
||||
r, err := NewReceiver(jiraReceiverJSON("https://acme.atlassian.net/", "KAN", "Task", true))
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "https://acme.atlassian.net", r.JiraConfigs[0].Site)
|
||||
assert.Equal(t, "https://acme.atlassian.net/rest/api/3", r.JiraConfigs[0].APIBaseURL(""))
|
||||
}
|
||||
|
||||
func TestJiraReceiverConfigValidation(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
json string
|
||||
}{
|
||||
{"missing site", `{"name":"j","jira_configs":[{"project":"KAN","issue_type":"Task","http_config":{"basic_auth":{"username":"e","password":"t"}}}]}`},
|
||||
{"http site", jiraReceiverJSON("http://acme.atlassian.net", "KAN", "Task", true)},
|
||||
{"non-cloud host", jiraReceiverJSON("https://jira.acme.com", "KAN", "Task", true)},
|
||||
{"lookalike host suffix", jiraReceiverJSON("https://www.iamnotatlassian.net", "KAN", "Task", true)},
|
||||
{"bare atlassian.net", jiraReceiverJSON("https://atlassian.net", "KAN", "Task", true)},
|
||||
{"missing project", jiraReceiverJSON("https://acme.atlassian.net", "", "Task", true)},
|
||||
{"missing issue_type", jiraReceiverJSON("https://acme.atlassian.net", "KAN", "", true)},
|
||||
{"missing basic auth", jiraReceiverJSON("https://acme.atlassian.net", "KAN", "Task", false)},
|
||||
{"invalid reopen_duration format", `{"name":"j","jira_configs":[{"site":"https://acme.atlassian.net","project":"KAN","issue_type":"Task","reopen_duration":"3days","http_config":{"basic_auth":{"username":"e","password":"t"}}}]}`},
|
||||
{"sub-minute reopen_duration", `{"name":"j","jira_configs":[{"site":"https://acme.atlassian.net","project":"KAN","issue_type":"Task","reopen_duration":"30s","http_config":{"basic_auth":{"username":"e","password":"t"}}}]}`},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
_, err := NewReceiver(c.json)
|
||||
assert.Error(t, err)
|
||||
})
|
||||
}
|
||||
}
|
||||
75
pkg/types/alertmanagertypes/jsmops.go
Normal file
75
pkg/types/alertmanagertypes/jsmops.go
Normal file
@@ -0,0 +1,75 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
)
|
||||
|
||||
// JSMOpsAPIBaseURL is the native JSM Ops integration-events gateway. It is a
|
||||
// single global host keyed by the integration API key (no region/cloud id in
|
||||
// the path). The trailing slash is required: the Opsgenie notifier appends
|
||||
// "v2/alerts..." to APIURL.Path with no separator.
|
||||
const JSMOpsAPIBaseURL = "https://api.atlassian.com/jsm/ops/integration/"
|
||||
|
||||
// JSM Ops speaks the Opsgenie alert API, so a JSM alert description takes the
|
||||
// same HTML subset and 15,000-char limit; message caps at 130. The templates
|
||||
// mirror Google Chat / Jira for a consistent default across channels.
|
||||
const (
|
||||
DefaultJSMOpsMessageTemplate = `[{{ .Status | toUpper }}{{ if eq .Status "firing" }}:{{ .Alerts.Firing | len }}{{ end }}] {{ .CommonLabels.alertname }}`
|
||||
|
||||
DefaultJSMOpsDescriptionTemplate = `{{ range .Alerts -}}
|
||||
**Alert:** {{ .Labels.alertname }}{{ if .Labels.severity }} ({{ .Labels.severity }}){{ end }}
|
||||
|
||||
{{ if .Annotations.summary }}**Summary:** {{ .Annotations.summary }}
|
||||
|
||||
{{ end }}{{ if .Annotations.description }}**Description:** {{ .Annotations.description }}
|
||||
|
||||
{{ end }}{{ if .GeneratorURL }}[View in SigNoz]({{ .GeneratorURL }})
|
||||
|
||||
{{ end }}{{ if .Annotations.related_logs }}[View related logs]({{ .Annotations.related_logs }})
|
||||
|
||||
{{ end }}{{ if .Annotations.related_traces }}[View related traces]({{ .Annotations.related_traces }})
|
||||
|
||||
{{ end }}{{ end }}`
|
||||
)
|
||||
|
||||
// JSMOpsReceiverConfig is the SigNoz Jira Service Management Ops receiver. It is
|
||||
// delivered by reusing the Opsgenie notifier (JSM Ops is the ex-Opsgenie alert
|
||||
// API): the notifier package maps these fields onto config.OpsGenieConfig with
|
||||
// APIURL pinned to JSMOpsAPIBaseURL.
|
||||
type JSMOpsReceiverConfig struct {
|
||||
config.NotifierConfig `yaml:",inline" json:",inline"`
|
||||
|
||||
HTTPConfig *commoncfg.HTTPClientConfig `yaml:"http_config,omitempty" json:"http_config,omitempty"`
|
||||
|
||||
APIKey config.Secret `yaml:"api_key,omitempty" json:"api_key,omitempty"`
|
||||
Message string `yaml:"message,omitempty" json:"message,omitempty"`
|
||||
Description string `yaml:"description,omitempty" json:"description,omitempty"`
|
||||
Priority string `yaml:"priority,omitempty" json:"priority,omitempty"`
|
||||
Tags string `yaml:"tags,omitempty" json:"tags,omitempty"`
|
||||
}
|
||||
|
||||
// send_resolved has no omitempty upstream, so a var default here is overwritten
|
||||
// by the yaml round-trip to the request value (false when omitted); the UI sends
|
||||
// it explicitly, defaulted on, so JSM alerts close on resolve.
|
||||
var DefaultJSMOpsReceiverConfig = JSMOpsReceiverConfig{
|
||||
NotifierConfig: config.NotifierConfig{
|
||||
VSendResolved: false,
|
||||
},
|
||||
Message: DefaultJSMOpsMessageTemplate,
|
||||
Description: DefaultJSMOpsDescriptionTemplate,
|
||||
Tags: "signoz",
|
||||
}
|
||||
|
||||
func (c *JSMOpsReceiverConfig) UnmarshalYAML(unmarshal func(any) error) error {
|
||||
*c = DefaultJSMOpsReceiverConfig
|
||||
type plain JSMOpsReceiverConfig
|
||||
if err := unmarshal((*plain)(c)); err != nil {
|
||||
return err
|
||||
}
|
||||
if c.APIKey == "" {
|
||||
return errors.New(errors.TypeInvalidInput, errors.CodeInvalidInput, "jsm ops api_key is required")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
70
pkg/types/alertmanagertypes/jsmops_test.go
Normal file
70
pkg/types/alertmanagertypes/jsmops_test.go
Normal file
@@ -0,0 +1,70 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestJSMOpsReceiverConfigDefaults(t *testing.T) {
|
||||
r, err := NewReceiver(`{"name":"jsm","jsmops_configs":[{"api_key":"key-123"}]}`)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, r.JSMOpsConfigs, 1)
|
||||
|
||||
c := r.JSMOpsConfigs[0]
|
||||
assert.Equal(t, "key-123", string(c.APIKey))
|
||||
assert.Equal(t, DefaultJSMOpsMessageTemplate, c.Message)
|
||||
assert.Equal(t, DefaultJSMOpsDescriptionTemplate, c.Description)
|
||||
assert.Equal(t, "signoz", c.Tags)
|
||||
assert.False(t, c.SendResolved()) // default off when omitted, like other channels
|
||||
|
||||
ch, err := NewChannelFromReceiver(r, "org-1")
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "jsmops", ch.Type)
|
||||
}
|
||||
|
||||
func TestJSMOpsReceiverConfigOverrides(t *testing.T) {
|
||||
r, err := NewReceiver(`{"name":"jsm","jsmops_configs":[{"api_key":"k","message":"m","description":"d","priority":"P1","tags":"a,b","send_resolved":true}]}`)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, r.JSMOpsConfigs, 1)
|
||||
|
||||
c := r.JSMOpsConfigs[0]
|
||||
assert.Equal(t, "m", c.Message)
|
||||
assert.Equal(t, "d", c.Description)
|
||||
assert.Equal(t, "P1", c.Priority)
|
||||
assert.Equal(t, "a,b", c.Tags)
|
||||
assert.True(t, c.SendResolved())
|
||||
}
|
||||
|
||||
func TestJSMOpsReceiverConfigSendResolved(t *testing.T) {
|
||||
withSendResolved := func(v bool) string {
|
||||
return fmt.Sprintf(`{"name":"jsm","jsmops_configs":[{"api_key":"k","send_resolved":%t}]}`, v)
|
||||
}
|
||||
on, err := NewReceiver(withSendResolved(true))
|
||||
require.NoError(t, err)
|
||||
require.Len(t, on.JSMOpsConfigs, 1)
|
||||
assert.True(t, on.JSMOpsConfigs[0].SendResolved())
|
||||
|
||||
off, err := NewReceiver(withSendResolved(false))
|
||||
require.NoError(t, err)
|
||||
require.Len(t, off.JSMOpsConfigs, 1)
|
||||
assert.False(t, off.JSMOpsConfigs[0].SendResolved())
|
||||
}
|
||||
|
||||
func TestJSMOpsReceiverConfigValidation(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
json string
|
||||
}{
|
||||
{"missing api_key", `{"name":"jsm","jsmops_configs":[{"message":"m"}]}`},
|
||||
{"empty api_key", `{"name":"jsm","jsmops_configs":[{"api_key":""}]}`},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
_, err := NewReceiver(c.json)
|
||||
assert.Error(t, err)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -23,6 +23,11 @@ import (
|
||||
type Receiver struct {
|
||||
*config.Receiver
|
||||
GoogleChatConfigs []*GoogleChatReceiverConfig `json:"googlechat_configs,omitempty" yaml:"googlechat_configs,omitempty"`
|
||||
// Shadows upstream's jira_configs so our custom notifier (rich ADF, deep-links,
|
||||
// lifecycle comments) handles it instead of upstream's plain Jira notifier.
|
||||
JiraConfigs []*JiraReceiverConfig `json:"jira_configs,omitempty" yaml:"jira_configs,omitempty"`
|
||||
// JSM Ops (ex-Opsgenie alert API); delivered by reusing the Opsgenie notifier.
|
||||
JSMOpsConfigs []*JSMOpsReceiverConfig `json:"jsmops_configs,omitempty" yaml:"jsmops_configs,omitempty"`
|
||||
}
|
||||
|
||||
// NewReceiver builds a Receiver from its JSON input, applying each notifier
|
||||
@@ -51,6 +56,22 @@ func NewReceiver(input string) (*Receiver, error) {
|
||||
receiver.GoogleChatConfigs[i] = defaulted
|
||||
}
|
||||
|
||||
for i, jc := range receiver.JiraConfigs {
|
||||
defaulted, err := defaultedNotifierConfig(jc)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
receiver.JiraConfigs[i] = defaulted
|
||||
}
|
||||
|
||||
for i, jc := range receiver.JSMOpsConfigs {
|
||||
defaulted, err := defaultedNotifierConfig(jc)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
receiver.JSMOpsConfigs[i] = defaulted
|
||||
}
|
||||
|
||||
return receiver, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -62,7 +62,7 @@ var (
|
||||
ResourceMetaResourcePipeline = NewResourceMetaResource(KindPipeline)
|
||||
ResourceMetaResourceUserPreference = NewResourceMetaResource(KindUserPreference)
|
||||
ResourceMetaResourceOrgPreference = NewResourceMetaResource(KindOrgPreference)
|
||||
ResourceMetaResourceQuickFilter = NewResourceMetaResource(KindQuickFilter)
|
||||
ResourceMetaResourceQuickFilter = NewResourceMetaResource(KindQuickFilter, VerbList, VerbRead, VerbUpdate)
|
||||
ResourceMetaResourceTTLSetting = NewResourceMetaResource(KindTTLSetting)
|
||||
ResourceMetaResourceRule = NewResourceMetaResource(KindRule)
|
||||
ResourceMetaResourcePlannedMaintenance = NewResourceMetaResource(KindPlannedMaintenance)
|
||||
|
||||
@@ -5,58 +5,58 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/uptrace/bun"
|
||||
)
|
||||
|
||||
type Signal struct {
|
||||
type Source struct {
|
||||
valuer.String
|
||||
}
|
||||
|
||||
func (enum *Signal) UnmarshalJSON(data []byte) error {
|
||||
func (enum *Source) UnmarshalJSON(data []byte) error {
|
||||
var str string
|
||||
if err := json.Unmarshal(data, &str); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
signal, err := NewSignal(str)
|
||||
source, err := NewSource(str)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
*enum = signal
|
||||
*enum = source
|
||||
return nil
|
||||
}
|
||||
|
||||
var (
|
||||
SignalTraces = Signal{valuer.NewString("traces")}
|
||||
SignalLogs = Signal{valuer.NewString("logs")}
|
||||
SignalApiMonitoring = Signal{valuer.NewString("api_monitoring")}
|
||||
SignalExceptions = Signal{valuer.NewString("exceptions")}
|
||||
SignalMeter = Signal{valuer.NewString("meter")}
|
||||
SignalAiObservability = Signal{valuer.NewString("ai_observability")}
|
||||
SourceTraces = Source{valuer.NewString("traces")}
|
||||
SourceLogs = Source{valuer.NewString("logs")}
|
||||
SourceApiMonitoring = Source{valuer.NewString("api_monitoring")}
|
||||
SourceExceptions = Source{valuer.NewString("exceptions")}
|
||||
SourceMeter = Source{valuer.NewString("meter")}
|
||||
SourceAiObservability = Source{valuer.NewString("ai_observability")}
|
||||
)
|
||||
|
||||
// NewSignal creates a Signal from a string.
|
||||
func NewSignal(s string) (Signal, error) {
|
||||
// NewSource creates a Source from a string.
|
||||
func NewSource(s string) (Source, error) {
|
||||
switch s {
|
||||
case "traces":
|
||||
return SignalTraces, nil
|
||||
return SourceTraces, nil
|
||||
case "logs":
|
||||
return SignalLogs, nil
|
||||
return SourceLogs, nil
|
||||
case "api_monitoring":
|
||||
return SignalApiMonitoring, nil
|
||||
return SourceApiMonitoring, nil
|
||||
case "exceptions":
|
||||
return SignalExceptions, nil
|
||||
return SourceExceptions, nil
|
||||
case "meter":
|
||||
return SignalMeter, nil
|
||||
return SourceMeter, nil
|
||||
case "ai_observability":
|
||||
return SignalAiObservability, nil
|
||||
return SourceAiObservability, nil
|
||||
default:
|
||||
return Signal{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid signal: %s", s)
|
||||
return Source{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid source: %s", s)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -65,33 +65,46 @@ type StorableQuickFilter struct {
|
||||
types.Identifiable
|
||||
OrgID valuer.UUID `bun:"org_id,type:text,notnull"`
|
||||
Filter string `bun:"filter,type:text,notnull"`
|
||||
Signal Signal `bun:"signal,type:text,notnull"`
|
||||
Source Source `bun:"source,type:text,notnull"`
|
||||
types.TimeAuditable
|
||||
}
|
||||
|
||||
type SignalFilters struct {
|
||||
Signal Signal `json:"signal"`
|
||||
Filters []v3.AttributeKey `json:"filters"`
|
||||
type SourceFilters struct {
|
||||
types.Identifiable
|
||||
types.TimeAuditable
|
||||
|
||||
OrgID valuer.UUID `json:"orgId"`
|
||||
Source Source `json:"source"`
|
||||
Filters []telemetrytypes.TelemetryFieldKey `json:"filters" required:"true" nullable:"false"`
|
||||
}
|
||||
|
||||
type UpdatableQuickFilters struct {
|
||||
Signal Signal `json:"signal"`
|
||||
Filters []v3.AttributeKey `json:"filters"`
|
||||
Filters []telemetrytypes.TelemetryFieldKey `json:"filters" required:"true" nullable:"false"`
|
||||
}
|
||||
|
||||
// NewStorableQuickFilter creates a new StorableQuickFilter after validation.
|
||||
func NewStorableQuickFilter(orgID valuer.UUID, signal Signal, filterJSON []byte) (*StorableQuickFilter, error) {
|
||||
if orgID.StringValue() == "" {
|
||||
func NewStorableQuickFilter(orgID valuer.UUID, source Source, filters []telemetrytypes.TelemetryFieldKey) (*StorableQuickFilter, error) {
|
||||
if orgID.IsZero() {
|
||||
return nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "orgID is required")
|
||||
}
|
||||
|
||||
if _, err := NewSignal(signal.StringValue()); err != nil {
|
||||
if _, err := NewSource(source.StringValue()); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var filters []v3.AttributeKey
|
||||
if err := json.Unmarshal(filterJSON, &filters); err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter JSON")
|
||||
if err := validateFilters(filters); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// A nil slice marshals to the JSON literal "null"; store an empty array so
|
||||
// reads never have to render a null filter list.
|
||||
if filters == nil {
|
||||
filters = []telemetrytypes.TelemetryFieldKey{}
|
||||
}
|
||||
|
||||
filterJSON, err := json.Marshal(filters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "error marshalling filters")
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
@@ -100,7 +113,7 @@ func NewStorableQuickFilter(orgID valuer.UUID, signal Signal, filterJSON []byte)
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Signal: signal,
|
||||
Source: source,
|
||||
Filter: string(filterJSON),
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: now,
|
||||
@@ -109,25 +122,21 @@ func NewStorableQuickFilter(orgID valuer.UUID, signal Signal, filterJSON []byte)
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Update updates an existing StorableQuickFilter with new filter data after validation.
|
||||
func (quickfilter *StorableQuickFilter) Update(filterJSON []byte) error {
|
||||
var filters []v3.AttributeKey
|
||||
if err := json.Unmarshal(filterJSON, &filters); err != nil {
|
||||
return errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid filter JSON")
|
||||
// NewSourceFiltersFromSource creates a SourceFilters with no filters for a source.
|
||||
func NewSourceFiltersFromSource(source Source) *SourceFilters {
|
||||
return &SourceFilters{
|
||||
Source: source,
|
||||
Filters: []telemetrytypes.TelemetryFieldKey{},
|
||||
}
|
||||
|
||||
quickfilter.Filter = string(filterJSON)
|
||||
quickfilter.UpdatedAt = time.Now()
|
||||
return nil
|
||||
}
|
||||
|
||||
// NewSignalFilterFromStorableQuickFilter converts a StorableQuickFilter to a SignalFilters object.
|
||||
func NewSignalFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFilter) (*SignalFilters, error) {
|
||||
// NewSourceFilterFromStorableQuickFilter converts a StorableQuickFilter to a SourceFilters object.
|
||||
func NewSourceFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFilter) (*SourceFilters, error) {
|
||||
if storableQuickFilter == nil {
|
||||
return nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "storableQuickFilter cannot be nil")
|
||||
}
|
||||
|
||||
var filters []v3.AttributeKey
|
||||
filters := []telemetrytypes.TelemetryFieldKey{}
|
||||
if storableQuickFilter.Filter != "" {
|
||||
err := json.Unmarshal([]byte(storableQuickFilter.Filter), &filters)
|
||||
if err != nil {
|
||||
@@ -135,178 +144,114 @@ func NewSignalFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFi
|
||||
}
|
||||
}
|
||||
|
||||
return &SignalFilters{
|
||||
Signal: storableQuickFilter.Signal,
|
||||
Filters: filters,
|
||||
// Stored filter JSON can be the literal "null" (a nil slice was upserted),
|
||||
// which unmarshals to nil; the API contract requires a non-null array.
|
||||
if filters == nil {
|
||||
filters = []telemetrytypes.TelemetryFieldKey{}
|
||||
}
|
||||
|
||||
return &SourceFilters{
|
||||
Identifiable: storableQuickFilter.Identifiable,
|
||||
OrgID: storableQuickFilter.OrgID,
|
||||
Source: storableQuickFilter.Source,
|
||||
Filters: filters,
|
||||
TimeAuditable: storableQuickFilter.TimeAuditable,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// NewDefaultQuickFilter generates default filters for all supported signals.
|
||||
// NewDefaultQuickFilter generates default filters for all supported sources.
|
||||
func NewDefaultQuickFilter(orgID valuer.UUID) ([]*StorableQuickFilter, error) {
|
||||
tracesFilters := []map[string]interface{}{
|
||||
{"key": "duration_nano", "dataType": "float64", "type": "tag"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "hasError", "dataType": "bool", "type": "tag"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "name", "dataType": "string", "type": "tag"},
|
||||
{"key": "rpc.method", "dataType": "string", "type": "tag"},
|
||||
{"key": "response_status_code", "dataType": "string", "type": "tag"},
|
||||
{"key": "http_host", "dataType": "string", "type": "tag"},
|
||||
{"key": "http.method", "dataType": "string", "type": "tag"},
|
||||
{"key": "http.route", "dataType": "string", "type": "tag"},
|
||||
{"key": "http_url", "dataType": "string", "type": "tag"},
|
||||
{"key": "trace_id", "dataType": "string", "type": "tag"},
|
||||
tracesFilters := []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "duration_nano", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeNumber},
|
||||
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "hasError", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeBool},
|
||||
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "name", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "rpc.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "response_status_code", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "http_host", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "http.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "http.route", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "http_url", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "trace_id", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
}
|
||||
|
||||
logsFilters := []map[string]interface{}{
|
||||
{"key": "severity_text", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
|
||||
logsFilters := []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "severity_text", FieldContext: telemetrytypes.FieldContextLog, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "host.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.cluster.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.deployment.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.namespace.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.pod.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
}
|
||||
|
||||
apiMonitoringFilters := []map[string]interface{}{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "rpc.method", "dataType": "string", "type": "tag"},
|
||||
apiMonitoringFilters := []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "rpc.method", FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
}
|
||||
|
||||
exceptionsFilters := []map[string]interface{}{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.deployment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.namespace.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.pod.name", "dataType": "string", "type": "resource"},
|
||||
exceptionsFilters := []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "host.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.cluster.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.deployment.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.namespace.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "k8s.pod.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
}
|
||||
|
||||
meterFilters := []map[string]interface{}{
|
||||
{"key": "deployment.environment", "dataType": "float64", "type": "Sum"},
|
||||
{"key": "service.name", "dataType": "float64", "type": "Sum"},
|
||||
{"key": "host.name", "dataType": "float64", "type": "Sum"},
|
||||
// Meter keys are label names with no context or datatype: the meter fields
|
||||
// API returns them as name+signal only, so the defaults mirror that shape.
|
||||
meterFilters := []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "deployment.environment", Signal: telemetrytypes.SignalMetrics},
|
||||
{Name: "service.name", Signal: telemetrytypes.SignalMetrics},
|
||||
{Name: "host.name", Signal: telemetrytypes.SignalMetrics},
|
||||
}
|
||||
|
||||
// AI observability (builder_ai_query trace explorer), ordered by expected
|
||||
// usage: env scoping, the LLM identity keys, then service and the rest.
|
||||
aiObservabilityFilters := []map[string]interface{}{
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": aiobservabilitytypes.GenAIOperationName, "dataType": "string", "type": "tag"},
|
||||
{"key": aiobservabilitytypes.GenAIProviderName, "dataType": "string", "type": "tag"},
|
||||
{"key": aiobservabilitytypes.GenAIRequestModel, "dataType": "string", "type": "tag"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": aiobservabilitytypes.GenAIToolName, "dataType": "string", "type": "tag"},
|
||||
{"key": aiobservabilitytypes.GenAIAgentName, "dataType": "string", "type": "tag"},
|
||||
aiObservabilityFilters := []telemetrytypes.TelemetryFieldKey{
|
||||
{Name: "deployment.environment", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: aiobservabilitytypes.GenAIOperationName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: aiobservabilitytypes.GenAIProviderName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: aiobservabilitytypes.GenAIRequestModel, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: "service.name", FieldContext: telemetrytypes.FieldContextResource, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: aiobservabilitytypes.GenAIToolName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
{Name: aiobservabilitytypes.GenAIAgentName, FieldContext: telemetrytypes.FieldContextAttribute, FieldDataType: telemetrytypes.FieldDataTypeString},
|
||||
}
|
||||
|
||||
tracesJSON, err := json.Marshal(tracesFilters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal traces filters")
|
||||
defaults := []struct {
|
||||
source Source
|
||||
filters []telemetrytypes.TelemetryFieldKey
|
||||
}{
|
||||
{SourceTraces, tracesFilters},
|
||||
{SourceLogs, logsFilters},
|
||||
{SourceApiMonitoring, apiMonitoringFilters},
|
||||
{SourceExceptions, exceptionsFilters},
|
||||
{SourceMeter, meterFilters},
|
||||
{SourceAiObservability, aiObservabilityFilters},
|
||||
}
|
||||
|
||||
logsJSON, err := json.Marshal(logsFilters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal logs filters")
|
||||
storableQuickFilters := make([]*StorableQuickFilter, 0, len(defaults))
|
||||
for _, def := range defaults {
|
||||
storableQuickFilter, err := NewStorableQuickFilter(orgID, def.source, def.filters)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
storableQuickFilters = append(storableQuickFilters, storableQuickFilter)
|
||||
}
|
||||
|
||||
apiMonitoringJSON, err := json.Marshal(apiMonitoringFilters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal api monitoring filters")
|
||||
}
|
||||
|
||||
exceptionsJSON, err := json.Marshal(exceptionsFilters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal exceptions filters")
|
||||
}
|
||||
|
||||
meterJSON, err := json.Marshal(meterFilters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal meter filters")
|
||||
}
|
||||
|
||||
aiObservabilityJSON, err := json.Marshal(aiObservabilityFilters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, errors.TypeInternal, errors.CodeInternal, "failed to marshal ai observability filters")
|
||||
}
|
||||
|
||||
timeRightNow := time.Now()
|
||||
|
||||
return []*StorableQuickFilter{
|
||||
{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: string(tracesJSON),
|
||||
Signal: SignalTraces,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: timeRightNow,
|
||||
UpdatedAt: timeRightNow,
|
||||
},
|
||||
},
|
||||
{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: string(logsJSON),
|
||||
Signal: SignalLogs,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: timeRightNow,
|
||||
UpdatedAt: timeRightNow,
|
||||
},
|
||||
},
|
||||
{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: string(apiMonitoringJSON),
|
||||
Signal: SignalApiMonitoring,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: timeRightNow,
|
||||
UpdatedAt: timeRightNow,
|
||||
},
|
||||
},
|
||||
{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: string(exceptionsJSON),
|
||||
Signal: SignalExceptions,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: timeRightNow,
|
||||
UpdatedAt: timeRightNow,
|
||||
},
|
||||
},
|
||||
{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: string(meterJSON),
|
||||
Signal: SignalMeter,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: timeRightNow,
|
||||
UpdatedAt: timeRightNow,
|
||||
},
|
||||
},
|
||||
{
|
||||
Identifiable: types.Identifiable{
|
||||
ID: valuer.GenerateUUID(),
|
||||
},
|
||||
OrgID: orgID,
|
||||
Filter: string(aiObservabilityJSON),
|
||||
Signal: SignalAiObservability,
|
||||
TimeAuditable: types.TimeAuditable{
|
||||
CreatedAt: timeRightNow,
|
||||
UpdatedAt: timeRightNow,
|
||||
},
|
||||
},
|
||||
}, nil
|
||||
return storableQuickFilters, nil
|
||||
}
|
||||
|
||||
func validateFilters(filters []telemetrytypes.TelemetryFieldKey) error {
|
||||
for _, filter := range filters {
|
||||
if filter.Name == "" {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "filter name is required")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -10,10 +10,10 @@ type QuickFilterStore interface {
|
||||
// Get retrieves all filters for an organization
|
||||
Get(ctx context.Context, orgID valuer.UUID) ([]*StorableQuickFilter, error)
|
||||
|
||||
// GetBySignal retrieves filters for a specific signal in an organization
|
||||
GetBySignal(ctx context.Context, orgID valuer.UUID, signal string) (*StorableQuickFilter, error)
|
||||
// GetBySource retrieves filters for a specific source in an organization
|
||||
GetBySource(ctx context.Context, orgID valuer.UUID, source string) (*StorableQuickFilter, error)
|
||||
|
||||
// Upsert inserts or updates filters for an organization and signal
|
||||
// Upsert inserts or updates filters for an organization and source
|
||||
Upsert(ctx context.Context, filter *StorableQuickFilter) error
|
||||
Create(ctx context.Context, filter []*StorableQuickFilter) error
|
||||
}
|
||||
|
||||
276
tests/integration/tests/quickfilter/01_quick_filter.py
Normal file
276
tests/integration/tests/quickfilter/01_quick_filter.py
Normal file
@@ -0,0 +1,276 @@
|
||||
from collections.abc import Callable
|
||||
from http import HTTPStatus
|
||||
|
||||
import requests
|
||||
from sqlalchemy import sql
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
|
||||
ALL_SOURCES = {
|
||||
"traces",
|
||||
"logs",
|
||||
"api_monitoring",
|
||||
"exceptions",
|
||||
"meter",
|
||||
"ai_observability",
|
||||
}
|
||||
|
||||
|
||||
def test_get_quick_filters_returns_defaults(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
data = response.json()["data"]
|
||||
assert {source_filters["source"] for source_filters in data} == ALL_SOURCES
|
||||
|
||||
for source_filters in data:
|
||||
assert source_filters["id"] != "00000000-0000-0000-0000-000000000000"
|
||||
assert source_filters["orgId"] != "00000000-0000-0000-0000-000000000000"
|
||||
assert source_filters["createdAt"] != ""
|
||||
assert source_filters["updatedAt"] != ""
|
||||
assert len(source_filters["filters"]) > 0
|
||||
for field_key in source_filters["filters"]:
|
||||
assert field_key["name"] != ""
|
||||
assert "fieldContext" in field_key
|
||||
assert "fieldDataType" in field_key
|
||||
assert "key" not in field_key
|
||||
|
||||
|
||||
def test_v1_get_serves_legacy_shape(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
data = response.json()["data"]
|
||||
assert {source_filters["signal"] for source_filters in data} == ALL_SOURCES
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters/traces"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
filters = response.json()["data"]["filters"]
|
||||
assert filters[0]["key"] == "duration_nano"
|
||||
assert filters[0]["type"] == "tag"
|
||||
assert filters[0]["dataType"] == "float64"
|
||||
assert all("name" not in legacy_filter for legacy_filter in filters)
|
||||
|
||||
|
||||
def test_v1_update_round_trips_to_v2(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters"),
|
||||
json={
|
||||
"signal": "exceptions",
|
||||
"filters": [
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "http.method", "dataType": "string", "type": "tag"},
|
||||
{"key": "code_line", "dataType": "int64", "type": "tag"},
|
||||
],
|
||||
},
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/exceptions"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
filters = response.json()["data"]["filters"]
|
||||
assert [(field_key["name"], field_key["fieldContext"]) for field_key in filters] == [
|
||||
("service.name", "resource"),
|
||||
("http.method", "attribute"),
|
||||
("code_line", "attribute"),
|
||||
]
|
||||
assert filters[2]["fieldDataType"] == "number"
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters/exceptions"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
filters = response.json()["data"]["filters"]
|
||||
assert [(legacy_filter["key"], legacy_filter["type"]) for legacy_filter in filters] == [
|
||||
("service.name", "resource"),
|
||||
("http.method", "tag"),
|
||||
("code_line", "tag"),
|
||||
]
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters"),
|
||||
json={
|
||||
"signal": "meter",
|
||||
"filters": [{"key": "host.name", "dataType": "string", "type": ""}],
|
||||
},
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/meter"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert [(field_key["name"], field_key["signal"]) for field_key in response.json()["data"]["filters"]] == [("host.name", "metrics")]
|
||||
|
||||
|
||||
def test_update_quick_filters_round_trip(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/logs"),
|
||||
json={
|
||||
"filters": [
|
||||
{
|
||||
"name": "k8s.pod.name",
|
||||
"fieldContext": "resource",
|
||||
"fieldDataType": "string",
|
||||
},
|
||||
{
|
||||
"name": "body.status",
|
||||
"fieldContext": "body",
|
||||
"fieldDataType": "string",
|
||||
},
|
||||
],
|
||||
},
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/logs"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
filters = response.json()["data"]["filters"]
|
||||
assert [field_key["name"] for field_key in filters] == [
|
||||
"k8s.pod.name",
|
||||
"body.status",
|
||||
]
|
||||
assert filters[0]["fieldContext"] == "resource"
|
||||
assert filters[1]["fieldContext"] == "body"
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/orgs/me/filters/logs"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert [(legacy_filter["key"], legacy_filter["type"]) for legacy_filter in response.json()["data"]["filters"]] == [
|
||||
("k8s.pod.name", "resource"),
|
||||
("body.status", ""),
|
||||
]
|
||||
|
||||
|
||||
def test_update_quick_filters_creates_row_for_source_without_one(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
with signoz.sqlstore.conn.connect() as conn:
|
||||
conn.execute(
|
||||
sql.text("DELETE FROM quick_filter WHERE source = :source"),
|
||||
{"source": "api_monitoring"},
|
||||
)
|
||||
conn.commit()
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/api_monitoring"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
data = response.json()["data"]
|
||||
assert data["source"] == "api_monitoring"
|
||||
assert data["filters"] == []
|
||||
assert data["id"] == "00000000-0000-0000-0000-000000000000"
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/api_monitoring"),
|
||||
json={
|
||||
"filters": [
|
||||
{
|
||||
"name": "service.name",
|
||||
"fieldContext": "resource",
|
||||
"fieldDataType": "string",
|
||||
},
|
||||
],
|
||||
},
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/quick_filters/api_monitoring"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
data = response.json()["data"]
|
||||
assert [field_key["name"] for field_key in data["filters"]] == ["service.name"]
|
||||
assert data["id"] != "00000000-0000-0000-0000-000000000000"
|
||||
|
||||
|
||||
def test_update_quick_filters_rejects_invalid_input(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
for source, invalid_body in [
|
||||
(
|
||||
"traces",
|
||||
{"filters": [{"key": "service.name", "dataType": "string", "type": "resource"}]},
|
||||
),
|
||||
("invalid", {"filters": []}),
|
||||
]:
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v2/quick_filters/{source}"),
|
||||
json=invalid_body,
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=2,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
|
||||
Reference in New Issue
Block a user