diff --git a/pkg/engine/internal/planner/logical/planner.go b/pkg/engine/internal/planner/logical/planner.go index 27b022dea46..a55cd9d289d 100644 --- a/pkg/engine/internal/planner/logical/planner.go +++ b/pkg/engine/internal/planner/logical/planner.go @@ -259,8 +259,25 @@ func walkRangeAggregation(e *syntax.RangeAggregationExpr, params logql.Params) ( builder = builder.Cast(unwrapIdentifier, unwrapOperation) } - rangeAggType := convertRangeAggregationType(e.Operation) - if rangeAggType == types.RangeAggregationTypeInvalid { + var rangeAggType types.RangeAggregationType + switch e.Operation { + case syntax.OpRangeTypeCount: + rangeAggType = types.RangeAggregationTypeCount + case syntax.OpRangeTypeSum: + rangeAggType = types.RangeAggregationTypeSum + //case syntax.OpRangeTypeMax: + // rangeAggType = types.RangeAggregationTypeMax + //case syntax.OpRangeTypeMin: + // rangeAggType = types.RangeAggregationTypeMin + //case syntax.OpRangeTypeBytesRate: + // rangeAggType = types.RangeAggregationTypeBytes // bytes_rate is implemented as bytes_over_time/$interval + case syntax.OpRangeTypeRate: + if e.Left.Unwrap != nil { + rangeAggType = types.RangeAggregationTypeSum // rate of an unwrap is implemented as sum_over_time/$interval + } else { + rangeAggType = types.RangeAggregationTypeCount // rate is implemented as count_over_time/$interval + } + default: return nil, errUnimplemented } @@ -268,6 +285,19 @@ func walkRangeAggregation(e *syntax.RangeAggregationExpr, params logql.Params) ( nil, rangeAggType, params.Start(), params.End(), params.Step(), rangeInterval, ) + switch e.Operation { + //case syntax.OpRangeTypeBytesRate: + // // bytes_rate is implemented as bytes_over_time/$interval + // builder = builder.BinOpRight(types.BinaryOpDiv, &Literal{ + // Literal: NewLiteral(rangeInterval.Seconds()), + // }) + case syntax.OpRangeTypeRate: + // rate is implemented as count_over_time/$interval + builder = builder.BinOpRight(types.BinaryOpDiv, &Literal{ + Literal: NewLiteral(rangeInterval.Seconds()), + }) + } + return builder.Value(), nil } @@ -426,21 +456,6 @@ func convertVectorAggregationType(op string) types.VectorAggregationType { } } -func convertRangeAggregationType(op string) types.RangeAggregationType { - switch op { - case syntax.OpRangeTypeCount: - return types.RangeAggregationTypeCount - case syntax.OpRangeTypeSum: - return types.RangeAggregationTypeSum - //case syntax.OpRangeTypeMax: - // return types.RangeAggregationTypeMax - //case syntax.OpRangeTypeMin: - // return types.RangeAggregationTypeMin - default: - return types.RangeAggregationTypeInvalid - } -} - func convertMatcherType(t labels.MatchType) types.BinaryOp { switch t { case labels.MatchEqual: diff --git a/pkg/engine/internal/planner/logical/planner_test.go b/pkg/engine/internal/planner/logical/planner_test.go index 03754dac476..8e54c2aaf19 100644 --- a/pkg/engine/internal/planner/logical/planner_test.go +++ b/pkg/engine/internal/planner/logical/planner_test.go @@ -193,6 +193,41 @@ RETURN %12 t.Logf("\n%s\n", sb.String()) }) + + t.Run(`rate metric query with nested math expression`, func(t *testing.T) { + q := &query{ + statement: `sum by (level) ((rate({cluster="prod"}[5m]) - 100) ^ 2)`, + start: 3600, + end: 7200, + interval: 5 * time.Minute, + } + + logicalPlan, err := BuildPlan(q) + require.NoError(t, err) + t.Logf("\n%s\n", logicalPlan.String()) + + expected := `%1 = EQ label.cluster "prod" +%2 = MAKETABLE [selector=%1, predicates=[], shard=0_of_1] +%3 = GTE builtin.timestamp 1970-01-01T00:55:00Z +%4 = SELECT %2 [predicate=%3] +%5 = LT builtin.timestamp 1970-01-01T02:00:00Z +%6 = SELECT %4 [predicate=%5] +%7 = RANGE_AGGREGATION %6 [operation=count, start_ts=1970-01-01T01:00:00Z, end_ts=1970-01-01T02:00:00Z, step=0s, range=5m0s] +%8 = DIV %7 300 +%9 = SUB %8 100 +%10 = POW %9 2 +%11 = VECTOR_AGGREGATION %10 [operation=sum, group_by=(ambiguous.level)] +%12 = LOGQL_COMPAT %11 +RETURN %12 +` + + require.Equal(t, expected, logicalPlan.String()) + + var sb strings.Builder + PrintTree(&sb, logicalPlan.Value()) + + t.Logf("\n%s\n", sb.String()) + }) } func TestCanExecuteQuery(t *testing.T) { @@ -284,8 +319,8 @@ func TestCanExecuteQuery(t *testing.T) { expected: true, }, { - // rate is not supported statement: `sum by (level) (rate({env="prod"}[1m]))`, + expected: true, }, { // max is not supported diff --git a/pkg/engine/internal/types/aggregations.go b/pkg/engine/internal/types/aggregations.go index 89a407a3575..6049b6b6cee 100644 --- a/pkg/engine/internal/types/aggregations.go +++ b/pkg/engine/internal/types/aggregations.go @@ -10,6 +10,7 @@ const ( RangeAggregationTypeSum // Represents sum_over_time range aggregation RangeAggregationTypeMax // Represents max_over_time range aggregation RangeAggregationTypeMin // Represents min_over_time range aggregation + RangeAggregationTypeBytes // Represents bytes_over_time range aggregation ) func (op RangeAggregationType) String() string { @@ -22,6 +23,8 @@ func (op RangeAggregationType) String() string { return "max" case RangeAggregationTypeMin: return "min" + case RangeAggregationTypeBytes: + return "bytes" default: return "invalid" }