Compare commits

..

3 Commits

Author SHA1 Message Date
Srikanth Chekuri
3e950fb7cc Merge branch 'main' into nv/promql-vector-labelset-err 2026-10-02 02:08:56 +05:30
Naman Verma
55763eded7 fix(promql): serve transpiled series without their synthetic __name__ 2026-10-01 17:53:57 +05:30
Naman Verma
7021fe71d9 test: add test to demonstrate failure 2026-10-01 15:08:04 +05:30
3 changed files with 90 additions and 77 deletions

View File

@@ -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:

View File

@@ -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))
}

View File

@@ -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)