Compare commits

...

23 Commits

Author SHA1 Message Date
nityanandagohain
f537154df9 Merge remote-tracking branch 'origin/main' into issue_5947 2026-09-01 19:46:16 +05:30
Nityananda Gohain
7c50fe3763 fix: quick filters old migration cleanup (#12746)
#### Description
This PR makes sure that the old migrations of quick filters are
decoupled from types as they should and not imported.

<!--Reference issues using `Closes #issue-number` to enable automatic
closure on merge. -->
Part of https://github.com/SigNoz/engineering-pod/issues/5947
2026-09-01 13:54:36 +00:00
Nikhil Mantri
160a1b018c feat(alert-channel-integrations): jira + jsm ops channel backend (#12478)
#### Description

Adds two Atlassian alert channels. Backend only — frontend is #12488;
channels are created via the API.

**Jira issues — `jira_configs`**

- A firing alert creates a Jira Cloud issue; when the alert resolves,
the issue is transitioned to done. A re-fire within 3 days reopens the
same issue instead of creating a new one. 3 days is default but can be
edited via frontend form.
- The issue body is rich **Atlassian Document Format (ADF)**: a status
panel, the rendered alert description, and deep-links back to SigNoz.
- Re-fires keep the issue in sync (summary and description are
refreshed), and every notification after the first — re-fire, resolve,
reopen — also posts a **comment** carrying the same rich ADF snapshot,
so the issue holds a full lifecycle timeline.
- Per-rule custom notification templates (title/body) are honored, same
as every other channel; multi-alert custom bodies render as
divider-separated sections.
- Auth is Atlassian email + API token; Atlassian **service accounts**
also work (routed via the `api.atlassian.com` gateway automatically —
the cloud id is resolved server-side and client-supplied values are
ignored). Jira Cloud only.

**JSM Ops alerts — `jsmops_configs`**

- A firing alert opens a JSM Operations alert (the ex-Opsgenie alert
product); resolve **closes** it. A fire after close opens a fresh alert
— there is no reopen window.
- Re-fires dedupe into the same alert and increment its count. The alert
description keeps the first-fire snapshot; the value-over-time story
lives in the notes.
- Every fire and the resolve appends a **note** to the alert. JSM Ops
notes support **plain text only** (they render neither HTML nor
markdown), so notes use a new plain-text renderer with links flattened
to `text (url)`.
- The alert description supports JSM's **HTML subset**, rendered from
the same markdown templates.
- Auth is the JSM integration API key. No region/site config needed.

**Also in this PR**

- Unit tests for both config types, both notifiers, and the new ADF +
plain-text renderers.
- OpenAPI spec regenerated (adds `jsmops_configs`).

#### Issues closed by this PR

Closes SigNoz/pulse-pod#168 · Discussion: SigNoz/pulse-pod#169

#### Screenshots / Screen Recordings

Jira alert issue:

<img width="1171" height="739" alt="Screenshot 2026-08-18 at 12 49
31 PM"
src="https://github.com/user-attachments/assets/1ec74757-5d4d-4534-a5ca-13bed07cce72"
/>

Jira issue comments as a timeline:

<img width="1034" height="746" alt="Screenshot 2026-08-18 at 12 50
32 PM"
src="https://github.com/user-attachments/assets/09afb687-b62b-4eb3-86ab-39489e38513e"
/>

JSM Ops alerts page look: 

<img width="1317" height="460" alt="Screenshot 2026-08-18 at 12 52
29 PM"
src="https://github.com/user-attachments/assets/dd5a60b5-ea6c-49b7-afde-f3b2c72b178a"
/>

JSM Ops alerts main body + comment timeline ( comments only support
plain text today ) :

<img width="1323" height="784" alt="Screenshot 2026-08-18 at 12 53
10 PM"
src="https://github.com/user-attachments/assets/871fb91c-307a-4c67-b2c4-bc9548506cbc"
/>

#### Additional Information

Notes for reviewers:

- Jira shadows upstream Alertmanager's `jira_configs` so our notifier
handles it instead of upstream's; this needs a small dedupe in
`PostableChannel.JSONSchema()` and leaves every other channel type
untouched.
- JSM Ops reuses the existing Opsgenie notifier; all new behaviour sits
behind a single `advancedFeatures` flag, so plain Opsgenie is unchanged
when it's off.
- `send_resolved` defaults off for both channels, so resolve-time
behaviour (Jira transition, JSM close + resolved note) needs it set on
the channel; the frontend will send it on by default.
- Notes are best-effort: a permanently-failed note (e.g. the first-fire
note racing JSM's asynchronous alert create) is dropped with a warning
instead of failing the whole notification. Nothing is lost — that first
datapoint is already in the alert body; retryable failures (429) still
retry.

---------

Co-authored-by: Naman Verma <naman.verma@signoz.io>
2026-09-01 13:03:48 +00:00
nityanandagohain
1c6d966afa Merge remote-tracking branch 'origin/main' into issue_5947 2026-09-01 14:50:33 +05:30
nityanandagohain
f22a18d0a7 fix: more cleanup 2026-09-01 14:49:03 +05:30
nityanandagohain
752899a7a4 fix: add to registry resource 2026-09-01 01:00:51 +05:30
nityanandagohain
eb4c570d53 fix: address comments 2026-09-01 00:35:40 +05:30
nityanandagohain
7ce405aff7 fix: updatablequickfilters updated 2026-08-31 15:40:03 +05:30
nityanandagohain
72618d1d83 Merge remote-tracking branch 'origin/main' into issue_5947 2026-08-31 15:31:19 +05:30
nityanandagohain
80b7edc22a fix: address comments 2026-08-31 15:28:58 +05:30
nityanandagohain
45a8bb424c Merge remote-tracking branch 'origin/main' into issue_5947 2026-08-28 20:28:22 +05:30
nityanandagohain
17afa7a3cf fix: address comments 2026-08-28 20:25:14 +05:30
nityanandagohain
6dc9bf7b16 Merge remote-tracking branch 'origin/main' into issue_5947 2026-08-27 14:13:30 +05:30
nityanandagohain
d15a1452f8 fix: minor fixes 2026-08-27 14:13:08 +05:30
nityanandagohain
2fe033ec78 fix: trigger build 2026-08-26 18:20:23 +05:30
nityanandagohain
02a2800f89 Merge remote-tracking branch 'origin/main' into issue_5947 2026-08-26 18:20:04 +05:30
nityanandagohain
915aa2eb70 Merge remote-tracking branch 'origin/issue_5947' into issue_5947 2026-08-26 16:44:26 +05:30
nityanandagohain
518caff0f2 Merge remote-tracking branch 'origin/main' into issue_5947 2026-08-26 16:44:09 +05:30
nityanandagohain
0d3ac28286 fix: update openapi spec 2026-08-26 16:41:50 +05:30
Nityananda Gohain
47beef07de Merge branch 'main' into issue_5947 2026-08-26 16:32:51 +05:30
nityanandagohain
5f2891bd6c fix: test cleanup 2026-08-26 16:32:01 +05:30
nityanandagohain
97f0e832ab fix: minor cleanup 2026-08-26 16:31:17 +05:30
nityanandagohain
770a8f7b0a fix: add quick filters v2 api to support TelemetryFieldKey 2026-08-26 11:12:59 +05:30
42 changed files with 4815 additions and 478 deletions

View File

@@ -50,6 +50,7 @@ jobs:
- logspipelines
- passwordauthn
- preference
- quickfilter
- querierlogs
- queriertraces
- queriermetrics

View File

@@ -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

View 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));
};

View File

@@ -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;
/**

View 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: &notify.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])
}

View 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)
})
}
}

View 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
}

View 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)
}

View File

@@ -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: &notify.Retrier{RetryCodes: []int{http.StatusTooManyRequests}},
templater: templater,
conf: c,
tmpl: t,
logger: l,
client: client,
retrier: &notify.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

View File

@@ -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])), &noteMsg))
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)

View File

@@ -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

View File

@@ -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
}

View 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
}

View File

@@ -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,
}
}

View 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])
}

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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)
}

View File

@@ -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)

View File

@@ -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

View File

@@ -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,
),
)
}

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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
}

View 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))]
}

View 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
}

View 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()
}

View 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(""))
}

View File

@@ -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 {

View 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())
}

View 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))
})
}
}

View File

@@ -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())
}

View File

@@ -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
}
}
}
}

View 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, "/"))
}

View 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)
})
}
}

View 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
}

View 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)
})
}
}

View File

@@ -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
}

View File

@@ -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)

View File

@@ -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
}

View File

@@ -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
}

View 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