Skip to content
49 changes: 32 additions & 17 deletions pkg/engine/internal/planner/logical/planner.go
Original file line number Diff line number Diff line change
Expand Up @@ -259,15 +259,45 @@ 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
}

builder = builder.RangeAggregation(
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
}

Expand Down Expand Up @@ -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:
Expand Down
37 changes: 36 additions & 1 deletion pkg/engine/internal/planner/logical/planner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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
Expand Down
3 changes: 3 additions & 0 deletions pkg/engine/internal/types/aggregations.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -22,6 +23,8 @@ func (op RangeAggregationType) String() string {
return "max"
case RangeAggregationTypeMin:
return "min"
case RangeAggregationTypeBytes:
return "bytes"
default:
return "invalid"
}
Expand Down