mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-31 16:40:44 +01:00
Compare commits
2 Commits
feat/licen
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4d015a927e | ||
|
|
abebea532b |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -58,6 +58,7 @@ jobs:
|
||||
- querierai
|
||||
- rawexportdata
|
||||
- promqlconformance
|
||||
- promapiconformance
|
||||
- querierauthz
|
||||
- role
|
||||
- rootuser
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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
|
||||
|
||||
396
frontend/src/api/generated/services/prometheus/index.ts
Normal file
396
frontend/src/api/generated/services/prometheus/index.ts
Normal file
@@ -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<PrometheusSuccessResponseSchemaDTO>({
|
||||
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<ReturnType<typeof prometheusQuery>>,
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
>(
|
||||
params: PrometheusQueryParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof prometheusQuery>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey = queryOptions?.queryKey ?? getPrometheusQueryQueryKey(params);
|
||||
|
||||
const queryFn: QueryFunction<Awaited<ReturnType<typeof prometheusQuery>>> = ({
|
||||
signal,
|
||||
}) => prometheusQuery(params, signal);
|
||||
|
||||
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof prometheusQuery>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type PrometheusQueryQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof prometheusQuery>>
|
||||
>;
|
||||
export type PrometheusQueryQueryError = ErrorType<
|
||||
PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO
|
||||
>;
|
||||
|
||||
/**
|
||||
* @summary Prometheus instant query
|
||||
*/
|
||||
|
||||
export function usePrometheusQuery<
|
||||
TData = Awaited<ReturnType<typeof prometheusQuery>>,
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
>(
|
||||
params: PrometheusQueryParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof prometheusQuery>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getPrometheusQueryQueryOptions(params, options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return { ...query, queryKey: queryOptions.queryKey };
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary Prometheus instant query
|
||||
*/
|
||||
export const invalidatePrometheusQuery = async (
|
||||
queryClient: QueryClient,
|
||||
params: PrometheusQueryParams,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
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<PrometheusSuccessResponseSchemaDTO>({
|
||||
url: `/prometheus/api/v1/query`,
|
||||
method: 'POST',
|
||||
params,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getPrometheusQueryPostMutationOptions = <
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryPost>>,
|
||||
TError,
|
||||
{ params: PrometheusQueryPostParams },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryPost>>,
|
||||
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<ReturnType<typeof prometheusQueryPost>>,
|
||||
{ params: PrometheusQueryPostParams }
|
||||
> = (props) => {
|
||||
const { params } = props ?? {};
|
||||
|
||||
return prometheusQueryPost(params);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type PrometheusQueryPostMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof prometheusQueryPost>>
|
||||
>;
|
||||
|
||||
export type PrometheusQueryPostMutationError = ErrorType<
|
||||
PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO
|
||||
>;
|
||||
|
||||
/**
|
||||
* @summary Prometheus instant query
|
||||
*/
|
||||
export const usePrometheusQueryPost = <
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryPost>>,
|
||||
TError,
|
||||
{ params: PrometheusQueryPostParams },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof prometheusQueryPost>>,
|
||||
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<PrometheusSuccessResponseSchemaDTO>({
|
||||
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<ReturnType<typeof prometheusQueryRange>>,
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
>(
|
||||
params: PrometheusQueryRangeParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryRange>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
) => {
|
||||
const { query: queryOptions } = options ?? {};
|
||||
|
||||
const queryKey =
|
||||
queryOptions?.queryKey ?? getPrometheusQueryRangeQueryKey(params);
|
||||
|
||||
const queryFn: QueryFunction<
|
||||
Awaited<ReturnType<typeof prometheusQueryRange>>
|
||||
> = ({ signal }) => prometheusQueryRange(params, signal);
|
||||
|
||||
return { queryKey, queryFn, ...queryOptions } as UseQueryOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryRange>>,
|
||||
TError,
|
||||
TData
|
||||
> & { queryKey: QueryKey };
|
||||
};
|
||||
|
||||
export type PrometheusQueryRangeQueryResult = NonNullable<
|
||||
Awaited<ReturnType<typeof prometheusQueryRange>>
|
||||
>;
|
||||
export type PrometheusQueryRangeQueryError = ErrorType<
|
||||
PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO
|
||||
>;
|
||||
|
||||
/**
|
||||
* @summary Prometheus range query
|
||||
*/
|
||||
|
||||
export function usePrometheusQueryRange<
|
||||
TData = Awaited<ReturnType<typeof prometheusQueryRange>>,
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
>(
|
||||
params: PrometheusQueryRangeParams,
|
||||
options?: {
|
||||
query?: UseQueryOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryRange>>,
|
||||
TError,
|
||||
TData
|
||||
>;
|
||||
},
|
||||
): UseQueryResult<TData, TError> & { queryKey: QueryKey } {
|
||||
const queryOptions = getPrometheusQueryRangeQueryOptions(params, options);
|
||||
|
||||
const query = useQuery(queryOptions) as UseQueryResult<TData, TError> & {
|
||||
queryKey: QueryKey;
|
||||
};
|
||||
|
||||
return { ...query, queryKey: queryOptions.queryKey };
|
||||
}
|
||||
|
||||
/**
|
||||
* @summary Prometheus range query
|
||||
*/
|
||||
export const invalidatePrometheusQueryRange = async (
|
||||
queryClient: QueryClient,
|
||||
params: PrometheusQueryRangeParams,
|
||||
options?: InvalidateOptions,
|
||||
): Promise<QueryClient> => {
|
||||
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<PrometheusSuccessResponseSchemaDTO>({
|
||||
url: `/prometheus/api/v1/query_range`,
|
||||
method: 'POST',
|
||||
params,
|
||||
signal,
|
||||
});
|
||||
};
|
||||
|
||||
export const getPrometheusQueryRangePostMutationOptions = <
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryRangePost>>,
|
||||
TError,
|
||||
{ params: PrometheusQueryRangePostParams },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryRangePost>>,
|
||||
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<ReturnType<typeof prometheusQueryRangePost>>,
|
||||
{ params: PrometheusQueryRangePostParams }
|
||||
> = (props) => {
|
||||
const { params } = props ?? {};
|
||||
|
||||
return prometheusQueryRangePost(params);
|
||||
};
|
||||
|
||||
return { mutationFn, ...mutationOptions };
|
||||
};
|
||||
|
||||
export type PrometheusQueryRangePostMutationResult = NonNullable<
|
||||
Awaited<ReturnType<typeof prometheusQueryRangePost>>
|
||||
>;
|
||||
|
||||
export type PrometheusQueryRangePostMutationError = ErrorType<
|
||||
PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO
|
||||
>;
|
||||
|
||||
/**
|
||||
* @summary Prometheus range query
|
||||
*/
|
||||
export const usePrometheusQueryRangePost = <
|
||||
TError = ErrorType<PrometheusErrorResponseSchemaDTO | RenderErrorResponseDTO>,
|
||||
TContext = unknown,
|
||||
>(options?: {
|
||||
mutation?: UseMutationOptions<
|
||||
Awaited<ReturnType<typeof prometheusQueryRangePost>>,
|
||||
TError,
|
||||
{ params: PrometheusQueryRangePostParams },
|
||||
TContext
|
||||
>;
|
||||
}): UseMutationResult<
|
||||
Awaited<ReturnType<typeof prometheusQueryRangePost>>,
|
||||
TError,
|
||||
{ params: PrometheusQueryRangePostParams },
|
||||
TContext
|
||||
> => {
|
||||
return useMutation(getPrometheusQueryRangePostMutationOptions(options));
|
||||
};
|
||||
@@ -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;
|
||||
};
|
||||
|
||||
102
pkg/apiserver/signozapiserver/prometheus.go
Normal file
102
pkg/apiserver/signozapiserver/prometheus.go
Normal file
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 != "" {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
|
||||
211
pkg/prometheus/handler.go
Normal file
211
pkg/prometheus/handler.go
Normal file
@@ -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)
|
||||
}
|
||||
164
pkg/prometheus/render.go
Normal file
164
pkg/prometheus/render.go
Normal file
@@ -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"`
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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 }{},
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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(
|
||||
|
||||
66
tests/fixtures/promqltestcorpus.py
vendored
Normal file
66
tests/fixtures/promqltestcorpus.py
vendored
Normal file
@@ -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
|
||||
@@ -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])
|
||||
37
tests/integration/tests/promapiconformance/conftest.py
Normal file
37
tests/integration/tests/promapiconformance/conftest.py
Normal file
@@ -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",
|
||||
},
|
||||
)
|
||||
@@ -2,21 +2,21 @@ import json
|
||||
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 @@ LEGS: list[tuple[str, dict | None]] = [
|
||||
("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}
|
||||
|
||||
Reference in New Issue
Block a user