mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-30 16:20:27 +01:00
Compare commits
4 Commits
chore/dash
...
promql-dro
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2c2c982ee0 | ||
|
|
c320fdbe63 | ||
|
|
e248419a27 | ||
|
|
a077ec03f4 |
1
.github/workflows/integrationci.yaml
vendored
1
.github/workflows/integrationci.yaml
vendored
@@ -58,6 +58,7 @@ jobs:
|
||||
- querierai
|
||||
- rawexportdata
|
||||
- promqlconformance
|
||||
- promapiconformance
|
||||
- querierauthz
|
||||
- role
|
||||
- rootuser
|
||||
|
||||
@@ -299,8 +299,21 @@ 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`.
|
||||
|
||||
Both paths enforce fetch budgets
|
||||
(`prometheus::clickhousev2::max_fetched_series` and
|
||||
`::max_fetched_samples`; 0 disables). The engine path counts matched series
|
||||
and scanned samples. The transpiled path counts buffered grid cells (series
|
||||
times grid width) across a plan's units, because transpiled results never
|
||||
pass the engine's sample limiter. A refusal is a typed invalid-input error.
|
||||
It pierces the engine's `promql.ErrStorage` wrapper
|
||||
(`prometheus.TypedStorageError`), so the APIs report a user error, not an
|
||||
internal one.
|
||||
|
||||
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 +328,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
|
||||
|
||||
43
pkg/instrumentation/promqlspanfilter.go
Normal file
43
pkg/instrumentation/promqlspanfilter.go
Normal file
@@ -0,0 +1,43 @@
|
||||
package instrumentation
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.opentelemetry.io/otel/trace/embedded"
|
||||
tracenoop "go.opentelemetry.io/otel/trace/noop"
|
||||
)
|
||||
|
||||
// promqlSpanFilter drops the promql engine's spans. The engine traces every
|
||||
// evaluation through the global tracer provider under an unnamed scope: one
|
||||
// span per timer plus one per AST node per query ("promqlExec",
|
||||
// "promqlInnerEval eval *promql.BinaryExpr", ...), which floods traces
|
||||
// without adding value. Filtered spans return a non-recording span that
|
||||
// keeps the parent's span context, so descendants (the ClickHouse query
|
||||
// spans) still attach to the surrounding span.
|
||||
type promqlSpanFilter struct {
|
||||
embedded.TracerProvider
|
||||
delegate trace.TracerProvider
|
||||
}
|
||||
|
||||
func (p promqlSpanFilter) Tracer(name string, opts ...trace.TracerOption) trace.Tracer {
|
||||
tracer := p.delegate.Tracer(name, opts...)
|
||||
if name != "" {
|
||||
return tracer
|
||||
}
|
||||
return promqlSpanFilterTracer{delegate: tracer, noop: tracenoop.NewTracerProvider().Tracer("")}
|
||||
}
|
||||
|
||||
type promqlSpanFilterTracer struct {
|
||||
embedded.Tracer
|
||||
delegate trace.Tracer
|
||||
noop trace.Tracer
|
||||
}
|
||||
|
||||
func (t promqlSpanFilterTracer) Start(ctx context.Context, spanName string, opts ...trace.SpanStartOption) (context.Context, trace.Span) {
|
||||
if strings.HasPrefix(spanName, "promql") {
|
||||
return t.noop.Start(ctx, spanName)
|
||||
}
|
||||
return t.delegate.Start(ctx, spanName, opts...)
|
||||
}
|
||||
40
pkg/instrumentation/promqlspanfilter_test.go
Normal file
40
pkg/instrumentation/promqlspanfilter_test.go
Normal file
@@ -0,0 +1,40 @@
|
||||
package instrumentation
|
||||
|
||||
import (
|
||||
"context"
|
||||
"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"
|
||||
)
|
||||
|
||||
func TestPromqlSpanFilter(t *testing.T) {
|
||||
recorder := tracetest.NewSpanRecorder()
|
||||
provider := promqlSpanFilter{delegate: sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(recorder))}
|
||||
|
||||
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")
|
||||
}
|
||||
@@ -108,7 +108,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(promqlSpanFilter{delegate: sdk.TracerProvider()})
|
||||
|
||||
return &SDK{
|
||||
sdk: sdk,
|
||||
|
||||
@@ -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 != "" {
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"math"
|
||||
"slices"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/prometheus"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
@@ -29,6 +30,7 @@ type client struct {
|
||||
settings factory.ScopedProviderSettings
|
||||
telemetryStore telemetrystore.TelemetryStore
|
||||
lookbackMs int64
|
||||
cfg prometheus.ClickhouseV2Config
|
||||
}
|
||||
|
||||
func newClient(settings factory.ScopedProviderSettings, telemetryStore telemetrystore.TelemetryStore, cfg prometheus.Config) *client {
|
||||
@@ -41,6 +43,7 @@ func newClient(settings factory.ScopedProviderSettings, telemetryStore telemetry
|
||||
settings: settings,
|
||||
telemetryStore: telemetryStore,
|
||||
lookbackMs: lookback.Milliseconds(),
|
||||
cfg: cfg.ClickhouseV2,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -77,6 +80,13 @@ func (c *client) selectSeries(ctx context.Context, query string, args []any) (*s
|
||||
if name := lset.Get(metricNameLabel); name != "" {
|
||||
names[name] = struct{}{}
|
||||
}
|
||||
if c.cfg.MaxFetchedSeries > 0 && len(lookup.fingerprints) > c.cfg.MaxFetchedSeries {
|
||||
return nil, errors.NewInvalidInputf(
|
||||
errors.CodeInvalidInput,
|
||||
"promql selector matched more than %d series; narrow the label matchers or raise prometheus::clickhousev2::max_fetched_series",
|
||||
c.cfg.MaxFetchedSeries,
|
||||
)
|
||||
}
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, err
|
||||
@@ -136,6 +146,8 @@ func (c *client) selectSamples(ctx context.Context, query string, args []any, lo
|
||||
first = true
|
||||
haveCurrent bool
|
||||
staleMarker = math.Float64frombits(promValue.StaleNaN)
|
||||
maxSamples = c.cfg.MaxFetchedSamples
|
||||
fetched int64
|
||||
unknownCount int
|
||||
)
|
||||
|
||||
@@ -144,6 +156,15 @@ func (c *client) selectSamples(ctx context.Context, query string, args []any, lo
|
||||
return nil, err
|
||||
}
|
||||
|
||||
fetched++
|
||||
if maxSamples > 0 && fetched > maxSamples {
|
||||
return nil, errors.NewInvalidInputf(
|
||||
errors.CodeInvalidInput,
|
||||
"promql query would fetch more than %d samples; narrow the selector or time range, or raise prometheus::clickhousev2::max_fetched_samples",
|
||||
maxSamples,
|
||||
)
|
||||
}
|
||||
|
||||
if first || fingerprint != prevFp {
|
||||
first = false
|
||||
prevFp = fingerprint
|
||||
|
||||
49
pkg/prometheus/clickhouseprometheusv2/client_test.go
Normal file
49
pkg/prometheus/clickhouseprometheusv2/client_test.go
Normal file
@@ -0,0 +1,49 @@
|
||||
package clickhouseprometheusv2
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
cmock "github.com/SigNoz/clickhouse-go-mock"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/prometheus/prometheus/model/labels"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
var samplesCols = []cmock.ColumnType{
|
||||
{Name: "fingerprint", Type: "UInt64"},
|
||||
{Name: "unix_milli", Type: "Int64"},
|
||||
{Name: "value", Type: "Float64"},
|
||||
{Name: "flags", Type: "UInt32"},
|
||||
}
|
||||
|
||||
func TestSelectSeriesBudget(t *testing.T) {
|
||||
c, store := newTestClient(t)
|
||||
c.cfg.MaxFetchedSeries = 1
|
||||
|
||||
store.Mock().ExpectQuery("SELECT fingerprint, any\\(labels\\)").WillReturnRows(cmock.NewRows(seriesCols, [][]any{
|
||||
{uint64(1), `{"__name__":"up","instance":"a"}`},
|
||||
{uint64(2), `{"__name__":"up","instance":"b"}`},
|
||||
}))
|
||||
|
||||
_, err := c.selectSeries(context.Background(), "SELECT fingerprint, any(labels) FROM t", nil)
|
||||
require.Error(t, err)
|
||||
assert.True(t, errors.Ast(err, errors.TypeInvalidInput), "budget refusal must be typed invalid input, got %v", err)
|
||||
}
|
||||
|
||||
func TestSelectSamplesBudget(t *testing.T) {
|
||||
c, store := newTestClient(t)
|
||||
c.cfg.MaxFetchedSamples = 2
|
||||
|
||||
store.Mock().ExpectQuery("SELECT fingerprint, unix_milli").WillReturnRows(cmock.NewRows(samplesCols, [][]any{
|
||||
{uint64(1), int64(1_700_000_000_000), 1.0, uint32(0)},
|
||||
{uint64(1), int64(1_700_000_060_000), 2.0, uint32(0)},
|
||||
{uint64(1), int64(1_700_000_120_000), 3.0, uint32(0)},
|
||||
}))
|
||||
|
||||
lookup := &seriesLookup{fingerprints: map[uint64]labels.Labels{1: labels.FromStrings("__name__", "up")}}
|
||||
_, err := c.selectSamples(context.Background(), "SELECT fingerprint, unix_milli, value, flags FROM t", nil, lookup)
|
||||
require.Error(t, err)
|
||||
assert.True(t, errors.Ast(err, errors.TypeInvalidInput), "budget refusal must be typed invalid input, got %v", err)
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"encoding/json"
|
||||
"math"
|
||||
"sort"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
@@ -88,12 +89,17 @@ 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). The units share one grid-cell budget:
|
||||
// transpiled results never pass through the engine's sample limiter, so
|
||||
// without it a large series-count x grid-width query would buffer
|
||||
// unbounded arrays — the OOM this provider exists to prevent.
|
||||
results := make([][]transpiledSeries, len(plan.units))
|
||||
var gridCells atomic.Int64
|
||||
eg, egCtx := errgroup.WithContext(ctx)
|
||||
for i, unit := range plan.units {
|
||||
eg.Go(func() error {
|
||||
res, err := e.executeUnit(egCtx, &unit.core, unit.grid)
|
||||
res, err := e.executeUnit(egCtx, &unit.core, unit.grid, &gridCells)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -133,7 +139,7 @@ type transpiledSeries struct {
|
||||
values []*float64
|
||||
}
|
||||
|
||||
func (e *executor) executeUnit(ctx context.Context, unit *coreUnit, grid gridContext) ([]transpiledSeries, error) {
|
||||
func (e *executor) executeUnit(ctx context.Context, unit *coreUnit, grid gridContext, gridCells *atomic.Int64) ([]transpiledSeries, error) {
|
||||
startMs, endMs, stepMs := grid.startMs, grid.endMs, grid.stepMs
|
||||
windowMs := unit.rangeMs
|
||||
if unit.kind == unitInstant {
|
||||
@@ -142,19 +148,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
|
||||
}
|
||||
@@ -190,6 +204,17 @@ func (e *executor) executeUnit(ctx context.Context, unit *coreUnit, grid gridCon
|
||||
if err := rows.Scan(targets...); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// One row buffers one grid array; series count times grid width is
|
||||
// the transpiled equivalent of fetched samples. Counted per row as
|
||||
// the arrays accumulate: without a series lookup there is no series
|
||||
// count to charge up front.
|
||||
if maxSamples := e.client.cfg.MaxFetchedSamples; maxSamples > 0 && gridCells.Add(int64(len(gridValues))) > maxSamples {
|
||||
return nil, errors.NewInvalidInputf(
|
||||
errors.CodeInvalidInput,
|
||||
"promql query would buffer more than %d output points; narrow the selector or time range, or raise prometheus::clickhousev2::max_fetched_samples",
|
||||
maxSamples,
|
||||
)
|
||||
}
|
||||
var lset labels.Labels
|
||||
if keyNames != nil {
|
||||
builder := labels.NewScratchBuilder(len(keyNames))
|
||||
|
||||
@@ -2,6 +2,7 @@ package clickhouseprometheusv2
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -32,6 +33,17 @@ var seriesCols = []cmock.ColumnType{
|
||||
{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 {
|
||||
t.Helper()
|
||||
expr, err := parser.NewParser(parser.Options{}).ParseExpr(q)
|
||||
@@ -534,6 +546,30 @@ func TestDisjointWindowLattice(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Transpiled results never pass the engine's sample limiter, so the grid
|
||||
// cells (series x grid width) must be budgeted as the arrays accumulate —
|
||||
// otherwise a wide query rebuilds the OOM this provider exists to prevent.
|
||||
func TestExecuteUnit_GridCellBudget(t *testing.T) {
|
||||
c, store := newTestClient(t)
|
||||
c.cfg.MaxFetchedSamples = 100
|
||||
e := &executor{client: c, parser: prometheus.NewParser()}
|
||||
|
||||
grid61 := make([]*float64, 61)
|
||||
store.Mock().ExpectQuery("timeSeriesRateToGrid").WithArgs(anyArgs(7)...).WillReturnRows(cmock.NewRows(
|
||||
[]cmock.ColumnType{{Name: "g0", Type: "String"}, {Name: "grid", Type: "Array(Nullable(Float64))"}},
|
||||
[][]any{{"api", grid61}, {"web", grid61}},
|
||||
))
|
||||
|
||||
plan, ok := classify(parse(t, `sum by (job) (rate(up[5m]))`), gridContext{startMs: 1_700_000_000_000, endMs: 1_700_003_600_000, stepMs: 60_000})
|
||||
require.True(t, ok)
|
||||
|
||||
// 2 series x 61 grid points = 122 cells > 100.
|
||||
var cells atomic.Int64
|
||||
_, err := e.executeUnit(context.Background(), &plan.units[0].core, plan.units[0].grid, &cells)
|
||||
require.Error(t, err)
|
||||
assert.True(t, errors.Ast(err, errors.TypeInvalidInput), "budget refusal must be typed invalid input, got %v", err)
|
||||
}
|
||||
|
||||
func TestTryExecuteRange_WindowedGateFallsBack(t *testing.T) {
|
||||
c, store := newTestClient(t)
|
||||
e := &executor{client: c, parser: prometheus.NewParser()}
|
||||
@@ -553,7 +589,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 +673,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")
|
||||
|
||||
@@ -13,6 +13,16 @@ type ActiveQueryTrackerConfig struct {
|
||||
MaxConcurrent int `mapstructure:"max_concurrent"`
|
||||
}
|
||||
|
||||
type ClickhouseV2Config struct {
|
||||
// MaxFetchedSeries caps the series one selector may match; 0 disables
|
||||
// the cap.
|
||||
MaxFetchedSeries int `mapstructure:"max_fetched_series"`
|
||||
|
||||
// MaxFetchedSamples caps the samples (engine path) or buffered grid
|
||||
// cells (transpiled path) one query may fetch; 0 disables the cap.
|
||||
MaxFetchedSamples int64 `mapstructure:"max_fetched_samples"`
|
||||
}
|
||||
|
||||
type Config struct {
|
||||
ActiveQueryTrackerConfig ActiveQueryTrackerConfig `mapstructure:"active_query_tracker"`
|
||||
|
||||
@@ -28,6 +38,9 @@ type Config struct {
|
||||
// ProviderName selects the storage provider: "clickhouse" (default) or
|
||||
// "clickhousev2".
|
||||
ProviderName string `mapstructure:"provider"`
|
||||
|
||||
// ClickhouseV2 configures the clickhousev2 provider.
|
||||
ClickhouseV2 ClickhouseV2Config `mapstructure:"clickhousev2"`
|
||||
}
|
||||
|
||||
func NewConfigFactory() factory.ConfigFactory {
|
||||
@@ -43,6 +56,10 @@ func newConfig() factory.Config {
|
||||
},
|
||||
Timeout: 2 * time.Minute,
|
||||
ProviderName: "clickhouse",
|
||||
ClickhouseV2: ClickhouseV2Config{
|
||||
MaxFetchedSeries: 500_000,
|
||||
MaxFetchedSamples: 50_000_000,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,6 +70,9 @@ func (c Config) Validate() error {
|
||||
if c.ProviderName != "" && c.ProviderName != "clickhouse" && c.ProviderName != "clickhousev2" {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "prometheus::provider must be one of [clickhouse, clickhousev2], got %q", c.ProviderName)
|
||||
}
|
||||
if c.ClickhouseV2.MaxFetchedSeries < 0 || c.ClickhouseV2.MaxFetchedSamples < 0 {
|
||||
return errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "prometheus::clickhousev2 limits must not be negative")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
30
pkg/prometheus/errors.go
Normal file
30
pkg/prometheus/errors.go
Normal file
@@ -0,0 +1,30 @@
|
||||
package prometheus
|
||||
|
||||
import (
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/prometheus/prometheus/promql"
|
||||
)
|
||||
|
||||
// TypedStorageError walks an engine execution error chain looking for a
|
||||
// SigNoz-typed invalid-input error raised by the storage layer (the fetch
|
||||
// budget refusals). Every wrapper level is stepped through by hand: Ast is a
|
||||
// bare type cast, not an unwrap — it misses a typed error behind the
|
||||
// engine's "expanding series: %w" — and promql.ErrStorage has no Unwrap
|
||||
// method at all, so a plain unwrap loop would stop at it.
|
||||
func TypedStorageError(execErr error) error {
|
||||
for e := execErr; e != nil; {
|
||||
if errors.Ast(e, errors.TypeInvalidInput) {
|
||||
return e
|
||||
}
|
||||
if es, ok := e.(promql.ErrStorage); ok {
|
||||
e = es.Err
|
||||
continue
|
||||
}
|
||||
u, ok := e.(interface{ Unwrap() error })
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
e = u.Unwrap()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
22
pkg/prometheus/errors_test.go
Normal file
22
pkg/prometheus/errors_test.go
Normal file
@@ -0,0 +1,22 @@
|
||||
package prometheus
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/prometheus/prometheus/promql"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestTypedStorageError(t *testing.T) {
|
||||
budget := errors.NewInvalidInputf(errors.CodeInvalidInput, "too many series")
|
||||
|
||||
// The engine wraps a storage error as expanding series: %w inside
|
||||
// promql.ErrStorage, which has no Unwrap method.
|
||||
wrapped := promql.ErrStorage{Err: fmt.Errorf("expanding series: %w", budget)}
|
||||
assert.Equal(t, budget, TypedStorageError(wrapped))
|
||||
|
||||
assert.Nil(t, TypedStorageError(promql.ErrStorage{Err: fmt.Errorf("connection refused")}))
|
||||
assert.Nil(t, TypedStorageError(nil))
|
||||
}
|
||||
9
pkg/prometheus/handler.go
Normal file
9
pkg/prometheus/handler.go
Normal file
@@ -0,0 +1,9 @@
|
||||
package prometheus
|
||||
|
||||
import "net/http"
|
||||
|
||||
type Handler interface {
|
||||
Query(http.ResponseWriter, *http.Request)
|
||||
|
||||
QueryRange(http.ResponseWriter, *http.Request)
|
||||
}
|
||||
266
pkg/prometheus/promapi/handler.go
Normal file
266
pkg/prometheus/promapi/handler.go
Normal file
@@ -0,0 +1,266 @@
|
||||
// Package promapi serves the Prometheus HTTP query API over a
|
||||
// prometheus.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.
|
||||
package promapi
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"math"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
promModel "github.com/prometheus/common/model"
|
||||
"github.com/prometheus/prometheus/promql"
|
||||
"github.com/prometheus/prometheus/promql/parser"
|
||||
"github.com/prometheus/prometheus/util/stats"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/prometheus"
|
||||
)
|
||||
|
||||
type handler struct {
|
||||
logger *slog.Logger
|
||||
prom prometheus.Prometheus
|
||||
}
|
||||
|
||||
func NewHandler(logger *slog.Logger, prom prometheus.Prometheus) prometheus.Handler {
|
||||
return &handler{logger: logger, prom: prom}
|
||||
}
|
||||
|
||||
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"`
|
||||
}
|
||||
|
||||
// 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.(prometheus.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:
|
||||
// A fetch-budget refusal is the storage-level twin of the
|
||||
// engine's own too-many-samples error, which upstream maps to
|
||||
// "execution", not "internal".
|
||||
if typed := prometheus.TypedStorageError(res.Err); typed != nil {
|
||||
h.respondError(ctx, w, errExec, typed)
|
||||
return
|
||||
}
|
||||
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
|
||||
}
|
||||
|
||||
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))
|
||||
}
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
@@ -43,6 +43,10 @@ var quotedMetricOutsideBracesPattern = regexp.MustCompile(`"([^"]+)"\s*\{`)
|
||||
// tryEnhancePromQLExecError attempts to convert a PromQL execution error into
|
||||
// a properly typed error. Returns nil if the error is not a recognized execution error.
|
||||
func tryEnhancePromQLExecError(execErr error) error {
|
||||
if typed := prometheus.TypedStorageError(execErr); typed != nil {
|
||||
return typed
|
||||
}
|
||||
|
||||
var eqc promql.ErrQueryCanceled
|
||||
var eqt promql.ErrQueryTimeout
|
||||
var es promql.ErrStorage
|
||||
|
||||
@@ -484,6 +484,9 @@ 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("/prometheus/api/v1/query_range", am.ViewAccess(aH.Signoz.Handlers.PrometheusHandler.QueryRange)).Methods(http.MethodGet, http.MethodPost)
|
||||
router.HandleFunc("/prometheus/api/v1/query", am.ViewAccess(aH.Signoz.Handlers.PrometheusHandler.Query)).Methods(http.MethodGet, http.MethodPost)
|
||||
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)
|
||||
|
||||
@@ -48,6 +48,8 @@ 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/prometheus/promapi"
|
||||
"github.com/SigNoz/signoz/pkg/querier"
|
||||
"github.com/SigNoz/signoz/pkg/ruler"
|
||||
"github.com/SigNoz/signoz/pkg/ruler/signozruler"
|
||||
@@ -81,6 +83,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
|
||||
@@ -101,6 +104,7 @@ func NewHandlers(
|
||||
zeusService zeus.Zeus,
|
||||
registryHandler factory.Handler,
|
||||
alertmanagerService alertmanager.Alertmanager,
|
||||
prometheusService prometheus.Prometheus,
|
||||
rulerService ruler.Ruler,
|
||||
statsAggregator statsreporter.Aggregator,
|
||||
) Handlers {
|
||||
@@ -129,6 +133,7 @@ func NewHandlers(
|
||||
CloudIntegrationHandler: implcloudintegration.NewHandler(modules.CloudIntegration),
|
||||
SpanMapperHandler: implspanmapper.NewHandler(modules.SpanMapper),
|
||||
AlertmanagerHandler: signozalertmanager.NewHandler(alertmanagerService),
|
||||
PrometheusHandler: promapi.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)
|
||||
|
||||
@@ -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