mirror of
https://github.com/SigNoz/signoz.git
synced 2026-10-02 00:00:42 +01:00
Compare commits
3 Commits
main
...
nv/promql-
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3e950fb7cc | ||
|
|
55763eded7 | ||
|
|
7021fe71d9 |
@@ -381,6 +381,12 @@ func toMatrix(series []transpiledSeries, startMs, stepMs int64) promql.Matrix {
|
||||
// lookback cannot resurrect the previous grid point. Each unit's synthetic
|
||||
// samples sit on its own grid: the query grid, or the subquery grid for
|
||||
// units inside subqueries.
|
||||
//
|
||||
// Synthetic series carry no __name__: substituted units all drop it, so
|
||||
// nameless matches the replaced expressions' output. A stamped name splits
|
||||
// or arms into per-unit series that a later name drop collides into the
|
||||
// duplicate-labelset error. hybridQuerier.Select resolves the selector
|
||||
// from the matcher, not from series labels.
|
||||
func (e *executor) executeHybrid(ctx context.Context, plan *transpilePlan, results [][]transpiledSeries) (promql.Matrix, error) {
|
||||
synthetic := make(map[string][]*series, len(plan.units))
|
||||
staleMarker := math.Float64frombits(promValue.StaleNaN)
|
||||
@@ -395,9 +401,7 @@ func (e *executor) executeHybrid(ctx context.Context, plan *transpilePlan, resul
|
||||
}
|
||||
list := make([]*series, 0, len(results[i]))
|
||||
for _, cs := range results[i] {
|
||||
builder := labels.NewBuilder(cs.lset)
|
||||
builder.Set(metricNameLabel, unit.name)
|
||||
s := &series{lset: builder.Labels()}
|
||||
s := &series{lset: cs.lset}
|
||||
s.ts = make([]int64, 0, gridLen)
|
||||
s.vs = make([]float64, 0, gridLen)
|
||||
for idx := 0; idx < gridLen; idx++ {
|
||||
@@ -440,65 +444,17 @@ func (e *executor) executeHybrid(ctx context.Context, plan *transpilePlan, resul
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Deep-copy before Close returns the result's slices to the engine pool,
|
||||
// and drop the synthetic __name__ that filter comparisons preserve.
|
||||
// Deep-copy before Close returns the result's slices to the engine pool.
|
||||
out := make(promql.Matrix, 0, len(matrix))
|
||||
for _, s := range matrix {
|
||||
lset := s.Metric
|
||||
if name := lset.Get(metricNameLabel); len(name) >= len(syntheticNamePrefix) && name[:len(syntheticNamePrefix)] == syntheticNamePrefix {
|
||||
builder := labels.NewBuilder(lset)
|
||||
builder.Del(metricNameLabel)
|
||||
lset = builder.Labels()
|
||||
}
|
||||
floats := make([]promql.FPoint, len(s.Floats))
|
||||
copy(floats, s.Floats)
|
||||
out = append(out, promql.Series{Metric: lset.Copy(), Floats: floats})
|
||||
}
|
||||
// The strip can leave twins: two units' outputs that only their
|
||||
// synthetic names told apart (e.g. -metric_a or -metric_b, both {}
|
||||
// once real names are dropped). The engine assembles its matrix by
|
||||
// labelset. It merges such temporally-disjoint elements into one
|
||||
// series. Reproduce that, with its duplicate error on same-timestamp
|
||||
// overlap.
|
||||
out, err = mergeMatrixByLabelset(out)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
out = append(out, promql.Series{Metric: s.Metric.Copy(), Floats: floats})
|
||||
}
|
||||
sort.Slice(out, func(i, j int) bool { return labels.Compare(out[i].Metric, out[j].Metric) < 0 })
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// mergeMatrixByLabelset merges series that share a labelset. It interleaves
|
||||
// their points in timestamp order. A timestamp present in both is the
|
||||
// engine's duplicate-labelset error.
|
||||
func mergeMatrixByLabelset(matrix promql.Matrix) (promql.Matrix, error) {
|
||||
index := make(map[uint64]int, len(matrix))
|
||||
out := matrix[:0]
|
||||
for _, s := range matrix {
|
||||
hash := s.Metric.Hash()
|
||||
idx, ok := index[hash]
|
||||
if ok && labels.Equal(out[idx].Metric, s.Metric) {
|
||||
merged := make([]promql.FPoint, 0, len(out[idx].Floats)+len(s.Floats))
|
||||
a, b := out[idx].Floats, s.Floats
|
||||
for len(a) > 0 && len(b) > 0 {
|
||||
switch {
|
||||
case a[0].T < b[0].T:
|
||||
merged, a = append(merged, a[0]), a[1:]
|
||||
case b[0].T < a[0].T:
|
||||
merged, b = append(merged, b[0]), b[1:]
|
||||
default:
|
||||
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "vector cannot contain metrics with the same labelset")
|
||||
}
|
||||
}
|
||||
out[idx].Floats = append(append(merged, a...), b...)
|
||||
continue
|
||||
}
|
||||
index[hash] = len(out)
|
||||
out = append(out, s)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func resultToMatrix(res *promql.Result) (promql.Matrix, error) {
|
||||
switch v := res.Value.(type) {
|
||||
case promql.Matrix:
|
||||
|
||||
@@ -14,7 +14,6 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
|
||||
"github.com/prometheus/prometheus/model/labels"
|
||||
"github.com/prometheus/prometheus/promql"
|
||||
"github.com/prometheus/prometheus/promql/parser"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -679,26 +678,3 @@ func TestMergeSameLabelsetSeries(t *testing.T) {
|
||||
require.Error(t, err, "two values on one evaluation timestamp is the engine's duplicate error")
|
||||
assert.True(t, errors.Ast(err, errors.TypeInvalidInput))
|
||||
}
|
||||
|
||||
// Hybrid twin case: stripping the synthetic __name__ can leave two engine
|
||||
// output series distinguishable only by those names (-metric_a or -metric_b:
|
||||
// both {} once real names are dropped). Pinned by conformance cases
|
||||
// name_label_dropping.test:137 and operators.test:1016.
|
||||
func TestMergeMatrixByLabelset(t *testing.T) {
|
||||
empty := labels.EmptyLabels()
|
||||
|
||||
out, err := mergeMatrixByLabelset(promql.Matrix{
|
||||
{Metric: empty, Floats: []promql.FPoint{{T: 0, F: -1}}},
|
||||
{Metric: empty, Floats: []promql.FPoint{{T: 600_000, F: -4}}},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, out, 1)
|
||||
assert.Equal(t, []promql.FPoint{{T: 0, F: -1}, {T: 600_000, F: -4}}, out[0].Floats)
|
||||
|
||||
_, err = mergeMatrixByLabelset(promql.Matrix{
|
||||
{Metric: empty, Floats: []promql.FPoint{{T: 0, F: -1}}},
|
||||
{Metric: empty, Floats: []promql.FPoint{{T: 0, F: -3}}},
|
||||
})
|
||||
require.Error(t, err)
|
||||
assert.True(t, errors.Ast(err, errors.TypeInvalidInput))
|
||||
}
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
from collections.abc import Callable
|
||||
from datetime import UTC, datetime
|
||||
from http import HTTPStatus
|
||||
from uuid import uuid4
|
||||
|
||||
import requests
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.metrics import Metrics
|
||||
from fixtures.querier import get_all_series, make_query_request
|
||||
|
||||
MINUTE_MS = 60_000
|
||||
QUERY_TIMEOUT = 30
|
||||
|
||||
|
||||
# Valid PromQL: the [30m:5m] subquery evaluates the or expression every 5m,
|
||||
# yielding sum(flicker) while the flicker metric has data and sum(steady)
|
||||
# after it stops. sum() drops __name__ from both, so every evaluation shares
|
||||
# one labelset and the result is a single series. query_range must return it —
|
||||
# not a "vector cannot contain metrics with the same labelset" error — and
|
||||
# must agree with /prometheus/api/v1/query on the same expression.
|
||||
def test_or_arms_merge_under_subquery_name_drop(
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
insert_metrics: Callable[[list[Metrics]], None],
|
||||
) -> None:
|
||||
now_ms = int(datetime.now(tz=UTC).timestamp() * 1000)
|
||||
base_ms = (now_ms // (5 * MINUTE_MS)) * (5 * MINUTE_MS) - 45 * MINUTE_MS
|
||||
|
||||
flicker = f"or_flicker_arm_{uuid4().hex[:8]}"
|
||||
steady = f"or_steady_arm_{uuid4().hex[:8]}"
|
||||
insert_metrics(
|
||||
[
|
||||
Metrics(
|
||||
metric_name=name,
|
||||
labels={"host": "server-01"},
|
||||
timestamp=datetime.fromtimestamp((base_ms + minute * MINUTE_MS) / 1000, tz=UTC),
|
||||
value=1.0,
|
||||
)
|
||||
# The flicker arm stops at minute 9, so every 30m subquery window
|
||||
# below sees it present at some 5m-aligned steps and absent (past
|
||||
# lookback) at others, with the steady arm filling the gaps.
|
||||
for name, minutes in ((flicker, range(10)), (steady, range(36)))
|
||||
for minute in minutes
|
||||
]
|
||||
)
|
||||
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
query = f"max_over_time((sum({flicker}) or sum({steady}))[30m:5m])"
|
||||
start_ms = base_ms + 20 * MINUTE_MS
|
||||
end_ms = base_ms + 35 * MINUTE_MS
|
||||
|
||||
# A range query is semantically the instant evaluation of the same
|
||||
# expression at each grid timestamp. Each one must yield the single
|
||||
# nameless series; together they are the reference for query_range.
|
||||
instant_values: dict[int, float] = {}
|
||||
for ts_ms in range(start_ms, end_ms + 1, 5 * MINUTE_MS):
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/prometheus/api/v1/query"),
|
||||
params={"query": query, "time": ts_ms / 1000},
|
||||
timeout=QUERY_TIMEOUT,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text[:300]
|
||||
body = response.json()
|
||||
assert body.get("status") == "success", body
|
||||
result = body["data"]["result"]
|
||||
assert [series["metric"] for series in result] == [{}], (ts_ms, result)
|
||||
instant_values[ts_ms] = float(result[0]["value"][1])
|
||||
assert set(instant_values.values()) == {1.0}, instant_values
|
||||
|
||||
# query_range over the same grid must agree point for point.
|
||||
spec = {"name": "A", "query": query, "step": 300}
|
||||
response = make_query_request(signoz, token, start_ms, end_ms, [{"type": "promql", "spec": spec}])
|
||||
assert response.status_code == HTTPStatus.OK, response.text[:300]
|
||||
series = get_all_series(response.json(), "A")
|
||||
assert len(series) == 1, f"both or arms must merge into one series: {series}"
|
||||
points = {point["timestamp"]: point["value"] for point in series[0].get("values") or []}
|
||||
assert points == instant_values, (points, instant_values)
|
||||
Reference in New Issue
Block a user