mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-30 23:30:40 +01:00
Compare commits
8 Commits
feat/alert
...
feat/ai-tr
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
61f6370f45 | ||
|
|
f8b22c0feb | ||
|
|
9e3a6bb35d | ||
|
|
47dd1fabf3 | ||
|
|
270988fb48 | ||
|
|
f5f019f61b | ||
|
|
39badeb591 | ||
|
|
ec05bfe755 |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -68,6 +68,7 @@ jobs:
|
||||
- semconvfamilies
|
||||
- serviceaccount
|
||||
- spanmapper
|
||||
- tracedetail
|
||||
- querier_json_body
|
||||
- querier_skip_resource_fingerprint
|
||||
- ttl
|
||||
|
||||
@@ -8928,30 +8928,6 @@ components:
|
||||
- kind
|
||||
- spec
|
||||
type: object
|
||||
RuletypesGettableRuleView:
|
||||
properties:
|
||||
createdAt:
|
||||
format: date-time
|
||||
type: string
|
||||
data:
|
||||
$ref: '#/components/schemas/RuletypesRuleViewData'
|
||||
id:
|
||||
type: string
|
||||
name:
|
||||
type: string
|
||||
orgId:
|
||||
type: string
|
||||
updatedAt:
|
||||
format: date-time
|
||||
type: string
|
||||
required:
|
||||
- id
|
||||
- name
|
||||
- data
|
||||
- orgId
|
||||
- createdAt
|
||||
- updatedAt
|
||||
type: object
|
||||
RuletypesGettableTestRule:
|
||||
properties:
|
||||
alertCount:
|
||||
@@ -9019,15 +8995,6 @@ components:
|
||||
- alertType
|
||||
- ruleType
|
||||
type: object
|
||||
RuletypesListableRuleViews:
|
||||
properties:
|
||||
views:
|
||||
items:
|
||||
$ref: '#/components/schemas/RuletypesGettableRuleView'
|
||||
type: array
|
||||
required:
|
||||
- views
|
||||
type: object
|
||||
RuletypesListableRules:
|
||||
properties:
|
||||
labels:
|
||||
@@ -9125,16 +9092,6 @@ components:
|
||||
- ruleType
|
||||
- condition
|
||||
type: object
|
||||
RuletypesPostableRuleView:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/RuletypesRuleViewData'
|
||||
name:
|
||||
type: string
|
||||
required:
|
||||
- name
|
||||
- data
|
||||
type: object
|
||||
RuletypesQueryType:
|
||||
enum:
|
||||
- builder
|
||||
@@ -9273,23 +9230,6 @@ components:
|
||||
- promql_rule
|
||||
- anomaly_rule
|
||||
type: string
|
||||
RuletypesRuleViewData:
|
||||
properties:
|
||||
order:
|
||||
$ref: '#/components/schemas/RuletypesListOrder'
|
||||
query:
|
||||
type: string
|
||||
sort:
|
||||
$ref: '#/components/schemas/RuletypesListSort'
|
||||
states:
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
version:
|
||||
type: string
|
||||
required:
|
||||
- version
|
||||
type: object
|
||||
RuletypesScheduleType:
|
||||
enum:
|
||||
- hourly
|
||||
@@ -9707,6 +9647,17 @@ components:
|
||||
required:
|
||||
- aggregations
|
||||
type: object
|
||||
SpantypesGettableTraceThread:
|
||||
properties:
|
||||
nextCursor:
|
||||
type: string
|
||||
spans:
|
||||
items:
|
||||
$ref: '#/components/schemas/SpantypesThreadSpan'
|
||||
type: array
|
||||
required:
|
||||
- spans
|
||||
type: object
|
||||
SpantypesGettableWaterfallTrace:
|
||||
properties:
|
||||
endTimestampMillis:
|
||||
@@ -10021,6 +9972,84 @@ components:
|
||||
nullable: true
|
||||
type: object
|
||||
type: object
|
||||
SpantypesThreadSpan:
|
||||
properties:
|
||||
attributes:
|
||||
additionalProperties: {}
|
||||
nullable: true
|
||||
type: object
|
||||
db_name:
|
||||
type: string
|
||||
db_operation:
|
||||
type: string
|
||||
duration_nano:
|
||||
minimum: 0
|
||||
type: integer
|
||||
events:
|
||||
items:
|
||||
$ref: '#/components/schemas/SpantypesEvent'
|
||||
nullable: true
|
||||
type: array
|
||||
external_http_method:
|
||||
type: string
|
||||
external_http_url:
|
||||
type: string
|
||||
flags:
|
||||
minimum: 0
|
||||
type: integer
|
||||
has_children:
|
||||
type: boolean
|
||||
has_error:
|
||||
type: boolean
|
||||
http_host:
|
||||
type: string
|
||||
http_method:
|
||||
type: string
|
||||
http_url:
|
||||
type: string
|
||||
is_remote:
|
||||
type: string
|
||||
kind_string:
|
||||
type: string
|
||||
level:
|
||||
minimum: 0
|
||||
type: integer
|
||||
name:
|
||||
type: string
|
||||
parent_span_id:
|
||||
type: string
|
||||
references:
|
||||
items:
|
||||
$ref: '#/components/schemas/SpantypesOtelSpanRef'
|
||||
type: array
|
||||
resource:
|
||||
additionalProperties:
|
||||
type: string
|
||||
nullable: true
|
||||
type: object
|
||||
response_status_code:
|
||||
type: string
|
||||
span_id:
|
||||
type: string
|
||||
status_code:
|
||||
type: integer
|
||||
status_code_string:
|
||||
type: string
|
||||
status_message:
|
||||
type: string
|
||||
sub_tree_node_count:
|
||||
minimum: 0
|
||||
type: integer
|
||||
time_unix:
|
||||
minimum: 0
|
||||
type: integer
|
||||
trace_id:
|
||||
type: string
|
||||
trace_state:
|
||||
type: string
|
||||
required:
|
||||
- references
|
||||
type: object
|
||||
SpantypesUpdatableSpanMapper:
|
||||
properties:
|
||||
config:
|
||||
@@ -15745,6 +15774,80 @@ paths:
|
||||
tags:
|
||||
- tracedetail
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/traces/{traceID}/thread:
|
||||
get:
|
||||
deprecated: false
|
||||
description: Returns the spans carrying gen_ai input or output messages in timestamp
|
||||
order. Pages are fetched with the returned nextCursor.
|
||||
operationId: GetTraceThread
|
||||
parameters:
|
||||
- in: query
|
||||
name: limit
|
||||
schema:
|
||||
type: integer
|
||||
- in: query
|
||||
name: cursor
|
||||
schema:
|
||||
type: string
|
||||
- in: path
|
||||
name: traceID
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/SpantypesGettableTraceThread'
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"404":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Not Found
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: Get thread view for a trace
|
||||
tags:
|
||||
- tracedetail
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/user/me:
|
||||
get:
|
||||
deprecated: true
|
||||
@@ -21207,235 +21310,6 @@ paths:
|
||||
tags:
|
||||
- users
|
||||
x-signoz-stability: alpha
|
||||
/api/v2/rule_views:
|
||||
get:
|
||||
deprecated: false
|
||||
description: Returns every saved view in the calling user's org. Saved views
|
||||
are shared org-wide.
|
||||
operationId: ListRuleViews
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/RuletypesListableRuleViews'
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"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:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: List rule saved views
|
||||
tags:
|
||||
- rules
|
||||
x-signoz-stability: alpha
|
||||
post:
|
||||
deprecated: false
|
||||
description: Persists the calling user's rule listing state (query, states,
|
||||
sort, order) as a named, reusable view shared across the org.
|
||||
operationId: CreateRuleView
|
||||
requestBody:
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RuletypesPostableRuleView'
|
||||
responses:
|
||||
"201":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/RuletypesGettableRuleView'
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: Created
|
||||
"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:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: Create rule saved view
|
||||
tags:
|
||||
- rules
|
||||
x-signoz-stability: alpha
|
||||
/api/v2/rule_views/{id}:
|
||||
delete:
|
||||
deprecated: false
|
||||
description: Removes a saved view. Saved views are shared org-wide. Deleting
|
||||
a non-existent view returns 404.
|
||||
operationId: DeleteRuleView
|
||||
parameters:
|
||||
- in: path
|
||||
name: id
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
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
|
||||
"404":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Not Found
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: Delete rule saved view
|
||||
tags:
|
||||
- rules
|
||||
x-signoz-stability: alpha
|
||||
put:
|
||||
deprecated: false
|
||||
description: Replaces a saved view's name and data. Saved views are shared org-wide.
|
||||
operationId: UpdateRuleView
|
||||
parameters:
|
||||
- in: path
|
||||
name: id
|
||||
required: true
|
||||
schema:
|
||||
type: string
|
||||
requestBody:
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RuletypesPostableRuleView'
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
$ref: '#/components/schemas/RuletypesGettableRuleView'
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"404":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Not Found
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: Update rule saved view
|
||||
tags:
|
||||
- rules
|
||||
x-signoz-stability: alpha
|
||||
/api/v2/rules:
|
||||
get:
|
||||
deprecated: true
|
||||
|
||||
@@ -19,9 +19,7 @@ import type {
|
||||
|
||||
import type {
|
||||
CreateRule201,
|
||||
CreateRuleView201,
|
||||
DeleteRuleByIDPathParameters,
|
||||
DeleteRuleViewPathParameters,
|
||||
GetRuleByID200,
|
||||
GetRuleByIDPathParameters,
|
||||
GetRuleHistoryFilterKeys200,
|
||||
@@ -42,7 +40,6 @@ import type {
|
||||
GetRuleHistoryTopContributors200,
|
||||
GetRuleHistoryTopContributorsParams,
|
||||
GetRuleHistoryTopContributorsPathParameters,
|
||||
ListRuleViews200,
|
||||
ListRules200,
|
||||
ListRulesV3200,
|
||||
ListRulesV3Params,
|
||||
@@ -50,11 +47,8 @@ import type {
|
||||
PatchRuleByIDPathParameters,
|
||||
RenderErrorResponseDTO,
|
||||
RuletypesPostableRuleDTO,
|
||||
RuletypesPostableRuleViewDTO,
|
||||
TestRule200,
|
||||
UpdateRuleByIDPathParameters,
|
||||
UpdateRuleView200,
|
||||
UpdateRuleViewPathParameters,
|
||||
} from '../sigNoz.schemas';
|
||||
|
||||
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
|
||||
@@ -80,351 +74,6 @@ const withQueryKey = <T extends object, K>(
|
||||
return result;
|
||||
};
|
||||
|
||||
/**
|
||||
* Returns every saved view in the calling user's org. Saved views are shared org-wide.
|
||||
* @summary List rule saved views
|
||||
*/
|
||||
export const listRuleViews = (signal?: AbortSignal) => {
|
||||
return GeneratedAPIInstance<ListRuleViews200>({
|
||||
url: `/api/v2/rule_views`,
|
||||
method: 'GET',
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getListRuleViewsQueryKey = () => {
|
||||
return [`/api/v2/rule_views`] as const;
|
||||
};
|
||||
|
||||
export const getListRuleViewsQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof listRuleViews>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listRuleViews>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
}) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey = queryOptions?.queryKey ?? getListRuleViewsQueryKey();
|
||||
|
||||
const queryFn: QueryFunction<Awaited<ReturnType<typeof listRuleViews>>> = ({
|
||||
signal,
|
||||
}) => listRuleViews(signal);
|
||||
|
||||
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listRuleViews>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type ListRuleViewsQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof listRuleViews>>
|
||||
>;
|
||||
export type ListRuleViewsQueryError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary List rule saved views
|
||||
*/
|
||||
|
||||
export function useListRuleViews<
|
||||
TData = Awaited<ReturnType<typeof listRuleViews>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listRuleViews>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
}): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getListRuleViewsQueryOptions(options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return withQueryKey(query, queryOptions.queryKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary List rule saved views
|
||||
*/
|
||||
export const invalidateListRuleViews = async (
|
||||
queryClient: QueryClient,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getListRuleViewsQueryKey() },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* Persists the calling user's rule listing state (query, states, sort, order) as a named, reusable view shared across the org.
|
||||
* @summary Create rule saved view
|
||||
*/
|
||||
export const createRuleView = (
|
||||
ruletypesPostableRuleViewDTO?: BodyType<RuletypesPostableRuleViewDTO>,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<CreateRuleView201>({
|
||||
url: `/api/v2/rule_views`,
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
data: ruletypesPostableRuleViewDTO,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getCreateRuleViewMutationOptions = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof createRuleView>>,
|
||||
TError,
|
||||
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof createRuleView>>,
|
||||
TError,
|
||||
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
|
||||
TContext
|
||||
> => {
|
||||
const mutationKey = ['createRuleView'];
|
||||
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 createRuleView>>,
|
||||
{ data?: BodyType<RuletypesPostableRuleViewDTO> }
|
||||
> = (props) => {
|
||||
const { data } = props ?? {};
|
||||
|
||||
return createRuleView(data);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type CreateRuleViewMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof createRuleView>>
|
||||
>;
|
||||
export type CreateRuleViewMutationBody =
|
||||
| BodyType<RuletypesPostableRuleViewDTO>
|
||||
| undefined;
|
||||
export type CreateRuleViewMutationError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Create rule saved view
|
||||
*/
|
||||
export const useCreateRuleView = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof createRuleView>>,
|
||||
TError,
|
||||
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof createRuleView>>,
|
||||
TError,
|
||||
{ data?: BodyType<RuletypesPostableRuleViewDTO> },
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getCreateRuleViewMutationOptions(options));
|
||||
};
|
||||
/**
|
||||
* Removes a saved view. Saved views are shared org-wide. Deleting a non-existent view returns 404.
|
||||
* @summary Delete rule saved view
|
||||
*/
|
||||
export const deleteRuleView = (
|
||||
{ id }: DeleteRuleViewPathParameters,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<void>({
|
||||
url: `/api/v2/rule_views/${id}`,
|
||||
method: 'DELETE',
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getDeleteRuleViewMutationOptions = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof deleteRuleView>>,
|
||||
TError,
|
||||
{ pathParams: DeleteRuleViewPathParameters },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof deleteRuleView>>,
|
||||
TError,
|
||||
{ pathParams: DeleteRuleViewPathParameters },
|
||||
TContext
|
||||
> => {
|
||||
const mutationKey = ['deleteRuleView'];
|
||||
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 deleteRuleView>>,
|
||||
{ pathParams: DeleteRuleViewPathParameters }
|
||||
> = (props) => {
|
||||
const { pathParams } = props ?? {};
|
||||
|
||||
return deleteRuleView(pathParams);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type DeleteRuleViewMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof deleteRuleView>>
|
||||
>;
|
||||
|
||||
export type DeleteRuleViewMutationError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Delete rule saved view
|
||||
*/
|
||||
export const useDeleteRuleView = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof deleteRuleView>>,
|
||||
TError,
|
||||
{ pathParams: DeleteRuleViewPathParameters },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof deleteRuleView>>,
|
||||
TError,
|
||||
{ pathParams: DeleteRuleViewPathParameters },
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getDeleteRuleViewMutationOptions(options));
|
||||
};
|
||||
/**
|
||||
* Replaces a saved view's name and data. Saved views are shared org-wide.
|
||||
* @summary Update rule saved view
|
||||
*/
|
||||
export const updateRuleView = (
|
||||
{ id }: UpdateRuleViewPathParameters,
|
||||
ruletypesPostableRuleViewDTO?: BodyType<RuletypesPostableRuleViewDTO>,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<UpdateRuleView200>({
|
||||
url: `/api/v2/rule_views/${id}`,
|
||||
method: 'PUT',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
data: ruletypesPostableRuleViewDTO,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getUpdateRuleViewMutationOptions = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof updateRuleView>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateRuleViewPathParameters;
|
||||
data?: BodyType<RuletypesPostableRuleViewDTO>;
|
||||
},
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof updateRuleView>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateRuleViewPathParameters;
|
||||
data?: BodyType<RuletypesPostableRuleViewDTO>;
|
||||
},
|
||||
TContext
|
||||
> => {
|
||||
const mutationKey = ['updateRuleView'];
|
||||
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 updateRuleView>>,
|
||||
{
|
||||
pathParams: UpdateRuleViewPathParameters;
|
||||
data?: BodyType<RuletypesPostableRuleViewDTO>;
|
||||
}
|
||||
> = (props) => {
|
||||
const { pathParams, data } = props ?? {};
|
||||
|
||||
return updateRuleView(pathParams, data);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type UpdateRuleViewMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof updateRuleView>>
|
||||
>;
|
||||
export type UpdateRuleViewMutationBody =
|
||||
| BodyType<RuletypesPostableRuleViewDTO>
|
||||
| undefined;
|
||||
export type UpdateRuleViewMutationError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Update rule saved view
|
||||
*/
|
||||
export const useUpdateRuleView = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof updateRuleView>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateRuleViewPathParameters;
|
||||
data?: BodyType<RuletypesPostableRuleViewDTO>;
|
||||
},
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof updateRuleView>>,
|
||||
TError,
|
||||
{
|
||||
pathParams: UpdateRuleViewPathParameters;
|
||||
data?: BodyType<RuletypesPostableRuleViewDTO>;
|
||||
},
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getUpdateRuleViewMutationOptions(options));
|
||||
};
|
||||
/**
|
||||
* This endpoint lists all alert rules with their current evaluation state. Deprecated: use ListRulesV3, which supports filtering, sorting and pagination.
|
||||
* @deprecated
|
||||
|
||||
@@ -10177,60 +10177,6 @@ export enum RuletypesEvaluationKindDTO {
|
||||
rolling = 'rolling',
|
||||
cumulative = 'cumulative',
|
||||
}
|
||||
export enum RuletypesListOrderDTO {
|
||||
asc = 'asc',
|
||||
desc = 'desc',
|
||||
}
|
||||
export enum RuletypesListSortDTO {
|
||||
updated_at = 'updated_at',
|
||||
created_at = 'created_at',
|
||||
name = 'name',
|
||||
state = 'state',
|
||||
severity = 'severity',
|
||||
}
|
||||
export interface RuletypesRuleViewDataDTO {
|
||||
order?: RuletypesListOrderDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
query?: string;
|
||||
sort?: RuletypesListSortDTO;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
states?: string[];
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
version: string;
|
||||
}
|
||||
|
||||
export interface RuletypesGettableRuleViewDTO {
|
||||
/**
|
||||
* @type string
|
||||
* @format date-time
|
||||
*/
|
||||
createdAt: string;
|
||||
data: RuletypesRuleViewDataDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
id: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
name: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
orgId: string;
|
||||
/**
|
||||
* @type string
|
||||
* @format date-time
|
||||
*/
|
||||
updatedAt: string;
|
||||
}
|
||||
|
||||
export interface RuletypesGettableTestRuleDTO {
|
||||
/**
|
||||
* @type integer
|
||||
@@ -10253,6 +10199,17 @@ export interface RuletypesLabelPairDTO {
|
||||
value: string;
|
||||
}
|
||||
|
||||
export enum RuletypesListOrderDTO {
|
||||
asc = 'asc',
|
||||
desc = 'desc',
|
||||
}
|
||||
export enum RuletypesListSortDTO {
|
||||
updated_at = 'updated_at',
|
||||
created_at = 'created_at',
|
||||
name = 'name',
|
||||
state = 'state',
|
||||
severity = 'severity',
|
||||
}
|
||||
export type RuletypesListableRuleDTOLabels = { [key: string]: string };
|
||||
|
||||
export enum RuletypesRuleTypeDTO {
|
||||
@@ -10304,13 +10261,6 @@ export interface RuletypesListableRuleDTO {
|
||||
updatedBy?: string;
|
||||
}
|
||||
|
||||
export interface RuletypesListableRuleViewsDTO {
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
views: RuletypesGettableRuleViewDTO[];
|
||||
}
|
||||
|
||||
export interface RuletypesListableRulesDTO {
|
||||
/**
|
||||
* @type array
|
||||
@@ -10479,14 +10429,6 @@ export interface RuletypesPostableRuleDTO {
|
||||
version?: string;
|
||||
}
|
||||
|
||||
export interface RuletypesPostableRuleViewDTO {
|
||||
data: RuletypesRuleViewDataDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
name: string;
|
||||
}
|
||||
|
||||
export type RuletypesRuleDTOAnnotations = { [key: string]: string };
|
||||
|
||||
export type RuletypesRuleDTOLabels = { [key: string]: string };
|
||||
@@ -11200,6 +11142,157 @@ export interface SpantypesOtelSpanRefDTO {
|
||||
traceId?: string;
|
||||
}
|
||||
|
||||
export type SpantypesThreadSpanDTOAttributesAnyOf = { [key: string]: unknown };
|
||||
|
||||
/**
|
||||
* @nullable
|
||||
*/
|
||||
export type SpantypesThreadSpanDTOAttributes =
|
||||
SpantypesThreadSpanDTOAttributesAnyOf | null;
|
||||
|
||||
export type SpantypesThreadSpanDTOResourceAnyOf = { [key: string]: string };
|
||||
|
||||
/**
|
||||
* @nullable
|
||||
*/
|
||||
export type SpantypesThreadSpanDTOResource =
|
||||
SpantypesThreadSpanDTOResourceAnyOf | null;
|
||||
|
||||
export interface SpantypesThreadSpanDTO {
|
||||
/**
|
||||
* @type object,null
|
||||
*/
|
||||
attributes?: SpantypesThreadSpanDTOAttributes;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
db_name?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
db_operation?: string;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
duration_nano?: number;
|
||||
/**
|
||||
* @type array,null
|
||||
*/
|
||||
events?: SpantypesEventDTO[] | null;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
external_http_method?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
external_http_url?: string;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
flags?: number;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
has_children?: boolean;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
has_error?: boolean;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
http_host?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
http_method?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
http_url?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
is_remote?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
kind_string?: string;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
level?: number;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
name?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
parent_span_id?: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
references: SpantypesOtelSpanRefDTO[];
|
||||
/**
|
||||
* @type object,null
|
||||
*/
|
||||
resource?: SpantypesThreadSpanDTOResource;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
response_status_code?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
span_id?: string;
|
||||
/**
|
||||
* @type integer
|
||||
*/
|
||||
status_code?: number;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status_code_string?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status_message?: string;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
sub_tree_node_count?: number;
|
||||
/**
|
||||
* @type integer
|
||||
* @minimum 0
|
||||
*/
|
||||
time_unix?: number;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
trace_id?: string;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
trace_state?: string;
|
||||
}
|
||||
|
||||
export interface SpantypesGettableTraceThreadDTO {
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
nextCursor?: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
spans: SpantypesThreadSpanDTO[];
|
||||
}
|
||||
|
||||
export type SpantypesWaterfallSpanDTOAttributesAnyOf = {
|
||||
[key: string]: unknown;
|
||||
};
|
||||
@@ -12873,6 +12966,30 @@ export type GetTraceAggregations200 = {
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type GetTraceThreadPathParameters = {
|
||||
traceID: string;
|
||||
};
|
||||
export type GetTraceThreadParams = {
|
||||
/**
|
||||
* @type integer
|
||||
* @description undefined
|
||||
*/
|
||||
limit?: number;
|
||||
/**
|
||||
* @type string
|
||||
* @description undefined
|
||||
*/
|
||||
cursor?: string;
|
||||
};
|
||||
|
||||
export type GetTraceThread200 = {
|
||||
data: SpantypesGettableTraceThreadDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListUserPreferences200 = {
|
||||
/**
|
||||
* @type array
|
||||
@@ -13774,36 +13891,6 @@ export type GetUsersByRoleID200 = {
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListRuleViews200 = {
|
||||
data: RuletypesListableRuleViewsDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type CreateRuleView201 = {
|
||||
data: RuletypesGettableRuleViewDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type DeleteRuleViewPathParameters = {
|
||||
id: string;
|
||||
};
|
||||
export type UpdateRuleViewPathParameters = {
|
||||
id: string;
|
||||
};
|
||||
export type UpdateRuleView200 = {
|
||||
data: RuletypesGettableRuleViewDTO;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListRules200 = {
|
||||
/**
|
||||
* @type array
|
||||
|
||||
@@ -4,11 +4,17 @@
|
||||
* * regenerate with 'pnpm generate:api'
|
||||
* SigNoz
|
||||
*/
|
||||
import { useMutation } from 'react-query';
|
||||
import { useMutation, useQuery } from 'react-query';
|
||||
import type {
|
||||
InvalidateOptions,
|
||||
MutationFunction,
|
||||
QueryClient,
|
||||
QueryFunction,
|
||||
QueryKey,
|
||||
UseMutationOptions,
|
||||
UseMutationResult,
|
||||
UseQueryOptions,
|
||||
UseQueryResult,
|
||||
} from 'react-query';
|
||||
|
||||
import type {
|
||||
@@ -16,6 +22,9 @@ import type {
|
||||
GetFlamegraphPathParameters,
|
||||
GetTraceAggregations200,
|
||||
GetTraceAggregationsPathParameters,
|
||||
GetTraceThread200,
|
||||
GetTraceThreadParams,
|
||||
GetTraceThreadPathParameters,
|
||||
GetWaterfallV4200,
|
||||
GetWaterfallV4PathParameters,
|
||||
RenderErrorResponseDTO,
|
||||
@@ -27,6 +36,26 @@ import type {
|
||||
import { GeneratedAPIInstance } from '../../../generatedAPIInstance';
|
||||
import type { ErrorType, BodyType } from '../../../generatedAPIInstance';
|
||||
|
||||
const withQueryKey = <T extends object, K>(
|
||||
query: T,
|
||||
queryKey: K,
|
||||
): T & { queryKey: K } => {
|
||||
const result = { queryKey } as T & { queryKey: K };
|
||||
for (const key of Object.keys(query)) {
|
||||
// The explicit queryKey always wins, matching the previous
|
||||
// `{ ...query, queryKey }` spread where it was set last.
|
||||
if (key === 'queryKey') {
|
||||
continue;
|
||||
}
|
||||
Object.defineProperty(result, key, {
|
||||
enumerable: true,
|
||||
configurable: true,
|
||||
get: () => (query as Record<string, unknown>)[key],
|
||||
});
|
||||
}
|
||||
return result;
|
||||
};
|
||||
|
||||
/**
|
||||
* Computes span aggregations grouped by requested field.
|
||||
* @summary Get aggregations for a trace
|
||||
@@ -127,6 +156,121 @@ export const useGetTraceAggregations = <
|
||||
> => {
|
||||
return useMutation(getGetTraceAggregationsMutationOptions(options));
|
||||
};
|
||||
/**
|
||||
* Returns the spans carrying gen_ai input or output messages in timestamp order. Pages are fetched with the returned nextCursor.
|
||||
* @summary Get thread view for a trace
|
||||
*/
|
||||
export const getTraceThread = (
|
||||
{ traceID }: GetTraceThreadPathParameters,
|
||||
params?: GetTraceThreadParams,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<GetTraceThread200>({
|
||||
url: `/api/v1/traces/${traceID}/thread`,
|
||||
method: 'GET',
|
||||
params,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getGetTraceThreadQueryKey = (
|
||||
{ traceID }: GetTraceThreadPathParameters,
|
||||
params?: GetTraceThreadParams,
|
||||
) => {
|
||||
return [
|
||||
`/api/v1/traces/${traceID}/thread`,
|
||||
...(params ? [params] : []),
|
||||
] as const;
|
||||
};
|
||||
|
||||
export const getGetTraceThreadQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof getTraceThread>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
{ traceID }: GetTraceThreadPathParameters,
|
||||
params?: GetTraceThreadParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getTraceThread>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey =
|
||||
queryOptions?.queryKey ?? getGetTraceThreadQueryKey({ traceID }, params);
|
||||
|
||||
const queryFn: QueryFunction<Awaited<ReturnType<typeof getTraceThread>>> = ({
|
||||
signal,
|
||||
}) => getTraceThread({ traceID }, params, signal);
|
||||
|
||||
return {
|
||||
queryKey,
|
||||
queryFn,
|
||||
enabled: traceID !== null && traceID !== undefined,
|
||||
...queryOptions,
|
||||
} as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getTraceThread>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type GetTraceThreadQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof getTraceThread>>
|
||||
>;
|
||||
export type GetTraceThreadQueryError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Get thread view for a trace
|
||||
*/
|
||||
|
||||
export function useGetTraceThread<
|
||||
TData = Awaited<ReturnType<typeof getTraceThread>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
{ traceID }: GetTraceThreadPathParameters,
|
||||
params?: GetTraceThreadParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof getTraceThread>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getGetTraceThreadQueryOptions(
|
||||
{ traceID },
|
||||
params,
|
||||
options,
|
||||
);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return withQueryKey(query, queryOptions.queryKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary Get thread view for a trace
|
||||
*/
|
||||
export const invalidateGetTraceThread = async (
|
||||
queryClient: QueryClient,
|
||||
{ traceID }: GetTraceThreadPathParameters,
|
||||
params?: GetTraceThreadParams,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getGetTraceThreadQueryKey({ traceID }, params) },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* Returns the flamegraph view of spans for a given trace ID.
|
||||
* @summary Get flamegraph view for a trace
|
||||
|
||||
@@ -132,64 +132,6 @@ func (provider *provider) addRulerRoutes(router *mux.Router) error {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v2/rule_views", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.ListRuleViews), handler.OpenAPIDef{
|
||||
ID: "ListRuleViews",
|
||||
Tags: []string{"rules"},
|
||||
Summary: "List rule saved views",
|
||||
Description: "Returns every saved view in the calling user's org. Saved views are shared org-wide.",
|
||||
Response: new(ruletypes.ListableRuleViews),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusOK,
|
||||
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
|
||||
})).Methods(http.MethodGet).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v2/rule_views", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.CreateRuleView), handler.OpenAPIDef{
|
||||
ID: "CreateRuleView",
|
||||
Tags: []string{"rules"},
|
||||
Summary: "Create rule saved view",
|
||||
Description: "Persists the calling user's rule listing state (query, states, sort, order) as a named, reusable view shared across the org.",
|
||||
Request: new(ruletypes.PostableRuleView),
|
||||
RequestContentType: "application/json",
|
||||
Response: new(ruletypes.GettableRuleView),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusCreated,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest},
|
||||
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
|
||||
})).Methods(http.MethodPost).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v2/rule_views/{id}", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.UpdateRuleView), handler.OpenAPIDef{
|
||||
ID: "UpdateRuleView",
|
||||
Tags: []string{"rules"},
|
||||
Summary: "Update rule saved view",
|
||||
Description: "Replaces a saved view's name and data. Saved views are shared org-wide.",
|
||||
Request: new(ruletypes.UpdatableRuleView),
|
||||
RequestContentType: "application/json",
|
||||
Response: new(ruletypes.GettableRuleView),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusOK,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
|
||||
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
|
||||
})).Methods(http.MethodPut).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v2/rule_views/{id}", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.DeleteRuleView), handler.OpenAPIDef{
|
||||
ID: "DeleteRuleView",
|
||||
Tags: []string{"rules"},
|
||||
Summary: "Delete rule saved view",
|
||||
Description: "Removes a saved view. Saved views are shared org-wide. Deleting a non-existent view returns 404.",
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusNoContent,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
|
||||
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
|
||||
})).Methods(http.MethodDelete).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v1/downtime_schedules", handler.New(provider.authzMiddleware.ViewAccess(provider.rulerHandler.ListDowntimeSchedules), handler.OpenAPIDef{
|
||||
ID: "ListDowntimeSchedules",
|
||||
Tags: []string{"downtimeschedules"},
|
||||
|
||||
@@ -67,5 +67,23 @@ func (provider *provider) addTraceDetailRoutes(router *mux.Router) error {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v1/traces/{traceID}/thread", handler.New(
|
||||
provider.authzMiddleware.ViewAccess(provider.traceDetailHandler.GetThread),
|
||||
handler.OpenAPIDef{
|
||||
ID: "GetTraceThread",
|
||||
Tags: []string{"tracedetail"},
|
||||
Summary: "Get thread view for a trace",
|
||||
Description: "Returns the spans carrying gen_ai input or output messages in timestamp order. Pages are fetched with the returned nextCursor.",
|
||||
RequestQuery: new(spantypes.QueryableThread),
|
||||
Response: new(spantypes.GettableTraceThread),
|
||||
ResponseContentType: "application/json",
|
||||
SuccessStatusCode: http.StatusOK,
|
||||
ErrorStatusCodes: []int{http.StatusBadRequest, http.StatusNotFound},
|
||||
SecuritySchemes: newSecuritySchemes(types.RoleViewer),
|
||||
},
|
||||
)).Methods(http.MethodGet).GetError(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -75,3 +75,25 @@ func (h *handler) GetFlamegraph(rw http.ResponseWriter, r *http.Request) {
|
||||
|
||||
render.Success(rw, http.StatusOK, result)
|
||||
}
|
||||
|
||||
func (h *handler) GetThread(rw http.ResponseWriter, r *http.Request) {
|
||||
req := new(spantypes.QueryableThread)
|
||||
if err := binding.Query.BindQuery(r.URL.Query(), req); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
query, err := spantypes.NewThreadQuery(req)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
result, err := h.module.GetThread(r.Context(), mux.Vars(r)["traceID"], query)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, result)
|
||||
}
|
||||
|
||||
@@ -173,6 +173,19 @@ func (m *module) getWindowedWaterfall(ctx context.Context, traceID, selectedSpan
|
||||
), nil
|
||||
}
|
||||
|
||||
func (m *module) GetThread(ctx context.Context, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error) {
|
||||
summary, err := m.store.GetTraceSummary(ctx, traceID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
spans, err := m.store.GetThreadSpans(ctx, traceID, summary, query.Cursor, query.Limit+1)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return spantypes.NewGettableTraceThread(traceID, spans, query.Limit), nil
|
||||
}
|
||||
|
||||
func (m *module) getFullFlamegraph(ctx context.Context, traceID string, summary *spantypes.TraceSummary, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error) {
|
||||
fullSpans, err := m.store.GetFlamegraphSpans(ctx, traceID, summary.Start, summary.End, nil)
|
||||
if err != nil {
|
||||
|
||||
@@ -11,12 +11,23 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/clickhousesql"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
"github.com/SigNoz/signoz/pkg/types/aiobservabilitytypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/spantypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
const colServiceName = `resource_string_service$$$$name` // $ gets escaped so $$$$ converts to $$.
|
||||
|
||||
var fullSpanColumns = []string{
|
||||
"duration_nano", "span_id", "has_error", "kind",
|
||||
colServiceName, "name",
|
||||
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
|
||||
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
|
||||
"flags", "is_remote", "trace_state", "status_code",
|
||||
"db_name", "db_operation", "http_method", "http_url", "http_host",
|
||||
"external_http_method", "external_http_url", "response_status_code", "links as references",
|
||||
}
|
||||
|
||||
func buildFieldExpr(fieldKey telemetrytypes.TelemetryFieldKey) (string, error) {
|
||||
switch fieldKey.FieldContext {
|
||||
case telemetrytypes.FieldContextResource:
|
||||
@@ -123,16 +134,8 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
|
||||
return []spantypes.StorableSpan{}, nil
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
sb.Select(
|
||||
"DISTINCT ON (span_id) timestamp",
|
||||
"duration_nano", "span_id", "has_error", "kind",
|
||||
colServiceName, "name",
|
||||
"attributes_string", "attributes_number", "attributes_bool", "resources_string",
|
||||
"events", "status_message", "status_code_string", "kind_string", "parent_span_id",
|
||||
"flags", "is_remote", "trace_state", "status_code",
|
||||
"db_name", "db_operation", "http_method", "http_url", "http_host",
|
||||
"external_http_method", "external_http_url", "response_status_code", "links as references",
|
||||
)
|
||||
sb.Select("DISTINCT ON (span_id) timestamp")
|
||||
sb.SelectMore(fullSpanColumns...)
|
||||
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
|
||||
ids := make([]any, len(spanIDs))
|
||||
for i, id := range spanIDs {
|
||||
@@ -155,6 +158,44 @@ func (s *traceStore) GetTraceSpansByIDs(ctx context.Context, traceID string, sta
|
||||
return spans, nil
|
||||
}
|
||||
|
||||
func (s *traceStore) GetThreadSpans(ctx context.Context, traceID string, summary *spantypes.TraceSummary, cursor *spantypes.ThreadCursor, limit int) ([]spantypes.StorableSpan, error) {
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
sb.Select("DISTINCT ON (span_id) timestamp")
|
||||
sb.SelectMore(fullSpanColumns...)
|
||||
sb.SelectMore("attributes")
|
||||
sb.From(fmt.Sprintf("%s.%s", spantypes.TraceDB, spantypes.TraceTable))
|
||||
sb.Where(
|
||||
sb.E("trace_id", traceID),
|
||||
sb.GE("ts_bucket_start", summary.Start.Unix()-1800),
|
||||
sb.LE("ts_bucket_start", summary.End.Unix()),
|
||||
// Reads only the JSON column; spans with messages only in the legacy maps are skipped.
|
||||
// todo(nitya): pick the column from the attribute evolution metadata.
|
||||
sb.Or(
|
||||
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIInputMessages))),
|
||||
sqlbuilder.Escape(fmt.Sprintf("attributes.%s IS NOT NULL", clickhousesql.Identifier(aiobservabilitytypes.GenAIOutputMessages))),
|
||||
),
|
||||
)
|
||||
if cursor != nil {
|
||||
// ClickHouse can't use an index for a tuple comparison, so the separate timestamp and
|
||||
// ts_bucket_start bounds are what skip the data before the cursor.
|
||||
sb.Where(
|
||||
sb.GE("ts_bucket_start", int64(cursor.TimeUnixNano/uint64(time.Second))-1800),
|
||||
sb.GE("timestamp", fmt.Sprintf("%d", cursor.TimeUnixNano)),
|
||||
sb.GT("(toUnixTimestamp64Nano(timestamp), span_id)", sqlbuilder.Tuple(cursor.TimeUnixNano, cursor.SpanID)),
|
||||
)
|
||||
}
|
||||
sb.OrderByAsc("timestamp")
|
||||
sb.OrderByAsc("span_id")
|
||||
sb.Limit(limit)
|
||||
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
|
||||
var spans []spantypes.StorableSpan
|
||||
if err := s.telemetryStore.ClickhouseDB().Select(ctx, &spans, query, args...); err != nil {
|
||||
return nil, errors.WrapInternalf(err, errors.CodeInternal, "error querying thread spans")
|
||||
}
|
||||
return spans, nil
|
||||
}
|
||||
|
||||
func (s *traceStore) GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]spantypes.StorableSpan, error) {
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
sb.Select(
|
||||
|
||||
@@ -13,6 +13,7 @@ type Handler interface {
|
||||
GetWaterfallV4(http.ResponseWriter, *http.Request)
|
||||
GetTraceAggregations(http.ResponseWriter, *http.Request)
|
||||
GetFlamegraph(http.ResponseWriter, *http.Request)
|
||||
GetThread(http.ResponseWriter, *http.Request)
|
||||
}
|
||||
|
||||
// Module defines the business logic for trace detail operations.
|
||||
@@ -20,4 +21,5 @@ type Module interface {
|
||||
GetWaterfallV4(ctx context.Context, traceID string, selectedSpanID string, uncollapsedSpans []string) (*spantypes.GettableWaterfallTrace, error)
|
||||
GetTraceAggregations(ctx context.Context, traceID string, req *spantypes.PostableTraceAggregations) (*spantypes.GettableTraceAggregations, error)
|
||||
GetFlamegraph(ctx context.Context, traceID string, selectedSpanID string, selectFields []telemetrytypes.TelemetryFieldKey) (*spantypes.GettableFlamegraphTrace, error)
|
||||
GetThread(ctx context.Context, traceID string, query *spantypes.ThreadQuery) (*spantypes.GettableTraceThread, error)
|
||||
}
|
||||
|
||||
@@ -566,24 +566,6 @@ func readAsRaw(rows driver.Rows, queryName string) (*qbtypes.RawData, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
// flattenJSONPaths flattens a decoded JSON document into dotted keys, overwriting existing keys in out.
|
||||
func flattenJSONPaths(prefix string, m map[string]any, out map[string]any) {
|
||||
for k, v := range m {
|
||||
key := k
|
||||
if prefix != "" {
|
||||
key = prefix + "." + k
|
||||
}
|
||||
switch child := v.(type) {
|
||||
case map[string]any:
|
||||
flattenJSONPaths(key, child, out)
|
||||
case telemetrystoretypes.JSONValue:
|
||||
flattenJSONPaths(key, child, out)
|
||||
default:
|
||||
out[key] = v
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// mergeSpanAttributeColumns merges (attributes_string, attributes_number, attributes_bool, resources_string) into
|
||||
// unified "attributes" and "resource" keys, and parses the stringified `events`
|
||||
// and `links` columns into structured slices. Raw DB columns are removed.
|
||||
@@ -598,7 +580,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
|
||||
resStr, hasRes := data["resources_string"]
|
||||
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {
|
||||
attributes := make(map[string]any)
|
||||
flattenJSONPaths("", attrJSON, attributes)
|
||||
attrJSON.FlattenInto("", attributes)
|
||||
if m, ok := attrStr.(map[string]string); ok {
|
||||
for k, v := range m {
|
||||
attributes[k] = v
|
||||
|
||||
@@ -36,7 +36,7 @@ func TestManager_ListRules_ValidatesParams(t *testing.T) {
|
||||
_, err = m.ListRules(context.Background(), &ruletypes.ListRulesParams{Limit: -1})
|
||||
require.ErrorContains(t, err, "invalid limit")
|
||||
|
||||
_, err = m.ListRules(context.Background(), &ruletypes.ListRulesParams{ListFilter: ruletypes.ListFilter{States: []string{"bogus"}}})
|
||||
_, err = m.ListRules(context.Background(), &ruletypes.ListRulesParams{States: []string{"bogus"}})
|
||||
require.ErrorContains(t, err, `invalid state "bogus"`)
|
||||
}
|
||||
|
||||
|
||||
@@ -12,11 +12,6 @@ type Handler interface {
|
||||
PatchRuleByID(http.ResponseWriter, *http.Request)
|
||||
TestRule(http.ResponseWriter, *http.Request)
|
||||
|
||||
ListRuleViews(http.ResponseWriter, *http.Request)
|
||||
CreateRuleView(http.ResponseWriter, *http.Request)
|
||||
UpdateRuleView(http.ResponseWriter, *http.Request)
|
||||
DeleteRuleView(http.ResponseWriter, *http.Request)
|
||||
|
||||
ListDowntimeSchedules(http.ResponseWriter, *http.Request)
|
||||
GetDowntimeScheduleByID(http.ResponseWriter, *http.Request)
|
||||
CreateDowntimeSchedule(http.ResponseWriter, *http.Request)
|
||||
|
||||
@@ -49,9 +49,4 @@ type Ruler interface {
|
||||
// TODO: expose downtime CRUD as methods on Ruler directly instead of leaking the
|
||||
// store interface. The handler should not call store methods directly.
|
||||
MaintenanceStore() alertmanagertypes.MaintenanceStore
|
||||
|
||||
CreateRuleView(ctx context.Context, orgID valuer.UUID, postable ruletypes.PostableRuleView) (*ruletypes.GettableRuleView, error)
|
||||
ListRuleViews(ctx context.Context, orgID valuer.UUID) (*ruletypes.ListableRuleViews, error)
|
||||
UpdateRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID, updatable ruletypes.UpdatableRuleView) (*ruletypes.GettableRuleView, error)
|
||||
DeleteRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) error
|
||||
}
|
||||
|
||||
@@ -1,93 +0,0 @@
|
||||
package sqlrulestore
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
ruletypes "github.com/SigNoz/signoz/pkg/types/ruletypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
)
|
||||
|
||||
func (r *rule) CreateRuleView(ctx context.Context, view *ruletypes.StorableRuleView) error {
|
||||
_, err := r.sqlstore.
|
||||
BunDBCtx(ctx).
|
||||
NewInsert().
|
||||
Model(view).
|
||||
Exec(ctx)
|
||||
if err != nil {
|
||||
return r.sqlstore.WrapAlreadyExistsErrf(err, errors.CodeAlreadyExists, "rule view with id %s already exists", view.ID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *rule) GetRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) (*ruletypes.StorableRuleView, error) {
|
||||
view := new(ruletypes.StorableRuleView)
|
||||
err := r.sqlstore.
|
||||
BunDBCtx(ctx).
|
||||
NewSelect().
|
||||
Model(view).
|
||||
Where("id = ?", id).
|
||||
Where("org_id = ?", orgID).
|
||||
Scan(ctx)
|
||||
if err != nil {
|
||||
return nil, r.sqlstore.WrapNotFoundErrf(err, ruletypes.ErrCodeRuleViewNotFound, "rule view with id %s doesn't exist", id)
|
||||
}
|
||||
return view, nil
|
||||
}
|
||||
|
||||
func (r *rule) ListRuleViews(ctx context.Context, orgID valuer.UUID) ([]*ruletypes.StorableRuleView, error) {
|
||||
views := make([]*ruletypes.StorableRuleView, 0)
|
||||
err := r.sqlstore.
|
||||
BunDBCtx(ctx).
|
||||
NewSelect().
|
||||
Model(&views).
|
||||
Where("org_id = ?", orgID).
|
||||
OrderExpr("updated_at DESC").
|
||||
Scan(ctx)
|
||||
if err != nil {
|
||||
return nil, errors.WrapInternalf(err, errors.CodeInternal, "couldn't list rule views")
|
||||
}
|
||||
return views, nil
|
||||
}
|
||||
|
||||
func (r *rule) UpdateRuleView(ctx context.Context, view *ruletypes.StorableRuleView) error {
|
||||
res, err := r.sqlstore.
|
||||
BunDBCtx(ctx).
|
||||
NewUpdate().
|
||||
Model(view).
|
||||
WherePK().
|
||||
Where("org_id = ?", view.OrgID).
|
||||
Exec(ctx)
|
||||
if err != nil {
|
||||
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't update rule view")
|
||||
}
|
||||
rows, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't read rule view update result")
|
||||
}
|
||||
if rows == 0 {
|
||||
return errors.Newf(errors.TypeNotFound, ruletypes.ErrCodeRuleViewNotFound, "rule view with id %s doesn't exist", view.ID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *rule) DeleteRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) error {
|
||||
res, err := r.sqlstore.
|
||||
BunDBCtx(ctx).
|
||||
NewDelete().
|
||||
Model(new(ruletypes.StorableRuleView)).
|
||||
Where("id = ?", id).
|
||||
Where("org_id = ?", orgID).
|
||||
Exec(ctx)
|
||||
if err != nil {
|
||||
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't delete rule view")
|
||||
}
|
||||
rows, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return errors.WrapInternalf(err, errors.CodeInternal, "couldn't read rule view delete result")
|
||||
}
|
||||
if rows == 0 {
|
||||
return errors.Newf(errors.TypeNotFound, ruletypes.ErrCodeRuleViewNotFound, "rule view with id %s doesn't exist", id)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -345,122 +345,3 @@ func (handler *handler) DeleteDowntimeScheduleByID(rw http.ResponseWriter, req *
|
||||
|
||||
render.Success(rw, http.StatusNoContent, nil)
|
||||
}
|
||||
|
||||
func (handler *handler) ListRuleViews(rw http.ResponseWriter, req *http.Request) {
|
||||
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
claims, err := authtypes.ClaimsFromContext(ctx)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
orgID, err := valuer.NewUUID(claims.OrgID)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
views, err := handler.ruler.ListRuleViews(ctx, orgID)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, views)
|
||||
}
|
||||
|
||||
func (handler *handler) CreateRuleView(rw http.ResponseWriter, req *http.Request) {
|
||||
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
claims, err := authtypes.ClaimsFromContext(ctx)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
orgID, err := valuer.NewUUID(claims.OrgID)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
var postable ruletypes.PostableRuleView
|
||||
if err := binding.JSON.BindBody(req.Body, &postable); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
view, err := handler.ruler.CreateRuleView(ctx, orgID, postable)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusCreated, view)
|
||||
}
|
||||
|
||||
func (handler *handler) UpdateRuleView(rw http.ResponseWriter, req *http.Request) {
|
||||
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
claims, err := authtypes.ClaimsFromContext(ctx)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
orgID, err := valuer.NewUUID(claims.OrgID)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
id, err := valuer.NewUUID(mux.Vars(req)["id"])
|
||||
if err != nil {
|
||||
render.Error(rw, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "id is not a valid uuid-v7"))
|
||||
return
|
||||
}
|
||||
|
||||
var updatable ruletypes.UpdatableRuleView
|
||||
if err := binding.JSON.BindBody(req.Body, &updatable); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
view, err := handler.ruler.UpdateRuleView(ctx, orgID, id, updatable)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusOK, view)
|
||||
}
|
||||
|
||||
func (handler *handler) DeleteRuleView(rw http.ResponseWriter, req *http.Request) {
|
||||
ctx, cancel := context.WithTimeout(req.Context(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
claims, err := authtypes.ClaimsFromContext(ctx)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
orgID, err := valuer.NewUUID(claims.OrgID)
|
||||
if err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
id, err := valuer.NewUUID(mux.Vars(req)["id"])
|
||||
if err != nil {
|
||||
render.Error(rw, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "id is not a valid uuid-v7"))
|
||||
return
|
||||
}
|
||||
|
||||
if err := handler.ruler.DeleteRuleView(ctx, orgID, id); err != nil {
|
||||
render.Error(rw, err)
|
||||
return
|
||||
}
|
||||
|
||||
render.Success(rw, http.StatusNoContent, nil)
|
||||
}
|
||||
|
||||
@@ -147,41 +147,3 @@ func (provider *provider) TestNotification(ctx context.Context, orgID valuer.UUI
|
||||
func (provider *provider) MaintenanceStore() alertmanagertypes.MaintenanceStore {
|
||||
return provider.manager.MaintenanceStore()
|
||||
}
|
||||
|
||||
func (provider *provider) CreateRuleView(ctx context.Context, orgID valuer.UUID, postable ruletypes.PostableRuleView) (*ruletypes.GettableRuleView, error) {
|
||||
if err := postable.Validate(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
storable := postable.ToStorableRuleView(orgID)
|
||||
if err := provider.ruleStore.CreateRuleView(ctx, storable); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return storable.ToGettableRuleView(), nil
|
||||
}
|
||||
|
||||
func (provider *provider) ListRuleViews(ctx context.Context, orgID valuer.UUID) (*ruletypes.ListableRuleViews, error) {
|
||||
storables, err := provider.ruleStore.ListRuleViews(ctx, orgID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &ruletypes.ListableRuleViews{Views: ruletypes.NewGettableRuleViewsFromStorableRuleViews(storables)}, nil
|
||||
}
|
||||
|
||||
func (provider *provider) UpdateRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID, updatable ruletypes.UpdatableRuleView) (*ruletypes.GettableRuleView, error) {
|
||||
if err := updatable.Validate(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
storable, err := provider.ruleStore.GetRuleView(ctx, orgID, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
storable.Update(updatable)
|
||||
if err := provider.ruleStore.UpdateRuleView(ctx, storable); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return storable.ToGettableRuleView(), nil
|
||||
}
|
||||
|
||||
func (provider *provider) DeleteRuleView(ctx context.Context, orgID valuer.UUID, id valuer.UUID) error {
|
||||
return provider.ruleStore.DeleteRuleView(ctx, orgID, id)
|
||||
}
|
||||
|
||||
@@ -257,7 +257,6 @@ func NewSQLMigrationProviderFactories(
|
||||
sqlmigration.NewAddCloudIntegrationTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddNotificationChannelTuplesFactory(sqlstore),
|
||||
sqlmigration.NewAddAIObservabilityQuickFiltersFactory(sqlstore),
|
||||
sqlmigration.NewAddRuleViewFactory(sqlstore, sqlschema),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -13,7 +13,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/sqlschema"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/types/ruletypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
@@ -50,6 +49,11 @@ type rule struct {
|
||||
OrgID string `bun:"org_id,type:text"`
|
||||
}
|
||||
|
||||
type routePolicyRuleData struct {
|
||||
PreferredChannels []string `json:"preferredChannels"`
|
||||
Labels map[string]string `json:"labels"`
|
||||
}
|
||||
|
||||
type addRoutePolicies struct {
|
||||
sqlstore sqlstore.SQLStore
|
||||
sqlschema sqlschema.SQLSchema
|
||||
@@ -187,20 +191,20 @@ func (migration *addRoutePolicies) migrateRulesToRoutePolicies(ctx context.Conte
|
||||
func (migration *addRoutePolicies) convertRulesToRoutes(rules []*rule, channelsByOrg map[string][]string) ([]*expressionRoute, error) {
|
||||
var routes []*expressionRoute
|
||||
for _, r := range rules {
|
||||
var gettableRule ruletypes.GettableRule
|
||||
if err := json.Unmarshal([]byte(r.Data), &gettableRule); err != nil {
|
||||
var ruleData routePolicyRuleData
|
||||
if err := json.Unmarshal([]byte(r.Data), &ruleData); err != nil {
|
||||
return nil, errors.NewInternalf(errors.CodeInternal, "failed to unmarshal rule data for rule ID %s: %v", r.ID, err)
|
||||
}
|
||||
|
||||
if len(gettableRule.PreferredChannels) == 0 {
|
||||
if len(ruleData.PreferredChannels) == 0 {
|
||||
channels, exists := channelsByOrg[r.OrgID]
|
||||
if !exists || len(channels) == 0 {
|
||||
continue
|
||||
}
|
||||
gettableRule.PreferredChannels = channels
|
||||
ruleData.PreferredChannels = channels
|
||||
}
|
||||
severity := "critical"
|
||||
if v, ok := gettableRule.Labels["severity"]; ok {
|
||||
if v, ok := ruleData.Labels["severity"]; ok {
|
||||
severity = v
|
||||
}
|
||||
expression := fmt.Sprintf(`%s == "%s" && %s == "%s"`, "threshold.name", severity, "ruleId", r.ID.String())
|
||||
@@ -218,7 +222,7 @@ func (migration *addRoutePolicies) convertRulesToRoutes(rules []*rule, channelsB
|
||||
},
|
||||
Expression: expression,
|
||||
ExpressionKind: "rule",
|
||||
Channels: gettableRule.PreferredChannels,
|
||||
Channels: ruleData.PreferredChannels,
|
||||
Name: r.ID.StringValue(),
|
||||
Enabled: true,
|
||||
OrgID: r.OrgID,
|
||||
|
||||
@@ -1,78 +0,0 @@
|
||||
package sqlmigration
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/sqlschema"
|
||||
"github.com/SigNoz/signoz/pkg/sqlstore"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
)
|
||||
|
||||
type addRuleView struct {
|
||||
sqlstore sqlstore.SQLStore
|
||||
sqlschema sqlschema.SQLSchema
|
||||
}
|
||||
|
||||
func NewAddRuleViewFactory(sqlstore sqlstore.SQLStore, sqlschema sqlschema.SQLSchema) factory.ProviderFactory[SQLMigration, Config] {
|
||||
return factory.NewProviderFactory(factory.MustNewName("add_rule_view"), func(ctx context.Context, ps factory.ProviderSettings, c Config) (SQLMigration, error) {
|
||||
return &addRuleView{
|
||||
sqlstore: sqlstore,
|
||||
sqlschema: sqlschema,
|
||||
}, nil
|
||||
})
|
||||
}
|
||||
|
||||
func (migration *addRuleView) Register(migrations *migrate.Migrations) error {
|
||||
return migrations.Register(migration.Up, migration.Down)
|
||||
}
|
||||
|
||||
func (migration *addRuleView) Up(ctx context.Context, db *bun.DB) error {
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
sqls := migration.sqlschema.Operator().CreateTable(&sqlschema.Table{
|
||||
Name: "rule_view",
|
||||
Columns: []*sqlschema.Column{
|
||||
{Name: "id", DataType: sqlschema.DataTypeText, Nullable: false},
|
||||
{Name: "name", DataType: sqlschema.DataTypeText, Nullable: false},
|
||||
{Name: "data", DataType: sqlschema.DataTypeText, Nullable: false},
|
||||
{Name: "org_id", DataType: sqlschema.DataTypeText, Nullable: false},
|
||||
{Name: "created_at", DataType: sqlschema.DataTypeTimestamp, Nullable: false},
|
||||
{Name: "updated_at", DataType: sqlschema.DataTypeTimestamp, Nullable: false},
|
||||
},
|
||||
PrimaryKeyConstraint: &sqlschema.PrimaryKeyConstraint{ColumnNames: []sqlschema.ColumnName{"id"}},
|
||||
ForeignKeyConstraints: []*sqlschema.ForeignKeyConstraint{
|
||||
{
|
||||
ReferencingColumnName: sqlschema.ColumnName("org_id"),
|
||||
ReferencedTableName: sqlschema.TableName("organizations"),
|
||||
ReferencedColumnName: sqlschema.ColumnName("id"),
|
||||
},
|
||||
},
|
||||
})
|
||||
|
||||
for _, sql := range sqls {
|
||||
if _, err := tx.ExecContext(ctx, string(sql)); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
if _, err := tx.NewCreateIndex().
|
||||
Table("rule_view").
|
||||
Column("org_id").
|
||||
Index("idx_rule_view_org_id").
|
||||
IfNotExists().
|
||||
Exec(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
func (migration *addRuleView) Down(_ context.Context, _ *bun.DB) error {
|
||||
return nil
|
||||
}
|
||||
@@ -364,12 +364,26 @@ func (provider *provider) gc(ctx context.Context, org *types.Organization) error
|
||||
}
|
||||
|
||||
func (provider *provider) flushLastObservedAt(ctx context.Context, org *types.Organization) error {
|
||||
accessTokenToLastObservedAt, err := provider.listLastObservedAtDesc(ctx, org.ID)
|
||||
tokens, err := provider.tokenStore.ListByOrgID(ctx, org.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := provider.tokenStore.UpdateLastObservedAtByAccessToken(ctx, accessTokenToLastObservedAt); err != nil {
|
||||
observedTokens := make([]*authtypes.StorableToken, 0, len(tokens))
|
||||
for _, token := range tokens {
|
||||
cachedLastObservedAt, ok := provider.lastObservedAtCache.Get(lastObservedAtCacheKey(token.AccessToken, token.UserID))
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
if err := token.UpdateLastObservedAt(cachedLastObservedAt); err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
observedTokens = append(observedTokens, token)
|
||||
}
|
||||
|
||||
if err := provider.tokenStore.UpdateLastObservedAt(ctx, observedTokens); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
@@ -232,15 +232,16 @@ func (store *store) ListByUserID(ctx context.Context, userID valuer.UUID) ([]*au
|
||||
return tokens, nil
|
||||
}
|
||||
|
||||
func (store *store) UpdateLastObservedAtByAccessToken(ctx context.Context, accessTokenToLastObservedAt []map[string]any) error {
|
||||
if len(accessTokenToLastObservedAt) == 0 {
|
||||
func (store *store) UpdateLastObservedAt(ctx context.Context, tokens []*authtypes.StorableToken) error {
|
||||
if len(tokens) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
values := store.
|
||||
sqlstore.
|
||||
BunDBCtx(ctx).
|
||||
NewValues(&accessTokenToLastObservedAt)
|
||||
NewValues(&tokens).
|
||||
Column("id", "last_observed_at", "updated_at")
|
||||
|
||||
_, err := store.
|
||||
sqlstore.
|
||||
@@ -250,8 +251,8 @@ func (store *store) UpdateLastObservedAtByAccessToken(ctx context.Context, acces
|
||||
Model((*authtypes.StorableToken)(nil)).
|
||||
TableExpr("update_cte").
|
||||
Set("last_observed_at = update_cte.last_observed_at").
|
||||
Where("auth_token.access_token = update_cte.access_token").
|
||||
Where("auth_token.user_id = update_cte.user_id").
|
||||
Set("updated_at = update_cte.updated_at").
|
||||
Where("auth_token.id = update_cte.id").
|
||||
Exec(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
62
pkg/types/alertmanagertypes/channel_email.go
Normal file
62
pkg/types/alertmanagertypes/channel_email.go
Normal file
@@ -0,0 +1,62 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"maps"
|
||||
"net/textproto"
|
||||
"slices"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
)
|
||||
|
||||
// ChannelEmailConfig carries no SMTP transport fields: the smarthost,
|
||||
// credentials and TLS settings come from the deployment's global config, so a
|
||||
// channel can only choose recipients and body.
|
||||
type ChannelEmailConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
To string `json:"to" required:"true"`
|
||||
HTML valuer.UnsetOrNonEmptyString `json:"html"`
|
||||
Headers map[string]string `json:"headers,omitempty"`
|
||||
}
|
||||
|
||||
func (c ChannelEmailConfig) Validate() error {
|
||||
if c.To == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.to is required for an email channel")
|
||||
}
|
||||
|
||||
// A read reports header names as textproto canonicalizes them, turning
|
||||
// "subject" into "Subject", so a name that is not already in that form is
|
||||
// rejected rather than answered with one the caller never sent.
|
||||
for _, header := range slices.Sorted(maps.Keys(c.Headers)) {
|
||||
if canonical := textproto.CanonicalMIMEHeaderKey(header); canonical != header {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.headers name %q must be written as %q", header, canonical)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelEmailConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
EmailConfigs: []*config.EmailConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultEmailConfig.VSendResolved)},
|
||||
To: c.To,
|
||||
HTML: c.HTML.StringValue(),
|
||||
Headers: c.Headers,
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelEmailConfigFromReceiver(_ string, receiver *Receiver) (ChannelSpec, error) {
|
||||
email := receiver.EmailConfigs[0]
|
||||
sendResolved := email.VSendResolved
|
||||
|
||||
return &ChannelEmailConfig{
|
||||
SendResolved: &sendResolved,
|
||||
To: email.To,
|
||||
HTML: valuer.UnsetIfEmpty(email.HTML),
|
||||
Headers: email.Headers,
|
||||
}, nil
|
||||
}
|
||||
@@ -5,10 +5,59 @@ import (
|
||||
"strings"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
)
|
||||
|
||||
type ChannelGoogleChatConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Text valuer.UnsetOrNonEmptyString `json:"text"`
|
||||
}
|
||||
|
||||
func (c ChannelGoogleChatConfig) Validate() error {
|
||||
if c.WebhookURL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for a googlechat channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelGoogleChatConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
webhookURL, err := parseSecretURL(c.WebhookURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
GoogleChatConfigs: []*GoogleChatReceiverConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultGoogleChatReceiverConfig.VSendResolved)},
|
||||
WebhookURL: webhookURL,
|
||||
Title: c.Title.StringValue(),
|
||||
Text: c.Text.StringValue(),
|
||||
}},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelGoogleChatConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
googlechat := receiver.GoogleChatConfigs[0]
|
||||
sendResolved := googlechat.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, googlechat.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelGoogleChatConfig{
|
||||
SendResolved: &sendResolved,
|
||||
WebhookURL: formatSecretURL(googlechat.WebhookURL),
|
||||
Title: valuer.UnsetIfEmpty(googlechat.Title),
|
||||
Text: valuer.UnsetIfEmpty(googlechat.Text),
|
||||
}, nil
|
||||
}
|
||||
|
||||
type GoogleChatReceiverConfig struct {
|
||||
config.NotifierConfig `yaml:",inline" json:",inline"`
|
||||
|
||||
@@ -6,10 +6,64 @@ import (
|
||||
"strings"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
)
|
||||
|
||||
type ChannelIncidentIOConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
URL string `json:"url" required:"true"`
|
||||
Token string `json:"token" required:"true" format:"password"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Metadata map[string]string `json:"metadata,omitempty"`
|
||||
}
|
||||
|
||||
func (c ChannelIncidentIOConfig) Validate() error {
|
||||
if c.URL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for an incidentio channel")
|
||||
}
|
||||
|
||||
if c.Token == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.token is required for an incidentio channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelIncidentIOConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
IncidentIOConfigs: []*IncidentIOReceiverConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultIncidentIOReceiverConfig.VSendResolved)},
|
||||
URL: c.URL,
|
||||
Token: config.Secret(c.Token),
|
||||
Title: c.Title.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Metadata: c.Metadata,
|
||||
}},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelIncidentIOConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
incidentio := receiver.IncidentIOConfigs[0]
|
||||
sendResolved := incidentio.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, incidentio.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelIncidentIOConfig{
|
||||
SendResolved: &sendResolved,
|
||||
URL: incidentio.URL,
|
||||
Token: string(incidentio.Token),
|
||||
Title: valuer.UnsetIfEmpty(incidentio.Title),
|
||||
Description: valuer.UnsetIfEmpty(incidentio.Description),
|
||||
Metadata: incidentio.Metadata,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// incidentIOEventsPathPrefix is the path of incident.io's HTTP alert source
|
||||
// endpoint (Alert Events V2 API). The full URL is per-source:
|
||||
// https://api.incident.io/v2/alert_events/http/<source_config_id>.
|
||||
@@ -7,11 +7,147 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
"github.com/prometheus/common/model"
|
||||
)
|
||||
|
||||
type ChannelJiraConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
// Site is the Jira Cloud base URL, https://<site>.atlassian.net. Only Jira
|
||||
// Cloud is supported; the REST base is derived from it.
|
||||
Site string `json:"site" required:"true"`
|
||||
Project string `json:"project" required:"true"`
|
||||
IssueType string `json:"issueType" required:"true"`
|
||||
Summary valuer.UnsetOrNonEmptyString `json:"summary"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Priority string `json:"priority"`
|
||||
Labels []string `json:"labels,omitempty"`
|
||||
ResolveTransition string `json:"resolveTransition"`
|
||||
ReopenTransition string `json:"reopenTransition"`
|
||||
ReopenDuration valuer.UnsetOrNonEmptyString `json:"reopenDuration"`
|
||||
WontFixResolution string `json:"wontFixResolution"`
|
||||
CustomFields map[string]any `json:"customFields,omitempty"`
|
||||
|
||||
Email string `json:"email" required:"true"`
|
||||
APIToken string `json:"apiToken" required:"true" format:"password"`
|
||||
}
|
||||
|
||||
func (c ChannelJiraConfig) Validate() error {
|
||||
for _, required := range []struct {
|
||||
value string
|
||||
field string
|
||||
}{
|
||||
{c.Site, "site"},
|
||||
{c.Project, "project"},
|
||||
{c.IssueType, "issueType"},
|
||||
{c.Email, "email"},
|
||||
{c.APIToken, "apiToken"},
|
||||
} {
|
||||
if required.value == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.%s is required for a jira channel", required.field)
|
||||
}
|
||||
}
|
||||
|
||||
if !c.ReopenDuration.IsZero() {
|
||||
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
|
||||
if err != nil {
|
||||
return errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q is not a valid duration", c.ReopenDuration)
|
||||
}
|
||||
|
||||
// A read reports the duration as model.Duration formats it, collapsing
|
||||
// "72h" into "3d", so a value that is not already in that form is rejected
|
||||
// rather than answered with one the caller never sent.
|
||||
if canonical := reopenDuration.String(); canonical != c.ReopenDuration.StringValue() {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q must be written as %q", c.ReopenDuration, canonical)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelJiraConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
// Seeded from upstream's default rather than a zero value: FollowRedirects
|
||||
// and EnableHTTP2 marshal unconditionally, so a zero value would persist them
|
||||
// as false and read back as a config ChannelJiraConfig cannot represent.
|
||||
httpConfig := commoncfg.DefaultHTTPClientConfig
|
||||
httpConfig.BasicAuth = &commoncfg.BasicAuth{
|
||||
Username: c.Email,
|
||||
Password: commoncfg.Secret(c.APIToken),
|
||||
}
|
||||
|
||||
jira := &JiraReceiverConfig{
|
||||
// JiraReceiverConfig seeds no send_resolved of its own, so unset means off.
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, false)},
|
||||
Site: c.Site,
|
||||
Project: c.Project,
|
||||
IssueType: c.IssueType,
|
||||
Summary: c.Summary.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Priority: c.Priority,
|
||||
Labels: c.Labels,
|
||||
ResolveTransition: c.ResolveTransition,
|
||||
ReopenTransition: c.ReopenTransition,
|
||||
WontFixResolution: c.WontFixResolution,
|
||||
CustomFields: c.CustomFields,
|
||||
HTTPConfig: &httpConfig,
|
||||
}
|
||||
|
||||
if !c.ReopenDuration.IsZero() {
|
||||
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
|
||||
if err != nil {
|
||||
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "parse reopenDuration %q", c.ReopenDuration)
|
||||
}
|
||||
jira.ReopenDuration = reopenDuration
|
||||
}
|
||||
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
JiraConfigs: []*JiraReceiverConfig{jira},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelJiraConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
jira := receiver.JiraConfigs[0]
|
||||
sendResolved := jira.VSendResolved
|
||||
|
||||
if err := rejectUnsupportedHTTPConfig(name, jira.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if jira.HTTPConfig != nil && jira.HTTPConfig.Authorization != nil {
|
||||
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization, which is not supported", name)
|
||||
}
|
||||
|
||||
if err := rejectHTTPBasicAuthBeyondPassword(name, jira.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
spec := &ChannelJiraConfig{
|
||||
SendResolved: &sendResolved,
|
||||
Site: jira.Site,
|
||||
Project: jira.Project,
|
||||
IssueType: jira.IssueType,
|
||||
Summary: valuer.UnsetIfEmpty(jira.Summary),
|
||||
Description: valuer.UnsetIfEmpty(jira.Description),
|
||||
Priority: jira.Priority,
|
||||
Labels: jira.Labels,
|
||||
ResolveTransition: jira.ResolveTransition,
|
||||
ReopenTransition: jira.ReopenTransition,
|
||||
ReopenDuration: valuer.UnsetIfEmpty(jira.ReopenDuration.String()),
|
||||
WontFixResolution: jira.WontFixResolution,
|
||||
CustomFields: jira.CustomFields,
|
||||
}
|
||||
|
||||
if jira.HTTPConfig != nil && jira.HTTPConfig.BasicAuth != nil {
|
||||
spec.Email = jira.HTTPConfig.BasicAuth.Username
|
||||
spec.APIToken = string(jira.HTTPConfig.BasicAuth.Password)
|
||||
}
|
||||
|
||||
return spec, nil
|
||||
}
|
||||
|
||||
const defaultJiraReopenDuration = model.Duration(3 * 24 * time.Hour)
|
||||
|
||||
// Service accounts authenticate against the api.atlassian.com gateway (keyed by
|
||||
@@ -2,10 +2,63 @@ package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
)
|
||||
|
||||
// ChannelJSMOpsConfig carries no API URL: JSM Ops is a single global gateway
|
||||
// keyed by the integration API key, which the notifier pins itself.
|
||||
type ChannelJSMOpsConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
APIKey string `json:"apiKey" required:"true" format:"password"`
|
||||
Message valuer.UnsetOrNonEmptyString `json:"message"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Priority string `json:"priority"`
|
||||
// Tags is the comma-separated list JSM Ops attaches to the alert.
|
||||
Tags valuer.UnsetOrNonEmptyString `json:"tags"`
|
||||
}
|
||||
|
||||
func (c ChannelJSMOpsConfig) Validate() error {
|
||||
if c.APIKey == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for a jsmops channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelJSMOpsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
JSMOpsConfigs: []*JSMOpsReceiverConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultJSMOpsReceiverConfig.VSendResolved)},
|
||||
APIKey: config.Secret(c.APIKey),
|
||||
Message: c.Message.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Priority: c.Priority,
|
||||
Tags: c.Tags.StringValue(),
|
||||
}},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelJSMOpsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
jsmops := receiver.JSMOpsConfigs[0]
|
||||
sendResolved := jsmops.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, jsmops.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelJSMOpsConfig{
|
||||
SendResolved: &sendResolved,
|
||||
APIKey: string(jsmops.APIKey),
|
||||
Message: valuer.UnsetIfEmpty(jsmops.Message),
|
||||
Description: valuer.UnsetIfEmpty(jsmops.Description),
|
||||
Priority: jsmops.Priority,
|
||||
Tags: valuer.UnsetIfEmpty(jsmops.Tags),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// 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
|
||||
55
pkg/types/alertmanagertypes/channel_msteams.go
Normal file
55
pkg/types/alertmanagertypes/channel_msteams.go
Normal file
@@ -0,0 +1,55 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
)
|
||||
|
||||
type ChannelMSTeamsConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Text valuer.UnsetOrNonEmptyString `json:"text"`
|
||||
}
|
||||
|
||||
func (c ChannelMSTeamsConfig) Validate() error {
|
||||
if c.WebhookURL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for an msteams channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelMSTeamsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
webhookURL, err := parseSecretURL(c.WebhookURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
MSTeamsV2Configs: []*config.MSTeamsV2Config{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultMSTeamsV2Config.VSendResolved)},
|
||||
WebhookURL: webhookURL,
|
||||
Title: c.Title.StringValue(),
|
||||
Text: c.Text.StringValue(),
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelMSTeamsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
msteams := receiver.MSTeamsV2Configs[0]
|
||||
sendResolved := msteams.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, msteams.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelMSTeamsConfig{
|
||||
SendResolved: &sendResolved,
|
||||
WebhookURL: formatSecretURL(msteams.WebhookURL),
|
||||
Title: valuer.UnsetIfEmpty(msteams.Title),
|
||||
Text: valuer.UnsetIfEmpty(msteams.Text),
|
||||
}, nil
|
||||
}
|
||||
71
pkg/types/alertmanagertypes/channel_opsgenie.go
Normal file
71
pkg/types/alertmanagertypes/channel_opsgenie.go
Normal file
@@ -0,0 +1,71 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
)
|
||||
|
||||
type ChannelOpsgenieConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
APIKey string `json:"apiKey" required:"true" format:"password"`
|
||||
APIURL string `json:"apiUrl"`
|
||||
Message valuer.UnsetOrNonEmptyString `json:"message"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Source valuer.UnsetOrNonEmptyString `json:"source"`
|
||||
Details map[string]string `json:"details,omitempty"`
|
||||
Priority string `json:"priority"`
|
||||
}
|
||||
|
||||
func (c ChannelOpsgenieConfig) Validate() error {
|
||||
if c.APIKey == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for an opsgenie channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelOpsgenieConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
var apiURL *config.URL
|
||||
if c.APIURL != "" {
|
||||
parsed, err := parseUpstreamURL(c.APIURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
apiURL = parsed
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
OpsGenieConfigs: []*config.OpsGenieConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultOpsGenieConfig.VSendResolved)},
|
||||
APIKey: config.Secret(c.APIKey),
|
||||
APIURL: apiURL,
|
||||
Message: c.Message.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Source: c.Source.StringValue(),
|
||||
Priority: c.Priority,
|
||||
Details: c.Details,
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelOpsgenieConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
opsgenie := receiver.OpsGenieConfigs[0]
|
||||
sendResolved := opsgenie.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, opsgenie.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelOpsgenieConfig{
|
||||
SendResolved: &sendResolved,
|
||||
APIKey: string(opsgenie.APIKey),
|
||||
APIURL: formatUpstreamURL(opsgenie.APIURL),
|
||||
Message: valuer.UnsetIfEmpty(opsgenie.Message),
|
||||
Description: valuer.UnsetIfEmpty(opsgenie.Description),
|
||||
Source: valuer.UnsetIfEmpty(opsgenie.Source),
|
||||
Priority: opsgenie.Priority,
|
||||
Details: opsgenie.Details,
|
||||
}, nil
|
||||
}
|
||||
92
pkg/types/alertmanagertypes/channel_pagerduty.go
Normal file
92
pkg/types/alertmanagertypes/channel_pagerduty.go
Normal file
@@ -0,0 +1,92 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
)
|
||||
|
||||
type ChannelPagerdutyConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
RoutingKey string `json:"routingKey" required:"true" format:"password"`
|
||||
URL string `json:"url"`
|
||||
Source valuer.UnsetOrNonEmptyString `json:"source"`
|
||||
Client valuer.UnsetOrNonEmptyString `json:"client"`
|
||||
ClientURL valuer.UnsetOrNonEmptyString `json:"clientUrl"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Severity string `json:"severity"`
|
||||
Component string `json:"component"`
|
||||
Group string `json:"group"`
|
||||
Class string `json:"class"`
|
||||
Details map[string]string `json:"details,omitempty"`
|
||||
}
|
||||
|
||||
func (c ChannelPagerdutyConfig) Validate() error {
|
||||
if c.RoutingKey == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.routingKey is required for a pagerduty channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelPagerdutyConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
var eventsURL *config.URL
|
||||
if c.URL != "" {
|
||||
parsed, err := parseUpstreamURL(c.URL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
eventsURL = parsed
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
PagerdutyConfigs: []*config.PagerdutyConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultPagerdutyConfig.VSendResolved)},
|
||||
RoutingKey: config.Secret(c.RoutingKey),
|
||||
URL: eventsURL,
|
||||
Source: c.Source.StringValue(),
|
||||
Client: c.Client.StringValue(),
|
||||
ClientURL: c.ClientURL.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Severity: c.Severity,
|
||||
Component: c.Component,
|
||||
Group: c.Group,
|
||||
Class: c.Class,
|
||||
Details: newUpstreamDetails(c.Details),
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelPagerdutyConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
pagerduty := receiver.PagerdutyConfigs[0]
|
||||
sendResolved := pagerduty.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, pagerduty.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var details map[string]string
|
||||
if len(pagerduty.Details) > 0 {
|
||||
extracted, err := extractStringDetails(name, pagerduty.Details)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
details = extracted
|
||||
}
|
||||
|
||||
return &ChannelPagerdutyConfig{
|
||||
SendResolved: &sendResolved,
|
||||
RoutingKey: string(pagerduty.RoutingKey),
|
||||
URL: formatUpstreamURL(pagerduty.URL),
|
||||
Source: valuer.UnsetIfEmpty(pagerduty.Source),
|
||||
Client: valuer.UnsetIfEmpty(pagerduty.Client),
|
||||
ClientURL: valuer.UnsetIfEmpty(pagerduty.ClientURL),
|
||||
Description: valuer.UnsetIfEmpty(pagerduty.Description),
|
||||
Severity: pagerduty.Severity,
|
||||
Component: pagerduty.Component,
|
||||
Group: pagerduty.Group,
|
||||
Class: pagerduty.Class,
|
||||
Details: details,
|
||||
}, nil
|
||||
}
|
||||
182
pkg/types/alertmanagertypes/channel_slack.go
Normal file
182
pkg/types/alertmanagertypes/channel_slack.go
Normal file
@@ -0,0 +1,182 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
)
|
||||
|
||||
type ChannelSlackConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
APIURL string `json:"apiUrl" required:"true" format:"password"`
|
||||
Channel string `json:"channel"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Text valuer.UnsetOrNonEmptyString `json:"text"`
|
||||
Color valuer.UnsetOrNonEmptyString `json:"color"`
|
||||
TitleLink valuer.UnsetOrNonEmptyString `json:"titleLink"`
|
||||
Pretext valuer.UnsetOrNonEmptyString `json:"pretext"`
|
||||
Fallback valuer.UnsetOrNonEmptyString `json:"fallback"`
|
||||
Footer valuer.UnsetOrNonEmptyString `json:"footer"`
|
||||
Fields []ChannelSlackField `json:"fields,omitempty"`
|
||||
Actions []ChannelSlackAction `json:"actions,omitempty"`
|
||||
}
|
||||
|
||||
type ChannelSlackField struct {
|
||||
Title string `json:"title" required:"true"`
|
||||
Value string `json:"value" required:"true"`
|
||||
Short *bool `json:"short,omitempty"`
|
||||
}
|
||||
|
||||
// ChannelSlackAction is a link button when URL is set, otherwise a message
|
||||
// button that needs Name. Upstream clears whichever side is not in use.
|
||||
type ChannelSlackAction struct {
|
||||
Type string `json:"type" required:"true"`
|
||||
Text string `json:"text" required:"true"`
|
||||
URL string `json:"url"`
|
||||
Style string `json:"style"`
|
||||
Name string `json:"name"`
|
||||
Value string `json:"value"`
|
||||
Confirm *ChannelSlackConfirmation `json:"confirm,omitempty"`
|
||||
}
|
||||
|
||||
type ChannelSlackConfirmation struct {
|
||||
Text string `json:"text" required:"true"`
|
||||
Title string `json:"title"`
|
||||
OkText string `json:"okText"`
|
||||
DismissText string `json:"dismissText"`
|
||||
}
|
||||
|
||||
func (c ChannelSlackConfig) Validate() error {
|
||||
if c.APIURL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiUrl is required for a slack channel")
|
||||
}
|
||||
|
||||
for i, field := range c.Fields {
|
||||
if field.Title == "" || field.Value == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.fields[%d] requires title and value", i)
|
||||
}
|
||||
}
|
||||
|
||||
for i, action := range c.Actions {
|
||||
if action.Type == "" || action.Text == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires type and text", i)
|
||||
}
|
||||
if action.URL == "" && action.Name == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires url or name", i)
|
||||
}
|
||||
if action.Confirm != nil && action.Confirm.Text == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d].confirm requires text", i)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelSlackConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
apiURL, err := parseSecretURL(c.APIURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
SlackConfigs: []*config.SlackConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultSlackConfig.VSendResolved)},
|
||||
APIURL: apiURL,
|
||||
Channel: c.Channel,
|
||||
Title: c.Title.StringValue(),
|
||||
Text: c.Text.StringValue(),
|
||||
Color: c.Color.StringValue(),
|
||||
TitleLink: c.TitleLink.StringValue(),
|
||||
Pretext: c.Pretext.StringValue(),
|
||||
Fallback: c.Fallback.StringValue(),
|
||||
Footer: c.Footer.StringValue(),
|
||||
Fields: newUpstreamSlackFields(c.Fields),
|
||||
Actions: newUpstreamSlackActions(c.Actions),
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelSlackConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
slack := receiver.SlackConfigs[0]
|
||||
sendResolved := slack.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, slack.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelSlackConfig{
|
||||
SendResolved: &sendResolved,
|
||||
APIURL: formatSecretURL(slack.APIURL),
|
||||
Channel: slack.Channel,
|
||||
Title: valuer.UnsetIfEmpty(slack.Title),
|
||||
Text: valuer.UnsetIfEmpty(slack.Text),
|
||||
Color: valuer.UnsetIfEmpty(slack.Color),
|
||||
TitleLink: valuer.UnsetIfEmpty(slack.TitleLink),
|
||||
Pretext: valuer.UnsetIfEmpty(slack.Pretext),
|
||||
Fallback: valuer.UnsetIfEmpty(slack.Fallback),
|
||||
Footer: valuer.UnsetIfEmpty(slack.Footer),
|
||||
Fields: newChannelSlackFields(slack.Fields),
|
||||
Actions: newChannelSlackActions(slack.Actions),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newUpstreamSlackFields(fields []ChannelSlackField) []*config.SlackField {
|
||||
if len(fields) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
upstream := make([]*config.SlackField, 0, len(fields))
|
||||
for _, field := range fields {
|
||||
upstream = append(upstream, &config.SlackField{Title: field.Title, Value: field.Value, Short: field.Short})
|
||||
}
|
||||
|
||||
return upstream
|
||||
}
|
||||
|
||||
func newChannelSlackFields(upstream []*config.SlackField) []ChannelSlackField {
|
||||
if len(upstream) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
fields := make([]ChannelSlackField, 0, len(upstream))
|
||||
for _, field := range upstream {
|
||||
fields = append(fields, ChannelSlackField{Title: field.Title, Value: field.Value, Short: field.Short})
|
||||
}
|
||||
|
||||
return fields
|
||||
}
|
||||
|
||||
func newUpstreamSlackActions(actions []ChannelSlackAction) []*config.SlackAction {
|
||||
if len(actions) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
upstream := make([]*config.SlackAction, 0, len(actions))
|
||||
for _, action := range actions {
|
||||
upstreamAction := &config.SlackAction{Type: action.Type, Text: action.Text, URL: action.URL, Style: action.Style, Name: action.Name, Value: action.Value}
|
||||
if action.Confirm != nil {
|
||||
upstreamAction.ConfirmField = &config.SlackConfirmationField{Text: action.Confirm.Text, Title: action.Confirm.Title, OkText: action.Confirm.OkText, DismissText: action.Confirm.DismissText}
|
||||
}
|
||||
upstream = append(upstream, upstreamAction)
|
||||
}
|
||||
|
||||
return upstream
|
||||
}
|
||||
|
||||
func newChannelSlackActions(upstream []*config.SlackAction) []ChannelSlackAction {
|
||||
if len(upstream) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
actions := make([]ChannelSlackAction, 0, len(upstream))
|
||||
for _, upstreamAction := range upstream {
|
||||
action := ChannelSlackAction{Type: upstreamAction.Type, Text: upstreamAction.Text, URL: upstreamAction.URL, Style: upstreamAction.Style, Name: upstreamAction.Name, Value: upstreamAction.Value}
|
||||
if upstreamAction.ConfirmField != nil {
|
||||
action.Confirm = &ChannelSlackConfirmation{Text: upstreamAction.ConfirmField.Text, Title: upstreamAction.ConfirmField.Title, OkText: upstreamAction.ConfirmField.OkText, DismissText: upstreamAction.ConfirmField.DismissText}
|
||||
}
|
||||
actions = append(actions, action)
|
||||
}
|
||||
|
||||
return actions
|
||||
}
|
||||
98
pkg/types/alertmanagertypes/channel_webhook.go
Normal file
98
pkg/types/alertmanagertypes/channel_webhook.go
Normal file
@@ -0,0 +1,98 @@
|
||||
package alertmanagertypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
)
|
||||
|
||||
// ChannelWebhookConfig splits apart the two authentication modes the legacy API
|
||||
// overloaded onto one password field, where an empty username meant the password
|
||||
// was really a bearer token. Username or Password may be set without the other,
|
||||
// as upstream allows, but not together with BearerToken.
|
||||
type ChannelWebhookConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
URL string `json:"url" required:"true" format:"password"`
|
||||
Username string `json:"username"`
|
||||
Password string `json:"password" format:"password"`
|
||||
BearerToken string `json:"bearerToken" format:"password"`
|
||||
}
|
||||
|
||||
func (c ChannelWebhookConfig) Validate() error {
|
||||
if c.URL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for a webhook channel")
|
||||
}
|
||||
|
||||
usesBasicAuth := c.Username != "" || c.Password != ""
|
||||
|
||||
if usesBasicAuth && c.BearerToken != "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.bearerToken cannot be combined with config.spec.username or config.spec.password")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelWebhookConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
webhook := &config.WebhookConfig{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultWebhookConfig.VSendResolved)},
|
||||
URL: config.SecretTemplateURL(c.URL),
|
||||
}
|
||||
|
||||
// Seeded from upstream's default rather than a zero value: FollowRedirects
|
||||
// and EnableHTTP2 marshal unconditionally, so a zero value would persist
|
||||
// them as false and read back as a config ChannelWebhookConfig cannot represent.
|
||||
switch {
|
||||
case c.Username != "" || c.Password != "":
|
||||
httpConfig := commoncfg.DefaultHTTPClientConfig
|
||||
httpConfig.BasicAuth = &commoncfg.BasicAuth{
|
||||
Username: c.Username,
|
||||
Password: commoncfg.Secret(c.Password),
|
||||
}
|
||||
webhook.HTTPConfig = &httpConfig
|
||||
case c.BearerToken != "":
|
||||
httpConfig := commoncfg.DefaultHTTPClientConfig
|
||||
httpConfig.Authorization = &commoncfg.Authorization{
|
||||
Type: bearerAuthorizationType,
|
||||
Credentials: commoncfg.Secret(c.BearerToken),
|
||||
}
|
||||
webhook.HTTPConfig = &httpConfig
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
WebhookConfigs: []*config.WebhookConfig{webhook},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelWebhookConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
upstream := receiver.WebhookConfigs[0]
|
||||
sendResolved := upstream.VSendResolved
|
||||
if err := rejectUnsupportedHTTPConfig(name, upstream.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rejectHTTPBasicAuthBeyondPassword(name, upstream.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rejectHTTPAuthorizationBeyondBearer(name, upstream.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
webhook := &ChannelWebhookConfig{
|
||||
SendResolved: &sendResolved,
|
||||
URL: string(upstream.URL),
|
||||
}
|
||||
|
||||
if upstream.HTTPConfig != nil {
|
||||
if basicAuth := upstream.HTTPConfig.BasicAuth; basicAuth != nil {
|
||||
webhook.Username = basicAuth.Username
|
||||
webhook.Password = string(basicAuth.Password)
|
||||
}
|
||||
if authorization := upstream.HTTPConfig.Authorization; authorization != nil {
|
||||
webhook.BearerToken = string(authorization.Credentials)
|
||||
}
|
||||
}
|
||||
|
||||
return webhook, nil
|
||||
}
|
||||
@@ -3,8 +3,6 @@ package alertmanagertypes
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"maps"
|
||||
"net/textproto"
|
||||
"net/url"
|
||||
"reflect"
|
||||
"slices"
|
||||
@@ -14,7 +12,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/prometheus/alertmanager/config"
|
||||
commoncfg "github.com/prometheus/common/config"
|
||||
"github.com/prometheus/common/model"
|
||||
"github.com/swaggest/jsonschema-go"
|
||||
)
|
||||
|
||||
@@ -205,808 +202,6 @@ type ChannelSpec interface {
|
||||
toUndefaultedReceiver(displayName string) (*Receiver, error)
|
||||
}
|
||||
|
||||
type ChannelSlackConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
APIURL string `json:"apiUrl" required:"true" format:"password"`
|
||||
Channel string `json:"channel"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Text valuer.UnsetOrNonEmptyString `json:"text"`
|
||||
Color valuer.UnsetOrNonEmptyString `json:"color"`
|
||||
TitleLink valuer.UnsetOrNonEmptyString `json:"titleLink"`
|
||||
Pretext valuer.UnsetOrNonEmptyString `json:"pretext"`
|
||||
Fallback valuer.UnsetOrNonEmptyString `json:"fallback"`
|
||||
Footer valuer.UnsetOrNonEmptyString `json:"footer"`
|
||||
Fields []ChannelSlackField `json:"fields,omitempty"`
|
||||
Actions []ChannelSlackAction `json:"actions,omitempty"`
|
||||
}
|
||||
|
||||
type ChannelSlackField struct {
|
||||
Title string `json:"title" required:"true"`
|
||||
Value string `json:"value" required:"true"`
|
||||
Short *bool `json:"short,omitempty"`
|
||||
}
|
||||
|
||||
// ChannelSlackAction is a link button when URL is set, otherwise a message
|
||||
// button that needs Name. Upstream clears whichever side is not in use.
|
||||
type ChannelSlackAction struct {
|
||||
Type string `json:"type" required:"true"`
|
||||
Text string `json:"text" required:"true"`
|
||||
URL string `json:"url"`
|
||||
Style string `json:"style"`
|
||||
Name string `json:"name"`
|
||||
Value string `json:"value"`
|
||||
Confirm *ChannelSlackConfirmation `json:"confirm,omitempty"`
|
||||
}
|
||||
|
||||
type ChannelSlackConfirmation struct {
|
||||
Text string `json:"text" required:"true"`
|
||||
Title string `json:"title"`
|
||||
OkText string `json:"okText"`
|
||||
DismissText string `json:"dismissText"`
|
||||
}
|
||||
|
||||
func (c ChannelSlackConfig) Validate() error {
|
||||
if c.APIURL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiUrl is required for a slack channel")
|
||||
}
|
||||
|
||||
for i, field := range c.Fields {
|
||||
if field.Title == "" || field.Value == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.fields[%d] requires title and value", i)
|
||||
}
|
||||
}
|
||||
|
||||
for i, action := range c.Actions {
|
||||
if action.Type == "" || action.Text == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires type and text", i)
|
||||
}
|
||||
if action.URL == "" && action.Name == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d] requires url or name", i)
|
||||
}
|
||||
if action.Confirm != nil && action.Confirm.Text == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.actions[%d].confirm requires text", i)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelSlackConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
apiURL, err := parseSecretURL(c.APIURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
SlackConfigs: []*config.SlackConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultSlackConfig.VSendResolved)},
|
||||
APIURL: apiURL,
|
||||
Channel: c.Channel,
|
||||
Title: c.Title.StringValue(),
|
||||
Text: c.Text.StringValue(),
|
||||
Color: c.Color.StringValue(),
|
||||
TitleLink: c.TitleLink.StringValue(),
|
||||
Pretext: c.Pretext.StringValue(),
|
||||
Fallback: c.Fallback.StringValue(),
|
||||
Footer: c.Footer.StringValue(),
|
||||
Fields: newUpstreamSlackFields(c.Fields),
|
||||
Actions: newUpstreamSlackActions(c.Actions),
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelSlackConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
slack := receiver.SlackConfigs[0]
|
||||
sendResolved := slack.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, slack.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelSlackConfig{
|
||||
SendResolved: &sendResolved,
|
||||
APIURL: formatSecretURL(slack.APIURL),
|
||||
Channel: slack.Channel,
|
||||
Title: valuer.UnsetIfEmpty(slack.Title),
|
||||
Text: valuer.UnsetIfEmpty(slack.Text),
|
||||
Color: valuer.UnsetIfEmpty(slack.Color),
|
||||
TitleLink: valuer.UnsetIfEmpty(slack.TitleLink),
|
||||
Pretext: valuer.UnsetIfEmpty(slack.Pretext),
|
||||
Fallback: valuer.UnsetIfEmpty(slack.Fallback),
|
||||
Footer: valuer.UnsetIfEmpty(slack.Footer),
|
||||
Fields: newChannelSlackFields(slack.Fields),
|
||||
Actions: newChannelSlackActions(slack.Actions),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newUpstreamSlackFields(fields []ChannelSlackField) []*config.SlackField {
|
||||
if len(fields) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
upstream := make([]*config.SlackField, 0, len(fields))
|
||||
for _, field := range fields {
|
||||
upstream = append(upstream, &config.SlackField{Title: field.Title, Value: field.Value, Short: field.Short})
|
||||
}
|
||||
|
||||
return upstream
|
||||
}
|
||||
|
||||
func newChannelSlackFields(upstream []*config.SlackField) []ChannelSlackField {
|
||||
if len(upstream) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
fields := make([]ChannelSlackField, 0, len(upstream))
|
||||
for _, field := range upstream {
|
||||
fields = append(fields, ChannelSlackField{Title: field.Title, Value: field.Value, Short: field.Short})
|
||||
}
|
||||
|
||||
return fields
|
||||
}
|
||||
|
||||
func newUpstreamSlackActions(actions []ChannelSlackAction) []*config.SlackAction {
|
||||
if len(actions) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
upstream := make([]*config.SlackAction, 0, len(actions))
|
||||
for _, action := range actions {
|
||||
upstreamAction := &config.SlackAction{Type: action.Type, Text: action.Text, URL: action.URL, Style: action.Style, Name: action.Name, Value: action.Value}
|
||||
if action.Confirm != nil {
|
||||
upstreamAction.ConfirmField = &config.SlackConfirmationField{Text: action.Confirm.Text, Title: action.Confirm.Title, OkText: action.Confirm.OkText, DismissText: action.Confirm.DismissText}
|
||||
}
|
||||
upstream = append(upstream, upstreamAction)
|
||||
}
|
||||
|
||||
return upstream
|
||||
}
|
||||
|
||||
func newChannelSlackActions(upstream []*config.SlackAction) []ChannelSlackAction {
|
||||
if len(upstream) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
actions := make([]ChannelSlackAction, 0, len(upstream))
|
||||
for _, upstreamAction := range upstream {
|
||||
action := ChannelSlackAction{Type: upstreamAction.Type, Text: upstreamAction.Text, URL: upstreamAction.URL, Style: upstreamAction.Style, Name: upstreamAction.Name, Value: upstreamAction.Value}
|
||||
if upstreamAction.ConfirmField != nil {
|
||||
action.Confirm = &ChannelSlackConfirmation{Text: upstreamAction.ConfirmField.Text, Title: upstreamAction.ConfirmField.Title, OkText: upstreamAction.ConfirmField.OkText, DismissText: upstreamAction.ConfirmField.DismissText}
|
||||
}
|
||||
actions = append(actions, action)
|
||||
}
|
||||
|
||||
return actions
|
||||
}
|
||||
|
||||
// ChannelEmailConfig carries no SMTP transport fields: the smarthost,
|
||||
// credentials and TLS settings come from the deployment's global config, so a
|
||||
// channel can only choose recipients and body.
|
||||
type ChannelEmailConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
To string `json:"to" required:"true"`
|
||||
HTML valuer.UnsetOrNonEmptyString `json:"html"`
|
||||
Headers map[string]string `json:"headers,omitempty"`
|
||||
}
|
||||
|
||||
func (c ChannelEmailConfig) Validate() error {
|
||||
if c.To == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.to is required for an email channel")
|
||||
}
|
||||
|
||||
// A read reports header names as textproto canonicalizes them, turning
|
||||
// "subject" into "Subject", so a name that is not already in that form is
|
||||
// rejected rather than answered with one the caller never sent.
|
||||
for _, header := range slices.Sorted(maps.Keys(c.Headers)) {
|
||||
if canonical := textproto.CanonicalMIMEHeaderKey(header); canonical != header {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.headers name %q must be written as %q", header, canonical)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelEmailConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
EmailConfigs: []*config.EmailConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultEmailConfig.VSendResolved)},
|
||||
To: c.To,
|
||||
HTML: c.HTML.StringValue(),
|
||||
Headers: c.Headers,
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelEmailConfigFromReceiver(_ string, receiver *Receiver) (ChannelSpec, error) {
|
||||
email := receiver.EmailConfigs[0]
|
||||
sendResolved := email.VSendResolved
|
||||
|
||||
return &ChannelEmailConfig{
|
||||
SendResolved: &sendResolved,
|
||||
To: email.To,
|
||||
HTML: valuer.UnsetIfEmpty(email.HTML),
|
||||
Headers: email.Headers,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// ChannelWebhookConfig splits apart the two authentication modes the legacy API
|
||||
// overloaded onto one password field, where an empty username meant the password
|
||||
// was really a bearer token. Username or Password may be set without the other,
|
||||
// as upstream allows, but not together with BearerToken.
|
||||
type ChannelWebhookConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
URL string `json:"url" required:"true" format:"password"`
|
||||
Username string `json:"username"`
|
||||
Password string `json:"password" format:"password"`
|
||||
BearerToken string `json:"bearerToken" format:"password"`
|
||||
}
|
||||
|
||||
func (c ChannelWebhookConfig) Validate() error {
|
||||
if c.URL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for a webhook channel")
|
||||
}
|
||||
|
||||
usesBasicAuth := c.Username != "" || c.Password != ""
|
||||
|
||||
if usesBasicAuth && c.BearerToken != "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.bearerToken cannot be combined with config.spec.username or config.spec.password")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelWebhookConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
webhook := &config.WebhookConfig{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultWebhookConfig.VSendResolved)},
|
||||
URL: config.SecretTemplateURL(c.URL),
|
||||
}
|
||||
|
||||
// Seeded from upstream's default rather than a zero value: FollowRedirects
|
||||
// and EnableHTTP2 marshal unconditionally, so a zero value would persist
|
||||
// them as false and read back as a config ChannelWebhookConfig cannot represent.
|
||||
switch {
|
||||
case c.Username != "" || c.Password != "":
|
||||
httpConfig := commoncfg.DefaultHTTPClientConfig
|
||||
httpConfig.BasicAuth = &commoncfg.BasicAuth{
|
||||
Username: c.Username,
|
||||
Password: commoncfg.Secret(c.Password),
|
||||
}
|
||||
webhook.HTTPConfig = &httpConfig
|
||||
case c.BearerToken != "":
|
||||
httpConfig := commoncfg.DefaultHTTPClientConfig
|
||||
httpConfig.Authorization = &commoncfg.Authorization{
|
||||
Type: bearerAuthorizationType,
|
||||
Credentials: commoncfg.Secret(c.BearerToken),
|
||||
}
|
||||
webhook.HTTPConfig = &httpConfig
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
WebhookConfigs: []*config.WebhookConfig{webhook},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelWebhookConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
upstream := receiver.WebhookConfigs[0]
|
||||
sendResolved := upstream.VSendResolved
|
||||
if err := rejectUnsupportedHTTPConfig(name, upstream.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rejectHTTPBasicAuthBeyondPassword(name, upstream.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := rejectHTTPAuthorizationBeyondBearer(name, upstream.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
webhook := &ChannelWebhookConfig{
|
||||
SendResolved: &sendResolved,
|
||||
URL: string(upstream.URL),
|
||||
}
|
||||
|
||||
if upstream.HTTPConfig != nil {
|
||||
if basicAuth := upstream.HTTPConfig.BasicAuth; basicAuth != nil {
|
||||
webhook.Username = basicAuth.Username
|
||||
webhook.Password = string(basicAuth.Password)
|
||||
}
|
||||
if authorization := upstream.HTTPConfig.Authorization; authorization != nil {
|
||||
webhook.BearerToken = string(authorization.Credentials)
|
||||
}
|
||||
}
|
||||
|
||||
return webhook, nil
|
||||
}
|
||||
|
||||
type ChannelPagerdutyConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
RoutingKey string `json:"routingKey" required:"true" format:"password"`
|
||||
URL string `json:"url"`
|
||||
Source valuer.UnsetOrNonEmptyString `json:"source"`
|
||||
Client valuer.UnsetOrNonEmptyString `json:"client"`
|
||||
ClientURL valuer.UnsetOrNonEmptyString `json:"clientUrl"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Severity string `json:"severity"`
|
||||
Component string `json:"component"`
|
||||
Group string `json:"group"`
|
||||
Class string `json:"class"`
|
||||
Details map[string]string `json:"details,omitempty"`
|
||||
}
|
||||
|
||||
func (c ChannelPagerdutyConfig) Validate() error {
|
||||
if c.RoutingKey == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.routingKey is required for a pagerduty channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelPagerdutyConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
var eventsURL *config.URL
|
||||
if c.URL != "" {
|
||||
parsed, err := parseUpstreamURL(c.URL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
eventsURL = parsed
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
PagerdutyConfigs: []*config.PagerdutyConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultPagerdutyConfig.VSendResolved)},
|
||||
RoutingKey: config.Secret(c.RoutingKey),
|
||||
URL: eventsURL,
|
||||
Source: c.Source.StringValue(),
|
||||
Client: c.Client.StringValue(),
|
||||
ClientURL: c.ClientURL.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Severity: c.Severity,
|
||||
Component: c.Component,
|
||||
Group: c.Group,
|
||||
Class: c.Class,
|
||||
Details: newUpstreamDetails(c.Details),
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelPagerdutyConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
pagerduty := receiver.PagerdutyConfigs[0]
|
||||
sendResolved := pagerduty.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, pagerduty.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var details map[string]string
|
||||
if len(pagerduty.Details) > 0 {
|
||||
extracted, err := extractStringDetails(name, pagerduty.Details)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
details = extracted
|
||||
}
|
||||
|
||||
return &ChannelPagerdutyConfig{
|
||||
SendResolved: &sendResolved,
|
||||
RoutingKey: string(pagerduty.RoutingKey),
|
||||
URL: formatUpstreamURL(pagerduty.URL),
|
||||
Source: valuer.UnsetIfEmpty(pagerduty.Source),
|
||||
Client: valuer.UnsetIfEmpty(pagerduty.Client),
|
||||
ClientURL: valuer.UnsetIfEmpty(pagerduty.ClientURL),
|
||||
Description: valuer.UnsetIfEmpty(pagerduty.Description),
|
||||
Severity: pagerduty.Severity,
|
||||
Component: pagerduty.Component,
|
||||
Group: pagerduty.Group,
|
||||
Class: pagerduty.Class,
|
||||
Details: details,
|
||||
}, nil
|
||||
}
|
||||
|
||||
type ChannelOpsgenieConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
APIKey string `json:"apiKey" required:"true" format:"password"`
|
||||
APIURL string `json:"apiUrl"`
|
||||
Message valuer.UnsetOrNonEmptyString `json:"message"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Source valuer.UnsetOrNonEmptyString `json:"source"`
|
||||
Details map[string]string `json:"details,omitempty"`
|
||||
Priority string `json:"priority"`
|
||||
}
|
||||
|
||||
func (c ChannelOpsgenieConfig) Validate() error {
|
||||
if c.APIKey == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for an opsgenie channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelOpsgenieConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
var apiURL *config.URL
|
||||
if c.APIURL != "" {
|
||||
parsed, err := parseUpstreamURL(c.APIURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
apiURL = parsed
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
OpsGenieConfigs: []*config.OpsGenieConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultOpsGenieConfig.VSendResolved)},
|
||||
APIKey: config.Secret(c.APIKey),
|
||||
APIURL: apiURL,
|
||||
Message: c.Message.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Source: c.Source.StringValue(),
|
||||
Priority: c.Priority,
|
||||
Details: c.Details,
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelOpsgenieConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
opsgenie := receiver.OpsGenieConfigs[0]
|
||||
sendResolved := opsgenie.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, opsgenie.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelOpsgenieConfig{
|
||||
SendResolved: &sendResolved,
|
||||
APIKey: string(opsgenie.APIKey),
|
||||
APIURL: formatUpstreamURL(opsgenie.APIURL),
|
||||
Message: valuer.UnsetIfEmpty(opsgenie.Message),
|
||||
Description: valuer.UnsetIfEmpty(opsgenie.Description),
|
||||
Source: valuer.UnsetIfEmpty(opsgenie.Source),
|
||||
Priority: opsgenie.Priority,
|
||||
Details: opsgenie.Details,
|
||||
}, nil
|
||||
}
|
||||
|
||||
type ChannelMSTeamsConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Text valuer.UnsetOrNonEmptyString `json:"text"`
|
||||
}
|
||||
|
||||
func (c ChannelMSTeamsConfig) Validate() error {
|
||||
if c.WebhookURL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for an msteams channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelMSTeamsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
webhookURL, err := parseSecretURL(c.WebhookURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Receiver{Receiver: &config.Receiver{
|
||||
Name: displayName,
|
||||
MSTeamsV2Configs: []*config.MSTeamsV2Config{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, config.DefaultMSTeamsV2Config.VSendResolved)},
|
||||
WebhookURL: webhookURL,
|
||||
Title: c.Title.StringValue(),
|
||||
Text: c.Text.StringValue(),
|
||||
}},
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func newChannelMSTeamsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
msteams := receiver.MSTeamsV2Configs[0]
|
||||
sendResolved := msteams.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, msteams.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelMSTeamsConfig{
|
||||
SendResolved: &sendResolved,
|
||||
WebhookURL: formatSecretURL(msteams.WebhookURL),
|
||||
Title: valuer.UnsetIfEmpty(msteams.Title),
|
||||
Text: valuer.UnsetIfEmpty(msteams.Text),
|
||||
}, nil
|
||||
}
|
||||
|
||||
type ChannelGoogleChatConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
WebhookURL string `json:"webhookUrl" required:"true" format:"password"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Text valuer.UnsetOrNonEmptyString `json:"text"`
|
||||
}
|
||||
|
||||
func (c ChannelGoogleChatConfig) Validate() error {
|
||||
if c.WebhookURL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.webhookUrl is required for a googlechat channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelGoogleChatConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
webhookURL, err := parseSecretURL(c.WebhookURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
GoogleChatConfigs: []*GoogleChatReceiverConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultGoogleChatReceiverConfig.VSendResolved)},
|
||||
WebhookURL: webhookURL,
|
||||
Title: c.Title.StringValue(),
|
||||
Text: c.Text.StringValue(),
|
||||
}},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelGoogleChatConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
googlechat := receiver.GoogleChatConfigs[0]
|
||||
sendResolved := googlechat.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, googlechat.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelGoogleChatConfig{
|
||||
SendResolved: &sendResolved,
|
||||
WebhookURL: formatSecretURL(googlechat.WebhookURL),
|
||||
Title: valuer.UnsetIfEmpty(googlechat.Title),
|
||||
Text: valuer.UnsetIfEmpty(googlechat.Text),
|
||||
}, nil
|
||||
}
|
||||
|
||||
type ChannelJiraConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
// Site is the Jira Cloud base URL, https://<site>.atlassian.net. Only Jira
|
||||
// Cloud is supported; the REST base is derived from it.
|
||||
Site string `json:"site" required:"true"`
|
||||
Project string `json:"project" required:"true"`
|
||||
IssueType string `json:"issueType" required:"true"`
|
||||
Summary valuer.UnsetOrNonEmptyString `json:"summary"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Priority string `json:"priority"`
|
||||
Labels []string `json:"labels,omitempty"`
|
||||
ResolveTransition string `json:"resolveTransition"`
|
||||
ReopenTransition string `json:"reopenTransition"`
|
||||
ReopenDuration valuer.UnsetOrNonEmptyString `json:"reopenDuration"`
|
||||
WontFixResolution string `json:"wontFixResolution"`
|
||||
CustomFields map[string]any `json:"customFields,omitempty"`
|
||||
|
||||
Email string `json:"email" required:"true"`
|
||||
APIToken string `json:"apiToken" required:"true" format:"password"`
|
||||
}
|
||||
|
||||
func (c ChannelJiraConfig) Validate() error {
|
||||
for _, required := range []struct {
|
||||
value string
|
||||
field string
|
||||
}{
|
||||
{c.Site, "site"},
|
||||
{c.Project, "project"},
|
||||
{c.IssueType, "issueType"},
|
||||
{c.Email, "email"},
|
||||
{c.APIToken, "apiToken"},
|
||||
} {
|
||||
if required.value == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.%s is required for a jira channel", required.field)
|
||||
}
|
||||
}
|
||||
|
||||
if !c.ReopenDuration.IsZero() {
|
||||
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
|
||||
if err != nil {
|
||||
return errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q is not a valid duration", c.ReopenDuration)
|
||||
}
|
||||
|
||||
// A read reports the duration as model.Duration formats it, collapsing
|
||||
// "72h" into "3d", so a value that is not already in that form is rejected
|
||||
// rather than answered with one the caller never sent.
|
||||
if canonical := reopenDuration.String(); canonical != c.ReopenDuration.StringValue() {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.reopenDuration %q must be written as %q", c.ReopenDuration, canonical)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelJiraConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
// Seeded from upstream's default rather than a zero value: FollowRedirects
|
||||
// and EnableHTTP2 marshal unconditionally, so a zero value would persist them
|
||||
// as false and read back as a config ChannelJiraConfig cannot represent.
|
||||
httpConfig := commoncfg.DefaultHTTPClientConfig
|
||||
httpConfig.BasicAuth = &commoncfg.BasicAuth{
|
||||
Username: c.Email,
|
||||
Password: commoncfg.Secret(c.APIToken),
|
||||
}
|
||||
|
||||
jira := &JiraReceiverConfig{
|
||||
// JiraReceiverConfig seeds no send_resolved of its own, so unset means off.
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, false)},
|
||||
Site: c.Site,
|
||||
Project: c.Project,
|
||||
IssueType: c.IssueType,
|
||||
Summary: c.Summary.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Priority: c.Priority,
|
||||
Labels: c.Labels,
|
||||
ResolveTransition: c.ResolveTransition,
|
||||
ReopenTransition: c.ReopenTransition,
|
||||
WontFixResolution: c.WontFixResolution,
|
||||
CustomFields: c.CustomFields,
|
||||
HTTPConfig: &httpConfig,
|
||||
}
|
||||
|
||||
if !c.ReopenDuration.IsZero() {
|
||||
reopenDuration, err := model.ParseDuration(c.ReopenDuration.StringValue())
|
||||
if err != nil {
|
||||
return nil, errors.WrapInvalidInputf(err, ErrCodeAlertmanagerChannelInvalid, "parse reopenDuration %q", c.ReopenDuration)
|
||||
}
|
||||
jira.ReopenDuration = reopenDuration
|
||||
}
|
||||
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
JiraConfigs: []*JiraReceiverConfig{jira},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelJiraConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
jira := receiver.JiraConfigs[0]
|
||||
sendResolved := jira.VSendResolved
|
||||
|
||||
if err := rejectUnsupportedHTTPConfig(name, jira.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if jira.HTTPConfig != nil && jira.HTTPConfig.Authorization != nil {
|
||||
return nil, errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "channel %q sets http_config.authorization, which is not supported", name)
|
||||
}
|
||||
|
||||
if err := rejectHTTPBasicAuthBeyondPassword(name, jira.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
spec := &ChannelJiraConfig{
|
||||
SendResolved: &sendResolved,
|
||||
Site: jira.Site,
|
||||
Project: jira.Project,
|
||||
IssueType: jira.IssueType,
|
||||
Summary: valuer.UnsetIfEmpty(jira.Summary),
|
||||
Description: valuer.UnsetIfEmpty(jira.Description),
|
||||
Priority: jira.Priority,
|
||||
Labels: jira.Labels,
|
||||
ResolveTransition: jira.ResolveTransition,
|
||||
ReopenTransition: jira.ReopenTransition,
|
||||
ReopenDuration: valuer.UnsetIfEmpty(jira.ReopenDuration.String()),
|
||||
WontFixResolution: jira.WontFixResolution,
|
||||
CustomFields: jira.CustomFields,
|
||||
}
|
||||
|
||||
if jira.HTTPConfig != nil && jira.HTTPConfig.BasicAuth != nil {
|
||||
spec.Email = jira.HTTPConfig.BasicAuth.Username
|
||||
spec.APIToken = string(jira.HTTPConfig.BasicAuth.Password)
|
||||
}
|
||||
|
||||
return spec, nil
|
||||
}
|
||||
|
||||
// ChannelJSMOpsConfig carries no API URL: JSM Ops is a single global gateway
|
||||
// keyed by the integration API key, which the notifier pins itself.
|
||||
type ChannelJSMOpsConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
APIKey string `json:"apiKey" required:"true" format:"password"`
|
||||
Message valuer.UnsetOrNonEmptyString `json:"message"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Priority string `json:"priority"`
|
||||
// Tags is the comma-separated list JSM Ops attaches to the alert.
|
||||
Tags valuer.UnsetOrNonEmptyString `json:"tags"`
|
||||
}
|
||||
|
||||
func (c ChannelJSMOpsConfig) Validate() error {
|
||||
if c.APIKey == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.apiKey is required for a jsmops channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelJSMOpsConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
JSMOpsConfigs: []*JSMOpsReceiverConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultJSMOpsReceiverConfig.VSendResolved)},
|
||||
APIKey: config.Secret(c.APIKey),
|
||||
Message: c.Message.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Priority: c.Priority,
|
||||
Tags: c.Tags.StringValue(),
|
||||
}},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelJSMOpsConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
jsmops := receiver.JSMOpsConfigs[0]
|
||||
sendResolved := jsmops.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, jsmops.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelJSMOpsConfig{
|
||||
SendResolved: &sendResolved,
|
||||
APIKey: string(jsmops.APIKey),
|
||||
Message: valuer.UnsetIfEmpty(jsmops.Message),
|
||||
Description: valuer.UnsetIfEmpty(jsmops.Description),
|
||||
Priority: jsmops.Priority,
|
||||
Tags: valuer.UnsetIfEmpty(jsmops.Tags),
|
||||
}, nil
|
||||
}
|
||||
|
||||
type ChannelIncidentIOConfig struct {
|
||||
SendResolved *bool `json:"sendResolved,omitempty"`
|
||||
URL string `json:"url" required:"true"`
|
||||
Token string `json:"token" required:"true" format:"password"`
|
||||
Title valuer.UnsetOrNonEmptyString `json:"title"`
|
||||
Description valuer.UnsetOrNonEmptyString `json:"description"`
|
||||
Metadata map[string]string `json:"metadata,omitempty"`
|
||||
}
|
||||
|
||||
func (c ChannelIncidentIOConfig) Validate() error {
|
||||
if c.URL == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.url is required for an incidentio channel")
|
||||
}
|
||||
|
||||
if c.Token == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeAlertmanagerChannelInvalid, "config.spec.token is required for an incidentio channel")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c ChannelIncidentIOConfig) toUndefaultedReceiver(displayName string) (*Receiver, error) {
|
||||
return &Receiver{
|
||||
Receiver: &config.Receiver{Name: displayName},
|
||||
IncidentIOConfigs: []*IncidentIOReceiverConfig{{
|
||||
NotifierConfig: config.NotifierConfig{VSendResolved: resolveSendResolved(c.SendResolved, DefaultIncidentIOReceiverConfig.VSendResolved)},
|
||||
URL: c.URL,
|
||||
Token: config.Secret(c.Token),
|
||||
Title: c.Title.StringValue(),
|
||||
Description: c.Description.StringValue(),
|
||||
Metadata: c.Metadata,
|
||||
}},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newChannelIncidentIOConfigFromReceiver(name string, receiver *Receiver) (ChannelSpec, error) {
|
||||
incidentio := receiver.IncidentIOConfigs[0]
|
||||
sendResolved := incidentio.VSendResolved
|
||||
|
||||
if err := rejectAnyHTTPAuth(name, incidentio.HTTPConfig); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &ChannelIncidentIOConfig{
|
||||
SendResolved: &sendResolved,
|
||||
URL: incidentio.URL,
|
||||
Token: string(incidentio.Token),
|
||||
Title: valuer.UnsetIfEmpty(incidentio.Title),
|
||||
Description: valuer.UnsetIfEmpty(incidentio.Description),
|
||||
Metadata: incidentio.Metadata,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// ════════════════════════════════════════════════════════════════════════
|
||||
// Helpers
|
||||
// ════════════════════════════════════════════════════════════════════════
|
||||
|
||||
@@ -258,6 +258,6 @@ type TokenStore interface {
|
||||
// Delete a token by userID.
|
||||
DeleteByUserID(context.Context, valuer.UUID) error
|
||||
|
||||
// Update last observed at by access token.
|
||||
UpdateLastObservedAtByAccessToken(context.Context, []map[string]any) error
|
||||
// Update last observed at of the given tokens.
|
||||
UpdateLastObservedAt(context.Context, []*StorableToken) error
|
||||
}
|
||||
|
||||
@@ -208,6 +208,35 @@ func NewGettableUnmappedModels(items []*UnmappedModel) *GettableUnmappedModels {
|
||||
}
|
||||
}
|
||||
|
||||
func (u *UpdatableLLMPricingRule) UnmarshalJSON(data []byte) error {
|
||||
type Alias UpdatableLLMPricingRule
|
||||
|
||||
var temp Alias
|
||||
if err := json.Unmarshal(data, &temp); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
*u = UpdatableLLMPricingRule(temp)
|
||||
return u.Validate()
|
||||
}
|
||||
|
||||
// Validate mirrors the collector's pattern check: at least one pattern, none
|
||||
// empty, all valid path.Match globs.
|
||||
func (u *UpdatableLLMPricingRule) Validate() error {
|
||||
if len(u.ModelPattern) == 0 {
|
||||
return errors.Newf(errors.TypeInvalidInput, ErrCodePricingRuleInvalidInput, "model %q: modelPattern must contain at least one pattern", u.Model)
|
||||
}
|
||||
for _, p := range u.ModelPattern {
|
||||
if p == "" {
|
||||
return errors.Newf(errors.TypeInvalidInput, ErrCodePricingRuleInvalidInput, "model %q: modelPattern must not contain an empty pattern", u.Model)
|
||||
}
|
||||
if _, err := path.Match(p, ""); err != nil {
|
||||
return errors.Newf(errors.TypeInvalidInput, ErrCodePricingRuleInvalidInput, "model %q: modelPattern %q is not a valid glob", u.Model, p)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func NewLLMPricingRuleFromUpdatable(u *UpdatableLLMPricingRule, orgID valuer.UUID, userEmail string, now time.Time) *LLMPricingRule {
|
||||
id := valuer.GenerateUUID()
|
||||
if u.ID != nil {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package llmpricingruletypes
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
@@ -126,3 +127,34 @@ func TestGenerateCollectorConfig_EmptyInputPassthrough(t *testing.T) {
|
||||
assert.Equal(t, in, out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdatableLLMPricingRuleUnmarshalJSON(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
pattern string
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "valid", pattern: `["gpt-4o*", "gpt-4o"]`},
|
||||
{name: "missing", pattern: ``, wantErr: true},
|
||||
{name: "null", pattern: `null`, wantErr: true},
|
||||
{name: "empty_list", pattern: `[]`, wantErr: true},
|
||||
{name: "empty_entry", pattern: `["gpt-4o*", ""]`, wantErr: true},
|
||||
{name: "bad_glob", pattern: `["gpt-["]`, wantErr: true},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
body := `{"modelName": "gpt-4o"}`
|
||||
if tc.pattern != "" {
|
||||
body = `{"modelName": "gpt-4o", "modelPattern": ` + tc.pattern + `}`
|
||||
}
|
||||
var req UpdatableLLMPricingRules
|
||||
err := json.Unmarshal([]byte(`{"rules": [`+body+`]}`), &req)
|
||||
if tc.wantErr {
|
||||
assert.Error(t, err)
|
||||
} else {
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -49,55 +49,66 @@ func (o ListOrder) IsValid() bool {
|
||||
return slices.ContainsFunc(o.Enum(), func(v any) bool { return v == o })
|
||||
}
|
||||
|
||||
// ListFilter is the rule listing state shared by the v3 list params and saved views.
|
||||
type ListFilter struct {
|
||||
Query string `query:"query" json:"query"`
|
||||
// gin cannot bind a slice of valuer enums; GetAlertStates converts these.
|
||||
States []string `query:"states" json:"states" nullable:"false"`
|
||||
Sort ListSort `query:"sort" json:"sort"`
|
||||
Order ListOrder `query:"order" json:"order"`
|
||||
type ListRulesParams struct {
|
||||
Query string `query:"query"`
|
||||
// gin cannot bind a slice of valuer enums; AlertStates converts these.
|
||||
States []string `query:"states"`
|
||||
Sort ListSort `query:"sort"`
|
||||
Order ListOrder `query:"order"`
|
||||
Limit int `query:"limit"`
|
||||
Offset int `query:"offset"`
|
||||
}
|
||||
|
||||
// Validate normalizes in place; zero sort/order get the defaults, nil states an empty slice.
|
||||
func (f *ListFilter) Validate() error {
|
||||
if f.States == nil {
|
||||
f.States = []string{}
|
||||
}
|
||||
|
||||
if n := utf8.RuneCountInString(f.Query); n > MaxListQueryLen {
|
||||
// Validate normalizes in place; an over-max limit is clamped, not rejected.
|
||||
func (p *ListRulesParams) Validate() error {
|
||||
if n := utf8.RuneCountInString(p.Query); n > MaxListQueryLen {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"query cannot be longer than %d characters, got %d", MaxListQueryLen, n)
|
||||
}
|
||||
|
||||
if _, err := f.GetAlertStates(); err != nil {
|
||||
if p.Sort.IsZero() {
|
||||
p.Sort = ListSortUpdatedAt
|
||||
} else if !p.Sort.IsValid() {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid sort %q, expected one of: `updated_at`, `created_at`, `name`, `state`, `severity`", p.Sort)
|
||||
}
|
||||
|
||||
if p.Order.IsZero() {
|
||||
p.Order = ListOrderDesc
|
||||
} else if !p.Order.IsValid() {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid order %q, expected `asc` or `desc`", p.Order)
|
||||
}
|
||||
|
||||
if p.Limit == 0 {
|
||||
p.Limit = DefaultListLimit
|
||||
} else if p.Limit < 0 {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid limit %d, must be a positive integer", p.Limit)
|
||||
} else if p.Limit > MaxListLimit {
|
||||
p.Limit = MaxListLimit
|
||||
}
|
||||
|
||||
if p.Offset < 0 {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid offset %d, must be a non-negative integer", p.Offset)
|
||||
}
|
||||
|
||||
if _, err := p.GetAlertStates(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if f.Sort.IsZero() {
|
||||
f.Sort = ListSortUpdatedAt
|
||||
} else if !f.Sort.IsValid() {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid sort %q, expected one of: `updated_at`, `created_at`, `name`, `state`, `severity`", f.Sort)
|
||||
}
|
||||
|
||||
if f.Order.IsZero() {
|
||||
f.Order = ListOrderDesc
|
||||
} else if !f.Order.IsValid() {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid order %q, expected `asc` or `desc`", f.Order)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetAlertStates parses States; empty means no state filtering.
|
||||
func (f *ListFilter) GetAlertStates() ([]AlertState, error) {
|
||||
if len(f.States) == 0 {
|
||||
func (p *ListRulesParams) GetAlertStates() ([]AlertState, error) {
|
||||
if len(p.States) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
states := make([]AlertState, 0, len(f.States))
|
||||
for _, raw := range f.States {
|
||||
states := make([]AlertState, 0, len(p.States))
|
||||
for _, raw := range p.States {
|
||||
state, err := parseAlertState(raw)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -116,32 +127,3 @@ func parseAlertState(raw string) (AlertState, error) {
|
||||
}
|
||||
return state, nil
|
||||
}
|
||||
|
||||
type ListRulesParams struct {
|
||||
ListFilter
|
||||
Limit int `query:"limit"`
|
||||
Offset int `query:"offset"`
|
||||
}
|
||||
|
||||
// Validate normalizes in place; an over-max limit is clamped, not rejected.
|
||||
func (p *ListRulesParams) Validate() error {
|
||||
if err := p.ListFilter.Validate(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if p.Limit == 0 {
|
||||
p.Limit = DefaultListLimit
|
||||
} else if p.Limit < 0 {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid limit %d, must be a positive integer", p.Limit)
|
||||
} else if p.Limit > MaxListLimit {
|
||||
p.Limit = MaxListLimit
|
||||
}
|
||||
|
||||
if p.Offset < 0 {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleListInvalid,
|
||||
"invalid offset %d, must be a non-negative integer", p.Offset)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -27,7 +27,7 @@ func TestListRulesParamsValidate(t *testing.T) {
|
||||
},
|
||||
{
|
||||
name: "ExplicitValues_Kept",
|
||||
params: ListRulesParams{ListFilter: ListFilter{Sort: ListSortSeverity, Order: ListOrderAsc}, Limit: 50, Offset: 100},
|
||||
params: ListRulesParams{Sort: ListSortSeverity, Order: ListOrderAsc, Limit: 50, Offset: 100},
|
||||
wantSort: ListSortSeverity,
|
||||
wantOrder: ListOrderAsc,
|
||||
wantLimit: 50,
|
||||
@@ -41,17 +41,17 @@ func TestListRulesParamsValidate(t *testing.T) {
|
||||
},
|
||||
{
|
||||
name: "InvalidState_Rejected",
|
||||
params: ListRulesParams{ListFilter: ListFilter{States: []string{"bogus"}}},
|
||||
params: ListRulesParams{States: []string{"bogus"}},
|
||||
wantErr: `invalid state "bogus"`,
|
||||
},
|
||||
{
|
||||
name: "InvalidSort_Rejected",
|
||||
params: ListRulesParams{ListFilter: ListFilter{Sort: ListSort{valuer.NewString("bogus")}}},
|
||||
params: ListRulesParams{Sort: ListSort{valuer.NewString("bogus")}},
|
||||
wantErr: "invalid sort",
|
||||
},
|
||||
{
|
||||
name: "InvalidOrder_Rejected",
|
||||
params: ListRulesParams{ListFilter: ListFilter{Order: ListOrder{valuer.NewString("bogus")}}},
|
||||
params: ListRulesParams{Order: ListOrder{valuer.NewString("bogus")}},
|
||||
wantErr: "invalid order",
|
||||
},
|
||||
{
|
||||
@@ -66,7 +66,7 @@ func TestListRulesParamsValidate(t *testing.T) {
|
||||
},
|
||||
{
|
||||
name: "OverLongQuery_Rejected",
|
||||
params: ListRulesParams{ListFilter: ListFilter{Query: strings.Repeat("a", MaxListQueryLen+1)}},
|
||||
params: ListRulesParams{Query: strings.Repeat("a", MaxListQueryLen+1)},
|
||||
wantErr: "query cannot be longer",
|
||||
},
|
||||
}
|
||||
@@ -112,7 +112,7 @@ func TestListRulesParamsAlertStates(t *testing.T) {
|
||||
|
||||
for _, tc := range testCases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
params := ListRulesParams{ListFilter: ListFilter{States: tc.states}}
|
||||
params := ListRulesParams{States: tc.states}
|
||||
states, err := params.GetAlertStates()
|
||||
if tc.wantErr != "" {
|
||||
require.Error(t, err)
|
||||
|
||||
@@ -65,10 +65,4 @@ type RuleStore interface {
|
||||
GetStoredRuleLabels(context.Context, string) ([]string, error)
|
||||
GetStoredRule(context.Context, valuer.UUID, valuer.UUID) (*StorableRule, error)
|
||||
GetStoredRulesByMetricName(context.Context, string, string) ([]RuleAlert, error)
|
||||
|
||||
CreateRuleView(context.Context, *StorableRuleView) error
|
||||
GetRuleView(context.Context, valuer.UUID, valuer.UUID) (*StorableRuleView, error)
|
||||
ListRuleViews(context.Context, valuer.UUID) ([]*StorableRuleView, error)
|
||||
UpdateRuleView(context.Context, *StorableRuleView) error
|
||||
DeleteRuleView(context.Context, valuer.UUID, valuer.UUID) error
|
||||
}
|
||||
|
||||
@@ -1,169 +0,0 @@
|
||||
package ruletypes
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/types"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/uptrace/bun"
|
||||
)
|
||||
|
||||
const (
|
||||
RuleViewSchemaVersion = "v1"
|
||||
MaxRuleViewNameLen = 64
|
||||
)
|
||||
|
||||
var (
|
||||
ErrCodeRuleViewInvalidInput = errors.MustNewCode("rule_view_invalid_input")
|
||||
ErrCodeRuleViewNotFound = errors.MustNewCode("rule_view_not_found")
|
||||
)
|
||||
|
||||
type StorableRuleView struct {
|
||||
bun.BaseModel `bun:"table:rule_view,alias:rule_view"`
|
||||
|
||||
types.Identifiable
|
||||
types.TimeAuditable
|
||||
|
||||
Name string `bun:"name,type:text,notnull"`
|
||||
Data storableRuleViewData `bun:"data,type:text,notnull"`
|
||||
OrgID valuer.UUID `bun:"org_id,type:text,notnull"`
|
||||
}
|
||||
|
||||
func (s *StorableRuleView) ToGettableRuleView() *GettableRuleView {
|
||||
return &GettableRuleView{
|
||||
ID: s.ID,
|
||||
Name: s.Name,
|
||||
Data: s.Data.toRuleViewData(),
|
||||
OrgID: s.OrgID,
|
||||
CreatedAt: s.CreatedAt,
|
||||
UpdatedAt: s.UpdatedAt,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *StorableRuleView) Update(updatable UpdatableRuleView) {
|
||||
s.Name = updatable.Name
|
||||
s.Data = newStorableRuleViewData(updatable.Data)
|
||||
s.UpdatedAt = time.Now()
|
||||
}
|
||||
|
||||
type GettableRuleView struct {
|
||||
ID valuer.UUID `json:"id" required:"true"`
|
||||
Name string `json:"name" required:"true"`
|
||||
Data RuleViewData `json:"data" required:"true"`
|
||||
OrgID valuer.UUID `json:"orgId" required:"true"`
|
||||
CreatedAt time.Time `json:"createdAt" required:"true"`
|
||||
UpdatedAt time.Time `json:"updatedAt" required:"true"`
|
||||
}
|
||||
|
||||
// RuleViewData holds the rule listing state (ListRulesParams minus pagination) a view replays.
|
||||
type RuleViewData struct {
|
||||
Version string `json:"version" required:"true"`
|
||||
ListFilter
|
||||
}
|
||||
|
||||
func (d *RuleViewData) Validate() error {
|
||||
if d.Version != RuleViewSchemaVersion {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput,
|
||||
"version must be %q, got %q", RuleViewSchemaVersion, d.Version)
|
||||
}
|
||||
return d.ListFilter.Validate()
|
||||
}
|
||||
|
||||
type PostableRuleView struct {
|
||||
Name string `json:"name" required:"true"`
|
||||
Data RuleViewData `json:"data" required:"true"`
|
||||
}
|
||||
|
||||
func (p *PostableRuleView) UnmarshalJSON(data []byte) error {
|
||||
dec := json.NewDecoder(bytes.NewReader(data))
|
||||
dec.DisallowUnknownFields()
|
||||
type alias PostableRuleView
|
||||
var tmp alias
|
||||
if err := dec.Decode(&tmp); err != nil {
|
||||
return errors.WrapInvalidInputf(err, ErrCodeRuleViewInvalidInput, "invalid saved view request body").WithAdditional(err.Error())
|
||||
}
|
||||
*p = PostableRuleView(tmp)
|
||||
return p.Validate()
|
||||
}
|
||||
|
||||
func (p *PostableRuleView) Validate() error {
|
||||
if err := validateRuleViewName(p.Name); err != nil {
|
||||
return err
|
||||
}
|
||||
return p.Data.Validate()
|
||||
}
|
||||
|
||||
func (p PostableRuleView) ToStorableRuleView(orgID valuer.UUID) *StorableRuleView {
|
||||
now := time.Now()
|
||||
return &StorableRuleView{
|
||||
Identifiable: types.Identifiable{ID: valuer.GenerateUUID()},
|
||||
TimeAuditable: types.TimeAuditable{CreatedAt: now, UpdatedAt: now},
|
||||
Name: p.Name,
|
||||
Data: newStorableRuleViewData(p.Data),
|
||||
OrgID: orgID,
|
||||
}
|
||||
}
|
||||
|
||||
type UpdatableRuleView = PostableRuleView
|
||||
|
||||
func NewGettableRuleViewsFromStorableRuleViews(storables []*StorableRuleView) []*GettableRuleView {
|
||||
views := make([]*GettableRuleView, 0, len(storables))
|
||||
for _, storable := range storables {
|
||||
views = append(views, storable.ToGettableRuleView())
|
||||
}
|
||||
return views
|
||||
}
|
||||
|
||||
type ListableRuleViews struct {
|
||||
Views []*GettableRuleView `json:"views" required:"true" nullable:"false"`
|
||||
}
|
||||
|
||||
// storableRuleViewData owns the persisted blob format; wire tag changes must not affect stored rows.
|
||||
type storableRuleViewData struct {
|
||||
Version string `json:"version"`
|
||||
Query string `json:"query"`
|
||||
States []string `json:"states"`
|
||||
Sort string `json:"sort"`
|
||||
Order string `json:"order"`
|
||||
}
|
||||
|
||||
func newStorableRuleViewData(data RuleViewData) storableRuleViewData {
|
||||
return storableRuleViewData{
|
||||
Version: data.Version,
|
||||
Query: data.Query,
|
||||
States: data.States,
|
||||
Sort: data.Sort.StringValue(),
|
||||
Order: data.Order.StringValue(),
|
||||
}
|
||||
}
|
||||
|
||||
func (d storableRuleViewData) toRuleViewData() RuleViewData {
|
||||
return RuleViewData{
|
||||
Version: d.Version,
|
||||
ListFilter: ListFilter{
|
||||
Query: d.Query,
|
||||
States: d.States,
|
||||
Sort: ListSort{valuer.NewString(d.Sort)},
|
||||
Order: ListOrder{valuer.NewString(d.Order)},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func validateRuleViewName(name string) error {
|
||||
if strings.TrimSpace(name) == "" {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput, "name is required")
|
||||
}
|
||||
if name != strings.TrimSpace(name) {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput, "name must not have leading or trailing whitespace")
|
||||
}
|
||||
if n := utf8.RuneCountInString(name); n > MaxRuleViewNameLen {
|
||||
return errors.NewInvalidInputf(ErrCodeRuleViewInvalidInput,
|
||||
"name must be at most %d characters, got %d", MaxRuleViewNameLen, n)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1,210 +0,0 @@
|
||||
package ruletypes
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestRuleViewDataValidate(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
data RuleViewData
|
||||
expectError bool
|
||||
}{
|
||||
{
|
||||
name: "AllFieldsSet_Valid",
|
||||
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Query: "name CONTAINS 'prod'", States: []string{"firing", "pending"}, Sort: ListSortName, Order: ListOrderAsc}},
|
||||
expectError: false,
|
||||
},
|
||||
{
|
||||
name: "ZeroStatesSortOrder_Valid",
|
||||
data: RuleViewData{Version: RuleViewSchemaVersion},
|
||||
expectError: false,
|
||||
},
|
||||
{
|
||||
name: "QueryOverCap_Rejected",
|
||||
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Query: strings.Repeat("x", MaxListQueryLen+1)}},
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "WrongVersion_Rejected",
|
||||
data: RuleViewData{Version: "v2"},
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "EmptyVersion_Rejected",
|
||||
data: RuleViewData{},
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "UnknownState_Rejected",
|
||||
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"exploding"}}},
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "UnknownSort_Rejected",
|
||||
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Sort: ListSort{valuer.NewString("bogus")}}},
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "UnknownOrder_Rejected",
|
||||
data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Order: ListOrder{valuer.NewString("sideways")}}},
|
||||
expectError: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
err := testCase.data.Validate()
|
||||
if testCase.expectError {
|
||||
assert.Error(t, err)
|
||||
} else {
|
||||
assert.NoError(t, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRuleViewDataValidateDefaults(t *testing.T) {
|
||||
data := RuleViewData{Version: RuleViewSchemaVersion}
|
||||
require.NoError(t, data.Validate())
|
||||
assert.Equal(t, ListSortUpdatedAt, data.Sort)
|
||||
assert.Equal(t, ListOrderDesc, data.Order)
|
||||
assert.Equal(t, []string{}, data.States)
|
||||
}
|
||||
|
||||
func TestPostableRuleViewUnmarshalJSON(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
body string
|
||||
expectError bool
|
||||
expectedErrMsg string
|
||||
expectedName string
|
||||
}{
|
||||
{
|
||||
name: "ValidBody_NameKeptAsIs",
|
||||
body: `{"name":"my view","data":{"version":"v1","query":"severity = 'critical'","states":["firing"],"sort":"name","order":"asc"}}`,
|
||||
expectError: false,
|
||||
expectedName: "my view",
|
||||
},
|
||||
{
|
||||
name: "NameSurroundingWhitespace_Rejected",
|
||||
body: `{"name":" my view ","data":{"version":"v1"}}`,
|
||||
expectError: true,
|
||||
expectedErrMsg: "name must not have leading or trailing whitespace",
|
||||
},
|
||||
{
|
||||
name: "UnknownField_Rejected",
|
||||
body: `{"name":"my view","data":{"version":"v1"},"extra":true}`,
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "BlankName_Rejected",
|
||||
body: `{"name":" ","data":{"version":"v1"}}`,
|
||||
expectError: true,
|
||||
expectedErrMsg: "name is required",
|
||||
},
|
||||
{
|
||||
name: "NameOverMaxLength_Rejected",
|
||||
body: `{"name":"` + strings.Repeat("x", MaxRuleViewNameLen+1) + `","data":{"version":"v1"}}`,
|
||||
expectError: true,
|
||||
expectedErrMsg: "name must be at most",
|
||||
},
|
||||
{
|
||||
name: "InvalidDataVersion_Rejected",
|
||||
body: `{"name":"my view","data":{"version":"v9"}}`,
|
||||
expectError: true,
|
||||
},
|
||||
{
|
||||
name: "InvalidState_Rejected",
|
||||
body: `{"name":"my view","data":{"version":"v1","states":["exploding"]}}`,
|
||||
expectError: true,
|
||||
expectedErrMsg: "invalid state",
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
var p PostableRuleView
|
||||
err := json.Unmarshal([]byte(testCase.body), &p)
|
||||
if testCase.expectError {
|
||||
assert.Error(t, err)
|
||||
if testCase.expectedErrMsg != "" {
|
||||
assert.ErrorContains(t, err, testCase.expectedErrMsg)
|
||||
}
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.expectedName, p.Name)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestToStorableRuleView(t *testing.T) {
|
||||
orgID := valuer.GenerateUUID()
|
||||
postable := PostableRuleView{
|
||||
Name: "my view",
|
||||
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"firing"}, Sort: ListSortName, Order: ListOrderAsc}},
|
||||
}
|
||||
|
||||
storable := postable.ToStorableRuleView(orgID)
|
||||
|
||||
assert.Equal(t, orgID, storable.OrgID)
|
||||
assert.Equal(t, "my view", storable.Name)
|
||||
assert.False(t, storable.ID.IsZero())
|
||||
assert.False(t, storable.CreatedAt.IsZero())
|
||||
assert.Equal(t, storable.CreatedAt, storable.UpdatedAt)
|
||||
|
||||
gettable := storable.ToGettableRuleView()
|
||||
assert.Equal(t, storable.ID, gettable.ID)
|
||||
assert.Equal(t, "my view", gettable.Name)
|
||||
assert.Equal(t, postable.Data, gettable.Data)
|
||||
assert.Equal(t, orgID, gettable.OrgID)
|
||||
assert.Equal(t, storable.CreatedAt, gettable.CreatedAt)
|
||||
}
|
||||
|
||||
func TestNewGettableRuleViewsFromStorableRuleViews(t *testing.T) {
|
||||
orgID := valuer.GenerateUUID()
|
||||
first := PostableRuleView{
|
||||
Name: "first view",
|
||||
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Sort: ListSortName, Order: ListOrderAsc}},
|
||||
}.ToStorableRuleView(orgID)
|
||||
second := PostableRuleView{
|
||||
Name: "second view",
|
||||
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"firing"}, Sort: ListSortState, Order: ListOrderDesc}},
|
||||
}.ToStorableRuleView(orgID)
|
||||
|
||||
views := NewGettableRuleViewsFromStorableRuleViews([]*StorableRuleView{first, second})
|
||||
|
||||
require.Len(t, views, 2)
|
||||
assert.Equal(t, "first view", views[0].Name)
|
||||
assert.Equal(t, second.Data.States, views[1].Data.States)
|
||||
assert.Equal(t, ListSortState, views[1].Data.Sort)
|
||||
}
|
||||
|
||||
func TestStorableRuleViewUpdate(t *testing.T) {
|
||||
orgID := valuer.GenerateUUID()
|
||||
storable := PostableRuleView{
|
||||
Name: "original",
|
||||
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{Sort: ListSortName, Order: ListOrderAsc}},
|
||||
}.ToStorableRuleView(orgID)
|
||||
createdAt := storable.CreatedAt
|
||||
|
||||
storable.Update(UpdatableRuleView{
|
||||
Name: "renamed",
|
||||
Data: RuleViewData{Version: RuleViewSchemaVersion, ListFilter: ListFilter{States: []string{"disabled"}, Sort: ListSortCreatedAt, Order: ListOrderDesc}},
|
||||
})
|
||||
|
||||
gettable := storable.ToGettableRuleView()
|
||||
assert.Equal(t, "renamed", gettable.Name)
|
||||
assert.Equal(t, []string{"disabled"}, gettable.Data.States)
|
||||
assert.Equal(t, ListSortCreatedAt, gettable.Data.Sort)
|
||||
assert.Equal(t, ListOrderDesc, gettable.Data.Order)
|
||||
assert.Equal(t, createdAt, storable.CreatedAt)
|
||||
assert.True(t, storable.UpdatedAt.After(createdAt))
|
||||
}
|
||||
@@ -36,6 +36,7 @@ type TraceStore interface {
|
||||
GetMinimalSpans(ctx context.Context, traceID string, start, end time.Time) ([]MinimalSpan, error)
|
||||
GetTraceSpansByIDs(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
|
||||
GetFlamegraphSpans(ctx context.Context, traceID string, start, end time.Time, spanIDs []string) ([]StorableSpan, error)
|
||||
GetThreadSpans(ctx context.Context, traceID string, summary *TraceSummary, cursor *ThreadCursor, limit int) ([]StorableSpan, error)
|
||||
|
||||
GetSpanCountByField(ctx context.Context, traceID string, summary *TraceSummary, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)
|
||||
GetSpanDurationByField(ctx context.Context, traceID string, summary *TraceSummary, fieldKey telemetrytypes.TelemetryFieldKey) (map[string]uint64, error)
|
||||
|
||||
113
pkg/types/spantypes/thread.go
Normal file
113
pkg/types/spantypes/thread.go
Normal file
@@ -0,0 +1,113 @@
|
||||
package spantypes
|
||||
|
||||
import (
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
)
|
||||
|
||||
const (
|
||||
threadDefaultLimit = 100
|
||||
threadMaxLimit = 1000
|
||||
)
|
||||
|
||||
var (
|
||||
ErrCodeThreadInvalidLimit = errors.MustNewCode("trace_thread_invalid_limit")
|
||||
ErrCodeThreadInvalidCursor = errors.MustNewCode("trace_thread_invalid_cursor")
|
||||
)
|
||||
|
||||
type QueryableThread struct {
|
||||
// Limit is the page size; 0 means 100.
|
||||
Limit int `query:"limit"`
|
||||
// Cursor is the nextCursor of the previous page; empty for the first page.
|
||||
Cursor string `query:"cursor"`
|
||||
}
|
||||
|
||||
type ThreadQuery struct {
|
||||
Limit int
|
||||
Cursor *ThreadCursor
|
||||
}
|
||||
|
||||
func NewThreadQuery(queryable *QueryableThread) (*ThreadQuery, error) {
|
||||
query := &ThreadQuery{Limit: queryable.Limit}
|
||||
if query.Limit < 0 {
|
||||
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidLimit, "limit cannot be negative, got %d", query.Limit)
|
||||
}
|
||||
if query.Limit == 0 {
|
||||
query.Limit = threadDefaultLimit
|
||||
}
|
||||
if query.Limit > threadMaxLimit {
|
||||
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidLimit, "limit cannot exceed %d, got %d", threadMaxLimit, query.Limit)
|
||||
}
|
||||
if queryable.Cursor != "" {
|
||||
cursor, err := DecodeThreadCursor(queryable.Cursor)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
query.Cursor = cursor
|
||||
}
|
||||
return query, nil
|
||||
}
|
||||
|
||||
// ThreadCursor is the (TimeUnixNano, SpanID) of the last span of a page.
|
||||
type ThreadCursor struct {
|
||||
TimeUnixNano uint64 `json:"t"`
|
||||
SpanID string `json:"s"`
|
||||
}
|
||||
|
||||
func (c ThreadCursor) Encode() string {
|
||||
data, _ := json.Marshal(c)
|
||||
return base64.RawURLEncoding.EncodeToString(data)
|
||||
}
|
||||
|
||||
func DecodeThreadCursor(cursor string) (*ThreadCursor, error) {
|
||||
data, err := base64.RawURLEncoding.DecodeString(cursor)
|
||||
if err != nil {
|
||||
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
|
||||
}
|
||||
c := new(ThreadCursor)
|
||||
if err := json.Unmarshal(data, c); err != nil {
|
||||
return nil, errors.WrapInvalidInputf(err, ErrCodeThreadInvalidCursor, "invalid cursor")
|
||||
}
|
||||
if c.SpanID == "" {
|
||||
return nil, errors.NewInvalidInputf(ErrCodeThreadInvalidCursor, "invalid cursor: missing span id")
|
||||
}
|
||||
return c, nil
|
||||
}
|
||||
|
||||
type GettableTraceThread struct {
|
||||
Spans []*ThreadSpan `json:"spans" required:"true" nullable:"false"`
|
||||
NextCursor string `json:"nextCursor,omitempty"`
|
||||
}
|
||||
|
||||
type ThreadSpan struct {
|
||||
WaterfallSpan
|
||||
}
|
||||
|
||||
// NewGettableTraceThread expects limit+1 spans; the extra one only signals a next page.
|
||||
func NewGettableTraceThread(traceID string, spans []StorableSpan, limit int) *GettableTraceThread {
|
||||
hasMore := len(spans) > limit
|
||||
if hasMore {
|
||||
spans = spans[:limit]
|
||||
}
|
||||
|
||||
out := make([]*ThreadSpan, len(spans))
|
||||
for i := range spans {
|
||||
out[i] = newThreadSpan(traceID, &spans[i])
|
||||
}
|
||||
|
||||
thread := &GettableTraceThread{Spans: out}
|
||||
if hasMore {
|
||||
last := spans[len(spans)-1]
|
||||
thread.NextCursor = ThreadCursor{TimeUnixNano: uint64(last.StartTime.UnixNano()), SpanID: last.SpanID}.Encode()
|
||||
}
|
||||
return thread
|
||||
}
|
||||
|
||||
func newThreadSpan(traceID string, storable *StorableSpan) *ThreadSpan {
|
||||
span := &ThreadSpan{WaterfallSpan: *storable.ToWaterfallSpan(traceID)}
|
||||
// client expects millis, as in the waterfall
|
||||
span.TimeUnix = span.TimeUnix / 1_000_000
|
||||
return span
|
||||
}
|
||||
102
pkg/types/spantypes/thread_test.go
Normal file
102
pkg/types/spantypes/thread_test.go
Normal file
@@ -0,0 +1,102 @@
|
||||
package spantypes
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestNewThreadQuery(t *testing.T) {
|
||||
cursor := ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
queryable QueryableThread
|
||||
want *ThreadQuery
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "ZeroLimit_UsesDefault", queryable: QueryableThread{}, want: &ThreadQuery{Limit: threadDefaultLimit}},
|
||||
{name: "PositiveLimit_Kept", queryable: QueryableThread{Limit: 25}, want: &ThreadQuery{Limit: 25}},
|
||||
{name: "MaxLimit_Kept", queryable: QueryableThread{Limit: threadMaxLimit}, want: &ThreadQuery{Limit: threadMaxLimit}},
|
||||
{name: "AboveMaxLimit_Rejected", queryable: QueryableThread{Limit: threadMaxLimit + 1}, wantErr: true},
|
||||
{name: "NegativeLimit_Rejected", queryable: QueryableThread{Limit: -1}, wantErr: true},
|
||||
{name: "Cursor_Decoded", queryable: QueryableThread{Limit: 10, Cursor: cursor.Encode()}, want: &ThreadQuery{Limit: 10, Cursor: &cursor}},
|
||||
{name: "InvalidCursor_Rejected", queryable: QueryableThread{Cursor: "not base64!"}, wantErr: true},
|
||||
}
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
got, err := NewThreadQuery(&testCase.queryable)
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.want, got)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDecodeThreadCursor(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
cursor string
|
||||
want *ThreadCursor
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "EncodedCursor_RoundTrips", cursor: ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}.Encode(), want: &ThreadCursor{TimeUnixNano: 1757500000123456789, SpanID: "f1fa1bc863e94dd0"}},
|
||||
{name: "NotBase64_Rejected", cursor: "not base64!", wantErr: true},
|
||||
{name: "NotJSON_Rejected", cursor: "bm90IGpzb24", wantErr: true},
|
||||
{name: "MissingSpanID_Rejected", cursor: "eyJ0IjogMX0", wantErr: true},
|
||||
}
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
got, err := DecodeThreadCursor(testCase.cursor)
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.want, got)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewGettableTraceThread(t *testing.T) {
|
||||
spans := []StorableSpan{
|
||||
{SpanID: "a", StartTime: time.Unix(1, 500_000_000)},
|
||||
{SpanID: "b", StartTime: time.Unix(2, 0)},
|
||||
{SpanID: "c", StartTime: time.Unix(3, 0)},
|
||||
}
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
spans []StorableSpan
|
||||
limit int
|
||||
wantSpanIDs []string
|
||||
wantTimeUnix []uint64
|
||||
wantNextCursor string
|
||||
}{
|
||||
{name: "MoreThanLimit_TrimsAndSetsCursor", spans: spans, limit: 2, wantSpanIDs: []string{"a", "b"}, wantTimeUnix: []uint64{1500, 2000}, wantNextCursor: ThreadCursor{TimeUnixNano: 2_000_000_000, SpanID: "b"}.Encode()},
|
||||
{name: "WithinLimit_NoCursor", spans: spans, limit: 3, wantSpanIDs: []string{"a", "b", "c"}, wantTimeUnix: []uint64{1500, 2000, 3000}},
|
||||
{name: "NoSpans_EmptyList", limit: 3, wantSpanIDs: []string{}, wantTimeUnix: []uint64{}},
|
||||
}
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
thread := NewGettableTraceThread("trace-1", testCase.spans, testCase.limit)
|
||||
require.NotNil(t, thread.Spans)
|
||||
spanIDs := make([]string, len(thread.Spans))
|
||||
timeUnix := make([]uint64, len(thread.Spans))
|
||||
for i, span := range thread.Spans {
|
||||
spanIDs[i] = span.SpanID
|
||||
timeUnix[i] = span.TimeUnix
|
||||
assert.Equal(t, "trace-1", span.TraceID)
|
||||
}
|
||||
assert.Equal(t, testCase.wantSpanIDs, spanIDs)
|
||||
assert.Equal(t, testCase.wantTimeUnix, timeUnix)
|
||||
assert.Equal(t, testCase.wantNextCursor, thread.NextCursor)
|
||||
})
|
||||
}
|
||||
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
@@ -93,35 +94,36 @@ type WaterfallSpan struct {
|
||||
|
||||
// StorableSpan is the ClickHouse scan struct for the v3 waterfall query.
|
||||
type StorableSpan struct {
|
||||
StartTime time.Time `ch:"timestamp"`
|
||||
DurationNano uint64 `ch:"duration_nano"`
|
||||
SpanID string `ch:"span_id"`
|
||||
HasError bool `ch:"has_error"`
|
||||
Kind int8 `ch:"kind"`
|
||||
ServiceName string `ch:"resource_string_service$$name"`
|
||||
Name string `ch:"name"`
|
||||
AttributesString map[string]string `ch:"attributes_string"`
|
||||
AttributesNumber map[string]float64 `ch:"attributes_number"`
|
||||
AttributesBool map[string]bool `ch:"attributes_bool"`
|
||||
ResourcesString map[string]string `ch:"resources_string"`
|
||||
Events []string `ch:"events"`
|
||||
StatusMessage string `ch:"status_message"`
|
||||
StatusCodeString string `ch:"status_code_string"`
|
||||
SpanKind string `ch:"kind_string"`
|
||||
ParentSpanID string `ch:"parent_span_id"`
|
||||
Flags uint32 `ch:"flags"`
|
||||
IsRemote string `ch:"is_remote"`
|
||||
TraceState string `ch:"trace_state"`
|
||||
StatusCode int16 `ch:"status_code"`
|
||||
DBName string `ch:"db_name"`
|
||||
DBOperation string `ch:"db_operation"`
|
||||
HTTPMethod string `ch:"http_method"`
|
||||
HTTPURL string `ch:"http_url"`
|
||||
HTTPHost string `ch:"http_host"`
|
||||
ExternalHTTPMethod string `ch:"external_http_method"`
|
||||
ExternalHTTPURL string `ch:"external_http_url"`
|
||||
ResponseStatusCode string `ch:"response_status_code"`
|
||||
References string `ch:"references"`
|
||||
StartTime time.Time `ch:"timestamp"`
|
||||
DurationNano uint64 `ch:"duration_nano"`
|
||||
SpanID string `ch:"span_id"`
|
||||
HasError bool `ch:"has_error"`
|
||||
Kind int8 `ch:"kind"`
|
||||
ServiceName string `ch:"resource_string_service$$name"`
|
||||
Name string `ch:"name"`
|
||||
AttributesString map[string]string `ch:"attributes_string"`
|
||||
AttributesNumber map[string]float64 `ch:"attributes_number"`
|
||||
AttributesBool map[string]bool `ch:"attributes_bool"`
|
||||
AttributesJSON telemetrystoretypes.JSONValue `ch:"attributes"`
|
||||
ResourcesString map[string]string `ch:"resources_string"`
|
||||
Events []string `ch:"events"`
|
||||
StatusMessage string `ch:"status_message"`
|
||||
StatusCodeString string `ch:"status_code_string"`
|
||||
SpanKind string `ch:"kind_string"`
|
||||
ParentSpanID string `ch:"parent_span_id"`
|
||||
Flags uint32 `ch:"flags"`
|
||||
IsRemote string `ch:"is_remote"`
|
||||
TraceState string `ch:"trace_state"`
|
||||
StatusCode int16 `ch:"status_code"`
|
||||
DBName string `ch:"db_name"`
|
||||
DBOperation string `ch:"db_operation"`
|
||||
HTTPMethod string `ch:"http_method"`
|
||||
HTTPURL string `ch:"http_url"`
|
||||
HTTPHost string `ch:"http_host"`
|
||||
ExternalHTTPMethod string `ch:"external_http_method"`
|
||||
ExternalHTTPURL string `ch:"external_http_url"`
|
||||
ResponseStatusCode string `ch:"response_status_code"`
|
||||
References string `ch:"references"`
|
||||
}
|
||||
|
||||
// MinimalSpan with only the fields needed to build the parent-child tree.
|
||||
@@ -277,8 +279,10 @@ func (item *StorableSpan) AttributeValue(name string) any {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Attributes flattens the JSON column first, so the legacy maps win on collision.
|
||||
func (item *StorableSpan) Attributes() map[string]any {
|
||||
attributes := make(map[string]any, len(item.AttributesString)+len(item.AttributesNumber)+len(item.AttributesBool))
|
||||
attributes := make(map[string]any, len(item.AttributesString)+len(item.AttributesNumber)+len(item.AttributesBool)+len(item.AttributesJSON))
|
||||
item.AttributesJSON.FlattenInto("", attributes)
|
||||
for k, v := range item.AttributesString {
|
||||
attributes[k] = v
|
||||
}
|
||||
|
||||
59
pkg/types/spantypes/waterfall_span_test.go
Normal file
59
pkg/types/spantypes/waterfall_span_test.go
Normal file
@@ -0,0 +1,59 @@
|
||||
package spantypes
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestStorableSpanAttributes(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
span StorableSpan
|
||||
wantAttrs map[string]any
|
||||
}{
|
||||
{
|
||||
name: "LegacyMapOnly_Kept",
|
||||
span: StorableSpan{AttributesString: map[string]string{
|
||||
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
|
||||
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
|
||||
}},
|
||||
wantAttrs: map[string]any{
|
||||
"gen_ai.input.messages": `[{"role":"user","parts":[{"type":"text","content":"hi"}]}]`,
|
||||
"gen_ai.output.messages": `[{"role":"assistant","parts":[{"type":"text","content":"hello"}],"finish_reason":"stop"}]`,
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "JSONColumn_FlattenedToDottedKeys",
|
||||
span: StorableSpan{AttributesJSON: telemetrystoretypes.JSONValue{
|
||||
"gen_ai": map[string]any{
|
||||
"input": map[string]any{"messages": `[{"role":"user","content":"hi"}]`},
|
||||
"request": map[string]any{"model": "gpt-4o"},
|
||||
},
|
||||
}},
|
||||
wantAttrs: map[string]any{
|
||||
"gen_ai.input.messages": `[{"role":"user","content":"hi"}]`,
|
||||
"gen_ai.request.model": "gpt-4o",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "LegacyMapWinsOverJSONColumn",
|
||||
span: StorableSpan{
|
||||
AttributesJSON: telemetrystoretypes.JSONValue{"gen_ai": map[string]any{"request": map[string]any{"model": "json"}}},
|
||||
AttributesString: map[string]string{"gen_ai.request.model": "map"},
|
||||
},
|
||||
wantAttrs: map[string]any{"gen_ai.request.model": "map"},
|
||||
},
|
||||
{
|
||||
name: "NoMessages_AttributesKept",
|
||||
span: StorableSpan{AttributesString: map[string]string{"http.method": "GET"}},
|
||||
wantAttrs: map[string]any{"http.method": "GET"},
|
||||
},
|
||||
}
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
assert.Equal(t, testCase.wantAttrs, testCase.span.Attributes())
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -35,3 +35,21 @@ func (v *JSONValue) Scan(src any) error {
|
||||
*v = decoded
|
||||
return nil
|
||||
}
|
||||
|
||||
// FlattenInto writes v into out under dotted keys, overwriting existing keys.
|
||||
func (v JSONValue) FlattenInto(prefix string, out map[string]any) {
|
||||
for k, value := range v {
|
||||
key := k
|
||||
if prefix != "" {
|
||||
key = prefix + "." + k
|
||||
}
|
||||
switch child := value.(type) {
|
||||
case map[string]any:
|
||||
JSONValue(child).FlattenInto(key, out)
|
||||
case JSONValue:
|
||||
child.FlattenInto(key, out)
|
||||
default:
|
||||
out[key] = value
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
30
tests/fixtures/alerts.py
vendored
30
tests/fixtures/alerts.py
vendored
@@ -131,36 +131,6 @@ def seed_alert_rules(
|
||||
return _seed_alert_rules
|
||||
|
||||
|
||||
@pytest.fixture(name="create_rule_view", scope="function")
|
||||
def create_rule_view(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> Callable[[dict], dict]:
|
||||
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
view_ids = []
|
||||
|
||||
def _create_rule_view(view: dict) -> dict:
|
||||
response = requests.post(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/rule_views"),
|
||||
json=view,
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=5,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.CREATED, f"Failed to create rule view, api returned {response.status_code} with response: {response.text}"
|
||||
created = response.json()["data"]
|
||||
view_ids.append(created["id"])
|
||||
return created
|
||||
|
||||
yield _create_rule_view
|
||||
# A view the test already deleted returns 404; only real failures are logged.
|
||||
for view_id in view_ids:
|
||||
response = requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"/api/v2/rule_views/{view_id}"),
|
||||
headers={"Authorization": f"Bearer {admin_token}"},
|
||||
timeout=5,
|
||||
)
|
||||
if response.status_code not in (HTTPStatus.NO_CONTENT, HTTPStatus.NOT_FOUND):
|
||||
logger.error("Error deleting rule view: %s", {"view_id": view_id, "response": response.text})
|
||||
|
||||
|
||||
def labels_to_map(labels: list[dict]) -> dict[str, str]:
|
||||
"""Converts the label list shape of the v2 rule history APIs to a plain map."""
|
||||
return {label["key"]["name"]: label["value"] for label in labels or []}
|
||||
|
||||
@@ -130,3 +130,17 @@ def test_bulk_sync(
|
||||
assert all(r["pricing"]["input"] == 5 for r in stored)
|
||||
|
||||
delete_all_llm_pricing_rules(signoz, token)
|
||||
|
||||
|
||||
def test_rejects_rule_without_pattern(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: types.Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
delete_all_llm_pricing_rules(signoz, token)
|
||||
|
||||
rules = zeus_rules(10)
|
||||
rules[1]["modelPattern"] = []
|
||||
assert upsert_llm_pricing_rules(signoz, token, rules).status_code == HTTPStatus.BAD_REQUEST
|
||||
assert list_llm_pricing_rules(signoz, token) == []
|
||||
|
||||
36
tests/integration/tests/passwordauthn/09_last_observed_at.py
Normal file
36
tests/integration/tests/passwordauthn/09_last_observed_at.py
Normal file
@@ -0,0 +1,36 @@
|
||||
import time
|
||||
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
|
||||
|
||||
|
||||
def test_last_observed_at_is_flushed(signoz: types.SigNoz, get_token: Callable[[str, str], str]) -> None:
|
||||
"""Verify the tokenizer GC persists the cached last observed at of a used token to the sql store."""
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/users/me"),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK
|
||||
|
||||
deadline = time.time() + 30
|
||||
while time.time() < deadline:
|
||||
with signoz.sqlstore.conn.connect() as conn:
|
||||
row = conn.execute(
|
||||
sql.text("SELECT last_observed_at FROM auth_token WHERE access_token = :access_token"),
|
||||
{"access_token": token},
|
||||
).fetchone()
|
||||
|
||||
if row is not None and row[0] is not None:
|
||||
return
|
||||
|
||||
time.sleep(1)
|
||||
|
||||
raise AssertionError("last_observed_at was not flushed to the sql store within 30s")
|
||||
33
tests/integration/tests/passwordauthn/conftest.py
Normal file
33
tests/integration/tests/passwordauthn/conftest.py
Normal file
@@ -0,0 +1,33 @@
|
||||
import pytest
|
||||
from testcontainers.core.container import Network
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.signoz import create_signoz
|
||||
|
||||
|
||||
@pytest.fixture(name="signoz", scope="package")
|
||||
def signoz_passwordauthn(
|
||||
network: Network,
|
||||
zeus: types.TestContainerDocker,
|
||||
gateway: types.TestContainerDocker,
|
||||
sqlstore: types.TestContainerSQL,
|
||||
clickhouse: types.TestContainerClickhouse,
|
||||
request: pytest.FixtureRequest,
|
||||
pytestconfig: pytest.Config,
|
||||
) -> types.SigNoz:
|
||||
"""
|
||||
Package-scoped fixture for SigNoz with a short tokenizer GC interval so the last observed at flush runs within a test.
|
||||
"""
|
||||
return create_signoz(
|
||||
network=network,
|
||||
zeus=zeus,
|
||||
gateway=gateway,
|
||||
sqlstore=sqlstore,
|
||||
clickhouse=clickhouse,
|
||||
request=request,
|
||||
pytestconfig=pytestconfig,
|
||||
cache_key="signoz-passwordauthn",
|
||||
env_overrides={
|
||||
"SIGNOZ_TOKENIZER_OPAQUE_GC_INTERVAL": "5s",
|
||||
},
|
||||
)
|
||||
@@ -1,269 +0,0 @@
|
||||
import uuid
|
||||
from collections.abc import Callable
|
||||
from http import HTTPStatus
|
||||
|
||||
import pytest
|
||||
import requests
|
||||
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.types import Operation, SigNoz
|
||||
|
||||
BASE_URL = "/api/v2/rule_views"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("body", "expected_code", "expected_message"),
|
||||
[
|
||||
({"data": {"version": "v1"}}, "rule_view_invalid_input", "name is required"),
|
||||
({"name": " ", "data": {"version": "v1"}}, "rule_view_invalid_input", "name is required"),
|
||||
(
|
||||
{"name": " Storage ", "data": {"version": "v1"}},
|
||||
"rule_view_invalid_input",
|
||||
"name must not have leading or trailing whitespace",
|
||||
),
|
||||
(
|
||||
{"name": "x" * 65, "data": {"version": "v1"}},
|
||||
"rule_view_invalid_input",
|
||||
"name must be at most 64 characters, got 65",
|
||||
),
|
||||
(
|
||||
{"name": "wrong-version", "data": {"version": "v2"}},
|
||||
"rule_view_invalid_input",
|
||||
'version must be "v1", got "v2"',
|
||||
),
|
||||
(
|
||||
{"name": "missing-version", "data": {}},
|
||||
"rule_view_invalid_input",
|
||||
'version must be "v1", got ""',
|
||||
),
|
||||
(
|
||||
{"name": "bad-state", "data": {"version": "v1", "states": ["exploding"]}},
|
||||
"rule_list_invalid",
|
||||
'invalid state "exploding"',
|
||||
),
|
||||
(
|
||||
{"name": "bad-sort", "data": {"version": "v1", "sort": "bogus"}},
|
||||
"rule_list_invalid",
|
||||
"invalid sort",
|
||||
),
|
||||
(
|
||||
{"name": "bad-order", "data": {"version": "v1", "order": "bogus"}},
|
||||
"rule_list_invalid",
|
||||
"invalid order",
|
||||
),
|
||||
(
|
||||
{"name": "long-query", "data": {"version": "v1", "query": "x" * 1025}},
|
||||
"rule_list_invalid",
|
||||
"query cannot be longer than 1024 characters",
|
||||
),
|
||||
(
|
||||
{"name": "rejects-unknown", "data": {"version": "v1"}, "unknownfield": "boom"},
|
||||
"rule_view_invalid_input",
|
||||
"invalid saved view request body",
|
||||
),
|
||||
],
|
||||
ids=[
|
||||
"missing_name",
|
||||
"blank_name",
|
||||
"whitespace_name",
|
||||
"name_too_long",
|
||||
"wrong_schema_version",
|
||||
"missing_version",
|
||||
"invalid_state",
|
||||
"invalid_sort",
|
||||
"invalid_order",
|
||||
"query_too_long",
|
||||
"unknown_field",
|
||||
],
|
||||
)
|
||||
def test_create_rejects_invalid_body(
|
||||
signoz: SigNoz,
|
||||
create_user_admin: Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
body: dict,
|
||||
expected_code: str,
|
||||
expected_message: str,
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.post(
|
||||
signoz.self.host_configs["8080"].get(BASE_URL),
|
||||
json=body,
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.BAD_REQUEST
|
||||
assert response.json()["error"]["code"] == expected_code
|
||||
assert expected_message in response.json()["error"]["message"]
|
||||
|
||||
|
||||
def test_update_rejects_malformed_id(
|
||||
signoz: SigNoz,
|
||||
create_user_admin: Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/not-a-uuid"),
|
||||
json={"name": "x", "data": {"version": "v1"}},
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.BAD_REQUEST
|
||||
|
||||
|
||||
def test_update_missing_view_returns_not_found(
|
||||
signoz: SigNoz,
|
||||
create_user_admin: Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{uuid.uuid4()}"),
|
||||
json={"name": "x", "data": {"version": "v1"}},
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.NOT_FOUND
|
||||
assert response.json()["error"]["code"] == "rule_view_not_found"
|
||||
|
||||
|
||||
def test_delete_rejects_malformed_id(
|
||||
signoz: SigNoz,
|
||||
create_user_admin: Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/not-a-uuid"),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.BAD_REQUEST
|
||||
|
||||
|
||||
def test_delete_missing_view_returns_not_found(
|
||||
signoz: SigNoz,
|
||||
create_user_admin: Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
response = requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{uuid.uuid4()}"),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
|
||||
assert response.status_code == HTTPStatus.NOT_FOUND
|
||||
assert response.json()["error"]["code"] == "rule_view_not_found"
|
||||
|
||||
|
||||
def test_rule_view_lifecycle(
|
||||
signoz: SigNoz,
|
||||
create_user_admin: Operation, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
create_rule_view: Callable[[dict], dict],
|
||||
):
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
|
||||
# List assertions filter on this test's names so foreign views never interfere.
|
||||
owned_names = {"Critical Prod", "Critical Staging", "Disabled"}
|
||||
|
||||
created = create_rule_view(
|
||||
{
|
||||
"name": "Critical Prod",
|
||||
"data": {
|
||||
"version": "v1",
|
||||
"query": "name CONTAINS 'prod' AND severity = 'critical'",
|
||||
"states": ["firing", "pending"],
|
||||
"sort": "name",
|
||||
"order": "asc",
|
||||
},
|
||||
}
|
||||
)
|
||||
view_id = created["id"]
|
||||
assert created["name"] == "Critical Prod"
|
||||
assert created["data"]["version"] == "v1"
|
||||
assert created["data"]["query"] == "name CONTAINS 'prod' AND severity = 'critical'"
|
||||
assert created["data"]["states"] == ["firing", "pending"]
|
||||
|
||||
# Omitted states, sort and order are normalized on save: [] and the list defaults, never null.
|
||||
disabled = create_rule_view({"name": "Disabled", "data": {"version": "v1"}})
|
||||
assert disabled["name"] == "Disabled"
|
||||
assert disabled["data"]["states"] == []
|
||||
assert disabled["data"]["sort"] == "updated_at"
|
||||
assert disabled["data"]["order"] == "desc"
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get(BASE_URL),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
views = [v for v in response.json()["data"]["views"] if v["name"] in owned_names]
|
||||
assert {v["name"] for v in views} == {"Critical Prod", "Disabled"}
|
||||
|
||||
response = requests.put(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{view_id}"),
|
||||
json={
|
||||
"name": "Critical Staging",
|
||||
"data": {
|
||||
"version": "v1",
|
||||
"query": "name CONTAINS 'staging'",
|
||||
"states": ["firing"],
|
||||
"sort": "created_at",
|
||||
"order": "desc",
|
||||
},
|
||||
},
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
updated = response.json()["data"]
|
||||
assert updated["id"] == view_id
|
||||
assert updated["name"] == "Critical Staging"
|
||||
assert updated["data"]["query"] == "name CONTAINS 'staging'"
|
||||
assert updated["data"]["states"] == ["firing"]
|
||||
assert updated["data"]["sort"] == "created_at"
|
||||
assert updated["data"]["order"] == "desc"
|
||||
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get(BASE_URL),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
listed = {v["name"]: v for v in response.json()["data"]["views"] if v["name"] in owned_names}
|
||||
assert set(listed) == {"Critical Staging", "Disabled"}
|
||||
assert listed["Critical Staging"]["data"]["query"] == "name CONTAINS 'staging'"
|
||||
|
||||
assert (
|
||||
requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{view_id}"),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
).status_code
|
||||
== HTTPStatus.NO_CONTENT
|
||||
)
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get(BASE_URL),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
)
|
||||
assert {v["name"] for v in response.json()["data"]["views"] if v["name"] in owned_names} == {"Disabled"}
|
||||
|
||||
assert (
|
||||
requests.delete(
|
||||
signoz.self.host_configs["8080"].get(f"{BASE_URL}/{view_id}"),
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=5,
|
||||
).status_code
|
||||
== HTTPStatus.NOT_FOUND
|
||||
)
|
||||
127
tests/integration/tests/tracedetail/01_thread.py
Normal file
127
tests/integration/tests/tracedetail/01_thread.py
Normal file
@@ -0,0 +1,127 @@
|
||||
import json
|
||||
from collections.abc import Callable
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
|
||||
import requests
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.traces import TraceIdGenerator, Traces, TracesKind
|
||||
|
||||
|
||||
def test_thread_returns_message_spans_in_order(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
) -> None:
|
||||
now = datetime.now(tz=UTC).replace(microsecond=0)
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
root_id, first_llm_id, tool_id, second_llm_id, third_llm_id = (TraceIdGenerator.span_id() for _ in range(5))
|
||||
resources = {"service.name": "tracedetail-thread"}
|
||||
first_input = json.dumps([{"role": "user", "parts": [{"type": "text", "content": "weather in Bangalore?"}]}])
|
||||
first_output = json.dumps([{"role": "assistant", "parts": [{"type": "tool_call", "id": "call_1", "name": "get_weather", "arguments": {"city": "Bangalore"}}], "finish_reason": "tool_call"}])
|
||||
second_input = json.dumps([{"role": "tool", "content": "sunny", "tool_call_id": "call_1"}])
|
||||
|
||||
insert_traces(
|
||||
[
|
||||
Traces(timestamp=now - timedelta(seconds=10), duration=timedelta(seconds=9), trace_id=trace_id, span_id=root_id, name="POST /chat", kind=TracesKind.SPAN_KIND_SERVER, resources=resources, attribute_write_mode="json_only"),
|
||||
Traces(
|
||||
timestamp=now - timedelta(seconds=8), trace_id=trace_id, span_id=first_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.input.messages": first_input, "gen_ai.output.messages": first_output}, attribute_write_mode="json_only"
|
||||
),
|
||||
Traces(timestamp=now - timedelta(seconds=6), trace_id=trace_id, span_id=tool_id, parent_span_id=root_id, name="execute_tool get_weather", resources=resources, attributes={"gen_ai.tool.name": "get_weather"}, attribute_write_mode="json_only"),
|
||||
Traces(timestamp=now - timedelta(seconds=4), trace_id=trace_id, span_id=second_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.input.messages": second_input}, attribute_write_mode="json_only"),
|
||||
Traces(timestamp=now - timedelta(seconds=2), trace_id=trace_id, span_id=third_llm_id, parent_span_id=root_id, name="chat gpt-4o", resources=resources, attributes={"gen_ai.request.model": "gpt-4o", "gen_ai.output.messages": "It is sunny in Bangalore."}, attribute_write_mode="json_only"),
|
||||
]
|
||||
)
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
response = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread"), headers={"Authorization": f"Bearer {token}"}, timeout=10)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
|
||||
thread = response.json()["data"]
|
||||
assert [span["span_id"] for span in thread["spans"]] == [first_llm_id, second_llm_id, third_llm_id]
|
||||
assert "nextCursor" not in thread
|
||||
|
||||
first, input_only, output_only = thread["spans"]
|
||||
assert first["time_unix"] == int((now - timedelta(seconds=8)).timestamp() * 1000)
|
||||
assert first["attributes"]["gen_ai.input.messages"] == first_input
|
||||
assert first["attributes"]["gen_ai.request.model"] == "gpt-4o"
|
||||
assert first["attributes"]["gen_ai.output.messages"] == first_output
|
||||
assert input_only["attributes"]["gen_ai.input.messages"] == second_input
|
||||
assert "gen_ai.output.messages" not in input_only["attributes"]
|
||||
assert "gen_ai.input.messages" not in output_only["attributes"]
|
||||
assert output_only["attributes"]["gen_ai.output.messages"] == "It is sunny in Bangalore."
|
||||
|
||||
|
||||
def test_thread_paginates_with_cursor(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
) -> None:
|
||||
now = datetime.now(tz=UTC).replace(microsecond=0)
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
span_ids = [TraceIdGenerator.span_id() for _ in range(3)]
|
||||
# identical timestamps on the last two exercise the span_id tie-break
|
||||
timestamps = [now - timedelta(seconds=6), now - timedelta(seconds=3), now - timedelta(seconds=3)]
|
||||
insert_traces(
|
||||
[
|
||||
Traces(timestamp=timestamp, trace_id=trace_id, span_id=span_id, name="chat gpt-4o", resources={"service.name": "tracedetail-thread-pages"}, attributes={"gen_ai.input.messages": json.dumps([{"role": "user", "content": span_id}])}, attribute_write_mode="json_only")
|
||||
for span_id, timestamp in zip(span_ids, timestamps, strict=True)
|
||||
]
|
||||
)
|
||||
expected_order = [span_ids[0], *sorted(span_ids[1:])]
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
url = signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread")
|
||||
headers = {"Authorization": f"Bearer {token}"}
|
||||
|
||||
first_page = requests.get(url, params={"limit": 2}, headers=headers, timeout=10)
|
||||
assert first_page.status_code == HTTPStatus.OK, first_page.text
|
||||
first = first_page.json()["data"]
|
||||
assert [span["span_id"] for span in first["spans"]] == expected_order[:2]
|
||||
assert first["nextCursor"]
|
||||
|
||||
second_page = requests.get(url, params={"limit": 2, "cursor": first["nextCursor"]}, headers=headers, timeout=10)
|
||||
assert second_page.status_code == HTTPStatus.OK, second_page.text
|
||||
second = second_page.json()["data"]
|
||||
assert [span["span_id"] for span in second["spans"]] == expected_order[2:]
|
||||
assert "nextCursor" not in second
|
||||
|
||||
|
||||
def test_thread_without_messages_is_empty(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
) -> None:
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
insert_traces([Traces(timestamp=datetime.now(tz=UTC) - timedelta(seconds=5), trace_id=trace_id, span_id=TraceIdGenerator.span_id(), name="GET /health", resources={"service.name": "tracedetail-thread-empty"}, attribute_write_mode="json_only")])
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
response = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread"), headers={"Authorization": f"Bearer {token}"}, timeout=10)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
assert response.json()["data"] == {"spans": []}
|
||||
|
||||
|
||||
def test_thread_rejects_invalid_requests(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
) -> None:
|
||||
trace_id = TraceIdGenerator.trace_id()
|
||||
insert_traces([Traces(timestamp=datetime.now(tz=UTC) - timedelta(seconds=5), trace_id=trace_id, span_id=TraceIdGenerator.span_id(), name="chat gpt-4o", resources={"service.name": "tracedetail-thread-invalid"}, attributes={"gen_ai.input.messages": "hi"}, attribute_write_mode="json_only")])
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
headers = {"Authorization": f"Bearer {token}"}
|
||||
url = signoz.self.host_configs["8080"].get(f"/api/v1/traces/{trace_id}/thread")
|
||||
|
||||
for params in ({"limit": -1}, {"limit": 1001}, {"cursor": "not-a-cursor"}):
|
||||
response = requests.get(url, params=params, headers=headers, timeout=10)
|
||||
assert response.status_code == HTTPStatus.BAD_REQUEST, f"{params}: {response.text}"
|
||||
|
||||
missing = requests.get(signoz.self.host_configs["8080"].get(f"/api/v1/traces/{TraceIdGenerator.trace_id()}/thread"), headers=headers, timeout=10)
|
||||
assert missing.status_code == HTTPStatus.NOT_FOUND, missing.text
|
||||
Reference in New Issue
Block a user