diff --git a/.github/workflows/integrationci.yaml b/.github/workflows/integrationci.yaml index 1f185eaba5d..ca80a12ad31 100644 --- a/.github/workflows/integrationci.yaml +++ b/.github/workflows/integrationci.yaml @@ -58,6 +58,7 @@ jobs: - querierai - rawexportdata - promqlconformance + - promapiconformance - querierauthz - role - rootuser diff --git a/docs/api/openapi.yml b/docs/api/openapi.yml index 79dd3f74a3e..ebd5aaf550f 100644 --- a/docs/api/openapi.yml +++ b/docs/api/openapi.yml @@ -6460,6 +6460,148 @@ components: type: object PreferencetypesValue: type: object + PrometheusErrorResponseSchema: + properties: + error: + type: string + errorType: + enum: + - bad_data + - execution + - canceled + - timeout + - internal + type: string + status: + enum: + - error + type: string + required: + - status + - errorType + - error + type: object + PrometheusMatrixDataSchema: + properties: + result: + items: + $ref: '#/components/schemas/PrometheusMatrixSeriesSchema' + nullable: true + type: array + resultType: + enum: + - matrix + type: string + required: + - resultType + - result + type: object + PrometheusMatrixSeriesSchema: + properties: + metric: + additionalProperties: + type: string + nullable: true + type: object + values: + items: + $ref: '#/components/schemas/PrometheusSamplePairSchema' + nullable: true + type: array + required: + - metric + - values + type: object + PrometheusQueryDataSchema: + oneOf: + - $ref: '#/components/schemas/PrometheusMatrixDataSchema' + - $ref: '#/components/schemas/PrometheusVectorDataSchema' + - $ref: '#/components/schemas/PrometheusScalarDataSchema' + - $ref: '#/components/schemas/PrometheusStringDataSchema' + type: object + PrometheusSamplePairSchema: + description: 'A [timestamp, value] pair: float unix seconds, then the string-encoded + sample value ("NaN", "+Inf", "-Inf" included).' + items: + oneOf: + - type: number + - type: string + maxItems: 2 + minItems: 2 + nullable: true + type: array + PrometheusScalarDataSchema: + properties: + result: + $ref: '#/components/schemas/PrometheusSamplePairSchema' + resultType: + enum: + - scalar + type: string + required: + - resultType + - result + type: object + PrometheusStringDataSchema: + properties: + result: + $ref: '#/components/schemas/PrometheusSamplePairSchema' + resultType: + enum: + - string + type: string + required: + - resultType + - result + type: object + PrometheusSuccessResponseSchema: + properties: + data: + $ref: '#/components/schemas/PrometheusQueryDataSchema' + infos: + items: + type: string + type: array + status: + enum: + - success + type: string + warnings: + items: + type: string + type: array + required: + - status + - data + type: object + PrometheusVectorDataSchema: + properties: + result: + items: + $ref: '#/components/schemas/PrometheusVectorSampleSchema' + nullable: true + type: array + resultType: + enum: + - vector + type: string + required: + - resultType + - result + type: object + PrometheusVectorSampleSchema: + properties: + metric: + additionalProperties: + type: string + nullable: true + type: object + value: + $ref: '#/components/schemas/PrometheusSamplePairSchema' + required: + - metric + - value + type: object PromotetypesPromotePath: properties: indexes: @@ -24811,6 +24953,374 @@ paths: summary: Replace variables tags: - querier + /prometheus/api/v1/query: + get: + description: 'Prometheus-compatible endpoint: the request and response contract + is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). + Parameters are accepted as URL query parameters or a form-encoded body, on + GET and POST alike.' + operationId: PrometheusQuery + parameters: + - description: PromQL expression. + in: query + name: query + required: true + schema: + description: PromQL expression. + type: string + - description: 'Evaluation timestamp: RFC3339 or float unix seconds. Defaults + to the server''s current time.' + in: query + name: time + schema: + description: 'Evaluation timestamp: RFC3339 or float unix seconds. Defaults + to the server''s current time.' + type: string + - description: 'Evaluation timeout: duration string or float seconds.' + in: query + name: timeout + schema: + description: 'Evaluation timeout: duration string or float seconds.' + type: string + - description: Any non-empty value includes query statistics in the response. + in: query + name: stats + schema: + description: Any non-empty value includes query statistics in the response. + type: string + responses: + "200": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusSuccessResponseSchema' + description: OK + "400": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + 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 + "422": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Unprocessable Entity + "500": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Internal Server Error + "503": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Service Unavailable + security: + - api_key: + - metrics:read + - tokenizer: + - metrics:read + summary: Prometheus instant query + tags: + - prometheus + post: + description: 'Prometheus-compatible endpoint: the request and response contract + is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). + Parameters are accepted as URL query parameters or a form-encoded body, on + GET and POST alike.' + operationId: PrometheusQueryPost + parameters: + - description: PromQL expression. + in: query + name: query + required: true + schema: + description: PromQL expression. + type: string + - description: 'Evaluation timestamp: RFC3339 or float unix seconds. Defaults + to the server''s current time.' + in: query + name: time + schema: + description: 'Evaluation timestamp: RFC3339 or float unix seconds. Defaults + to the server''s current time.' + type: string + - description: 'Evaluation timeout: duration string or float seconds.' + in: query + name: timeout + schema: + description: 'Evaluation timeout: duration string or float seconds.' + type: string + - description: Any non-empty value includes query statistics in the response. + in: query + name: stats + schema: + description: Any non-empty value includes query statistics in the response. + type: string + responses: + "200": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusSuccessResponseSchema' + description: OK + "400": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + 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 + "422": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Unprocessable Entity + "500": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Internal Server Error + "503": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Service Unavailable + security: + - api_key: + - metrics:read + - tokenizer: + - metrics:read + summary: Prometheus instant query + tags: + - prometheus + /prometheus/api/v1/query_range: + get: + description: 'Prometheus-compatible endpoint: the request and response contract + is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). + Parameters are accepted as URL query parameters or a form-encoded body, on + GET and POST alike.' + operationId: PrometheusQueryRange + parameters: + - description: PromQL expression. + in: query + name: query + required: true + schema: + description: PromQL expression. + type: string + - description: 'Range start: RFC3339 or float unix seconds.' + in: query + name: start + required: true + schema: + description: 'Range start: RFC3339 or float unix seconds.' + type: string + - description: 'Range end: RFC3339 or float unix seconds.' + in: query + name: end + required: true + schema: + description: 'Range end: RFC3339 or float unix seconds.' + type: string + - description: 'Resolution step: duration string or float seconds.' + in: query + name: step + required: true + schema: + description: 'Resolution step: duration string or float seconds.' + type: string + - description: 'Evaluation timeout: duration string or float seconds.' + in: query + name: timeout + schema: + description: 'Evaluation timeout: duration string or float seconds.' + type: string + - description: Any non-empty value includes query statistics in the response. + in: query + name: stats + schema: + description: Any non-empty value includes query statistics in the response. + type: string + responses: + "200": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusSuccessResponseSchema' + description: OK + "400": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + 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 + "422": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Unprocessable Entity + "500": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Internal Server Error + "503": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Service Unavailable + security: + - api_key: + - metrics:read + - tokenizer: + - metrics:read + summary: Prometheus range query + tags: + - prometheus + post: + description: 'Prometheus-compatible endpoint: the request and response contract + is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). + Parameters are accepted as URL query parameters or a form-encoded body, on + GET and POST alike.' + operationId: PrometheusQueryRangePost + parameters: + - description: PromQL expression. + in: query + name: query + required: true + schema: + description: PromQL expression. + type: string + - description: 'Range start: RFC3339 or float unix seconds.' + in: query + name: start + required: true + schema: + description: 'Range start: RFC3339 or float unix seconds.' + type: string + - description: 'Range end: RFC3339 or float unix seconds.' + in: query + name: end + required: true + schema: + description: 'Range end: RFC3339 or float unix seconds.' + type: string + - description: 'Resolution step: duration string or float seconds.' + in: query + name: step + required: true + schema: + description: 'Resolution step: duration string or float seconds.' + type: string + - description: 'Evaluation timeout: duration string or float seconds.' + in: query + name: timeout + schema: + description: 'Evaluation timeout: duration string or float seconds.' + type: string + - description: Any non-empty value includes query statistics in the response. + in: query + name: stats + schema: + description: Any non-empty value includes query statistics in the response. + type: string + responses: + "200": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusSuccessResponseSchema' + description: OK + "400": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + 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 + "422": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Unprocessable Entity + "500": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Internal Server Error + "503": + content: + application/json: + schema: + $ref: '#/components/schemas/PrometheusErrorResponseSchema' + description: Service Unavailable + security: + - api_key: + - metrics:read + - tokenizer: + - metrics:read + summary: Prometheus range query + tags: + - prometheus servers: - description: The fully qualified URL to the SigNoz APIServer. url: https://{host}:{port}{base_path} diff --git a/docs/contributing/prometheus.md b/docs/contributing/prometheus.md index 7f4573ac1d0..f9a9d72ea38 100644 --- a/docs/contributing/prometheus.md +++ b/docs/contributing/prometheus.md @@ -299,8 +299,11 @@ substituted. One subtlety makes it exact: we write stale markers at absent grid points. Without them, the engine's lookback would resurrect a point from up to `lookback` earlier. The marker encodes "absent here" the way the engine itself encodes it. Units evaluate concurrently. Each unit is one -series lookup plus one grid statement. A step of 0 is an instant query: a -single evaluation at `end`. +grid statement: the group-key join resolves the matchers, and the samples +primary key takes the metric name straight from the selector. Only a +selector without a static `__name__` runs the series lookup first, to learn +the concrete metric names. A step of 0 is an instant query: a single +evaluation at `end`. A note on the window sliver: when the window is narrower than the step, the grid windows cover only `window/step` of the timeline. A sample in a gap @@ -315,8 +318,9 @@ selectors and `last_over_time` transpile at window < step too. ## Series lookup -Both paths resolve matchers the same way, once per selector -(`selectSeries`). The series tables hold one row per (fingerprint, bucket) +The engine path resolves matchers once per selector (`selectSeries`); the +transpiled path builds the same conditions into its group-key join. Both +read the same tables. The series tables hold one row per (fingerprint, bucket) at 1h/6h/1d/1w granularities. The shared schema package (`pkg/telemetryschema/metricstelemetryschema`) picks the table whose bucket fits the window. It rounds the window start down to the bucket boundary, so diff --git a/frontend/src/api/generated/services/prometheus/index.ts b/frontend/src/api/generated/services/prometheus/index.ts new file mode 100644 index 00000000000..938dfc2026e --- /dev/null +++ b/frontend/src/api/generated/services/prometheus/index.ts @@ -0,0 +1,396 @@ +/** + * ! 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 { + PrometheusErrorResponseSchemaDTO, + PrometheusQueryParams, + PrometheusQueryPostParams, + PrometheusQueryRangeParams, + PrometheusQueryRangePostParams, + PrometheusSuccessResponseSchemaDTO, + RenderErrorResponseDTO, +} from '../sigNoz.schemas'; + +import { GeneratedAPIInstance } from '../../../generatedAPIInstance'; +import type { ErrorType } from '../../../generatedAPIInstance'; + +/** + * Prometheus-compatible endpoint: the request and response contract is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). Parameters are accepted as URL query parameters or a form-encoded body, on GET and POST alike. + * @summary Prometheus instant query + */ +export const prometheusQuery = ( + params: PrometheusQueryParams, + signal?: AbortSignal, +) => { + return GeneratedAPIInstance({ + url: `/prometheus/api/v1/query`, + method: 'GET', + params, + signal, + }); +}; + +export const getPrometheusQueryQueryKey = (params?: PrometheusQueryParams) => { + return [`/prometheus/api/v1/query`, ...(params ? [params] : [])] as const; +}; + +export const getPrometheusQueryQueryOptions = < + TData = Awaited>, + TError = ErrorType, +>( + params: PrometheusQueryParams, + options?: { + query?: UseQueryOptions< + Awaited>, + TError, + TData + >; + }, +) => { + const { query: queryOptions } = options ?? {}; + + const queryKey = queryOptions?.queryKey ?? getPrometheusQueryQueryKey(params); + + const queryFn: QueryFunction>> = ({ + signal, + }) => prometheusQuery(params, signal); + + return { queryKey, queryFn, ...queryOptions } as UseQueryOptions< + Awaited>, + TError, + TData + > & { queryKey: QueryKey }; +}; + +export type PrometheusQueryQueryResult = NonNullable< + Awaited> +>; +export type PrometheusQueryQueryError = ErrorType< + PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO +>; + +/** + * @summary Prometheus instant query + */ + +export function usePrometheusQuery< + TData = Awaited>, + TError = ErrorType, +>( + params: PrometheusQueryParams, + options?: { + query?: UseQueryOptions< + Awaited>, + TError, + TData + >; + }, +): UseQueryResult & { queryKey: QueryKey } { + const queryOptions = getPrometheusQueryQueryOptions(params, options); + + const query = useQuery(queryOptions) as UseQueryResult & { + queryKey: QueryKey; + }; + + return { ...query, queryKey: queryOptions.queryKey }; +} + +/** + * @summary Prometheus instant query + */ +export const invalidatePrometheusQuery = async ( + queryClient: QueryClient, + params: PrometheusQueryParams, + options?: InvalidateOptions, +): Promise => { + await queryClient.invalidateQueries( + { queryKey: getPrometheusQueryQueryKey(params) }, + options, + ); + + return queryClient; +}; + +/** + * Prometheus-compatible endpoint: the request and response contract is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). Parameters are accepted as URL query parameters or a form-encoded body, on GET and POST alike. + * @summary Prometheus instant query + */ +export const prometheusQueryPost = ( + params: PrometheusQueryPostParams, + signal?: AbortSignal, +) => { + return GeneratedAPIInstance({ + url: `/prometheus/api/v1/query`, + method: 'POST', + params, + signal, + }); +}; + +export const getPrometheusQueryPostMutationOptions = < + TError = ErrorType, + TContext = unknown, +>(options?: { + mutation?: UseMutationOptions< + Awaited>, + TError, + { params: PrometheusQueryPostParams }, + TContext + >; +}): UseMutationOptions< + Awaited>, + TError, + { params: PrometheusQueryPostParams }, + TContext +> => { + const mutationKey = ['prometheusQueryPost']; + 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>, + { params: PrometheusQueryPostParams } + > = (props) => { + const { params } = props ?? {}; + + return prometheusQueryPost(params); + }; + + return { mutationFn, ...mutationOptions }; +}; + +export type PrometheusQueryPostMutationResult = NonNullable< + Awaited> +>; + +export type PrometheusQueryPostMutationError = ErrorType< + PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO +>; + +/** + * @summary Prometheus instant query + */ +export const usePrometheusQueryPost = < + TError = ErrorType, + TContext = unknown, +>(options?: { + mutation?: UseMutationOptions< + Awaited>, + TError, + { params: PrometheusQueryPostParams }, + TContext + >; +}): UseMutationResult< + Awaited>, + TError, + { params: PrometheusQueryPostParams }, + TContext +> => { + return useMutation(getPrometheusQueryPostMutationOptions(options)); +}; +/** + * Prometheus-compatible endpoint: the request and response contract is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). Parameters are accepted as URL query parameters or a form-encoded body, on GET and POST alike. + * @summary Prometheus range query + */ +export const prometheusQueryRange = ( + params: PrometheusQueryRangeParams, + signal?: AbortSignal, +) => { + return GeneratedAPIInstance({ + url: `/prometheus/api/v1/query_range`, + method: 'GET', + params, + signal, + }); +}; + +export const getPrometheusQueryRangeQueryKey = ( + params?: PrometheusQueryRangeParams, +) => { + return [ + `/prometheus/api/v1/query_range`, + ...(params ? [params] : []), + ] as const; +}; + +export const getPrometheusQueryRangeQueryOptions = < + TData = Awaited>, + TError = ErrorType, +>( + params: PrometheusQueryRangeParams, + options?: { + query?: UseQueryOptions< + Awaited>, + TError, + TData + >; + }, +) => { + const { query: queryOptions } = options ?? {}; + + const queryKey = + queryOptions?.queryKey ?? getPrometheusQueryRangeQueryKey(params); + + const queryFn: QueryFunction< + Awaited> + > = ({ signal }) => prometheusQueryRange(params, signal); + + return { queryKey, queryFn, ...queryOptions } as UseQueryOptions< + Awaited>, + TError, + TData + > & { queryKey: QueryKey }; +}; + +export type PrometheusQueryRangeQueryResult = NonNullable< + Awaited> +>; +export type PrometheusQueryRangeQueryError = ErrorType< + PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO +>; + +/** + * @summary Prometheus range query + */ + +export function usePrometheusQueryRange< + TData = Awaited>, + TError = ErrorType, +>( + params: PrometheusQueryRangeParams, + options?: { + query?: UseQueryOptions< + Awaited>, + TError, + TData + >; + }, +): UseQueryResult & { queryKey: QueryKey } { + const queryOptions = getPrometheusQueryRangeQueryOptions(params, options); + + const query = useQuery(queryOptions) as UseQueryResult & { + queryKey: QueryKey; + }; + + return { ...query, queryKey: queryOptions.queryKey }; +} + +/** + * @summary Prometheus range query + */ +export const invalidatePrometheusQueryRange = async ( + queryClient: QueryClient, + params: PrometheusQueryRangeParams, + options?: InvalidateOptions, +): Promise => { + await queryClient.invalidateQueries( + { queryKey: getPrometheusQueryRangeQueryKey(params) }, + options, + ); + + return queryClient; +}; + +/** + * Prometheus-compatible endpoint: the request and response contract is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). Parameters are accepted as URL query parameters or a form-encoded body, on GET and POST alike. + * @summary Prometheus range query + */ +export const prometheusQueryRangePost = ( + params: PrometheusQueryRangePostParams, + signal?: AbortSignal, +) => { + return GeneratedAPIInstance({ + url: `/prometheus/api/v1/query_range`, + method: 'POST', + params, + signal, + }); +}; + +export const getPrometheusQueryRangePostMutationOptions = < + TError = ErrorType, + TContext = unknown, +>(options?: { + mutation?: UseMutationOptions< + Awaited>, + TError, + { params: PrometheusQueryRangePostParams }, + TContext + >; +}): UseMutationOptions< + Awaited>, + TError, + { params: PrometheusQueryRangePostParams }, + TContext +> => { + const mutationKey = ['prometheusQueryRangePost']; + 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>, + { params: PrometheusQueryRangePostParams } + > = (props) => { + const { params } = props ?? {}; + + return prometheusQueryRangePost(params); + }; + + return { mutationFn, ...mutationOptions }; +}; + +export type PrometheusQueryRangePostMutationResult = NonNullable< + Awaited> +>; + +export type PrometheusQueryRangePostMutationError = ErrorType< + PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO +>; + +/** + * @summary Prometheus range query + */ +export const usePrometheusQueryRangePost = < + TError = ErrorType, + TContext = unknown, +>(options?: { + mutation?: UseMutationOptions< + Awaited>, + TError, + { params: PrometheusQueryRangePostParams }, + TContext + >; +}): UseMutationResult< + Awaited>, + TError, + { params: PrometheusQueryRangePostParams }, + TContext +> => { + return useMutation(getPrometheusQueryRangePostMutationOptions(options)); +}; diff --git a/frontend/src/api/generated/services/sigNoz.schemas.ts b/frontend/src/api/generated/services/sigNoz.schemas.ts index 19404b3b037..0453f403cbc 100644 --- a/frontend/src/api/generated/services/sigNoz.schemas.ts +++ b/frontend/src/api/generated/services/sigNoz.schemas.ts @@ -7974,6 +7974,164 @@ export interface PreferencetypesUpdatablePreferenceDTO { value?: unknown; } +export enum PrometheusErrorResponseSchemaDTOErrorType { + bad_data = 'bad_data', + execution = 'execution', + canceled = 'canceled', + timeout = 'timeout', + internal = 'internal', +} +export enum PrometheusErrorResponseSchemaDTOStatus { + error = 'error', +} +export interface PrometheusErrorResponseSchemaDTO { + /** + * @type string + */ + error: string; + /** + * @enum bad_data,execution,canceled,timeout,internal + * @type string + */ + errorType: PrometheusErrorResponseSchemaDTOErrorType; + /** + * @enum error + * @type string + */ + status: PrometheusErrorResponseSchemaDTOStatus; +} + +export enum PrometheusMatrixDataSchemaDTOResultType { + matrix = 'matrix', +} +export type PrometheusSamplePairSchemaDTOItem = number | string; + +/** + * A [timestamp, value] pair: float unix seconds, then the string-encoded sample value ("NaN", "+Inf", "-Inf" included). + * @minItems 2 + * @maxItems 2 + * @nullable + */ +export type PrometheusSamplePairSchemaDTO = + | PrometheusSamplePairSchemaDTOItem[] + | null; + +export type PrometheusMatrixSeriesSchemaDTOMetricAnyOf = { + [key: string]: string; +}; + +/** + * @nullable + */ +export type PrometheusMatrixSeriesSchemaDTOMetric = + PrometheusMatrixSeriesSchemaDTOMetricAnyOf | null; + +export interface PrometheusMatrixSeriesSchemaDTO { + /** + * @type object,null + */ + metric: PrometheusMatrixSeriesSchemaDTOMetric; + /** + * @type array,null + */ + values: (PrometheusSamplePairSchemaDTO | null)[] | null; +} + +export interface PrometheusMatrixDataSchemaDTO { + /** + * @type array,null + */ + result: PrometheusMatrixSeriesSchemaDTO[] | null; + /** + * @enum matrix + * @type string + */ + resultType: PrometheusMatrixDataSchemaDTOResultType; +} + +export type PrometheusVectorSampleSchemaDTOMetricAnyOf = { + [key: string]: string; +}; + +/** + * @nullable + */ +export type PrometheusVectorSampleSchemaDTOMetric = + PrometheusVectorSampleSchemaDTOMetricAnyOf | null; + +export interface PrometheusVectorSampleSchemaDTO { + /** + * @type object,null + */ + metric: PrometheusVectorSampleSchemaDTOMetric; + value: PrometheusSamplePairSchemaDTO | null; +} + +export enum PrometheusVectorDataSchemaDTOResultType { + vector = 'vector', +} +export interface PrometheusVectorDataSchemaDTO { + /** + * @type array,null + */ + result: PrometheusVectorSampleSchemaDTO[] | null; + /** + * @enum vector + * @type string + */ + resultType: PrometheusVectorDataSchemaDTOResultType; +} + +export enum PrometheusScalarDataSchemaDTOResultType { + scalar = 'scalar', +} +export interface PrometheusScalarDataSchemaDTO { + result: PrometheusSamplePairSchemaDTO | null; + /** + * @enum scalar + * @type string + */ + resultType: PrometheusScalarDataSchemaDTOResultType; +} + +export enum PrometheusStringDataSchemaDTOResultType { + string = 'string', +} +export interface PrometheusStringDataSchemaDTO { + result: PrometheusSamplePairSchemaDTO | null; + /** + * @enum string + * @type string + */ + resultType: PrometheusStringDataSchemaDTOResultType; +} + +export type PrometheusQueryDataSchemaDTO = + | PrometheusMatrixDataSchemaDTO + | PrometheusVectorDataSchemaDTO + | PrometheusScalarDataSchemaDTO + | PrometheusStringDataSchemaDTO; + +export enum PrometheusSuccessResponseSchemaDTOStatus { + success = 'success', +} +export interface PrometheusSuccessResponseSchemaDTO { + data: PrometheusQueryDataSchemaDTO; + /** + * @type array + */ + infos?: string[]; + /** + * @enum success + * @type string + */ + status: PrometheusSuccessResponseSchemaDTOStatus; + /** + * @type array + */ + warnings?: string[]; +} + export interface PromotetypesWrappedIndexDTO { fieldDataType?: TelemetrytypesFieldDataTypeDTO; /** @@ -12471,3 +12629,115 @@ export type ReplaceVariables200 = { */ status: string; }; + +export type PrometheusQueryParams = { + /** + * @type string + * @description PromQL expression. + */ + query: string; + /** + * @type string + * @description Evaluation timestamp: RFC3339 or float unix seconds. Defaults to the server's current time. + */ + time?: string; + /** + * @type string + * @description Evaluation timeout: duration string or float seconds. + */ + timeout?: string; + /** + * @type string + * @description Any non-empty value includes query statistics in the response. + */ + stats?: string; +}; + +export type PrometheusQueryPostParams = { + /** + * @type string + * @description PromQL expression. + */ + query: string; + /** + * @type string + * @description Evaluation timestamp: RFC3339 or float unix seconds. Defaults to the server's current time. + */ + time?: string; + /** + * @type string + * @description Evaluation timeout: duration string or float seconds. + */ + timeout?: string; + /** + * @type string + * @description Any non-empty value includes query statistics in the response. + */ + stats?: string; +}; + +export type PrometheusQueryRangeParams = { + /** + * @type string + * @description PromQL expression. + */ + query: string; + /** + * @type string + * @description Range start: RFC3339 or float unix seconds. + */ + start: string; + /** + * @type string + * @description Range end: RFC3339 or float unix seconds. + */ + end: string; + /** + * @type string + * @description Resolution step: duration string or float seconds. + */ + step: string; + /** + * @type string + * @description Evaluation timeout: duration string or float seconds. + */ + timeout?: string; + /** + * @type string + * @description Any non-empty value includes query statistics in the response. + */ + stats?: string; +}; + +export type PrometheusQueryRangePostParams = { + /** + * @type string + * @description PromQL expression. + */ + query: string; + /** + * @type string + * @description Range start: RFC3339 or float unix seconds. + */ + start: string; + /** + * @type string + * @description Range end: RFC3339 or float unix seconds. + */ + end: string; + /** + * @type string + * @description Resolution step: duration string or float seconds. + */ + step: string; + /** + * @type string + * @description Evaluation timeout: duration string or float seconds. + */ + timeout?: string; + /** + * @type string + * @description Any non-empty value includes query statistics in the response. + */ + stats?: string; +}; diff --git a/pkg/apiserver/signozapiserver/prometheus.go b/pkg/apiserver/signozapiserver/prometheus.go new file mode 100644 index 00000000000..bdd6c12942d --- /dev/null +++ b/pkg/apiserver/signozapiserver/prometheus.go @@ -0,0 +1,102 @@ +package signozapiserver + +import ( + "net/http" + "strings" + + "github.com/SigNoz/signoz/pkg/errors" + "github.com/SigNoz/signoz/pkg/http/handler" + "github.com/SigNoz/signoz/pkg/http/render" + "github.com/SigNoz/signoz/pkg/prometheus" + "github.com/SigNoz/signoz/pkg/querybuilder" + "github.com/SigNoz/signoz/pkg/types/authtypes" + "github.com/SigNoz/signoz/pkg/types/coretypes" + "github.com/gorilla/mux" + openapi "github.com/swaggest/openapi-go" +) + +// prometheusOpenAPIHandler skips the default handler wrapper: that wraps +// every response in the house envelope, and these endpoints follow +// Prometheus' wire contract, described by the prometheus package's *Schema +// types. +type prometheusOpenAPIHandler struct { + handlerFunc http.HandlerFunc + id string + summary string + params any +} + +func (h *prometheusOpenAPIHandler) ServeHTTP(rw http.ResponseWriter, req *http.Request) { + h.handlerFunc.ServeHTTP(rw, req) +} + +func (h *prometheusOpenAPIHandler) ServeOpenAPI(opCtx openapi.OperationContext) { + // One route serves GET and POST; operation IDs must stay unique. + id := h.id + if strings.EqualFold(opCtx.Method(), http.MethodPost) { + id += "Post" + } + opCtx.SetID(id) + opCtx.SetTags("prometheus") + opCtx.SetSummary(h.summary) + opCtx.SetDescription("Prometheus-compatible endpoint: the request and response contract is the upstream Prometheus HTTP API (https://prometheus.io/docs/prometheus/latest/querying/api/). Parameters are accepted as URL query parameters or a form-encoded body, on GET and POST alike.") + + for _, scheme := range newScopedSecuritySchemes([]string{coretypes.ResourceTelemetryResourceMetrics.Scope(coretypes.VerbRead)}) { + opCtx.AddSecurity(scheme.Name, scheme.Scopes...) + } + + opCtx.AddReqStructure(h.params) + + opCtx.AddRespStructure( + prometheus.SuccessResponseSchema{}, + openapi.WithContentType("application/json"), + openapi.WithHTTPStatus(http.StatusOK), + ) + for _, statusCode := range []int{http.StatusBadRequest, http.StatusUnprocessableEntity, http.StatusServiceUnavailable, http.StatusInternalServerError} { + opCtx.AddRespStructure( + prometheus.ErrorResponseSchema{}, + openapi.WithContentType("application/json"), + openapi.WithHTTPStatus(statusCode), + ) + } + // The auth middleware answers before the handler and uses the house + // envelope, not Prometheus'. + for _, statusCode := range []int{http.StatusUnauthorized, http.StatusForbidden} { + opCtx.AddRespStructure( + render.ErrorResponse{Status: render.StatusError.String(), Error: &errors.JSON{}}, + openapi.WithContentType("application/json"), + openapi.WithHTTPStatus(statusCode), + ) + } +} + +func (h *prometheusOpenAPIHandler) ResourceDefs() []handler.ResourceDef { + return []handler.ResourceDef{handler.TelemetryResourceDef{ + Verb: coretypes.VerbRead, + Category: coretypes.ActionCategoryDataAccess, + Selector: querybuilder.TelemetrySelector, + Resources: querybuilder.PromQLResources, + }} +} + +func (provider *provider) addPrometheusRoutes(router *mux.Router) error { + if err := router.Handle("/prometheus/api/v1/query", &prometheusOpenAPIHandler{ + handlerFunc: provider.authzMiddleware.CheckResources(provider.prometheusHandler.Query, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName), + id: "PrometheusQuery", + summary: "Prometheus instant query", + params: new(prometheus.QueryParamsSchema), + }).Methods(http.MethodGet, http.MethodPost).GetError(); err != nil { + return err + } + + if err := router.Handle("/prometheus/api/v1/query_range", &prometheusOpenAPIHandler{ + handlerFunc: provider.authzMiddleware.CheckResources(provider.prometheusHandler.QueryRange, authtypes.SigNozAdminRoleName, authtypes.SigNozEditorRoleName, authtypes.SigNozViewerRoleName), + id: "PrometheusQueryRange", + summary: "Prometheus range query", + params: new(prometheus.QueryRangeParamsSchema), + }).Methods(http.MethodGet, http.MethodPost).GetError(); err != nil { + return err + } + + return nil +} diff --git a/pkg/apiserver/signozapiserver/provider.go b/pkg/apiserver/signozapiserver/provider.go index c3b79b858bc..9ffde5159bd 100644 --- a/pkg/apiserver/signozapiserver/provider.go +++ b/pkg/apiserver/signozapiserver/provider.go @@ -32,6 +32,7 @@ import ( "github.com/SigNoz/signoz/pkg/modules/spanmapper" "github.com/SigNoz/signoz/pkg/modules/tracedetail" "github.com/SigNoz/signoz/pkg/modules/user" + "github.com/SigNoz/signoz/pkg/prometheus" "github.com/SigNoz/signoz/pkg/querier" "github.com/SigNoz/signoz/pkg/ruler" "github.com/SigNoz/signoz/pkg/statsreporter" @@ -75,6 +76,7 @@ type provider struct { ruleStateHistoryHandler rulestatehistory.Handler spanMapperHandler spanmapper.Handler alertmanagerHandler alertmanager.Handler + prometheusHandler prometheus.Handler traceDetailHandler tracedetail.Handler rulerHandler ruler.Handler llmPricingRuleHandler llmpricingrule.Handler @@ -113,6 +115,7 @@ func NewFactory( ruleStateHistoryHandler rulestatehistory.Handler, spanMapperHandler spanmapper.Handler, alertmanagerHandler alertmanager.Handler, + prometheusHandler prometheus.Handler, llmPricingRuleHandler llmpricingrule.Handler, traceDetailHandler tracedetail.Handler, rulerHandler ruler.Handler, @@ -154,6 +157,7 @@ func NewFactory( ruleStateHistoryHandler, spanMapperHandler, alertmanagerHandler, + prometheusHandler, llmPricingRuleHandler, traceDetailHandler, rulerHandler, @@ -197,6 +201,7 @@ func newProvider( ruleStateHistoryHandler rulestatehistory.Handler, spanMapperHandler spanmapper.Handler, alertmanagerHandler alertmanager.Handler, + prometheusHandler prometheus.Handler, llmPricingRuleHandler llmpricingrule.Handler, traceDetailHandler tracedetail.Handler, rulerHandler ruler.Handler, @@ -239,6 +244,7 @@ func newProvider( ruleStateHistoryHandler: ruleStateHistoryHandler, spanMapperHandler: spanMapperHandler, alertmanagerHandler: alertmanagerHandler, + prometheusHandler: prometheusHandler, traceDetailHandler: traceDetailHandler, rulerHandler: rulerHandler, llmPricingRuleHandler: llmPricingRuleHandler, @@ -340,6 +346,10 @@ func (provider *provider) AddToRouter(router *mux.Router) error { return err } + if err := provider.addPrometheusRoutes(router); err != nil { + return err + } + if err := provider.addServiceAccountRoutes(router); err != nil { return err } diff --git a/pkg/instrumentation/sdk.go b/pkg/instrumentation/sdk.go index 1b252a5b3c7..867b11b171a 100644 --- a/pkg/instrumentation/sdk.go +++ b/pkg/instrumentation/sdk.go @@ -6,6 +6,7 @@ import ( "github.com/SigNoz/signoz/pkg/errors" "github.com/SigNoz/signoz/pkg/factory" + "github.com/SigNoz/signoz/pkg/instrumentation/tracehandler" "github.com/SigNoz/signoz/pkg/version" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/collectors" @@ -108,7 +109,7 @@ func New(ctx context.Context, cfg Config, build version.Build, serviceName strin } // Set the global tracer provider to the sdk tracer provider so that external packages can use this - otel.SetTracerProvider(sdk.TracerProvider()) + otel.SetTracerProvider(tracehandler.New(sdk.TracerProvider(), tracehandler.NewPromQL())) return &SDK{ sdk: sdk, diff --git a/pkg/instrumentation/tracehandler/promql.go b/pkg/instrumentation/tracehandler/promql.go new file mode 100644 index 00000000000..239d21e1b3d --- /dev/null +++ b/pkg/instrumentation/tracehandler/promql.go @@ -0,0 +1,27 @@ +package tracehandler + +import ( + "context" + "strings" + + "go.opentelemetry.io/otel/trace" + tracenoop "go.opentelemetry.io/otel/trace/noop" +) + +// TODO(srikanthccv): replace with the tracer scope filter (per-scope +// "enabled") when the otel-go trace SDK ships it +// (https://github.com/open-telemetry/opentelemetry-go/issues/8411). +func NewPromQL() Wrapper { + noop := tracenoop.NewTracerProvider().Tracer("") + return WrapperFunc(func(scope string, next StartFunc) StartFunc { + if scope != "" { + return next + } + return func(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, trace.Span) { + if strings.HasPrefix(spanName, "promql") { + return noop.Start(ctx, spanName) + } + return next(ctx, spanName, opts...) + } + }) +} diff --git a/pkg/instrumentation/tracehandler/tracehandler.go b/pkg/instrumentation/tracehandler/tracehandler.go new file mode 100644 index 00000000000..7b64a75322c --- /dev/null +++ b/pkg/instrumentation/tracehandler/tracehandler.go @@ -0,0 +1,53 @@ +package tracehandler + +import ( + "context" + + "go.opentelemetry.io/otel/trace" + "go.opentelemetry.io/otel/trace/embedded" +) + +// StartFunc is to trace.Tracer.Start as loghandler.LogHandlerFunc is to +// loghandler.LogHandler. +type StartFunc func(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, trace.Span) + +// Wrapper is an interface implemented by all trace handlers. scope is the +// instrumentation scope name of the tracer being wrapped; a wrapper that +// does not apply to a scope returns next unchanged. +type Wrapper interface { + Wrap(scope string, next StartFunc) StartFunc +} + +type WrapperFunc func(scope string, next StartFunc) StartFunc + +func (m WrapperFunc) Wrap(scope string, next StartFunc) StartFunc { + return m(scope, next) +} + +type provider struct { + embedded.TracerProvider + base trace.TracerProvider + wrappers []Wrapper +} + +func New(base trace.TracerProvider, wrappers ...Wrapper) trace.TracerProvider { + return &provider{base: base, wrappers: wrappers} +} + +func (p *provider) Tracer(name string, opts ...trace.TracerOption) trace.Tracer { + base := p.base.Tracer(name, opts...) + start := StartFunc(base.Start) + for i := len(p.wrappers) - 1; i >= 0; i-- { + start = p.wrappers[i].Wrap(name, start) + } + return &tracer{start: start} +} + +type tracer struct { + embedded.Tracer + start StartFunc +} + +func (t *tracer) Start(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, trace.Span) { + return t.start(ctx, spanName, opts...) +} diff --git a/pkg/instrumentation/tracehandler/tracehandler_test.go b/pkg/instrumentation/tracehandler/tracehandler_test.go new file mode 100644 index 00000000000..96c0b379438 --- /dev/null +++ b/pkg/instrumentation/tracehandler/tracehandler_test.go @@ -0,0 +1,87 @@ +package tracehandler + +import ( + "context" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" + "go.opentelemetry.io/otel/trace" + tracenoop "go.opentelemetry.io/otel/trace/noop" +) + +func TestWrappersChainInOrder(t *testing.T) { + recorder := tracetest.NewSpanRecorder() + var order []string + observer := func(name string) Wrapper { + return WrapperFunc(func(_ string, next StartFunc) StartFunc { + return func(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, trace.Span) { + order = append(order, name) + return next(ctx, spanName, opts...) + } + }) + } + provider := New(sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(recorder)), observer("first"), observer("second")) + + _, span := provider.Tracer("test").Start(context.Background(), "op") + span.End() + + assert.Equal(t, []string{"first", "second"}, order) + require.Len(t, recorder.Ended(), 1) + assert.Equal(t, "op", recorder.Ended()[0].Name()) +} + +func TestScopedWrapperSkipsOtherScopes(t *testing.T) { + recorder := tracetest.NewSpanRecorder() + noop := tracenoop.NewTracerProvider().Tracer("") + dropAnonymous := WrapperFunc(func(scope string, next StartFunc) StartFunc { + if scope != "" { + return next + } + return func(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, trace.Span) { + return noop.Start(ctx, spanName) + } + }) + provider := New(sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(recorder)), dropAnonymous) + + _, dropped := provider.Tracer("").Start(context.Background(), "anon") + dropped.End() + _, kept := provider.Tracer("named").Start(context.Background(), "op") + kept.End() + + require.Len(t, recorder.Ended(), 1) + assert.Equal(t, "op", recorder.Ended()[0].Name()) +} + +func TestPromQL(t *testing.T) { + recorder := tracetest.NewSpanRecorder() + provider := New(sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(recorder)), NewPromQL()) + + ctx, root := provider.Tracer("http").Start(context.Background(), "GET /api") + + engineCtx, engineSpan := provider.Tracer("").Start(ctx, "promqlInnerEval eval *promql.BinaryExpr") + assert.False(t, engineSpan.IsRecording(), "promql engine spans must not record") + assert.Equal(t, root.SpanContext().SpanID(), engineSpan.SpanContext().SpanID(), "the filtered span must keep the parent's span context") + + _, child := provider.Tracer("clickhouse").Start(engineCtx, "clickhouse.query") + child.End() + + _, other := provider.Tracer("").Start(ctx, "http.request") + other.End() + root.End() + + var names []string + var childParent string + for _, span := range recorder.Ended() { + names = append(names, span.Name()) + if span.Name() == "clickhouse.query" { + childParent = span.Parent().SpanID().String() + } + } + require.ElementsMatch(t, []string{"clickhouse.query", "http.request", "GET /api"}, names) + assert.Equal(t, root.SpanContext().SpanID().String(), childParent, "descendants of a filtered span must attach to the surrounding span") + assert.False(t, strings.HasPrefix(recorder.Ended()[0].Name(), "promql")) +} diff --git a/pkg/prometheus/clickhouseprometheusv2/capture.go b/pkg/prometheus/clickhouseprometheusv2/capture.go index 118716841b3..c6413de90f5 100644 --- a/pkg/prometheus/clickhouseprometheusv2/capture.go +++ b/pkg/prometheus/clickhouseprometheusv2/capture.go @@ -73,8 +73,8 @@ func (c *captureQuerier) LabelNames(context.Context, *storage.LabelHints, ...*la } // metricNamesFromMatchers extracts the statically known metric name, if any. -// The live path derives names from the matched series; the capture path has -// no execution results, so only a __name__ equality contributes. +// Only a __name__ equality contributes; a regex selector needs a series +// lookup to learn the concrete names. func metricNamesFromMatchers(matchers []*labels.Matcher) []string { for _, m := range matchers { if m.Name == metricNameLabel && m.Type == labels.MatchEqual && m.Value != "" { diff --git a/pkg/prometheus/clickhouseprometheusv2/transpiler_exec.go b/pkg/prometheus/clickhouseprometheusv2/transpiler_exec.go index 87fea487d6f..77007bded64 100644 --- a/pkg/prometheus/clickhouseprometheusv2/transpiler_exec.go +++ b/pkg/prometheus/clickhouseprometheusv2/transpiler_exec.go @@ -88,7 +88,8 @@ func (e *executor) TryExecuteRange(ctx context.Context, qs string, start, end ti } // Evaluate every unit concurrently on its own grid (the query grid, or a - // subquery grid); each is one series lookup plus one grid query. + // subquery grid); each is one grid query (see executeUnit for when a + // series lookup precedes it). results := make([][]transpiledSeries, len(plan.units)) eg, egCtx := errgroup.WithContext(ctx) for i, unit := range plan.units { @@ -142,19 +143,27 @@ func (e *executor) executeUnit(ctx context.Context, unit *coreUnit, grid gridCon dataStart := startMs - unit.offsetMs - windowMs dataEnd := endMs - unit.offsetMs - seriesQuery, seriesArgs, err := buildSeriesQuery(dataStart, dataEnd, unit.matchers) - if err != nil { - return nil, err - } - lookup, err := e.client.selectSeries(ctx, seriesQuery, seriesArgs) - if err != nil { - return nil, err - } - if len(lookup.fingerprints) == 0 { - return nil, nil + // The group-key join resolves the matchers on its own, so the unit + // statement only needs concrete metric names for the samples + // primary-key prefix. A selector without a static __name__ learns them + // through the series lookup; every other selector skips the roundtrip. + metricNames := metricNamesFromMatchers(unit.matchers) + if metricNames == nil { + seriesQuery, seriesArgs, err := buildSeriesQuery(dataStart, dataEnd, unit.matchers) + if err != nil { + return nil, err + } + lookup, err := e.client.selectSeries(ctx, seriesQuery, seriesArgs) + if err != nil { + return nil, err + } + if len(lookup.fingerprints) == 0 { + return nil, nil + } + metricNames = lookup.metricNames } - query, args, err := buildUnitSQL(unit, lookup.metricNames, dataStart, dataEnd, startMs, endMs, stepMs, e.client.lookbackMs) + query, args, err := buildUnitSQL(unit, metricNames, dataStart, dataEnd, startMs, endMs, stepMs, e.client.lookbackMs) if err != nil { return nil, err } diff --git a/pkg/prometheus/clickhouseprometheusv2/transpiler_test.go b/pkg/prometheus/clickhouseprometheusv2/transpiler_test.go index a042ae5cd4f..a14e7c79239 100644 --- a/pkg/prometheus/clickhouseprometheusv2/transpiler_test.go +++ b/pkg/prometheus/clickhouseprometheusv2/transpiler_test.go @@ -27,9 +27,15 @@ func newTestClient(t *testing.T) (*client, *telemetrystoretest.Provider) { return newClient(settings, store, prometheus.Config{}), store } -var seriesCols = []cmock.ColumnType{ - {Name: "fingerprint", Type: "UInt64"}, - {Name: "labels", Type: "String"}, +var unitCols = []cmock.ColumnType{ + {Name: "gkey", Type: "String"}, + {Name: "grid", Type: "Array(Nullable(Float64))"}, +} + +// anyArgs matches a bound-argument list by count alone: the mock treats a +// nil expected argument as a wildcard. +func anyArgs(n int) []any { + return make([]any, n) } func parse(t *testing.T, q string) parser.Expr { @@ -553,7 +559,7 @@ func TestTryExecuteRange_WindowedGateFallsBack(t *testing.T) { // 1m range at 5m step: the windows are disjoint slivers — no // divisibility or width requirement, so this transpiles. - store.Mock().ExpectQuery("SELECT fingerprint, any\\(labels\\)").WithArgs("up", int64(1_699_999_200_000), int64(1_700_003_600_000)).WillReturnRows(cmock.NewRows(seriesCols, [][]any{})) + store.Mock().ExpectQuery("FROM signoz_metrics\\.distributed_samples_v4").WithArgs(anyArgs(9)...).WillReturnRows(cmock.NewRows(unitCols, [][]any{})) _, ok, err = e.TryExecuteRange(context.Background(), `avg_over_time(up[1m])`, start, end, 5*time.Minute) require.NoError(t, err) assert.True(t, ok, "range below step is the disjoint form and must transpile") @@ -637,12 +643,12 @@ func TestTryExecuteRange_LastStyleWindowBelowStepTranspiles(t *testing.T) { start := time.UnixMilli(1_700_000_000_000) end := time.UnixMilli(1_700_003_600_000) - store.Mock().ExpectQuery("SELECT fingerprint, any\\(labels\\)").WithArgs("up", int64(1_699_999_200_000), int64(1_700_003_600_000)).WillReturnRows(cmock.NewRows(seriesCols, [][]any{})) + store.Mock().ExpectQuery("timeSeriesLastToGrid").WithArgs(anyArgs(10)...).WillReturnRows(cmock.NewRows(unitCols, [][]any{})) _, ok, err := e.TryExecuteRange(context.Background(), `sum by (pod) (up)`, start, end, time.Hour) require.NoError(t, err) assert.True(t, ok, "instant selection at step > lookback must transpile") - store.Mock().ExpectQuery("SELECT fingerprint, any\\(labels\\)").WithArgs("up", int64(1_699_999_200_000), int64(1_700_003_600_000)).WillReturnRows(cmock.NewRows(seriesCols, [][]any{})) + store.Mock().ExpectQuery("timeSeriesLastToGrid").WithArgs(anyArgs(9)...).WillReturnRows(cmock.NewRows(unitCols, [][]any{})) _, ok, err = e.TryExecuteRange(context.Background(), `last_over_time(up[10m])`, start, end, time.Hour) require.NoError(t, err) assert.True(t, ok, "last_over_time at range < step must transpile") diff --git a/pkg/prometheus/handler.go b/pkg/prometheus/handler.go new file mode 100644 index 00000000000..a5022b45c4c --- /dev/null +++ b/pkg/prometheus/handler.go @@ -0,0 +1,211 @@ +package prometheus + +import ( + "context" + "log/slog" + "math" + "net/http" + "strconv" + "time" + + promModel "github.com/prometheus/common/model" + "github.com/prometheus/prometheus/promql" + "github.com/prometheus/prometheus/util/stats" + + "github.com/SigNoz/signoz/pkg/errors" +) + +// Handler serves the Prometheus HTTP query API over a Prometheus provider: +// /query and /query_range in the shape of Prometheus' /api/v1 endpoints +// (https://prometheus.io/docs/prometheus/latest/querying/api/), intended to +// be mounted under a distinguishing prefix (/prometheus/api/v1) so +// PromQL-only endpoints are separate from the SigNoz query APIs. The request +// and response contracts follow Prometheus: form-encoded GET/POST params, +// {"status":"success","data":{resultType,result}} on success and +// {"status":"error","errorType","error"} with Prometheus' status codes on +// failure — so Prometheus-compatible clients can point at the prefix. The +// wire shapes are documented as OpenAPI schemas in render.go. +type Handler interface { + Query(http.ResponseWriter, *http.Request) + + QueryRange(http.ResponseWriter, *http.Request) +} + +type handler struct { + logger *slog.Logger + prom Prometheus +} + +func NewHandler(logger *slog.Logger, prom Prometheus) Handler { + return &handler{logger: logger, prom: prom} +} + +// QueryRange evaluates an expression over a grid: query, start, end, step, +// and optional timeout/stats params, all in Prometheus' formats. +func (h *handler) QueryRange(w http.ResponseWriter, r *http.Request) { + start, err := parseTime(r.FormValue("start")) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + end, err := parseTime(r.FormValue("end")) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + if end.Before(start) { + h.respondError(r.Context(), w, errBadData, errors.NewInvalidInputf(errors.CodeInvalidInput, "end timestamp must not be before start time")) + return + } + step, err := parseDuration(r.FormValue("step")) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + if step <= 0 { + h.respondError(r.Context(), w, errBadData, errors.NewInvalidInputf(errors.CodeInvalidInput, "zero or negative query resolution step widths are not accepted. Try a positive integer")) + return + } + // The engine materializes every point of every series; an unbounded + // grid is an unbounded allocation. 11,000 points covers 60s resolution + // for a week or 1h resolution for a year. + if end.Sub(start)/step > 11000 { + h.respondError(r.Context(), w, errBadData, errors.NewInvalidInputf(errors.CodeInvalidInput, "exceeded maximum resolution of 11,000 points per timeseries. Try decreasing the query resolution (?step=XX)")) + return + } + + ctx, cancel, err := h.contextWithTimeout(r) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + defer cancel() + + if h.tryRangeExecutor(ctx, w, r, start, end, step) { + return + } + + qry, err := h.prom.Engine().NewRangeQuery(ctx, h.prom.Storage(), nil, r.FormValue("query"), start, end, step) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + h.exec(ctx, w, r, qry) +} + +// tryRangeExecutor serves the query the way a RangeExecutor provider is +// designed to serve: evaluated inside the datastore when the shape allows. +// It reports whether the response was written. +func (h *handler) tryRangeExecutor(ctx context.Context, w http.ResponseWriter, r *http.Request, start, end time.Time, step time.Duration) bool { + re, ok := h.prom.(RangeExecutor) + if !ok { + return false + } + matrix, served, err := re.TryExecuteRange(ctx, r.FormValue("query"), start, end, step) + if err != nil { + h.respondError(ctx, w, errExec, err) + return true + } + if !served { + return false + } + h.respond(ctx, w, &queryData{ResultType: matrix.Type(), Result: matrix}, nil, nil) + return true +} + +// Query evaluates an expression at a single instant: query and optional +// time/timeout/stats params. A missing time evaluates at the server's now, +// as in Prometheus. +func (h *handler) Query(w http.ResponseWriter, r *http.Request) { + ts := time.Now() + if t := r.FormValue("time"); t != "" { + var err error + ts, err = parseTime(t) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + } + + ctx, cancel, err := h.contextWithTimeout(r) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + defer cancel() + + qry, err := h.prom.Engine().NewInstantQuery(ctx, h.prom.Storage(), nil, r.FormValue("query"), ts) + if err != nil { + h.respondError(r.Context(), w, errBadData, err) + return + } + h.exec(ctx, w, r, qry) +} + +func (h *handler) exec(ctx context.Context, w http.ResponseWriter, r *http.Request, qry promql.Query) { + defer qry.Close() + res := qry.Exec(ctx) + if res.Err != nil { + h.logger.ErrorContext(ctx, "error evaluating promql query", errors.Attr(res.Err)) + switch res.Err.(type) { + case promql.ErrQueryCanceled: + h.respondError(ctx, w, errCanceled, res.Err) + case promql.ErrQueryTimeout: + h.respondError(ctx, w, errTimeout, res.Err) + case promql.ErrStorage: + h.respondError(ctx, w, errInternal, res.Err) + default: + h.respondError(ctx, w, errExec, res.Err) + } + return + } + + data := &queryData{ResultType: res.Value.Type(), Result: res.Value} + if r.FormValue("stats") != "" { + data.Stats = stats.NewQueryStats(qry.Stats()) + } + warnings, infos := res.Warnings.AsStrings(r.FormValue("query"), 10, 10) + h.respond(ctx, w, data, warnings, infos) +} + +func (h *handler) contextWithTimeout(r *http.Request) (context.Context, context.CancelFunc, error) { + ctx := r.Context() + if to := r.FormValue("timeout"); to != "" { + timeout, err := parseDuration(to) + if err != nil { + return nil, nil, err + } + ctx, cancel := context.WithTimeout(ctx, timeout) + return ctx, cancel, nil + } + ctx, cancel := context.WithCancel(ctx) + return ctx, cancel, nil +} + +// parseTime accepts Prometheus' time formats: float unix seconds or RFC3339. +func parseTime(s string) (time.Time, error) { + if t, err := strconv.ParseFloat(s, 64); err == nil { + sec, ns := math.Modf(t) + return time.Unix(int64(sec), int64(ns*float64(time.Second))), nil + } + if t, err := time.Parse(time.RFC3339Nano, s); err == nil { + return t, nil + } + return time.Time{}, errors.NewInvalidInputf(errors.CodeInvalidInput, "cannot parse %q to a valid timestamp", s) +} + +// parseDuration accepts Prometheus' duration formats: float seconds or a +// duration string like 5m. +func parseDuration(s string) (time.Duration, error) { + if d, err := strconv.ParseFloat(s, 64); err == nil { + ts := d * float64(time.Second) + if ts > float64(math.MaxInt64) || ts < float64(math.MinInt64) { + return 0, errors.NewInvalidInputf(errors.CodeInvalidInput, "cannot parse %q to a valid duration. It overflows int64", s) + } + return time.Duration(ts), nil + } + if d, err := promModel.ParseDuration(s); err == nil { + return time.Duration(d), nil + } + return 0, errors.NewInvalidInputf(errors.CodeInvalidInput, "cannot parse %q to a valid duration", s) +} diff --git a/pkg/prometheus/render.go b/pkg/prometheus/render.go new file mode 100644 index 00000000000..868491b62fe --- /dev/null +++ b/pkg/prometheus/render.go @@ -0,0 +1,164 @@ +package prometheus + +import ( + "context" + "encoding/json" + "net/http" + + "github.com/prometheus/prometheus/promql/parser" + "github.com/prometheus/prometheus/util/stats" + "github.com/swaggest/jsonschema-go" + + "github.com/SigNoz/signoz/pkg/errors" +) + +// This file is the single description of the Prometheus API wire shapes: +// the runtime envelope the handler encodes, and the *Schema types that +// document the same shapes in the generated OpenAPI spec. The contract is +// upstream's (https://prometheus.io/docs/prometheus/latest/querying/api/); +// the schemas describe it, they do not define it. + +type errorType string + +const ( + errBadData errorType = "bad_data" + errExec errorType = "execution" + errCanceled errorType = "canceled" + errTimeout errorType = "timeout" + errInternal errorType = "internal" +) + +type queryData struct { + ResultType parser.ValueType `json:"resultType"` + Result parser.Value `json:"result"` + Stats stats.QueryStats `json:"stats,omitempty"` +} + +type response struct { + Status string `json:"status"` + Data *queryData `json:"data,omitempty"` + ErrorType errorType `json:"errorType,omitempty"` + Error string `json:"error,omitempty"` + Warnings []string `json:"warnings,omitempty"` + Infos []string `json:"infos,omitempty"` +} + +func (h *handler) respond(ctx context.Context, w http.ResponseWriter, data *queryData, warnings, infos []string) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + if err := json.NewEncoder(w).Encode(&response{Status: "success", Data: data, Warnings: warnings, Infos: infos}); err != nil { + h.logger.ErrorContext(ctx, "error writing prometheus api response", errors.Attr(err)) + } +} + +// respondError follows Prometheus' status-code mapping: bad_data 400, +// execution 422, canceled/timeout 503, internal 500. +func (h *handler) respondError(ctx context.Context, w http.ResponseWriter, typ errorType, err error) { + code := http.StatusInternalServerError + switch typ { + case errBadData: + code = http.StatusBadRequest + case errExec: + code = http.StatusUnprocessableEntity + case errCanceled, errTimeout: + code = http.StatusServiceUnavailable + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(code) + if encErr := json.NewEncoder(w).Encode(&response{Status: "error", ErrorType: typ, Error: err.Error()}); encErr != nil { + h.logger.ErrorContext(ctx, "error writing prometheus api error response", errors.Attr(encErr)) + } +} + +// The endpoints accept parameters as URL query params or a form-encoded +// body, on GET and POST alike. +type QueryParamsSchema struct { + Query string `query:"query" required:"true" description:"PromQL expression."` + Time string `query:"time" description:"Evaluation timestamp: RFC3339 or float unix seconds. Defaults to the server's current time."` + Timeout string `query:"timeout" description:"Evaluation timeout: duration string or float seconds."` + Stats string `query:"stats" description:"Any non-empty value includes query statistics in the response."` +} + +type QueryRangeParamsSchema struct { + Query string `query:"query" required:"true" description:"PromQL expression."` + Start string `query:"start" required:"true" description:"Range start: RFC3339 or float unix seconds."` + End string `query:"end" required:"true" description:"Range end: RFC3339 or float unix seconds."` + Step string `query:"step" required:"true" description:"Resolution step: duration string or float seconds."` + Timeout string `query:"timeout" description:"Evaluation timeout: duration string or float seconds."` + Stats string `query:"stats" description:"Any non-empty value includes query statistics in the response."` +} + +type SuccessResponseSchema struct { + Status string `json:"status" enum:"success" required:"true"` + Data QueryDataSchema `json:"data" required:"true"` + Warnings []string `json:"warnings,omitempty"` + Infos []string `json:"infos,omitempty"` +} + +// QueryDataSchema is the result union, discriminated by resultType. +type QueryDataSchema struct{} + +var _ jsonschema.OneOfExposer = QueryDataSchema{} + +func (QueryDataSchema) JSONSchemaOneOf() []interface{} { + return []interface{}{MatrixDataSchema{}, VectorDataSchema{}, ScalarDataSchema{}, StringDataSchema{}} +} + +type MatrixDataSchema struct { + ResultType string `json:"resultType" enum:"matrix" required:"true"` + Result []MatrixSeriesSchema `json:"result" required:"true"` +} + +type MatrixSeriesSchema struct { + Metric map[string]string `json:"metric" required:"true"` + Values []SamplePairSchema `json:"values" required:"true"` +} + +type VectorDataSchema struct { + ResultType string `json:"resultType" enum:"vector" required:"true"` + Result []VectorSampleSchema `json:"result" required:"true"` +} + +type VectorSampleSchema struct { + Metric map[string]string `json:"metric" required:"true"` + Value SamplePairSchema `json:"value" required:"true"` +} + +type ScalarDataSchema struct { + ResultType string `json:"resultType" enum:"scalar" required:"true"` + Result SamplePairSchema `json:"result" required:"true"` +} + +type StringDataSchema struct { + ResultType string `json:"resultType" enum:"string" required:"true"` + Result SamplePairSchema `json:"result" required:"true"` +} + +// SamplePairSchema is the positional [timestamp, value] pair: a float of +// unix seconds, then the value as a string ("NaN", "+Inf" and "-Inf" +// included). Struct reflection cannot express a positional array, so the +// schema is authored by hand. +type SamplePairSchema struct{} + +var _ jsonschema.Exposer = SamplePairSchema{} + +func (SamplePairSchema) JSONSchema() (jsonschema.Schema, error) { + item := jsonschema.Schema{} + item.WithOneOf( + (&jsonschema.Schema{}).WithType(jsonschema.Number.Type()).ToSchemaOrBool(), + (&jsonschema.Schema{}).WithType(jsonschema.String.Type()).ToSchemaOrBool(), + ) + s := jsonschema.Schema{} + s.WithType(jsonschema.Array.Type()) + s.WithMinItems(2) + s.WithMaxItems(2) + s.WithItems(*(&jsonschema.Items{}).WithSchemaOrBool(item.ToSchemaOrBool())) + s.WithDescription(`A [timestamp, value] pair: float unix seconds, then the string-encoded sample value ("NaN", "+Inf", "-Inf" included).`) + return s, nil +} + +type ErrorResponseSchema struct { + Status string `json:"status" enum:"error" required:"true"` + ErrorType string `json:"errorType" enum:"bad_data,execution,canceled,timeout,internal" required:"true"` + Error string `json:"error" required:"true"` +} diff --git a/pkg/query-service/app/http_handler.go b/pkg/query-service/app/http_handler.go index ddf71110c2d..e0a963d51e2 100644 --- a/pkg/query-service/app/http_handler.go +++ b/pkg/query-service/app/http_handler.go @@ -387,6 +387,7 @@ func (aH *APIHandler) Respond(w http.ResponseWriter, data interface{}) { func (aH *APIHandler) RegisterRoutes(router *mux.Router, am *middleware.AuthZ) { router.HandleFunc("/api/v1/query_range", am.ViewAccess(aH.queryRangeMetrics)).Methods(http.MethodGet) router.HandleFunc("/api/v1/query", am.ViewAccess(aH.queryMetrics)).Methods(http.MethodGet) + router.HandleFunc("/api/v1/rules", am.ViewAccess(aH.listRules)).Methods(http.MethodGet) router.HandleFunc("/api/v1/rules/{id}", am.ViewAccess(aH.getRule)).Methods(http.MethodGet) router.HandleFunc("/api/v1/rules", am.EditAccess(aH.createRule)).Methods(http.MethodPost) diff --git a/pkg/querybuilder/query_range_resources.go b/pkg/querybuilder/query_range_resources.go index c92284b0c41..447047c2564 100644 --- a/pkg/querybuilder/query_range_resources.go +++ b/pkg/querybuilder/query_range_resources.go @@ -75,6 +75,16 @@ func queryRangeVariables(body []byte) (map[string]qbtypes.VariableItem, error) { return variables, nil } +// PromQLResources is the resource set of a bare PromQL query: metrics on +// the promql wildcard, the same ID resourcesForQuery assigns to a PromQL +// query inside a composite — one grant covers both entry points. +func PromQLResources(coretypes.ExtractorContext) ([]coretypes.ResourceWithID, error) { + return []coretypes.ResourceWithID{{ + Resource: coretypes.ResourceTelemetryResourceMetrics, + ID: qbtypes.QueryTypePromQL.StringValue() + "/" + coretypes.WildCardSelectorString, + }}, nil +} + func resourcesForQuery(query gjson.Result, variables map[string]qbtypes.VariableItem) ([]coretypes.ResourceWithID, error) { queryType := query.Get("type").String() typeWildcard := queryType + "/" + coretypes.WildCardSelectorString diff --git a/pkg/signoz/handler.go b/pkg/signoz/handler.go index e56018db876..80b471ae20e 100644 --- a/pkg/signoz/handler.go +++ b/pkg/signoz/handler.go @@ -50,6 +50,7 @@ import ( "github.com/SigNoz/signoz/pkg/modules/tracedetail/impltracedetail" "github.com/SigNoz/signoz/pkg/modules/tracefunnel" "github.com/SigNoz/signoz/pkg/modules/tracefunnel/impltracefunnel" + "github.com/SigNoz/signoz/pkg/prometheus" "github.com/SigNoz/signoz/pkg/querier" "github.com/SigNoz/signoz/pkg/ruler" "github.com/SigNoz/signoz/pkg/ruler/signozruler" @@ -84,6 +85,7 @@ type Handlers struct { RuleStateHistory rulestatehistory.Handler SpanMapperHandler spanmapper.Handler AlertmanagerHandler alertmanager.Handler + PrometheusHandler prometheus.Handler TraceDetail tracedetail.Handler RulerHandler ruler.Handler LLMPricingRuleHandler llmpricingrule.Handler @@ -104,6 +106,7 @@ func NewHandlers( zeusService zeus.Zeus, registryHandler factory.Handler, alertmanagerService alertmanager.Alertmanager, + prometheusService prometheus.Prometheus, rulerService ruler.Ruler, statsAggregator statsreporter.Aggregator, ) Handlers { @@ -133,6 +136,7 @@ func NewHandlers( CloudIntegrationHandler: implcloudintegration.NewHandler(modules.CloudIntegration), SpanMapperHandler: implspanmapper.NewHandler(modules.SpanMapper), AlertmanagerHandler: signozalertmanager.NewHandler(alertmanagerService), + PrometheusHandler: prometheus.NewHandler(providerSettings.Logger, prometheusService), TraceDetail: impltracedetail.NewHandler(modules.TraceDetail), RulerHandler: signozruler.NewHandler(rulerService), LLMPricingRuleHandler: impllmpricingrule.NewHandler(modules.LLMPricingRule), diff --git a/pkg/signoz/handler_test.go b/pkg/signoz/handler_test.go index 60ac147c4fb..5d88a12b9e0 100644 --- a/pkg/signoz/handler_test.go +++ b/pkg/signoz/handler_test.go @@ -63,7 +63,7 @@ func TestNewHandlers(t *testing.T) { querierHandler := querier.NewHandler(providerSettings, nil, nil) registryHandler := factory.NewHandler(nil) - handlers := NewHandlers(modules, providerSettings, nil, querierHandler, nil, nil, nil, nil, nil, nil, nil, registryHandler, alertmanager, nil, nil) + handlers := NewHandlers(modules, providerSettings, nil, querierHandler, nil, nil, nil, nil, nil, nil, nil, registryHandler, alertmanager, nil, nil, nil) reflectVal := reflect.ValueOf(handlers) for i := 0; i < reflectVal.NumField(); i++ { f := reflectVal.Field(i) diff --git a/pkg/signoz/openapi.go b/pkg/signoz/openapi.go index 52a325f132c..5e6b3d8fe86 100644 --- a/pkg/signoz/openapi.go +++ b/pkg/signoz/openapi.go @@ -37,6 +37,7 @@ import ( "github.com/SigNoz/signoz/pkg/modules/spanmapper" "github.com/SigNoz/signoz/pkg/modules/tracedetail" "github.com/SigNoz/signoz/pkg/modules/user" + "github.com/SigNoz/signoz/pkg/prometheus" "github.com/SigNoz/signoz/pkg/querier" "github.com/SigNoz/signoz/pkg/ruler" "github.com/SigNoz/signoz/pkg/statsreporter" @@ -88,6 +89,7 @@ func NewOpenAPI(ctx context.Context, instrumentation instrumentation.Instrumenta struct{ rulestatehistory.Handler }{}, struct{ spanmapper.Handler }{}, struct{ alertmanager.Handler }{}, + struct{ prometheus.Handler }{}, struct{ llmpricingrule.Handler }{}, struct{ tracedetail.Handler }{}, struct{ ruler.Handler }{}, diff --git a/pkg/signoz/provider.go b/pkg/signoz/provider.go index 68cc7ab4a2d..e5dcbe0f093 100644 --- a/pkg/signoz/provider.go +++ b/pkg/signoz/provider.go @@ -343,6 +343,7 @@ func NewAPIServerProviderFactories(orgGetter organization.Getter, authz authz.Au handlers.RuleStateHistory, handlers.SpanMapperHandler, handlers.AlertmanagerHandler, + handlers.PrometheusHandler, handlers.LLMPricingRuleHandler, handlers.TraceDetail, handlers.RulerHandler, diff --git a/pkg/signoz/signoz.go b/pkg/signoz/signoz.go index 6224365f215..95f4a74487b 100644 --- a/pkg/signoz/signoz.go +++ b/pkg/signoz/signoz.go @@ -617,7 +617,7 @@ func New( // Initialize all handlers for the modules registryHandler := factory.NewHandler(registry) - handlers := NewHandlers(modules, providerSettings, analytics, querierHandler, licensing, global, flagger, gateway, telemetryMetadataStore, authz, zeus, registryHandler, alertmanager, rulerInstance, statsAggregator) + handlers := NewHandlers(modules, providerSettings, analytics, querierHandler, licensing, global, flagger, gateway, telemetryMetadataStore, authz, zeus, registryHandler, alertmanager, prometheus, rulerInstance, statsAggregator) // Initialize the API server (after registry so it can access service health) apiserverInstance, err := factory.NewProviderFromNamedMap( diff --git a/pkg/types/preferencetypes/name.go b/pkg/types/preferencetypes/name.go index 5b0643ae25c..d256e855610 100644 --- a/pkg/types/preferencetypes/name.go +++ b/pkg/types/preferencetypes/name.go @@ -23,6 +23,7 @@ var ( NameSpanDetailsPreviewAttributes = Name{valuer.NewString("span_details_preview_attributes")} NameSpanDetailsColorByAttribute = Name{valuer.NewString("span_details_color_by_attribute")} NameSpanPercentileResourceAttributes = Name{valuer.NewString("span_percentile_resource_attributes")} + NameLogDetailsPinnedAttributes = Name{valuer.NewString("log_details_pinned_attributes")} ) type Name struct{ valuer.String } @@ -45,6 +46,7 @@ func NewName(name string) (Name, error) { NameSpanDetailsPreviewAttributes.StringValue(), NameSpanDetailsColorByAttribute.StringValue(), NameSpanPercentileResourceAttributes.StringValue(), + NameLogDetailsPinnedAttributes.StringValue(), }, name, ) diff --git a/pkg/types/preferencetypes/preference.go b/pkg/types/preferencetypes/preference.go index df91acb777c..e65d31010a0 100644 --- a/pkg/types/preferencetypes/preference.go +++ b/pkg/types/preferencetypes/preference.go @@ -190,6 +190,15 @@ func NewAvailablePreference() map[Name]Preference { AllowedValues: []string{}, Value: MustNewValue([]any{}, ValueTypeArray), }, + NameLogDetailsPinnedAttributes: { + Name: NameLogDetailsPinnedAttributes, + Description: "List of pinned attributes in log details drawer.", + ValueType: ValueTypeArray, + DefaultValue: MustNewValue([]any{}, ValueTypeArray), + AllowedScopes: []Scope{ScopeUser}, + AllowedValues: []string{}, + Value: MustNewValue([]any{}, ValueTypeArray), + }, } } diff --git a/tests/fixtures/promqltestcorpus.py b/tests/fixtures/promqltestcorpus.py new file mode 100644 index 00000000000..87168a66d72 --- /dev/null +++ b/tests/fixtures/promqltestcorpus.py @@ -0,0 +1,66 @@ +import json +import math +import os +from collections.abc import Callable +from datetime import UTC, datetime, timedelta + +from fixtures.metrics import Metrics + +TESTDATA_DIR = os.path.join(os.path.dirname(__file__), "..", "integration", "testdata", "promqltestcorpus") +CORPUS_FILE = os.path.join(TESTDATA_DIR, "corpus.json") + +# Datasets sit on disjoint time windows (2h gaps, far beyond the 5m lookback) +# so one bulk ingest serves every case without cross-talk. +ISOLATION_GAP_MS = 2 * 3600 * 1000 +SPECIALS = {"NaN": math.nan, "Inf": math.inf, "-Inf": -math.inf} + + +def ingest_promqltest_corpus(insert_metrics: Callable[[list[Metrics]], None]) -> tuple[dict, dict[int, int]]: + """Loads the frozen corpus, lays its datasets end to end on the timeline + (newest last, ending safely in the past), ingests every sample, and + returns (corpus, dataset base timestamps). + + Dataset bases are hour-aligned: registration rows are hour-bucketed, so + behavior depends on where samples fall relative to hour boundaries, and + exact known-divergences enforcement needs identical placement every run.""" + with open(CORPUS_FILE, encoding="utf-8") as f: + corpus = json.load(f) + + cases_by_dataset: dict[int, list[dict]] = {} + for case in corpus["cases"]: + cases_by_dataset.setdefault(case["dataset"], []).append(case) + + spans = {} + for ds in corpus["datasets"]: + sample_max = max((s["samples"][-1][0] for s in ds["series"] if s["samples"]), default=0) + case_max = max((c["end_ms"] for c in cases_by_dataset.get(ds["id"], [])), default=0) + spans[ds["id"]] = max(sample_max, case_max) + corpus["meta"]["lookback_ms"] + + hour_ms = 3_600_000 + advances = {ds["id"]: -(-(spans[ds["id"]] + ISOLATION_GAP_MS) // hour_ms) * hour_ms for ds in corpus["datasets"]} + total = sum(advances.values()) + now = datetime.now(tz=UTC).replace(second=0, microsecond=0) + cursor = (int((now - timedelta(hours=1)).timestamp() * 1000) - total) // hour_ms * hour_ms + + bases: dict[int, int] = {} + metrics: list[Metrics] = [] + for ds in corpus["datasets"]: + bases[ds["id"]] = cursor + for series in ds["series"]: + labels = dict(series["labels"]) + metric_name = labels.pop("__name__") + for off_ms, raw in series["samples"]: + stale = raw == "stale" + metrics.append( + Metrics( + metric_name=metric_name, + labels=labels, + timestamp=datetime.fromtimestamp((cursor + off_ms) / 1000, tz=UTC), + value=0.0 if stale else (SPECIALS[raw] if isinstance(raw, str) else float(raw)), + flags=1 if stale else 0, + ) + ) + cursor += advances[ds["id"]] + + insert_metrics(metrics) + return corpus, bases diff --git a/tests/integration/tests/promapiconformance/01_prometheus_api_corpus.py b/tests/integration/tests/promapiconformance/01_prometheus_api_corpus.py new file mode 100644 index 00000000000..6a2c90aa106 --- /dev/null +++ b/tests/integration/tests/promapiconformance/01_prometheus_api_corpus.py @@ -0,0 +1,138 @@ +import json +import math +from collections.abc import Callable +from http import HTTPStatus + +import requests + +from fixtures import types +from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD +from fixtures.metrics import Metrics +from fixtures.promqltestcorpus import ingest_promqltest_corpus + +# The same frozen corpus the promqlconformance package replays through +# /api/v5/query_range, here replayed against the /prometheus/api/v1 endpoints +# with clickhousev2 as the serving provider (see conftest.py) — the two paths +# nothing else exercises. Range cases go to query_range, where a +# RangeExecutor provider serves transpiled statements when the shape allows. +# Instant cases go to /query with a real `time` parameter, so they need no +# grid encoding. +# +# Prometheus API sample values are strings, "NaN"/"+Inf"/"-Inf" included. +SPECIALS = {"NaN": math.nan, "Inf": math.inf, "+Inf": math.inf, "-Inf": -math.inf} +QUERY_TIMEOUT = 30 + + +def test_prometheus_api_corpus( + signoz: types.SigNoz, + create_user_admin: None, # pylint: disable=unused-argument + get_token: Callable[[str, str], str], + insert_metrics: Callable[[list[Metrics]], None], +) -> None: + corpus, bases = ingest_promqltest_corpus(insert_metrics) + token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD) + + failures: list[str] = [] + for case in corpus["cases"]: + # instant-coarse variants encode an instant eval as a coarse-step + # range because the v5 API cannot run true instants. This API can: + # the [base] form of the same eval goes through /query below, and the + # transpiled coarse-step serving the encoding exercises is covered + # (and its known divergences ledgered) by promqlconformance's + # clickhousev2 leg. + if case["variant"] == "instant-coarse": + continue + + base = bases[case["dataset"]] + start_ms = base + case["start_ms"] + end_ms = base + case["end_ms"] + step_s = max(1, case["step_ms"] // 1000) + case_id = f"{case['source']}[{case['variant']}]" + + if case["instant"]: + path, params = "/prometheus/api/v1/query", {"query": case["expr"], "time": end_ms / 1000} + else: + path, params = ( + "/prometheus/api/v1/query_range", + { + "query": case["expr"], + "start": start_ms / 1000, + "end": end_ms / 1000, + "step": step_s, + }, + ) + response = requests.get( + signoz.self.host_configs["8080"].get(path), + params=params, + timeout=QUERY_TIMEOUT, + headers={"authorization": f"Bearer {token}"}, + ) + if response.status_code != HTTPStatus.OK: + failures.append(f"{case_id}: HTTP {response.status_code} for {case['expr']!r}: {response.text[:200]}") + continue + body = response.json() + if body.get("status") != "success": + failures.append(f"{case_id}: status {body.get('status')!r} for {case['expr']!r}: {json.dumps(body)[:200]}") + continue + + result_type, result = body["data"]["resultType"], body["data"]["result"] + actual: dict[tuple, dict[int, float]] = {} + if result_type == "matrix": + for series in result: + points = {round(float(ts) * 1000): SPECIALS[v] if v in SPECIALS else float(v) for ts, v in series.get("values") or []} + actual[tuple(sorted((series.get("metric") or {}).items()))] = points + elif result_type == "vector": + for series in result: + ts, v = series["value"] + actual[tuple(sorted((series.get("metric") or {}).items()))] = {round(float(ts) * 1000): SPECIALS[v] if v in SPECIALS else float(v)} + elif result_type == "scalar": + ts, v = result + actual[()] = {round(float(ts) * 1000): SPECIALS[v] if v in SPECIALS else float(v)} + + expected: dict[tuple, dict[int, float]] = {} + for res in case["expected"]: + points = {base + off_ms: SPECIALS[v] if isinstance(v, str) else float(v) for off_ms, v in res["points"]} + expected[tuple(sorted(res["labels"].items()))] = points + + if set(actual) != set(expected): + missing = set(expected) - set(actual) + extra = set(actual) - set(expected) + failures.append(f"{case_id}: series mismatch for {case['expr']!r} (missing={sorted(missing)[:3]} extra={sorted(extra)[:3]})") + continue + + mismatch = None + for lset, exp_points in expected.items(): + act_points = actual[lset] + if set(act_points) != set(exp_points): + mismatch = f"{case_id}: timestamp mismatch for {case['expr']!r} series {dict(lset)} (expected {len(exp_points)} points, got {len(act_points)})" + break + for ts, exp_v in exp_points.items(): + act_v = act_points[ts] + if math.isnan(act_v) or math.isnan(exp_v): + close = math.isnan(act_v) and math.isnan(exp_v) + elif math.isinf(act_v) or math.isinf(exp_v): + close = act_v == exp_v + elif act_v == exp_v: + close = True + else: + # Expected values carry the v5 API's rounding (>=1: three + # decimal places; <1: three significant digits); this API + # returns raw floats. One rounding quantum covers the + # largest possible rounding difference. + scale = max(abs(act_v), abs(exp_v)) + if scale >= 1: + quantum = max(1e-3, scale * 1e-9) + else: + quantum = 10 ** (math.floor(math.log10(scale)) - 2) + close = abs(act_v - exp_v) <= quantum + 1e-12 + if not close: + mismatch = f"{case_id}: value mismatch for {case['expr']!r} series {dict(lset)} at {ts}: expected {exp_v}, got {act_v}" + break + if mismatch: + break + if mismatch: + failures.append(mismatch) + + for f_line in failures: + print("DIVERGED", f_line) + assert not failures, f"{len(failures)} corpus cases diverged:\n" + "\n".join(failures[:25]) diff --git a/tests/integration/tests/promapiconformance/conftest.py b/tests/integration/tests/promapiconformance/conftest.py new file mode 100644 index 00000000000..d60097bb503 --- /dev/null +++ b/tests/integration/tests/promapiconformance/conftest.py @@ -0,0 +1,37 @@ +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_promapi_v2( + network: Network, + migrator: types.Operation, # pylint: disable=unused-argument + zeus: types.TestContainerDocker, + gateway: types.TestContainerDocker, + sqlstore: types.TestContainerSQL, + clickhouse: types.TestContainerClickhouse, + request: pytest.FixtureRequest, + pytestconfig: pytest.Config, +) -> types.SigNoz: + """ + SigNoz with clickhousev2 as the serving prometheus provider. The corpus + replays against the /prometheus/api/v1 endpoints, so this package covers + the two paths nothing else serves: v2 as the provider (range queries + transpile when the shape allows), and the Prometheus HTTP API contract. + """ + return create_signoz( + network=network, + zeus=zeus, + gateway=gateway, + sqlstore=sqlstore, + clickhouse=clickhouse, + request=request, + pytestconfig=pytestconfig, + cache_key="signoz-promapi-v2", + env_overrides={ + "SIGNOZ_PROMETHEUS_PROVIDER": "clickhousev2", + }, + ) diff --git a/tests/integration/tests/promqlconformance/01_upstream_corpus.py b/tests/integration/tests/promqlconformance/01_upstream_corpus.py index 5311c3b97e8..d31c7cbf3b5 100644 --- a/tests/integration/tests/promqlconformance/01_upstream_corpus.py +++ b/tests/integration/tests/promqlconformance/01_upstream_corpus.py @@ -2,21 +2,21 @@ import math import os from collections.abc import Callable -from datetime import UTC, datetime, timedelta from http import HTTPStatus from fixtures import types from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD from fixtures.metrics import Metrics +from fixtures.promqltestcorpus import ingest_promqltest_corpus from fixtures.querier import get_all_series, make_query_request TESTDATA_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "testdata") -# Frozen corpus extracted from Prometheus' own promql/promqltest testdata by -# scripts/promqltestcorpus (upstream load scripts + the vendored reference engine). -# Unlike live-vs-live parity suites, the oracle is this committed file, so the suite -# keeps working when the serving path itself is the thing being changed — the one -# situation where comparing two live paths against each other is blind. -CORPUS_FILE = os.path.join(TESTDATA_DIR, "promqltestcorpus", "corpus.json") +# The corpus (see fixtures/promqltestcorpus.py) is frozen from Prometheus' own +# promql/promqltest testdata by scripts/promqltestcorpus (upstream load scripts +# + the vendored reference engine). Unlike live-vs-live parity suites, the +# oracle is a committed file, so the suite keeps working when the serving path +# itself is the thing being changed — the one situation where comparing two +# live paths against each other is blind. # One ledger per leg, enforced exactly in both directions. The default leg's # ledger is empty and pinned there; the clickhousev2 ledger is the rollout @@ -40,9 +40,6 @@ ("clickhousev2", {"X-SigNoz-PromQL-Provider": "clickhousev2"}), ] -# Datasets sit on disjoint time windows (2h gaps, far beyond the 5m lookback) so -# one bulk ingest serves every case without cross-talk. -ISOLATION_GAP_MS = 2 * 3600 * 1000 SPECIALS = {"NaN": math.nan, "Inf": math.inf, "-Inf": -math.inf} @@ -52,51 +49,7 @@ def test_upstream_promqltest_corpus( get_token: Callable[[str, str], str], insert_metrics: Callable[[list[Metrics]], None], ) -> None: - with open(CORPUS_FILE, encoding="utf-8") as f: - corpus = json.load(f) - - cases_by_dataset: dict[int, list[dict]] = {} - for case in corpus["cases"]: - cases_by_dataset.setdefault(case["dataset"], []).append(case) - - # Lay datasets end to end on the timeline, newest last, ending safely in - # the past; spans are per-dataset so the whole corpus stays within days. - spans = {} - for ds in corpus["datasets"]: - sample_max = max((s["samples"][-1][0] for s in ds["series"] if s["samples"]), default=0) - case_max = max((c["end_ms"] for c in cases_by_dataset.get(ds["id"], [])), default=0) - spans[ds["id"]] = max(sample_max, case_max) + corpus["meta"]["lookback_ms"] - - # Hour-aligned dataset bases: registration rows are hour-bucketed, so - # behavior depends on where samples fall relative to hour boundaries — - # the exact known-divergences enforcement needs that identical every run. - hour_ms = 3_600_000 - advances = {ds["id"]: -(-(spans[ds["id"]] + ISOLATION_GAP_MS) // hour_ms) * hour_ms for ds in corpus["datasets"]} - total = sum(advances.values()) - now = datetime.now(tz=UTC).replace(second=0, microsecond=0) - cursor = (int((now - timedelta(hours=1)).timestamp() * 1000) - total) // hour_ms * hour_ms - - bases: dict[int, int] = {} - metrics: list[Metrics] = [] - for ds in corpus["datasets"]: - bases[ds["id"]] = cursor - for series in ds["series"]: - labels = dict(series["labels"]) - metric_name = labels.pop("__name__") - for off_ms, raw in series["samples"]: - stale = raw == "stale" - metrics.append( - Metrics( - metric_name=metric_name, - labels=labels, - timestamp=datetime.fromtimestamp((cursor + off_ms) / 1000, tz=UTC), - value=0.0 if stale else (SPECIALS[raw] if isinstance(raw, str) else float(raw)), - flags=1 if stale else 0, - ) - ) - cursor += advances[ds["id"]] - - insert_metrics(metrics) + corpus, bases = ingest_promqltest_corpus(insert_metrics) token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD) failures: dict[str, list[str]] = {leg: [] for leg, _ in LEGS}