mirror of
https://github.com/SigNoz/signoz.git
synced 2026-10-02 16:20:56 +01:00
Compare commits
28 Commits
feat/scatt
...
ns/trace-a
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6f08f2c22a | ||
|
|
cffa1b2049 | ||
|
|
dd912ae453 | ||
|
|
d7e3a18b6a | ||
|
|
886efd7672 | ||
|
|
f80d6e541c | ||
|
|
7612cc134e | ||
|
|
bcaaa6548c | ||
|
|
0212b1d635 | ||
|
|
da70042e43 | ||
|
|
70602730dd | ||
|
|
3df607857e | ||
|
|
b685455b87 | ||
|
|
d19a578a4a | ||
|
|
cc134ba19e | ||
|
|
e66b954ee9 | ||
|
|
4b9e1103cf | ||
|
|
dd232e3b8a | ||
|
|
0c2e537ed8 | ||
|
|
6ec91554a4 | ||
|
|
d1ea0b2596 | ||
|
|
828d05b8cc | ||
|
|
47dd1fabf3 | ||
|
|
270988fb48 | ||
|
|
39badeb591 | ||
|
|
ec05bfe755 | ||
|
|
ed1bf7ab89 | ||
|
|
6e979c8318 |
@@ -7617,6 +7617,8 @@ components:
|
||||
type: object
|
||||
PromotetypesPromotePath:
|
||||
properties:
|
||||
context:
|
||||
type: string
|
||||
indexes:
|
||||
items:
|
||||
$ref: '#/components/schemas/PromotetypesWrappedIndex'
|
||||
@@ -7625,6 +7627,12 @@ components:
|
||||
type: string
|
||||
promote:
|
||||
type: boolean
|
||||
signal:
|
||||
type: string
|
||||
required:
|
||||
- signal
|
||||
- context
|
||||
- path
|
||||
type: object
|
||||
PromotetypesWrappedIndex:
|
||||
properties:
|
||||
@@ -13100,110 +13108,6 @@ paths:
|
||||
tags:
|
||||
- llmpricingrules
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/logs/promote_paths:
|
||||
get:
|
||||
deprecated: false
|
||||
description: This endpoints promotes and indexes paths
|
||||
operationId: ListPromotedAndIndexedPaths
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
items:
|
||||
$ref: '#/components/schemas/PromotetypesPromotePath'
|
||||
nullable: true
|
||||
type: array
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: Promote and index paths
|
||||
tags:
|
||||
- logs
|
||||
x-signoz-stability: alpha
|
||||
post:
|
||||
deprecated: false
|
||||
description: This endpoints promotes and indexes paths
|
||||
operationId: HandlePromoteAndIndexPaths
|
||||
requestBody:
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
items:
|
||||
$ref: '#/components/schemas/PromotetypesPromotePath'
|
||||
nullable: true
|
||||
type: array
|
||||
responses:
|
||||
"201":
|
||||
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:
|
||||
- EDITOR
|
||||
- tokenizer:
|
||||
- EDITOR
|
||||
summary: Promote and index paths
|
||||
tags:
|
||||
- logs
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/org/preferences:
|
||||
get:
|
||||
deprecated: false
|
||||
@@ -13375,6 +13279,133 @@ paths:
|
||||
tags:
|
||||
- preferences
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/promoted_path:
|
||||
get:
|
||||
deprecated: false
|
||||
description: This endpoint lists the promoted paths of every JSON column, each
|
||||
annotated with its signal and context. The signal, context, promoted and indexes
|
||||
query parameters filter the listing.
|
||||
operationId: ListPromotedPaths
|
||||
parameters:
|
||||
- in: query
|
||||
name: signal
|
||||
schema:
|
||||
type: string
|
||||
- in: query
|
||||
name: context
|
||||
schema:
|
||||
type: string
|
||||
- in: query
|
||||
name: promoted
|
||||
schema:
|
||||
nullable: true
|
||||
type: boolean
|
||||
- in: query
|
||||
name: indexes
|
||||
schema:
|
||||
nullable: true
|
||||
type: boolean
|
||||
responses:
|
||||
"200":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
properties:
|
||||
data:
|
||||
items:
|
||||
$ref: '#/components/schemas/PromotetypesPromotePath'
|
||||
nullable: true
|
||||
type: array
|
||||
status:
|
||||
type: string
|
||||
required:
|
||||
- status
|
||||
- data
|
||||
type: object
|
||||
description: OK
|
||||
"400":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Bad Request
|
||||
"401":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Unauthorized
|
||||
"403":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Forbidden
|
||||
"500":
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: '#/components/schemas/RenderErrorResponse'
|
||||
description: Internal Server Error
|
||||
security:
|
||||
- api_key:
|
||||
- VIEWER
|
||||
- tokenizer:
|
||||
- VIEWER
|
||||
summary: List promoted paths
|
||||
tags:
|
||||
- promote
|
||||
x-signoz-stability: alpha
|
||||
post:
|
||||
deprecated: false
|
||||
description: This endpoint promotes paths of JSON columns to their promoted
|
||||
columns. Each path names its promotion domain with its signal and context,
|
||||
e.g. traces/attribute.
|
||||
operationId: PromotePaths
|
||||
requestBody:
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
items:
|
||||
$ref: '#/components/schemas/PromotetypesPromotePath'
|
||||
nullable: true
|
||||
type: array
|
||||
responses:
|
||||
"201":
|
||||
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:
|
||||
- EDITOR
|
||||
- tokenizer:
|
||||
- EDITOR
|
||||
summary: Promote paths
|
||||
tags:
|
||||
- promote
|
||||
x-signoz-stability: alpha
|
||||
/api/v1/roles:
|
||||
get:
|
||||
deprecated: false
|
||||
|
||||
@@ -4,23 +4,15 @@
|
||||
* * regenerate with 'pnpm generate:api'
|
||||
* SigNoz
|
||||
*/
|
||||
import { useMutation, useQuery } from 'react-query';
|
||||
import { useMutation } from 'react-query';
|
||||
import type {
|
||||
InvalidateOptions,
|
||||
MutationFunction,
|
||||
QueryClient,
|
||||
QueryFunction,
|
||||
QueryKey,
|
||||
UseMutationOptions,
|
||||
UseMutationResult,
|
||||
UseQueryOptions,
|
||||
UseQueryResult,
|
||||
} from 'react-query';
|
||||
|
||||
import type {
|
||||
HandleExportRawDataPOSTParams,
|
||||
ListPromotedAndIndexedPaths200,
|
||||
PromotetypesPromotePathDTO,
|
||||
Querybuildertypesv5QueryRangeRequestDTO,
|
||||
RenderErrorResponseDTO,
|
||||
} from '../sigNoz.schemas';
|
||||
@@ -28,26 +20,6 @@ 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;
|
||||
};
|
||||
|
||||
/**
|
||||
* This endpoints allows complex query exporting raw data for traces and logs
|
||||
* @summary Export raw data
|
||||
@@ -149,175 +121,3 @@ export const useHandleExportRawDataPOST = <
|
||||
> => {
|
||||
return useMutation(getHandleExportRawDataPOSTMutationOptions(options));
|
||||
};
|
||||
/**
|
||||
* This endpoints promotes and indexes paths
|
||||
* @summary Promote and index paths
|
||||
*/
|
||||
export const listPromotedAndIndexedPaths = (signal?: AbortSignal) => {
|
||||
return GeneratedAPIInstance<ListPromotedAndIndexedPaths200>({
|
||||
url: `/api/v1/logs/promote_paths`,
|
||||
method: 'GET',
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getListPromotedAndIndexedPathsQueryKey = () => {
|
||||
return [`/api/v1/logs/promote_paths`] as const;
|
||||
};
|
||||
|
||||
export const getListPromotedAndIndexedPathsQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
}) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey =
|
||||
queryOptions?.queryKey ?? getListPromotedAndIndexedPathsQueryKey();
|
||||
|
||||
const queryFn: QueryFunction<
|
||||
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>
|
||||
> = ({ signal }) => listPromotedAndIndexedPaths(signal);
|
||||
|
||||
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type ListPromotedAndIndexedPathsQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>
|
||||
>;
|
||||
export type ListPromotedAndIndexedPathsQueryError =
|
||||
ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Promote and index paths
|
||||
*/
|
||||
|
||||
export function useListPromotedAndIndexedPaths<
|
||||
TData = Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listPromotedAndIndexedPaths>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
}): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getListPromotedAndIndexedPathsQueryOptions(options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return withQueryKey(query, queryOptions.queryKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary Promote and index paths
|
||||
*/
|
||||
export const invalidateListPromotedAndIndexedPaths = async (
|
||||
queryClient: QueryClient,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getListPromotedAndIndexedPathsQueryKey() },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* This endpoints promotes and indexes paths
|
||||
* @summary Promote and index paths
|
||||
*/
|
||||
export const handlePromoteAndIndexPaths = (
|
||||
promotetypesPromotePathDTONull?: BodyType<
|
||||
PromotetypesPromotePathDTO[] | null
|
||||
> | null,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<void>({
|
||||
url: `/api/v1/logs/promote_paths`,
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
data: promotetypesPromotePathDTONull,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getHandlePromoteAndIndexPathsMutationOptions = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
> => {
|
||||
const mutationKey = ['handlePromoteAndIndexPaths'];
|
||||
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 handlePromoteAndIndexPaths>>,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> }
|
||||
> = (props) => {
|
||||
const { data } = props ?? {};
|
||||
|
||||
return handlePromoteAndIndexPaths(data);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type HandlePromoteAndIndexPathsMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>
|
||||
>;
|
||||
export type HandlePromoteAndIndexPathsMutationBody =
|
||||
| BodyType<PromotetypesPromotePathDTO[] | null>
|
||||
| undefined;
|
||||
export type HandlePromoteAndIndexPathsMutationError =
|
||||
ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Promote and index paths
|
||||
*/
|
||||
export const useHandlePromoteAndIndexPaths = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof handlePromoteAndIndexPaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getHandlePromoteAndIndexPathsMutationOptions(options));
|
||||
};
|
||||
|
||||
232
frontend/src/api/generated/services/promote/index.ts
Normal file
232
frontend/src/api/generated/services/promote/index.ts
Normal file
@@ -0,0 +1,232 @@
|
||||
/**
|
||||
* ! Do not edit manually
|
||||
* * The file has been auto-generated using Orval for SigNoz
|
||||
* * regenerate with 'pnpm generate:api'
|
||||
* SigNoz
|
||||
*/
|
||||
import { useMutation, useQuery } from 'react-query';
|
||||
import type {
|
||||
InvalidateOptions,
|
||||
MutationFunction,
|
||||
QueryClient,
|
||||
QueryFunction,
|
||||
QueryKey,
|
||||
UseMutationOptions,
|
||||
UseMutationResult,
|
||||
UseQueryOptions,
|
||||
UseQueryResult,
|
||||
} from 'react-query';
|
||||
|
||||
import type {
|
||||
ListPromotedPaths200,
|
||||
ListPromotedPathsParams,
|
||||
PromotetypesPromotePathDTO,
|
||||
RenderErrorResponseDTO,
|
||||
} from '../sigNoz.schemas';
|
||||
|
||||
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;
|
||||
};
|
||||
|
||||
/**
|
||||
* This endpoint lists the promoted paths of every JSON column, each annotated with its signal and context. The signal, context, promoted and indexes query parameters filter the listing.
|
||||
* @summary List promoted paths
|
||||
*/
|
||||
export const listPromotedPaths = (
|
||||
params?: ListPromotedPathsParams,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<ListPromotedPaths200>({
|
||||
url: `/api/v1/promoted_path`,
|
||||
method: 'GET',
|
||||
params,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getListPromotedPathsQueryKey = (
|
||||
params?: ListPromotedPathsParams,
|
||||
) => {
|
||||
return [`/api/v1/promoted_path`, ...(params ? [params] : [])] as const;
|
||||
};
|
||||
|
||||
export const getListPromotedPathsQueryOptions = <
|
||||
TData = Awaited<ReturnType<typeof listPromotedPaths>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
params?: ListPromotedPathsParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listPromotedPaths>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey =
|
||||
queryOptions?.queryKey ?? getListPromotedPathsQueryKey(params);
|
||||
|
||||
const queryFn: QueryFunction<
|
||||
Awaited<ReturnType<typeof listPromotedPaths>>
|
||||
> = ({ signal }) => listPromotedPaths(params, signal);
|
||||
|
||||
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listPromotedPaths>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type ListPromotedPathsQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof listPromotedPaths>>
|
||||
>;
|
||||
export type ListPromotedPathsQueryError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary List promoted paths
|
||||
*/
|
||||
|
||||
export function useListPromotedPaths<
|
||||
TData = Awaited<ReturnType<typeof listPromotedPaths>>,
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
>(
|
||||
params?: ListPromotedPathsParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof listPromotedPaths>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getListPromotedPathsQueryOptions(params, options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return withQueryKey(query, queryOptions.queryKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary List promoted paths
|
||||
*/
|
||||
export const invalidateListPromotedPaths = async (
|
||||
queryClient: QueryClient,
|
||||
params?: ListPromotedPathsParams,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
await queryClient.invalidateQueries(
|
||||
{ queryKey: getListPromotedPathsQueryKey(params) },
|
||||
options,
|
||||
);
|
||||
|
||||
return queryClient;
|
||||
};
|
||||
|
||||
/**
|
||||
* This endpoint promotes paths of JSON columns to their promoted columns. Each path names its promotion domain with its signal and context, e.g. traces/attribute.
|
||||
* @summary Promote paths
|
||||
*/
|
||||
export const promotePaths = (
|
||||
promotetypesPromotePathDTONull?: BodyType<
|
||||
PromotetypesPromotePathDTO[] | null
|
||||
> | null,
|
||||
signal?: AbortSignal,
|
||||
) => {
|
||||
return GeneratedAPIInstance<void>({
|
||||
url: `/api/v1/promoted_path`,
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
data: promotetypesPromotePathDTONull,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getPromotePathsMutationOptions = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof promotePaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof promotePaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
> => {
|
||||
const mutationKey = ['promotePaths'];
|
||||
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 promotePaths>>,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> }
|
||||
> = (props) => {
|
||||
const { data } = props ?? {};
|
||||
|
||||
return promotePaths(data);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type PromotePathsMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof promotePaths>>
|
||||
>;
|
||||
export type PromotePathsMutationBody =
|
||||
| BodyType<PromotetypesPromotePathDTO[] | null>
|
||||
| undefined;
|
||||
export type PromotePathsMutationError = ErrorType<RenderErrorResponseDTO>;
|
||||
|
||||
/**
|
||||
* @summary Promote paths
|
||||
*/
|
||||
export const usePromotePaths = <
|
||||
TError = ErrorType<RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof promotePaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof promotePaths>>,
|
||||
TError,
|
||||
{ data?: BodyType<PromotetypesPromotePathDTO[] | null> },
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getPromotePathsMutationOptions(options));
|
||||
};
|
||||
@@ -9369,6 +9369,10 @@ export interface PromotetypesWrappedIndexDTO {
|
||||
}
|
||||
|
||||
export interface PromotetypesPromotePathDTO {
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
context: string;
|
||||
/**
|
||||
* @type array
|
||||
*/
|
||||
@@ -9376,11 +9380,15 @@ export interface PromotetypesPromotePathDTO {
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
path?: string;
|
||||
path: string;
|
||||
/**
|
||||
* @type boolean
|
||||
*/
|
||||
promote?: boolean;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
signal: string;
|
||||
}
|
||||
|
||||
export interface Querybuildertypesv5AggregationMetaDTO {
|
||||
@@ -12470,17 +12478,6 @@ export type ListUnmappedLLMModels200 = {
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListPromotedAndIndexedPaths200 = {
|
||||
/**
|
||||
* @type array,null
|
||||
*/
|
||||
data: PromotetypesPromotePathDTO[] | null;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListOrgPreferences200 = {
|
||||
/**
|
||||
* @type array
|
||||
@@ -12506,6 +12503,40 @@ export type GetOrgPreference200 = {
|
||||
export type UpdateOrgPreferencePathParameters = {
|
||||
name: string;
|
||||
};
|
||||
export type ListPromotedPathsParams = {
|
||||
/**
|
||||
* @type string
|
||||
* @description undefined
|
||||
*/
|
||||
signal?: string;
|
||||
/**
|
||||
* @type string
|
||||
* @description undefined
|
||||
*/
|
||||
context?: string;
|
||||
/**
|
||||
* @type boolean,null
|
||||
* @description undefined
|
||||
*/
|
||||
promoted?: boolean | null;
|
||||
/**
|
||||
* @type boolean,null
|
||||
* @description undefined
|
||||
*/
|
||||
indexes?: boolean | null;
|
||||
};
|
||||
|
||||
export type ListPromotedPaths200 = {
|
||||
/**
|
||||
* @type array,null
|
||||
*/
|
||||
data: PromotetypesPromotePathDTO[] | null;
|
||||
/**
|
||||
* @type string
|
||||
*/
|
||||
status: string;
|
||||
};
|
||||
|
||||
export type ListRoles200 = {
|
||||
/**
|
||||
* @type array
|
||||
|
||||
@@ -47,4 +47,5 @@ export enum LOCALSTORAGE {
|
||||
DASHBOARDS_LIST_VIEWS = 'DASHBOARDS_LIST_VIEWS',
|
||||
DASHBOARD_V2_PANEL_COLUMN_WIDTHS = 'DASHBOARD_V2_PANEL_COLUMN_WIDTHS',
|
||||
LLM_ATTRIBUTE_MAPPING_TEST_SPAN = 'LLM_ATTRIBUTE_MAPPING_TEST_SPAN',
|
||||
SAVED_VIEW_ENABLED = 'SAVED_VIEW_ENABLED',
|
||||
}
|
||||
|
||||
@@ -53,6 +53,10 @@
|
||||
z-index: 0;
|
||||
background: var(--l1-background);
|
||||
|
||||
// Column so the bottom strip sits under the scrolling content, not inside it.
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
|
||||
&.full-screen-content {
|
||||
width: 100%;
|
||||
}
|
||||
@@ -70,7 +74,9 @@
|
||||
|
||||
.chat-support-gateway {
|
||||
position: fixed;
|
||||
bottom: 20px;
|
||||
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
|
||||
// UI belongs in the bounded layout, not in another offset here.
|
||||
bottom: calc(20px + var(--bottom-strip-height, 0px));
|
||||
right: 20px;
|
||||
z-index: 1000;
|
||||
|
||||
|
||||
@@ -43,6 +43,7 @@ import { USER_PREFERENCES } from 'constants/userPreferences';
|
||||
import AIAssistantModal from 'container/AIAssistant/AIAssistantModal';
|
||||
import AIAssistantPanel from 'container/AIAssistant/AIAssistantPanel';
|
||||
import { useAIAssistantStore } from 'container/AIAssistant/store/useAIAssistantStore';
|
||||
import BottomStrip from 'container/BottomStrip';
|
||||
import SideNav from 'container/SideNav';
|
||||
import TopNav from 'container/TopNav';
|
||||
import dayjs from 'dayjs';
|
||||
@@ -51,6 +52,7 @@ import { useIsDarkMode } from 'hooks/useDarkMode';
|
||||
import { useGetTenantLicense } from 'hooks/useGetTenantLicense';
|
||||
import { useIsAIAssistantEnabled } from 'hooks/useIsAIAssistantEnabled';
|
||||
import { useNotifications } from 'hooks/useNotifications';
|
||||
import { useSavedViewEnabled } from 'hooks/useSavedViewEnabled';
|
||||
import useTabVisibility from 'hooks/useTabFocus';
|
||||
import history from 'lib/history';
|
||||
import { isNull } from 'lodash-es';
|
||||
@@ -402,6 +404,7 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
|
||||
}, [pathname]);
|
||||
|
||||
const isToDisplayLayout = isLoggedIn;
|
||||
const isSavedViewEnabled = useSavedViewEnabled();
|
||||
|
||||
const routeKey = useMemo(() => getRouteKey(pathname), [pathname]);
|
||||
const pageTitle = t(routeKey);
|
||||
@@ -868,6 +871,10 @@ function AppLayout(props: AppLayoutProps): JSX.Element {
|
||||
</OverlayScrollbar>
|
||||
</LayoutContent>
|
||||
</Sentry.ErrorBoundary>
|
||||
|
||||
{isSavedViewEnabled && isToDisplayLayout && !renderFullScreen && (
|
||||
<BottomStrip />
|
||||
)}
|
||||
</div>
|
||||
|
||||
{isLoggedIn && isAIAssistantEnabled && (
|
||||
|
||||
@@ -12,8 +12,12 @@ export const Layout = styled(LayoutComponent)`
|
||||
}
|
||||
`;
|
||||
|
||||
// Takes the height left in `.app-content` after the bottom strip.
|
||||
// `min-height: 0` is not needed right now, overlayscrollbars already sets
|
||||
// `overflow: auto` here. Kept so this does not break if that goes away.
|
||||
export const LayoutContent = styled(LayoutComponent.Content)`
|
||||
height: 100%;
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
&::-webkit-scrollbar {
|
||||
width: 0.1rem;
|
||||
}
|
||||
|
||||
36
frontend/src/container/BottomStrip/BottomStrip.module.scss
Normal file
36
frontend/src/container/BottomStrip/BottomStrip.module.scss
Normal file
@@ -0,0 +1,36 @@
|
||||
.strip {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: space-between;
|
||||
gap: var(--spacing-6);
|
||||
|
||||
flex-shrink: 0;
|
||||
height: var(--bottom-strip-height);
|
||||
padding: 0 var(--spacing-6);
|
||||
|
||||
background: var(--l2-background);
|
||||
border-top: 1px solid var(--l2-border);
|
||||
|
||||
font-family: var(--font-family-sf-mono, monospace);
|
||||
|
||||
// Above page content, below the body-portalled overlays that are meant to
|
||||
// cover the strip.
|
||||
position: relative;
|
||||
z-index: 1;
|
||||
}
|
||||
|
||||
.left,
|
||||
.right {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
gap: var(--spacing-6);
|
||||
min-width: 0;
|
||||
}
|
||||
|
||||
// Temporary placeholder for the left slot. Replaced later.
|
||||
.version {
|
||||
color: var(--l2-foreground);
|
||||
white-space: nowrap;
|
||||
overflow: hidden;
|
||||
text-overflow: ellipsis;
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
import { render } from 'tests/test-utils';
|
||||
|
||||
import BottomStrip, {
|
||||
BOTTOM_STRIP_HEIGHT,
|
||||
BOTTOM_STRIP_HEIGHT_VAR,
|
||||
BOTTOM_STRIP_ON_CLASS,
|
||||
} from '..';
|
||||
|
||||
describe('BottomStrip', () => {
|
||||
it('publishes the body class and height property while mounted', () => {
|
||||
const { unmount } = render(<BottomStrip />);
|
||||
|
||||
expect(document.body.classList.contains(BOTTOM_STRIP_ON_CLASS)).toBe(true);
|
||||
expect(document.body.style.getPropertyValue(BOTTOM_STRIP_HEIGHT_VAR)).toBe(
|
||||
`${BOTTOM_STRIP_HEIGHT}px`,
|
||||
);
|
||||
|
||||
unmount();
|
||||
|
||||
expect(document.body.classList.contains(BOTTOM_STRIP_ON_CLASS)).toBe(false);
|
||||
expect(document.body.style.getPropertyValue(BOTTOM_STRIP_HEIGHT_VAR)).toBe(
|
||||
'',
|
||||
);
|
||||
});
|
||||
|
||||
// The string is whatever the Go build injected, so it is rendered untouched —
|
||||
// same as SideNav. Release tags carry the "v", local builds do not.
|
||||
it.each([['v0.134.67'], ['main-64f1c2a']])(
|
||||
'renders the build version %p exactly as given',
|
||||
(version) => {
|
||||
const { getByTestId } = render(<BottomStrip />, undefined, {
|
||||
appContextOverrides: {
|
||||
versionData: { version, ee: 'Y', setupCompleted: true },
|
||||
},
|
||||
});
|
||||
|
||||
expect(getByTestId('bottom-strip-version')).toHaveTextContent(version);
|
||||
},
|
||||
);
|
||||
|
||||
it('renders the strip without a version when none is available', () => {
|
||||
const { getByTestId, queryByTestId } = render(<BottomStrip />, undefined, {
|
||||
appContextOverrides: { versionData: null },
|
||||
});
|
||||
|
||||
expect(getByTestId('bottom-strip')).toBeInTheDocument();
|
||||
expect(queryByTestId('bottom-strip-version')).not.toBeInTheDocument();
|
||||
});
|
||||
});
|
||||
42
frontend/src/container/BottomStrip/index.tsx
Normal file
42
frontend/src/container/BottomStrip/index.tsx
Normal file
@@ -0,0 +1,42 @@
|
||||
import { useLayoutEffect } from 'react';
|
||||
import { useAppContext } from 'providers/App/App';
|
||||
|
||||
import styles from './BottomStrip.module.scss';
|
||||
|
||||
export const BOTTOM_STRIP_HEIGHT = 24;
|
||||
|
||||
export const BOTTOM_STRIP_ON_CLASS = 'bottom-strip-on';
|
||||
export const BOTTOM_STRIP_HEIGHT_VAR = '--bottom-strip-height';
|
||||
|
||||
function BottomStrip(): JSX.Element {
|
||||
const { versionData } = useAppContext();
|
||||
const version = versionData?.version?.trim();
|
||||
|
||||
useLayoutEffect(() => {
|
||||
document.body.classList.add(BOTTOM_STRIP_ON_CLASS);
|
||||
document.body.style.setProperty(
|
||||
BOTTOM_STRIP_HEIGHT_VAR,
|
||||
`${BOTTOM_STRIP_HEIGHT}px`,
|
||||
);
|
||||
|
||||
return (): void => {
|
||||
document.body.classList.remove(BOTTOM_STRIP_ON_CLASS);
|
||||
document.body.style.removeProperty(BOTTOM_STRIP_HEIGHT_VAR);
|
||||
};
|
||||
}, []);
|
||||
|
||||
return (
|
||||
<div className={styles.strip} data-testid="bottom-strip">
|
||||
<div className={styles.left}>
|
||||
{version && (
|
||||
<span className={styles.version} data-testid="bottom-strip-version">
|
||||
{version}
|
||||
</span>
|
||||
)}
|
||||
</div>
|
||||
<div className={styles.right} />
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
export default BottomStrip;
|
||||
@@ -1,6 +1,8 @@
|
||||
.create-alert-v2-footer {
|
||||
position: fixed;
|
||||
bottom: 0;
|
||||
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
|
||||
// UI belongs in the bounded layout, not in another offset here.
|
||||
bottom: var(--bottom-strip-height, 0px);
|
||||
left: 63px;
|
||||
right: 0;
|
||||
background-color: var(--l1-background);
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
.explorer-options-container {
|
||||
position: fixed;
|
||||
bottom: 0px;
|
||||
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
|
||||
// UI belongs in the bounded layout, not in another offset here.
|
||||
bottom: var(--bottom-strip-height, 0px);
|
||||
left: calc(50% + 240px);
|
||||
transform: translate(calc(-50% - 120px), 0);
|
||||
transition: left 0.2s linear;
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
.explorer-option-droppable-container {
|
||||
position: fixed;
|
||||
bottom: 0;
|
||||
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
|
||||
// UI belongs in the bounded layout, not in another offset here.
|
||||
bottom: var(--bottom-strip-height, 0px);
|
||||
width: -webkit-fill-available;
|
||||
height: 24px;
|
||||
display: flex;
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
.home-container {
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
min-height: 100vh;
|
||||
overflow-y: auto;
|
||||
height: 100%;
|
||||
width: 100%;
|
||||
|
||||
@@ -1,7 +1,4 @@
|
||||
.licenses-page {
|
||||
max-height: 100vh;
|
||||
overflow: hidden;
|
||||
|
||||
.licenses-page-header {
|
||||
border-bottom: 1px solid var(--l1-border);
|
||||
background: var(--l1-background);
|
||||
@@ -32,7 +29,6 @@
|
||||
|
||||
.licenses-page-content {
|
||||
flex: 1;
|
||||
height: calc(100vh - 48px);
|
||||
background: var(--l1-background);
|
||||
padding: 10px 8px;
|
||||
overflow-y: auto;
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
gap: 1rem;
|
||||
height: calc(100vh - 62px);
|
||||
flex: 1;
|
||||
min-height: 400px;
|
||||
}
|
||||
|
||||
|
||||
@@ -181,7 +181,9 @@
|
||||
|
||||
.ant-pagination {
|
||||
position: fixed;
|
||||
bottom: 0;
|
||||
// Lifted above the bottom strip. Don't extend this pattern — new
|
||||
// fixed-bottom UI belongs in the bounded layout, not in another offset here.
|
||||
bottom: var(--bottom-strip-height, 0px);
|
||||
width: calc(100% - 54px);
|
||||
background: var(--l1-background);
|
||||
padding: 16px;
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
gap: 1rem;
|
||||
height: calc(100vh - 62px);
|
||||
flex: 1;
|
||||
min-height: 400px;
|
||||
padding-top: var(--spacing-8);
|
||||
}
|
||||
|
||||
@@ -1,7 +1,4 @@
|
||||
.version-container {
|
||||
max-height: 100vh;
|
||||
overflow: hidden;
|
||||
|
||||
.version-page-header {
|
||||
border-bottom: 1px solid var(--l1-border);
|
||||
background: var(--l1-background);
|
||||
|
||||
11
frontend/src/hooks/useSavedViewEnabled.ts
Normal file
11
frontend/src/hooks/useSavedViewEnabled.ts
Normal file
@@ -0,0 +1,11 @@
|
||||
import getLocalStorageKey from 'api/browser/localstorage/get';
|
||||
import { LOCALSTORAGE } from 'constants/localStorage';
|
||||
import { useState } from 'react';
|
||||
|
||||
export function useSavedViewEnabled(): boolean {
|
||||
const [isEnabled] = useState(
|
||||
() => getLocalStorageKey(LOCALSTORAGE.SAVED_VIEW_ENABLED) === 'true',
|
||||
);
|
||||
|
||||
return isEnabled;
|
||||
}
|
||||
@@ -1,4 +1,29 @@
|
||||
.alerts-container {
|
||||
// Hands the page height down to the active tab so its content can bound itself
|
||||
// instead of guessing with 100vh. Child combinators only, nested Tabs
|
||||
// (Configuration) must not be caught.
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
|
||||
> .ant-tabs-content-holder {
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
|
||||
> .ant-tabs-content {
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
|
||||
> .ant-tabs-tabpane-active {
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
.top-level-tab.periscope-tab {
|
||||
padding: 2px 0;
|
||||
}
|
||||
@@ -40,5 +65,9 @@
|
||||
|
||||
.alert-rules-container {
|
||||
margin-top: 10px;
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,7 +2,9 @@
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
position: fixed;
|
||||
bottom: 0;
|
||||
// Lifted above the bottom strip. Don't extend this pattern — new fixed-bottom
|
||||
// UI belongs in the bounded layout, not in another offset here.
|
||||
bottom: var(--bottom-strip-height, 0px);
|
||||
left: 0;
|
||||
width: 100%;
|
||||
z-index: 100;
|
||||
|
||||
@@ -1,7 +1,4 @@
|
||||
.support-page-container {
|
||||
max-height: 100vh;
|
||||
overflow: hidden;
|
||||
|
||||
.support-page-header {
|
||||
border-bottom: 1px solid var(--l1-border);
|
||||
background: var(--l1-background);
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
.root {
|
||||
height: calc(100vh);
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
}
|
||||
|
||||
@@ -1,13 +1,24 @@
|
||||
.traces-funnel-details {
|
||||
display: flex;
|
||||
// 45px -> height of the tab bar
|
||||
height: calc(100vh - 45px);
|
||||
height: 100%;
|
||||
|
||||
&__steps-config {
|
||||
flex-shrink: 0;
|
||||
width: 600px;
|
||||
border-right: 1px solid var(--l1-border);
|
||||
// Positioning context for the absolute .steps-footer.
|
||||
position: relative;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
|
||||
// Scoped here so the modal usage of FunnelConfiguration on trace details
|
||||
// stays in normal flow.
|
||||
.funnel-configuration {
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
}
|
||||
}
|
||||
&__steps-results {
|
||||
width: 100%;
|
||||
|
||||
@@ -4,14 +4,17 @@
|
||||
flex-direction: column;
|
||||
justify-content: flex-start;
|
||||
&.funnel-details-page {
|
||||
height: calc(
|
||||
100vh - 170px
|
||||
); // 64px bottom bar + 61px configuration header + 45px page navbar
|
||||
flex: 1;
|
||||
min-height: 0;
|
||||
// .steps-footer is absolute against the config column, so its 64px is
|
||||
// reserved rather than laid out.
|
||||
margin-bottom: 64px;
|
||||
overflow: auto;
|
||||
}
|
||||
}
|
||||
|
||||
&__header {
|
||||
flex-shrink: 0;
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: space-between;
|
||||
|
||||
@@ -10,11 +10,11 @@ import (
|
||||
)
|
||||
|
||||
func (provider *provider) addPromoteRoutes(router *mux.Router) error {
|
||||
if err := router.Handle("/api/v1/logs/promote_paths", handler.New(provider.authzMiddleware.EditAccess(provider.promoteHandler.HandlePromoteAndIndexPaths), handler.OpenAPIDef{
|
||||
ID: "HandlePromoteAndIndexPaths",
|
||||
Tags: []string{"logs"},
|
||||
Summary: "Promote and index paths",
|
||||
Description: "This endpoints promotes and indexes paths",
|
||||
if err := router.Handle("/api/v1/promoted_path", handler.New(provider.authzMiddleware.EditAccess(provider.promoteHandler.PromotePaths), handler.OpenAPIDef{
|
||||
ID: "PromotePaths",
|
||||
Tags: []string{"promote"},
|
||||
Summary: "Promote paths",
|
||||
Description: "This endpoint promotes paths of JSON columns to their promoted columns. Each path names its promotion domain with its signal and context, e.g. traces/attribute.",
|
||||
Request: new([]*promotetypes.PromotePath),
|
||||
RequestContentType: "application/json",
|
||||
Response: nil,
|
||||
@@ -26,12 +26,13 @@ func (provider *provider) addPromoteRoutes(router *mux.Router) error {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := router.Handle("/api/v1/logs/promote_paths", handler.New(provider.authzMiddleware.ViewAccess(provider.promoteHandler.ListPromotedAndIndexedPaths), handler.OpenAPIDef{
|
||||
ID: "ListPromotedAndIndexedPaths",
|
||||
Tags: []string{"logs"},
|
||||
Summary: "Promote and index paths",
|
||||
Description: "This endpoints promotes and indexes paths",
|
||||
if err := router.Handle("/api/v1/promoted_path", handler.New(provider.authzMiddleware.ViewAccess(provider.promoteHandler.ListPromotedPaths), handler.OpenAPIDef{
|
||||
ID: "ListPromotedPaths",
|
||||
Tags: []string{"promote"},
|
||||
Summary: "List promoted paths",
|
||||
Description: "This endpoint lists the promoted paths of every JSON column, each annotated with its signal and context. The signal, context, promoted and indexes query parameters filter the listing.",
|
||||
Request: nil,
|
||||
RequestQuery: new(promotetypes.ListPromotedPathsFilters),
|
||||
RequestContentType: "",
|
||||
Response: new([]*promotetypes.PromotePath),
|
||||
ResponseContentType: "",
|
||||
|
||||
@@ -19,7 +19,7 @@ func NewHandler(module promote.Module) promote.Handler {
|
||||
return &handler{module: module}
|
||||
}
|
||||
|
||||
func (h *handler) HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Request) {
|
||||
func (h *handler) PromotePaths(w http.ResponseWriter, r *http.Request) {
|
||||
// TODO(Nitya): Use in multi tenant setup
|
||||
_, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
@@ -33,7 +33,7 @@ func (h *handler) HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Requ
|
||||
return
|
||||
}
|
||||
|
||||
err = h.module.PromoteAndIndexPaths(r.Context(), req...)
|
||||
err = h.module.PromotePaths(r.Context(), req...)
|
||||
if err != nil {
|
||||
render.Error(w, err)
|
||||
return
|
||||
@@ -42,7 +42,7 @@ func (h *handler) HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Requ
|
||||
render.Success(w, http.StatusCreated, nil)
|
||||
}
|
||||
|
||||
func (h *handler) ListPromotedAndIndexedPaths(w http.ResponseWriter, r *http.Request) {
|
||||
func (h *handler) ListPromotedPaths(w http.ResponseWriter, r *http.Request) {
|
||||
// TODO(Nitya): Use in multi tenant setup
|
||||
_, err := authtypes.ClaimsFromContext(r.Context())
|
||||
if err != nil {
|
||||
@@ -50,7 +50,17 @@ func (h *handler) ListPromotedAndIndexedPaths(w http.ResponseWriter, r *http.Req
|
||||
return
|
||||
}
|
||||
|
||||
paths, err := h.module.ListPromotedAndIndexedPaths(r.Context())
|
||||
var filters promotetypes.ListPromotedPathsFilters
|
||||
if err := binding.Query.BindQuery(r.URL.Query(), &filters); err != nil {
|
||||
render.Error(w, err)
|
||||
return
|
||||
}
|
||||
if err := filters.Validate(); err != nil {
|
||||
render.Error(w, err)
|
||||
return
|
||||
}
|
||||
|
||||
paths, err := h.module.ListPromotedPaths(r.Context(), filters)
|
||||
if err != nil {
|
||||
render.Error(w, err)
|
||||
return
|
||||
|
||||
@@ -2,14 +2,11 @@ package implpromote
|
||||
|
||||
import (
|
||||
"context"
|
||||
"maps"
|
||||
"slices"
|
||||
"strings"
|
||||
|
||||
schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/modules/promote"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
|
||||
@@ -31,46 +28,77 @@ func NewModule(metadataStore telemetrytypes.MetadataStore, telemetrystore teleme
|
||||
return &module{metadataStore: metadataStore, telemetryStore: telemetrystore}
|
||||
}
|
||||
|
||||
func (m *module) ListPromotedAndIndexedPaths(ctx context.Context) ([]promotetypes.PromotePath, error) {
|
||||
func (m *module) ListPromotedPaths(ctx context.Context, filters promotetypes.ListPromotedPathsFilters) ([]promotetypes.PromotePath, error) {
|
||||
response := make([]promotetypes.PromotePath, 0)
|
||||
for _, target := range promotetypes.Targets() {
|
||||
if !filters.MatchesTarget(target) {
|
||||
continue
|
||||
}
|
||||
paths, err := m.listPromotedPaths(ctx, target)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, path := range paths {
|
||||
if filters.MatchesPath(path) {
|
||||
response = append(response, path)
|
||||
}
|
||||
}
|
||||
}
|
||||
return response, nil
|
||||
}
|
||||
|
||||
func (m *module) listPromotedPaths(ctx context.Context, target promotetypes.Target) ([]promotetypes.PromotePath, error) {
|
||||
promotedPaths, err := m.metadataStore.GetPromotedPaths(ctx, target.Entry)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
response := make([]promotetypes.PromotePath, 0, len(promotedPaths))
|
||||
for path := range promotedPaths {
|
||||
response = append(response, promotetypes.PromotePath{
|
||||
Signal: target.Entry.Signal.StringValue(),
|
||||
Context: target.Entry.FieldContext.StringValue(),
|
||||
Path: target.RequiredPathPrefix + path,
|
||||
Promote: true,
|
||||
})
|
||||
}
|
||||
|
||||
if !target.IndexesSupported {
|
||||
return response, nil
|
||||
}
|
||||
|
||||
indexes, err := m.metadataStore.ListLogsJSONIndexes(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// aggr keys are full sub-column paths: index.BaseColumn carries the
|
||||
// column prefix and index.Name the bare path.
|
||||
aggr := map[string][]promotetypes.WrappedIndex{}
|
||||
for _, index := range indexes {
|
||||
aggr[index.Name] = append(aggr[index.Name], promotetypes.WrappedIndex{
|
||||
fullPath := index.BaseColumn + index.Name
|
||||
aggr[fullPath] = append(aggr[fullPath], promotetypes.WrappedIndex{
|
||||
FieldDataType: index.FieldDataType,
|
||||
Type: index.IndexType,
|
||||
Granularity: index.Granularity,
|
||||
})
|
||||
}
|
||||
promotedPaths, err := m.listPromotedPaths(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
response := []promotetypes.PromotePath{}
|
||||
for _, path := range promotedPaths {
|
||||
fullPath := logstelemetryschema.BodyPromotedColumnPrefix + path
|
||||
path = telemetrytypes.BodyJSONStringSearchPrefix + path
|
||||
item := promotetypes.PromotePath{
|
||||
Path: path,
|
||||
Promote: true,
|
||||
}
|
||||
indexes, ok := aggr[fullPath]
|
||||
if ok {
|
||||
item.Indexes = indexes
|
||||
for i := range response {
|
||||
fullPath := target.PromotedColumnPrefix() + strings.TrimPrefix(response[i].Path, target.RequiredPathPrefix)
|
||||
if indexes, ok := aggr[fullPath]; ok {
|
||||
response[i].Indexes = indexes
|
||||
delete(aggr, fullPath)
|
||||
}
|
||||
response = append(response, item)
|
||||
}
|
||||
|
||||
// add the paths that are not promoted but have indexes
|
||||
for path, indexes := range aggr {
|
||||
path := strings.TrimPrefix(path, logstelemetryschema.BodyV2ColumnPrefix)
|
||||
path = telemetrytypes.BodyJSONStringSearchPrefix + path
|
||||
for fullPath, indexes := range aggr {
|
||||
path := strings.TrimPrefix(fullPath, target.BaseColumnPrefix())
|
||||
path = strings.TrimPrefix(path, target.PromotedColumnPrefix())
|
||||
path = target.RequiredPathPrefix + path
|
||||
response = append(response, promotetypes.PromotePath{
|
||||
Signal: target.Entry.Signal.StringValue(),
|
||||
Context: target.Entry.FieldContext.StringValue(),
|
||||
Path: path,
|
||||
Indexes: indexes,
|
||||
})
|
||||
@@ -78,68 +106,42 @@ func (m *module) ListPromotedAndIndexedPaths(ctx context.Context) ([]promotetype
|
||||
return response, nil
|
||||
}
|
||||
|
||||
func (m *module) listPromotedPaths(ctx context.Context) ([]string, error) {
|
||||
paths, err := m.metadataStore.GetPromotedPaths(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return slices.Collect(maps.Keys(paths)), nil
|
||||
}
|
||||
|
||||
// PromotePaths inserts provided JSON paths into the promoted paths table for logs queries.
|
||||
func (m *module) PromotePaths(ctx context.Context, paths []string) error {
|
||||
func (m *module) PromotePaths(ctx context.Context, paths ...*promotetypes.PromotePath) error {
|
||||
if len(paths) == 0 {
|
||||
return errors.NewInvalidInputf(errors.CodeInvalidInput, "paths cannot be empty")
|
||||
}
|
||||
|
||||
return m.metadataStore.PromotePaths(ctx, paths...)
|
||||
}
|
||||
|
||||
// createIndexes creates string ngram + token filter indexes on JSON path subcolumns for LIKE queries.
|
||||
func (m *module) createIndexes(ctx context.Context, indexes []schemamigrator.Index) error {
|
||||
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
|
||||
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalLogs.StringValue(),
|
||||
instrumentationtypes.CodeNamespace: "promote",
|
||||
instrumentationtypes.CodeFunctionName: "createIndexes",
|
||||
})
|
||||
if len(indexes) == 0 {
|
||||
return nil
|
||||
byTarget := map[promotetypes.Target][]*promotetypes.PromotePath{}
|
||||
targets := []promotetypes.Target{}
|
||||
for _, path := range paths {
|
||||
target, err := path.Target()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := path.ValidateAndSetDefaults(target); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, ok := byTarget[target]; !ok {
|
||||
targets = append(targets, target)
|
||||
}
|
||||
byTarget[target] = append(byTarget[target], path)
|
||||
}
|
||||
|
||||
for _, index := range indexes {
|
||||
alterStmt := schemamigrator.AlterTableAddIndex{
|
||||
Database: logstelemetryschema.DBName,
|
||||
Table: logstelemetryschema.LogsV2LocalTableName,
|
||||
Index: index,
|
||||
}
|
||||
op := alterStmt.OnCluster(m.telemetryStore.Cluster())
|
||||
if err := m.telemetryStore.ClickhouseDB().Exec(ctx, op.ToSQL()); err != nil {
|
||||
return errors.WrapInternalf(err, CodeFailedToCreateIndex, "failed to create index")
|
||||
for _, target := range targets {
|
||||
if err := m.promotePaths(ctx, target, byTarget[target]...); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// PromoteAndIndexPaths handles promoting paths and creating indexes in one call.
|
||||
func (m *module) PromoteAndIndexPaths(
|
||||
ctx context.Context,
|
||||
paths ...*promotetypes.PromotePath,
|
||||
) error {
|
||||
if len(paths) == 0 {
|
||||
return errors.NewInvalidInputf(errors.CodeInvalidInput, "paths cannot be empty")
|
||||
}
|
||||
|
||||
func (m *module) promotePaths(ctx context.Context, target promotetypes.Target, paths ...*promotetypes.PromotePath) error {
|
||||
pathsStr := []string{}
|
||||
// validate the paths
|
||||
for _, path := range paths {
|
||||
if err := path.ValidateAndSetDefaults(); err != nil {
|
||||
return err
|
||||
}
|
||||
pathsStr = append(pathsStr, path.Path)
|
||||
}
|
||||
|
||||
existingPromotedPaths, err := m.metadataStore.GetPromotedPaths(ctx, pathsStr...)
|
||||
existingPromotedPaths, err := m.metadataStore.GetPromotedPaths(ctx, target.Entry, pathsStr...)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -153,10 +155,10 @@ func (m *module) PromoteAndIndexPaths(
|
||||
}
|
||||
}
|
||||
if len(it.Indexes) > 0 {
|
||||
parentColumn := logstelemetryschema.LogsV2BodyV2Column
|
||||
parentColumn := target.BaseColumn
|
||||
// if the path is already promoted or is being promoted, add it to the promoted column
|
||||
if _, promoted := existingPromotedPaths[it.Path]; promoted || it.Promote {
|
||||
parentColumn = logstelemetryschema.LogsV2BodyPromotedColumn
|
||||
parentColumn = target.PromotedColumn()
|
||||
}
|
||||
|
||||
for _, index := range it.Indexes {
|
||||
@@ -182,17 +184,42 @@ func (m *module) PromoteAndIndexPaths(
|
||||
}
|
||||
|
||||
if len(toInsert) > 0 {
|
||||
err := m.PromotePaths(ctx, toInsert)
|
||||
err := m.metadataStore.PromotePaths(ctx, target.Entry, toInsert...)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
if len(indexes) > 0 {
|
||||
if err := m.createIndexes(ctx, indexes); err != nil {
|
||||
if err := m.createIndexes(ctx, target, indexes); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *module) createIndexes(ctx context.Context, target promotetypes.Target, indexes []schemamigrator.Index) error {
|
||||
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
|
||||
instrumentationtypes.TelemetrySignal: target.Entry.Signal.StringValue(),
|
||||
instrumentationtypes.CodeNamespace: "promote",
|
||||
instrumentationtypes.CodeFunctionName: "createIndexes",
|
||||
})
|
||||
if len(indexes) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
for _, index := range indexes {
|
||||
alterStmt := schemamigrator.AlterTableAddIndex{
|
||||
Database: target.DBName,
|
||||
Table: target.LocalTableName,
|
||||
Index: index,
|
||||
}
|
||||
op := alterStmt.OnCluster(m.telemetryStore.Cluster())
|
||||
if err := m.telemetryStore.ClickhouseDB().Exec(ctx, op.ToSQL()); err != nil {
|
||||
return errors.WrapInternalf(err, CodeFailedToCreateIndex, "failed to create index")
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
325
pkg/modules/promote/implpromote/module_test.go
Normal file
325
pkg/modules/promote/implpromote/module_test.go
Normal file
@@ -0,0 +1,325 @@
|
||||
package implpromote
|
||||
|
||||
import (
|
||||
"context"
|
||||
"regexp"
|
||||
"testing"
|
||||
|
||||
sqlmock "github.com/DATA-DOG/go-sqlmock"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
|
||||
"github.com/SigNoz/signoz/pkg/types/promotetypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestPromotePaths(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
paths []*promotetypes.PromotePath
|
||||
promoteTwice bool
|
||||
wantErr bool
|
||||
wantPromoted []string
|
||||
}{
|
||||
{
|
||||
name: "PromotesNewAttributes_Idempotent",
|
||||
paths: []*promotetypes.PromotePath{
|
||||
{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true},
|
||||
{Signal: "traces", Context: "attribute", Path: "span.operation", Promote: true},
|
||||
},
|
||||
promoteTwice: true,
|
||||
wantPromoted: []string{"http.method", "span.operation"},
|
||||
},
|
||||
{
|
||||
name: "MixedDomains_RecordedPerTarget",
|
||||
paths: []*promotetypes.PromotePath{
|
||||
{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true},
|
||||
{Signal: "logs", Context: "body", Path: "body.user.name", Promote: true},
|
||||
},
|
||||
wantPromoted: []string{"http.method", "user.name"},
|
||||
},
|
||||
{
|
||||
name: "NonPromoteEntries_NotRecorded",
|
||||
paths: []*promotetypes.PromotePath{{Signal: "traces", Context: "attribute", Path: "http.method"}},
|
||||
},
|
||||
{
|
||||
name: "InvalidSignal_Rejected",
|
||||
paths: []*promotetypes.PromotePath{{Signal: "events", Context: "attribute", Path: "http.method", Promote: true}},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "UnsupportedDomain_Rejected",
|
||||
paths: []*promotetypes.PromotePath{{Signal: "metrics", Context: "attribute", Path: "http.method", Promote: true}},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "ColumnPrefixedPath_Rejected",
|
||||
paths: []*promotetypes.PromotePath{{Signal: "traces", Context: "attribute", Path: "attributes.http.method", Promote: true}},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "EmptyPath_Rejected",
|
||||
paths: []*promotetypes.PromotePath{{Signal: "traces", Context: "attribute", Path: "", Promote: true}},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "EmptyRequest_Rejected",
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "PromotesBodyPath_PrefixStripped",
|
||||
paths: []*promotetypes.PromotePath{{Signal: "logs", Context: "body", Path: "body.user.name", Promote: true}},
|
||||
wantPromoted: []string{"user.name"},
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
store := telemetrytypestest.NewMockMetadataStore()
|
||||
m := NewModule(store, nil)
|
||||
|
||||
err := m.PromotePaths(ctx, testCase.paths...)
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
assert.Empty(t, store.PromotedPathsMap)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
|
||||
require.Len(t, store.PromotedPathsMap, len(testCase.wantPromoted))
|
||||
for _, path := range testCase.wantPromoted {
|
||||
assert.True(t, store.PromotedPathsMap[path], path)
|
||||
}
|
||||
|
||||
if testCase.promoteTwice {
|
||||
require.NoError(t, m.PromotePaths(ctx, testCase.paths...))
|
||||
assert.Len(t, store.PromotedPathsMap, len(testCase.wantPromoted))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestPromotePathsCreatesIndexes(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
promoted map[string]bool
|
||||
path *promotetypes.PromotePath
|
||||
wantDDLColumn string
|
||||
}{
|
||||
{
|
||||
name: "NewPromotion_IndexesPromotedColumn",
|
||||
path: &promotetypes.PromotePath{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.user.name",
|
||||
Promote: true,
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
|
||||
},
|
||||
},
|
||||
wantDDLColumn: "dynamicElement(body_promoted.user.name",
|
||||
},
|
||||
{
|
||||
name: "AlreadyPromoted_IndexesPromotedColumn",
|
||||
promoted: map[string]bool{"user.name": true},
|
||||
path: &promotetypes.PromotePath{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.user.name",
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
|
||||
},
|
||||
},
|
||||
wantDDLColumn: "dynamicElement(body_promoted.user.name",
|
||||
},
|
||||
{
|
||||
name: "UnpromotedPath_IndexesBaseColumn",
|
||||
path: &promotetypes.PromotePath{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.user.name",
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
|
||||
},
|
||||
},
|
||||
wantDDLColumn: "dynamicElement(body_v2.user.name",
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
ts := telemetrystoretest.New(telemetrystore.Config{}, sqlmock.QueryMatcherRegexp)
|
||||
store := telemetrytypestest.NewMockMetadataStore()
|
||||
if testCase.promoted != nil {
|
||||
store.PromotedPathsMap = testCase.promoted
|
||||
}
|
||||
m := NewModule(store, ts)
|
||||
|
||||
ts.Mock().ExpectExec("ADD INDEX (.+)" + regexp.QuoteMeta(testCase.wantDDLColumn)).WillReturnError(nil)
|
||||
require.NoError(t, m.PromotePaths(ctx, testCase.path))
|
||||
assert.NoError(t, ts.Mock().ExpectationsWereMet())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestListPromotedPaths(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
trueValue := true
|
||||
falseValue := false
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
filters promotetypes.ListPromotedPathsFilters
|
||||
promoted map[string]bool
|
||||
indexes []telemetrytypes.TelemetryFieldKeySkipIndex
|
||||
wantPaths []promotetypes.PromotePath
|
||||
}{
|
||||
{
|
||||
name: "PromotedPaths_EveryDomainAnnotated",
|
||||
promoted: map[string]bool{"http.method": true},
|
||||
wantPaths: []promotetypes.PromotePath{
|
||||
{Signal: "logs", Context: "body", Path: "body.http.method", Promote: true},
|
||||
{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "SignalFilter_SkipsOtherDomains",
|
||||
filters: promotetypes.ListPromotedPathsFilters{Signal: "traces"},
|
||||
promoted: map[string]bool{"http.method": true},
|
||||
wantPaths: []promotetypes.PromotePath{{Signal: "traces", Context: "attribute", Path: "http.method", Promote: true}},
|
||||
},
|
||||
{
|
||||
name: "ContextFilter_SkipsOtherDomains",
|
||||
filters: promotetypes.ListPromotedPathsFilters{Context: "body"},
|
||||
promoted: map[string]bool{"http.method": true},
|
||||
wantPaths: []promotetypes.PromotePath{{Signal: "logs", Context: "body", Path: "body.http.method", Promote: true}},
|
||||
},
|
||||
{
|
||||
name: "IndexedPaths_MergedForSupportingDomains",
|
||||
promoted: map[string]bool{"user.name": true},
|
||||
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
|
||||
{
|
||||
Name: "user.name",
|
||||
FieldContext: telemetrytypes.FieldContextBody,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
BaseColumn: "body_promoted.",
|
||||
IndexType: "ngrambf_v1(4, 1024, 2, 0)",
|
||||
Granularity: 1,
|
||||
},
|
||||
{
|
||||
Name: "request.duration",
|
||||
FieldContext: telemetrytypes.FieldContextBody,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeFloat64,
|
||||
BaseColumn: "body_v2.",
|
||||
IndexType: "minmax",
|
||||
Granularity: 1,
|
||||
},
|
||||
},
|
||||
wantPaths: []promotetypes.PromotePath{
|
||||
{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.user.name",
|
||||
Promote: true,
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
|
||||
},
|
||||
},
|
||||
{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.request.duration",
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeFloat64, Type: "minmax", Granularity: 1},
|
||||
},
|
||||
},
|
||||
{
|
||||
Signal: "traces",
|
||||
Context: "attribute",
|
||||
Path: "user.name",
|
||||
Promote: true,
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "PromotedFalseFilter_IndexOnlyPaths",
|
||||
filters: promotetypes.ListPromotedPathsFilters{Promoted: &falseValue},
|
||||
promoted: map[string]bool{"user.name": true},
|
||||
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
|
||||
{
|
||||
Name: "request.duration",
|
||||
FieldContext: telemetrytypes.FieldContextBody,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeFloat64,
|
||||
BaseColumn: "body_v2.",
|
||||
IndexType: "minmax",
|
||||
Granularity: 1,
|
||||
},
|
||||
},
|
||||
wantPaths: []promotetypes.PromotePath{
|
||||
{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.request.duration",
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeFloat64, Type: "minmax", Granularity: 1},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "IndexesTrueFilter_PathsWithIndexes",
|
||||
filters: promotetypes.ListPromotedPathsFilters{Indexes: &trueValue},
|
||||
promoted: map[string]bool{"user.name": true},
|
||||
indexes: []telemetrytypes.TelemetryFieldKeySkipIndex{
|
||||
{
|
||||
Name: "user.name",
|
||||
FieldContext: telemetrytypes.FieldContextBody,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
BaseColumn: "body_promoted.",
|
||||
IndexType: "ngrambf_v1(4, 1024, 2, 0)",
|
||||
Granularity: 1,
|
||||
},
|
||||
},
|
||||
wantPaths: []promotetypes.PromotePath{
|
||||
{
|
||||
Signal: "logs",
|
||||
Context: "body",
|
||||
Path: "body.user.name",
|
||||
Promote: true,
|
||||
Indexes: []promotetypes.WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
store := telemetrytypestest.NewMockMetadataStore()
|
||||
store.PromotedPathsMap = testCase.promoted
|
||||
store.LogsJSONIndexes = testCase.indexes
|
||||
m := NewModule(store, nil)
|
||||
|
||||
paths, err := m.ListPromotedPaths(ctx, testCase.filters)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, paths, len(testCase.wantPaths))
|
||||
|
||||
byDomainPath := map[string]promotetypes.PromotePath{}
|
||||
for _, path := range paths {
|
||||
byDomainPath[path.Signal+"/"+path.Context+"/"+path.Path] = path
|
||||
}
|
||||
for _, want := range testCase.wantPaths {
|
||||
key := want.Signal + "/" + want.Context + "/" + want.Path
|
||||
require.Contains(t, byDomainPath, key)
|
||||
assert.Equal(t, want, byDomainPath[key])
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -8,11 +8,11 @@ import (
|
||||
)
|
||||
|
||||
type Module interface {
|
||||
ListPromotedAndIndexedPaths(ctx context.Context) ([]promotetypes.PromotePath, error)
|
||||
PromoteAndIndexPaths(ctx context.Context, paths ...*promotetypes.PromotePath) error
|
||||
ListPromotedPaths(ctx context.Context, filters promotetypes.ListPromotedPathsFilters) ([]promotetypes.PromotePath, error)
|
||||
PromotePaths(ctx context.Context, paths ...*promotetypes.PromotePath) error
|
||||
}
|
||||
|
||||
type Handler interface {
|
||||
HandlePromoteAndIndexPaths(w http.ResponseWriter, r *http.Request)
|
||||
ListPromotedAndIndexedPaths(w http.ResponseWriter, r *http.Request)
|
||||
PromotePaths(w http.ResponseWriter, r *http.Request)
|
||||
ListPromotedPaths(w http.ResponseWriter, r *http.Request)
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -16,6 +16,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/promotetypes"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/huandu/go-sqlbuilder"
|
||||
)
|
||||
@@ -33,6 +34,8 @@ var (
|
||||
CodeFailedToAppendPath = errors.MustNewCode("failed_to_append_path_promoted_paths")
|
||||
)
|
||||
|
||||
var logsBodyPromotedEntry = promotetypes.NewLogsBodyTarget().Entry
|
||||
|
||||
// enrichJSONKeys enriches body-context keys with promoted path info, indexes,
|
||||
// and JSON access plans. parentTypeCache contains parent array types (ArrayJSON/ArrayDynamic)
|
||||
// pre-fetched in the main UNION query.
|
||||
@@ -67,7 +70,7 @@ func (t *telemetryMetaStore) enrichJSONKeys(ctx context.Context, selectors []*te
|
||||
}
|
||||
|
||||
// fetch promoted paths
|
||||
promoted, err := t.GetPromotedPaths(ctx, paths...)
|
||||
promoted, err := t.GetPromotedPaths(ctx, logsBodyPromotedEntry, paths...)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -157,7 +160,7 @@ func buildListLogsJSONIndexesQuery(cluster string, filters ...string) (string, [
|
||||
}
|
||||
|
||||
func (t *telemetryMetaStore) ListLogsJSONIndexes(ctx context.Context, filters ...string) ([]telemetrytypes.TelemetryFieldKeySkipIndex, error) {
|
||||
ctx = withTelemetryContext(ctx, "ListLogsJSONIndexes")
|
||||
ctx = withTelemetryContext(ctx, telemetrytypes.SignalLogs, "ListLogsJSONIndexes")
|
||||
query, args := buildListLogsJSONIndexesQuery(t.telemetrystore.Cluster(), filters...)
|
||||
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
|
||||
if err != nil {
|
||||
@@ -215,14 +218,14 @@ func (t *telemetryMetaStore) ListLogsJSONIndexes(ctx context.Context, filters ..
|
||||
|
||||
// TODO(Piyush): Remove this if not used in future.
|
||||
func (t *telemetryMetaStore) ListJSONValues(ctx context.Context, path string, limit int) (*telemetrytypes.TelemetryFieldValues, bool, error) {
|
||||
ctx = withTelemetryContext(ctx, "ListJSONValues")
|
||||
ctx = withTelemetryContext(ctx, telemetrytypes.SignalLogs, "ListJSONValues")
|
||||
path = CleanPathPrefixes(path)
|
||||
|
||||
if strings.Contains(path, telemetrytypes.ArraySep) || strings.Contains(path, telemetrytypes.ArrayAnyIndex) {
|
||||
return nil, false, errors.NewInvalidInputf(errors.CodeInvalidInput, "array paths are not supported")
|
||||
}
|
||||
|
||||
promoted, err := t.IsPathPromoted(ctx, path)
|
||||
promoted, err := t.isPathPromoted(ctx, logsBodyPromotedEntry, path)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
@@ -376,13 +379,12 @@ func derefValue(v any) any {
|
||||
return val.Interface()
|
||||
}
|
||||
|
||||
// IsPathPromoted checks if a specific path is promoted (Column Evolution table: field_name for logs body).
|
||||
func (t *telemetryMetaStore) IsPathPromoted(ctx context.Context, path string) (bool, error) {
|
||||
ctx = withTelemetryContext(ctx, "IsPathPromoted")
|
||||
func (t *telemetryMetaStore) isPathPromoted(ctx context.Context, entry telemetrytypes.EvolutionEntry, path string) (bool, error) {
|
||||
ctx = withTelemetryContext(ctx, entry.Signal, "isPathPromoted")
|
||||
split := strings.Split(path, telemetrytypes.ArraySep)
|
||||
pathSegment := split[0]
|
||||
query := fmt.Sprintf("SELECT 1 FROM %s.%s WHERE signal = ? AND column_name = ? AND field_context = ? AND field_name = ? LIMIT 1", DBName, PromotedPathsTableName)
|
||||
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, telemetrytypes.SignalLogs, logstelemetryschema.LogsV2BodyPromotedColumn, telemetrytypes.FieldContextBody, pathSegment)
|
||||
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, entry.Signal, entry.ColumnName, entry.FieldContext, pathSegment)
|
||||
if err != nil {
|
||||
return false, errors.WrapInternalf(err, CodeFailCheckPathPromoted, "failed to check if path %s is promoted", path)
|
||||
}
|
||||
@@ -391,14 +393,13 @@ func (t *telemetryMetaStore) IsPathPromoted(ctx context.Context, path string) (b
|
||||
return rows.Next(), nil
|
||||
}
|
||||
|
||||
// GetPromotedPaths returns promoted paths from the Column Evolution table (field_name for logs body).
|
||||
func (t *telemetryMetaStore) GetPromotedPaths(ctx context.Context, paths ...string) (map[string]bool, error) {
|
||||
ctx = withTelemetryContext(ctx, "GetPromotedPaths")
|
||||
func (t *telemetryMetaStore) GetPromotedPaths(ctx context.Context, entry telemetrytypes.EvolutionEntry, paths ...string) (map[string]bool, error) {
|
||||
ctx = withTelemetryContext(ctx, entry.Signal, "GetPromotedPaths")
|
||||
sb := sqlbuilder.Select("field_name").From(fmt.Sprintf("%s.%s", DBName, PromotedPathsTableName))
|
||||
conditions := []string{
|
||||
sb.Equal("signal", telemetrytypes.SignalLogs),
|
||||
sb.Equal("column_name", logstelemetryschema.LogsV2BodyPromotedColumn),
|
||||
sb.Equal("field_context", telemetrytypes.FieldContextBody),
|
||||
sb.Equal("signal", entry.Signal),
|
||||
sb.Equal("column_name", entry.ColumnName),
|
||||
sb.Equal("field_context", entry.FieldContext),
|
||||
sb.NotEqual("field_name", "__all__"),
|
||||
}
|
||||
if len(paths) > 0 {
|
||||
@@ -438,9 +439,8 @@ func CleanPathPrefixes(path string) string {
|
||||
return path
|
||||
}
|
||||
|
||||
// PromotePaths inserts promoted paths into the Column Evolution table (same schema as signoz-otel-collector metadata_migrations).
|
||||
func (t *telemetryMetaStore) PromotePaths(ctx context.Context, paths ...string) error {
|
||||
ctx = withTelemetryContext(ctx, "PromotePaths")
|
||||
func (t *telemetryMetaStore) PromotePaths(ctx context.Context, entry telemetrytypes.EvolutionEntry, paths ...string) error {
|
||||
ctx = withTelemetryContext(ctx, entry.Signal, "PromotePaths")
|
||||
batch, err := t.telemetrystore.ClickhouseDB().PrepareBatch(ctx,
|
||||
fmt.Sprintf("INSERT INTO %s.%s (signal, column_name, column_type, field_context, field_name, version, release_time) VALUES", DBName,
|
||||
PromotedPathsTableName))
|
||||
@@ -454,7 +454,7 @@ func (t *telemetryMetaStore) PromotePaths(ctx context.Context, paths ...string)
|
||||
if trimmed == "" {
|
||||
continue
|
||||
}
|
||||
if err := batch.Append(telemetrytypes.SignalLogs, logstelemetryschema.LogsV2BodyPromotedColumn, "JSON()", telemetrytypes.FieldContextBody, trimmed, 0, releaseTime); err != nil {
|
||||
if err := batch.Append(entry.Signal, entry.ColumnName, entry.ColumnType, entry.FieldContext, trimmed, entry.Version, releaseTime); err != nil {
|
||||
_ = batch.Abort()
|
||||
return errors.WrapInternalf(err, CodeFailedToAppendPath, "failed to append path")
|
||||
}
|
||||
@@ -466,9 +466,9 @@ func (t *telemetryMetaStore) PromotePaths(ctx context.Context, paths ...string)
|
||||
return nil
|
||||
}
|
||||
|
||||
func withTelemetryContext(ctx context.Context, functionName string) context.Context {
|
||||
func withTelemetryContext(ctx context.Context, signal telemetrytypes.Signal, functionName string) context.Context {
|
||||
return ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
|
||||
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalLogs.StringValue(),
|
||||
instrumentationtypes.TelemetrySignal: signal.StringValue(),
|
||||
instrumentationtypes.CodeNamespace: "metadata",
|
||||
instrumentationtypes.CodeFunctionName: functionName,
|
||||
})
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
103
pkg/types/promotetypes/target.go
Normal file
103
pkg/types/promotetypes/target.go
Normal file
@@ -0,0 +1,103 @@
|
||||
package promotetypes
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/tracestelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
// Target identifies a promotion domain.
|
||||
type Target struct {
|
||||
Entry telemetrytypes.EvolutionEntry // evolution row template; FieldName and ReleaseTime are set per write
|
||||
DBName string // index DDL database, used only when IndexesSupported
|
||||
LocalTableName string // index DDL local table, used only when IndexesSupported
|
||||
BaseColumn string // column holding every path; indexes for unpromoted paths are created on it
|
||||
RequiredPathPrefix string // prefix API paths must carry, stripped before storing; empty for bare names
|
||||
IndexesSupported bool
|
||||
}
|
||||
|
||||
func (t Target) PromotedColumn() string { return t.Entry.ColumnName }
|
||||
|
||||
func (t Target) BaseColumnPrefix() string { return t.BaseColumn + "." }
|
||||
|
||||
func (t Target) PromotedColumnPrefix() string { return t.PromotedColumn() + "." }
|
||||
|
||||
func NewTarget(entry telemetrytypes.EvolutionEntry, dbName, localTableName, baseColumn, requiredPathPrefix string, indexesSupported bool) Target {
|
||||
return Target{
|
||||
Entry: entry,
|
||||
DBName: dbName,
|
||||
LocalTableName: localTableName,
|
||||
BaseColumn: baseColumn,
|
||||
RequiredPathPrefix: requiredPathPrefix,
|
||||
IndexesSupported: indexesSupported,
|
||||
}
|
||||
}
|
||||
|
||||
// NewLogsBodyTarget returns the logs body domain (body_v2 -> body_promoted).
|
||||
func NewLogsBodyTarget() Target {
|
||||
return NewTarget(
|
||||
telemetrytypes.EvolutionEntry{
|
||||
Signal: telemetrytypes.SignalLogs,
|
||||
ColumnName: logstelemetryschema.LogsV2BodyPromotedColumn,
|
||||
ColumnType: "JSON()",
|
||||
FieldContext: telemetrytypes.FieldContextBody,
|
||||
},
|
||||
logstelemetryschema.DBName,
|
||||
logstelemetryschema.LogsV2LocalTableName,
|
||||
logstelemetryschema.LogsV2BodyV2Column,
|
||||
telemetrytypes.BodyJSONStringSearchPrefix,
|
||||
true,
|
||||
)
|
||||
}
|
||||
|
||||
// NewTracesAttributesTarget returns the spans attributes domain (attributes
|
||||
// -> attributes_promoted).
|
||||
func NewTracesAttributesTarget() Target {
|
||||
return NewTarget(
|
||||
telemetrytypes.EvolutionEntry{
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
ColumnName: tracestelemetryschema.SpanAttributesPromotedColumn,
|
||||
ColumnType: "JSON()",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
},
|
||||
tracestelemetryschema.DBName,
|
||||
tracestelemetryschema.SpanIndexV3LocalTableName,
|
||||
tracestelemetryschema.SpanAttributesColumn,
|
||||
"",
|
||||
false,
|
||||
)
|
||||
}
|
||||
|
||||
func NewTargetFromText(signal, context string) (Target, error) {
|
||||
parsedSignal, ok := telemetrytypes.SignalFromText(signal)
|
||||
if !ok {
|
||||
return Target{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid signal: %s", signal)
|
||||
}
|
||||
parsedContext, ok := telemetrytypes.FieldContextFromText(context)
|
||||
if !ok {
|
||||
return Target{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid context: %s", context)
|
||||
}
|
||||
target, ok := TargetFor(parsedSignal, parsedContext)
|
||||
if !ok {
|
||||
return Target{}, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "promotion is not supported for %s %s", parsedSignal.StringValue(), parsedContext.StringValue())
|
||||
}
|
||||
return target, nil
|
||||
}
|
||||
|
||||
func Targets() []Target {
|
||||
return []Target{
|
||||
NewLogsBodyTarget(),
|
||||
NewTracesAttributesTarget(),
|
||||
}
|
||||
}
|
||||
|
||||
func TargetFor(signal telemetrytypes.Signal, context telemetrytypes.FieldContext) (Target, bool) {
|
||||
for _, target := range Targets() {
|
||||
if target.Entry.Signal.StringValue() == signal.StringValue() &&
|
||||
target.Entry.FieldContext.StringValue() == context.StringValue() {
|
||||
return target, true
|
||||
}
|
||||
}
|
||||
return Target{}, false
|
||||
}
|
||||
@@ -3,7 +3,6 @@ package promotetypes
|
||||
import (
|
||||
"strings"
|
||||
|
||||
"github.com/SigNoz/signoz-otel-collector/constants"
|
||||
"github.com/SigNoz/signoz-otel-collector/pkg/keycheck"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
@@ -17,13 +16,62 @@ type WrappedIndex struct {
|
||||
}
|
||||
|
||||
type PromotePath struct {
|
||||
Path string `json:"path"`
|
||||
Signal string `json:"signal" required:"true"`
|
||||
Context string `json:"context" required:"true"`
|
||||
Path string `json:"path" required:"true"`
|
||||
Promote bool `json:"promote,omitempty"`
|
||||
|
||||
Indexes []WrappedIndex `json:"indexes,omitempty"`
|
||||
}
|
||||
|
||||
func (i *PromotePath) ValidateAndSetDefaults() error {
|
||||
func (i *PromotePath) Target() (Target, error) {
|
||||
return NewTargetFromText(i.Signal, i.Context)
|
||||
}
|
||||
|
||||
type ListPromotedPathsFilters struct {
|
||||
Signal string `query:"signal" json:"signal"`
|
||||
Context string `query:"context" json:"context"`
|
||||
Promoted *bool `query:"promoted" json:"promoted"`
|
||||
Indexes *bool `query:"indexes" json:"indexes"`
|
||||
}
|
||||
|
||||
// Validate checks the signal and context words are known; the pair need not
|
||||
// name a supported domain.
|
||||
func (f *ListPromotedPathsFilters) Validate() error {
|
||||
if f.Signal != "" {
|
||||
if _, ok := telemetrytypes.SignalFromText(f.Signal); !ok {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid signal: %s", f.Signal)
|
||||
}
|
||||
}
|
||||
if f.Context != "" {
|
||||
if _, ok := telemetrytypes.FieldContextFromText(f.Context); !ok {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid context: %s", f.Context)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *ListPromotedPathsFilters) MatchesTarget(target Target) bool {
|
||||
if f.Signal != "" && f.Signal != target.Entry.Signal.StringValue() {
|
||||
return false
|
||||
}
|
||||
if f.Context != "" && f.Context != target.Entry.FieldContext.StringValue() {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (f *ListPromotedPathsFilters) MatchesPath(path PromotePath) bool {
|
||||
if f.Promoted != nil && *f.Promoted != path.Promote {
|
||||
return false
|
||||
}
|
||||
if f.Indexes != nil && *f.Indexes != (len(path.Indexes) > 0) {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (i *PromotePath) ValidateAndSetDefaults(target Target) error {
|
||||
if i.Path == "" {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "path is required")
|
||||
}
|
||||
@@ -36,22 +84,26 @@ func (i *PromotePath) ValidateAndSetDefaults() error {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "array paths can not be promoted or indexed")
|
||||
}
|
||||
|
||||
if strings.HasPrefix(i.Path, constants.BodyV2ColumnPrefix) || strings.HasPrefix(i.Path, constants.BodyPromotedColumnPrefix) {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "`%s`, `%s` don't add these prefixes to the path", constants.BodyV2ColumnPrefix, constants.BodyPromotedColumnPrefix)
|
||||
if strings.HasPrefix(i.Path, target.BaseColumnPrefix()) || strings.HasPrefix(i.Path, target.PromotedColumnPrefix()) {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "`%s`, `%s` don't add these prefixes to the path", target.BaseColumnPrefix(), target.PromotedColumnPrefix())
|
||||
}
|
||||
|
||||
if !strings.HasPrefix(i.Path, telemetrytypes.BodyJSONStringSearchPrefix) {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "path must start with `body.`")
|
||||
if target.RequiredPathPrefix != "" {
|
||||
if !strings.HasPrefix(i.Path, target.RequiredPathPrefix) {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "path must start with `%s`", target.RequiredPathPrefix)
|
||||
}
|
||||
i.Path = strings.TrimPrefix(i.Path, target.RequiredPathPrefix)
|
||||
}
|
||||
|
||||
// remove the "body." prefix from the path
|
||||
i.Path = strings.TrimPrefix(i.Path, telemetrytypes.BodyJSONStringSearchPrefix)
|
||||
|
||||
isCardinal := keycheck.IsCardinal(i.Path)
|
||||
if isCardinal {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "cardinal paths can not be promoted or indexed")
|
||||
}
|
||||
|
||||
if len(i.Indexes) > 0 && !target.IndexesSupported {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "indexes are not supported for %s %s", target.Entry.Signal.StringValue(), target.Entry.FieldContext.StringValue())
|
||||
}
|
||||
|
||||
for idx, index := range i.Indexes {
|
||||
if index.Type == "" {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "index type is required")
|
||||
|
||||
333
pkg/types/promotetypes/types_test.go
Normal file
333
pkg/types/promotetypes/types_test.go
Normal file
@@ -0,0 +1,333 @@
|
||||
package promotetypes
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestPromotePathTarget(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
path *PromotePath
|
||||
want Target
|
||||
wantErr bool
|
||||
}{
|
||||
{
|
||||
name: "LogsBody_Resolved",
|
||||
path: &PromotePath{Signal: "logs", Context: "body", Path: "body.user.name"},
|
||||
want: NewLogsBodyTarget(),
|
||||
},
|
||||
{
|
||||
name: "TracesAttribute_Resolved",
|
||||
path: &PromotePath{Signal: "traces", Context: "attribute", Path: "http.method"},
|
||||
want: NewTracesAttributesTarget(),
|
||||
},
|
||||
{
|
||||
name: "InvalidSignal_Rejected",
|
||||
path: &PromotePath{Signal: "events", Context: "attribute", Path: "http.method"},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "InvalidContext_Rejected",
|
||||
path: &PromotePath{Signal: "logs", Context: "span", Path: "user.name"},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "UnsupportedDomain_Rejected",
|
||||
path: &PromotePath{Signal: "metrics", Context: "attribute", Path: "http.method"},
|
||||
wantErr: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
target, err := testCase.path.Target()
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.want, target)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateAndSetDefaultsLogsBody(t *testing.T) {
|
||||
target := NewLogsBodyTarget()
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
path *PromotePath
|
||||
wantErr bool
|
||||
wantPath string
|
||||
wantJSONDataType telemetrytypes.JSONDataType
|
||||
}{
|
||||
{
|
||||
name: "ValidPath_BodyPrefixStripped",
|
||||
path: &PromotePath{Path: "body.user.name", Promote: true},
|
||||
wantPath: "user.name",
|
||||
},
|
||||
{
|
||||
name: "PathWithoutBodyPrefix_Rejected",
|
||||
path: &PromotePath{Path: "user.name", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "BodyV2PrefixedPath_Rejected",
|
||||
path: &PromotePath{Path: "body_v2.user.name", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "BodyPromotedPrefixedPath_Rejected",
|
||||
path: &PromotePath{Path: "body_promoted.user.name", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "EmptyPath_Rejected",
|
||||
path: &PromotePath{Path: "", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "SpacedPath_Rejected",
|
||||
path: &PromotePath{Path: "body.my path", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "ArrayIndexPath_Rejected",
|
||||
path: &PromotePath{Path: "body.users[].id", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "ArrayWildcardPath_Rejected",
|
||||
path: &PromotePath{Path: "body.users[*].id", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "CardinalPath_Rejected",
|
||||
path: &PromotePath{Path: "body.request.550e8400-e29b-41d4-a716-446655440000", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "ValidIndex_JSONDataTypeDefaulted",
|
||||
path: &PromotePath{
|
||||
Path: "body.user.name",
|
||||
Indexes: []WrappedIndex{
|
||||
{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1},
|
||||
},
|
||||
},
|
||||
wantPath: "user.name",
|
||||
wantJSONDataType: telemetrytypes.String,
|
||||
},
|
||||
{
|
||||
name: "UnsupportedColumnTypeIndex_Rejected",
|
||||
path: &PromotePath{
|
||||
Path: "body.user.active",
|
||||
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeBool, Type: "minmax", Granularity: 1}},
|
||||
},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "IndexWithoutType_Rejected",
|
||||
path: &PromotePath{
|
||||
Path: "body.user.name",
|
||||
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeString, Granularity: 1}},
|
||||
},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "IndexWithoutGranularity_Rejected",
|
||||
path: &PromotePath{
|
||||
Path: "body.user.name",
|
||||
Indexes: []WrappedIndex{{FieldDataType: telemetrytypes.FieldDataTypeString, Type: "minmax"}},
|
||||
},
|
||||
wantErr: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
err := testCase.path.ValidateAndSetDefaults(target)
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.wantPath, testCase.path.Path)
|
||||
if testCase.wantJSONDataType != (telemetrytypes.JSONDataType{}) {
|
||||
require.Len(t, testCase.path.Indexes, 1)
|
||||
assert.Equal(t, testCase.wantJSONDataType, testCase.path.Indexes[0].JSONDataType)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateAndSetDefaultsTracesAttributes(t *testing.T) {
|
||||
target := NewTracesAttributesTarget()
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
path *PromotePath
|
||||
wantErr bool
|
||||
wantPath string
|
||||
}{
|
||||
{
|
||||
name: "BareAttributeName_KeptAsIs",
|
||||
path: &PromotePath{Path: "http.method", Promote: true},
|
||||
wantPath: "http.method",
|
||||
},
|
||||
{
|
||||
name: "AttributesPrefixedPath_Rejected",
|
||||
path: &PromotePath{Path: "attributes.http.method", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "AttributesPromotedPrefixedPath_Rejected",
|
||||
path: &PromotePath{Path: "attributes_promoted.http.method", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "EmptyPath_Rejected",
|
||||
path: &PromotePath{Path: "", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "SpacedPath_Rejected",
|
||||
path: &PromotePath{Path: "my attr", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "ArrayIndexPath_Rejected",
|
||||
path: &PromotePath{Path: "tags[].id", Promote: true},
|
||||
wantErr: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
err := testCase.path.ValidateAndSetDefaults(target)
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, testCase.wantPath, testCase.path.Path)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestListPromotedPathsFiltersValidate(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
filters ListPromotedPathsFilters
|
||||
wantErr bool
|
||||
}{
|
||||
{
|
||||
name: "EmptyFilters_Valid",
|
||||
filters: ListPromotedPathsFilters{},
|
||||
},
|
||||
{
|
||||
name: "KnownSignalAndContext_Valid",
|
||||
filters: ListPromotedPathsFilters{Signal: "traces", Context: "attribute"},
|
||||
},
|
||||
{
|
||||
name: "InvalidSignal_Rejected",
|
||||
filters: ListPromotedPathsFilters{Signal: "events"},
|
||||
wantErr: true,
|
||||
},
|
||||
{
|
||||
name: "InvalidContext_Rejected",
|
||||
filters: ListPromotedPathsFilters{Context: "json"},
|
||||
wantErr: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
err := testCase.filters.Validate()
|
||||
if testCase.wantErr {
|
||||
assert.Error(t, err)
|
||||
return
|
||||
}
|
||||
require.NoError(t, err)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestListPromotedPathsFiltersMatch(t *testing.T) {
|
||||
trueValue := true
|
||||
falseValue := false
|
||||
|
||||
testCases := []struct {
|
||||
name string
|
||||
filters ListPromotedPathsFilters
|
||||
target Target
|
||||
path PromotePath
|
||||
wantTarget bool
|
||||
wantPath bool
|
||||
}{
|
||||
{
|
||||
name: "EmptyFilters_MatchEverything",
|
||||
filters: ListPromotedPathsFilters{},
|
||||
target: NewTracesAttributesTarget(),
|
||||
path: PromotePath{Path: "http.method", Promote: true},
|
||||
wantTarget: true,
|
||||
wantPath: true,
|
||||
},
|
||||
{
|
||||
name: "SignalFilter_MatchesSameSignal",
|
||||
filters: ListPromotedPathsFilters{Signal: "traces"},
|
||||
target: NewTracesAttributesTarget(),
|
||||
wantTarget: true,
|
||||
},
|
||||
{
|
||||
name: "SignalFilter_SkipsOtherSignals",
|
||||
filters: ListPromotedPathsFilters{Signal: "traces"},
|
||||
target: NewLogsBodyTarget(),
|
||||
wantTarget: false,
|
||||
},
|
||||
{
|
||||
name: "ContextFilter_SkipsOtherContexts",
|
||||
filters: ListPromotedPathsFilters{Context: "body"},
|
||||
target: NewTracesAttributesTarget(),
|
||||
wantTarget: false,
|
||||
},
|
||||
{
|
||||
name: "PromotedFalseFilter_MatchesUnpromotedPath",
|
||||
filters: ListPromotedPathsFilters{Promoted: &falseValue},
|
||||
path: PromotePath{Path: "request.duration"},
|
||||
wantPath: true,
|
||||
},
|
||||
{
|
||||
name: "PromotedFalseFilter_SkipsPromotedPath",
|
||||
filters: ListPromotedPathsFilters{Promoted: &falseValue},
|
||||
path: PromotePath{Path: "http.method", Promote: true},
|
||||
wantPath: false,
|
||||
},
|
||||
{
|
||||
name: "IndexesTrueFilter_MatchesIndexedPath",
|
||||
filters: ListPromotedPathsFilters{Indexes: &trueValue},
|
||||
path: PromotePath{Path: "user.name", Indexes: []WrappedIndex{{Type: "minmax"}}},
|
||||
wantPath: true,
|
||||
},
|
||||
{
|
||||
name: "IndexesTrueFilter_SkipsUnindexedPath",
|
||||
filters: ListPromotedPathsFilters{Indexes: &trueValue},
|
||||
path: PromotePath{Path: "http.method", Promote: true},
|
||||
wantPath: false,
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
if testCase.target.Entry.Signal.StringValue() != "" {
|
||||
assert.Equal(t, testCase.wantTarget, testCase.filters.MatchesTarget(testCase.target))
|
||||
}
|
||||
if testCase.path.Path != "" {
|
||||
assert.Equal(t, testCase.wantPath, testCase.filters.MatchesPath(testCase.path))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -22,3 +22,18 @@ func (Signal) Enum() []any {
|
||||
SignalUnspecified,
|
||||
}
|
||||
}
|
||||
|
||||
// SignalFromText resolves a signal word to its Signal; ok is false for an
|
||||
// unknown word.
|
||||
func SignalFromText(text string) (Signal, bool) {
|
||||
s := Signal{valuer.NewString(text)}
|
||||
switch s {
|
||||
case SignalTraces:
|
||||
return SignalTraces, true
|
||||
case SignalLogs:
|
||||
return SignalLogs, true
|
||||
case SignalMetrics:
|
||||
return SignalMetrics, true
|
||||
}
|
||||
return Signal{}, false
|
||||
}
|
||||
|
||||
@@ -37,11 +37,13 @@ type MetadataStore interface {
|
||||
// ListLogsJSONIndexes lists the JSON indexes for the logs table.
|
||||
ListLogsJSONIndexes(ctx context.Context, filters ...string) ([]TelemetryFieldKeySkipIndex, error)
|
||||
|
||||
// ListPromotedPaths lists the promoted paths.
|
||||
GetPromotedPaths(ctx context.Context, paths ...string) (map[string]bool, error)
|
||||
// GetPromotedPaths lists the promoted paths recorded in the column
|
||||
// evolution table for the entry's signal, column and field context.
|
||||
GetPromotedPaths(ctx context.Context, entry EvolutionEntry, paths ...string) (map[string]bool, error)
|
||||
|
||||
// PromotePaths promotes the paths.
|
||||
PromotePaths(ctx context.Context, paths ...string) error
|
||||
// PromotePaths records promoted paths in the column evolution table as
|
||||
// rows templated by entry; FieldName and ReleaseTime are set per path.
|
||||
PromotePaths(ctx context.Context, entry EvolutionEntry, paths ...string) error
|
||||
|
||||
// GetFirstSeenFromMetricMetadata gets the first seen timestamp for a metric metadata lookup key.
|
||||
GetFirstSeenFromMetricMetadata(ctx context.Context, lookupKeys []MetricMetadataLookupKey) (map[MetricMetadataLookupKey]int64, error)
|
||||
|
||||
@@ -361,7 +361,7 @@ func (m *MockMetadataStore) SetTemporality(metricName string, temporality metric
|
||||
}
|
||||
|
||||
// PromotePaths promotes the paths.
|
||||
func (m *MockMetadataStore) PromotePaths(ctx context.Context, paths ...string) error {
|
||||
func (m *MockMetadataStore) PromotePaths(_ context.Context, _ telemetrytypes.EvolutionEntry, paths ...string) error {
|
||||
for _, path := range paths {
|
||||
m.PromotedPathsMap[path] = true
|
||||
}
|
||||
@@ -369,7 +369,7 @@ func (m *MockMetadataStore) PromotePaths(ctx context.Context, paths ...string) e
|
||||
}
|
||||
|
||||
// GetPromotedPaths returns the promoted paths.
|
||||
func (m *MockMetadataStore) GetPromotedPaths(ctx context.Context, paths ...string) (map[string]bool, error) {
|
||||
func (m *MockMetadataStore) GetPromotedPaths(_ context.Context, _ telemetrytypes.EvolutionEntry, _ ...string) (map[string]bool, error) {
|
||||
return m.PromotedPathsMap, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -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",
|
||||
},
|
||||
)
|
||||
Reference in New Issue
Block a user