package meterstatementbuilder import ( "context" "fmt" "log/slog" "github.com/SigNoz/signoz/pkg/errors" "github.com/SigNoz/signoz/pkg/factory" "github.com/SigNoz/signoz/pkg/flagger" "github.com/SigNoz/signoz/pkg/querybuilder" "github.com/SigNoz/signoz/pkg/statementbuilder" "github.com/SigNoz/signoz/pkg/statementbuilder/metricsstatementbuilder" "github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema" "github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema" "github.com/SigNoz/signoz/pkg/types/metrictypes" qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5" "github.com/SigNoz/signoz/pkg/types/telemetrytypes" "github.com/SigNoz/signoz/pkg/valuer" "github.com/huandu/go-sqlbuilder" ) type meterQueryStatementBuilder struct { logger *slog.Logger metadataStore telemetrytypes.MetadataStore fm qbtypes.FieldMapper cb qbtypes.ConditionBuilder metricsStatementBuilder *metricsstatementbuilder.StatementBuilder } var _ qbtypes.StatementBuilder[qbtypes.MetricAggregation] = (*meterQueryStatementBuilder)(nil) // NewFactory returns a provider factory for the meter statement builder. Its New // reuses the metrics FieldMapper/ConditionBuilder and delegates the final SELECT // to a metrics statement builder built via the metrics factory. func NewFactory( metadataStore telemetrytypes.MetadataStore, fl flagger.Flagger, ) factory.ProviderFactory[qbtypes.StatementBuilder[qbtypes.MetricAggregation], statementbuilder.Config] { return factory.NewProviderFactory( factory.MustNewName("meter"), func(ctx context.Context, settings factory.ProviderSettings, cfg statementbuilder.Config) (qbtypes.StatementBuilder[qbtypes.MetricAggregation], error) { metricsStatementBuilder, err := metricsstatementbuilder.NewFactory(metadataStore, fl).New(ctx, settings, cfg) if err != nil { return nil, err } fm := metricstelemetryschema.NewFieldMapper() cb := metricstelemetryschema.NewConditionBuilder(fm) return NewMeterQueryStatementBuilder(settings, metadataStore, fm, cb, metricsStatementBuilder), nil }, ) } func NewMeterQueryStatementBuilder( settings factory.ProviderSettings, metadataStore telemetrytypes.MetadataStore, fieldMapper qbtypes.FieldMapper, conditionBuilder qbtypes.ConditionBuilder, metricsStatementBuilder *metricsstatementbuilder.StatementBuilder, ) *meterQueryStatementBuilder { metricsSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema") return &meterQueryStatementBuilder{ logger: metricsSettings.Logger(), metadataStore: metadataStore, fm: fieldMapper, cb: conditionBuilder, metricsStatementBuilder: metricsStatementBuilder, } } func (b *meterQueryStatementBuilder) Build( ctx context.Context, orgID valuer.UUID, start uint64, end uint64, _ qbtypes.RequestType, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], variables map[string]qbtypes.VariableItem, ) (*qbtypes.Statement, error) { keySelectors := metricsstatementbuilder.GetKeySelectors(query) keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors) if err != nil { return nil, err } start, end = querybuilder.AdjustedMetricTimeRange(start, end, uint64(query.StepInterval.Seconds()), query) return b.buildPipelineStatement(ctx, orgID, start, end, query, keys, variables) } func (b *meterQueryStatementBuilder) buildPipelineStatement( ctx context.Context, orgID valuer.UUID, start, end uint64, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], keys map[string][]*telemetrytypes.TelemetryFieldKey, variables map[string]qbtypes.VariableItem, ) (*qbtypes.Statement, error) { var ( cteFragments []string cteArgs [][]any ) if qbtypes.CanShortCircuitDelta(query.Aggregations[0]) { // spatial_aggregation_cte directly for certain delta queries if frag, args, err := b.buildTemporalAggDeltaFastPath(ctx, orgID, start, end, query, keys, variables); err != nil { return nil, err } else if frag != "" { cteFragments = append(cteFragments, frag) cteArgs = append(cteArgs, args) } } else { // temporal_aggregation_cte if frag, args, err := b.buildTemporalAggregationCTE(ctx, orgID, start, end, query, keys, variables); err != nil { return nil, err } else if frag != "" { cteFragments = append(cteFragments, frag) cteArgs = append(cteArgs, args) } // spatial_aggregation_cte if frag, args, err := b.buildSpatialAggregationCTE(ctx, start, end, query, keys); err != nil { return nil, err } else if frag != "" { cteFragments = append(cteFragments, frag) cteArgs = append(cteArgs, args) } } // final SELECT return b.metricsStatementBuilder.BuildFinalSelect(cteFragments, cteArgs, query) } func (b *meterQueryStatementBuilder) buildTemporalAggDeltaFastPath( ctx context.Context, orgID valuer.UUID, start, end uint64, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], keys map[string][]*telemetrytypes.TelemetryFieldKey, variables map[string]qbtypes.VariableItem, ) (string, []any, error) { var filterWhere querybuilder.PreparedWhereClause var err error stepSec := int64(query.StepInterval.Seconds()) sb := sqlbuilder.NewSelectBuilder() sb.SelectMore(fmt.Sprintf( "toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts", stepSec, )) for _, g := range query.GroupBy { col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys) if err != nil { return "", nil, err } sb.SelectMore(col) } tbl := metertelemetryschema.WhichSamplesTableToUse(start, end, query.Aggregations[0].Type, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints) aggCol, err := metertelemetryschema.AggregationColumnForSamplesTable(start, end, query.Aggregations[0].Type, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints) if err != nil { return "", nil, err } if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate { aggCol = fmt.Sprintf("%s/%d", aggCol, stepSec) } sb.SelectMore(fmt.Sprintf("%s AS value", aggCol)) sb.From(fmt.Sprintf("%s.%s AS points", metertelemetryschema.DBName, tbl)) sb.Where( sb.In("metric_name", query.Aggregations[0].MetricName), sb.GTE("unix_milli", start), sb.LT("unix_milli", end), ) if query.Filter != nil && query.Filter.Expression != "" { filterWhere, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{ Context: ctx, OrgID: orgID, Logger: b.logger, FieldMapper: b.fm, ConditionBuilder: b.cb, FieldKeys: keys, FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"}, Variables: variables, StartNs: start, EndNs: end, }) if err != nil { return "", nil, err } } if !filterWhere.IsEmpty() { sb.AddWhereClause(filterWhere.WhereClause) } if query.Aggregations[0].Temporality != metrictypes.Unknown { sb.Where(sb.ILike("temporality", query.Aggregations[0].Temporality.StringValue())) } sb.GroupBy("ts") sb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...) q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse) return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, nil } func (b *meterQueryStatementBuilder) buildTemporalAggregationCTE( ctx context.Context, orgID valuer.UUID, start, end uint64, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], keys map[string][]*telemetrytypes.TelemetryFieldKey, variables map[string]qbtypes.VariableItem, ) (string, []any, error) { if query.Aggregations[0].Temporality == metrictypes.Delta { return b.buildTemporalAggDelta(ctx, orgID, start, end, query, keys, variables) } return b.buildTemporalAggCumulativeOrUnspecified(ctx, orgID, start, end, query, keys, variables) } func (b *meterQueryStatementBuilder) buildTemporalAggDelta( ctx context.Context, orgID valuer.UUID, start, end uint64, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], keys map[string][]*telemetrytypes.TelemetryFieldKey, variables map[string]qbtypes.VariableItem, ) (string, []any, error) { var filterWhere querybuilder.PreparedWhereClause var err error stepSec := int64(query.StepInterval.Seconds()) sb := sqlbuilder.NewSelectBuilder() sb.Select("fingerprint") sb.SelectMore(fmt.Sprintf( "toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts", stepSec, )) for _, g := range query.GroupBy { col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys) if err != nil { return "", nil, err } sb.SelectMore(col) } tbl := metertelemetryschema.WhichSamplesTableToUse(start, end, query.Aggregations[0].Type, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints) aggCol, err := metertelemetryschema.AggregationColumnForSamplesTable(start, end, query.Aggregations[0].Type, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints) if err != nil { return "", nil, err } if query.Aggregations[0].TimeAggregation == metrictypes.TimeAggregationRate { aggCol = fmt.Sprintf("%s/%d", aggCol, stepSec) } sb.SelectMore(fmt.Sprintf("%s AS per_series_value", aggCol)) sb.From(fmt.Sprintf("%s.%s AS points", metertelemetryschema.DBName, tbl)) sb.Where( sb.In("metric_name", query.Aggregations[0].MetricName), sb.GTE("unix_milli", start), sb.LT("unix_milli", end), ) if query.Filter != nil && query.Filter.Expression != "" { filterWhere, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{ Context: ctx, OrgID: orgID, Logger: b.logger, FieldMapper: b.fm, ConditionBuilder: b.cb, FieldKeys: keys, FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"}, Variables: variables, StartNs: start, EndNs: end, }) if err != nil { return "", nil, err } } if !filterWhere.IsEmpty() { sb.AddWhereClause(filterWhere.WhereClause) } if query.Aggregations[0].Temporality != metrictypes.Unknown { sb.Where(sb.ILike("temporality", query.Aggregations[0].Temporality.StringValue())) } sb.GroupBy("fingerprint", "ts") sb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...) sb.OrderBy("fingerprint", "ts") q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse) return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil } func (b *meterQueryStatementBuilder) buildTemporalAggCumulativeOrUnspecified( ctx context.Context, orgID valuer.UUID, start, end uint64, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], keys map[string][]*telemetrytypes.TelemetryFieldKey, variables map[string]qbtypes.VariableItem, ) (string, []any, error) { var filterWhere querybuilder.PreparedWhereClause var err error stepSec := int64(query.StepInterval.Seconds()) baseSb := sqlbuilder.NewSelectBuilder() baseSb.Select("fingerprint") baseSb.SelectMore(fmt.Sprintf( "toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(%d)) AS ts", stepSec, )) for _, g := range query.GroupBy { col, err := b.fm.ColumnExpressionFor(ctx, orgID, start, end, &g.TelemetryFieldKey, telemetrytypes.FieldDataTypeString, keys) if err != nil { return "", nil, err } baseSb.SelectMore(col) } tbl := metertelemetryschema.WhichSamplesTableToUse(start, end, query.Aggregations[0].Type, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints) aggCol, err := metertelemetryschema.AggregationColumnForSamplesTable(start, end, query.Aggregations[0].Type, query.Aggregations[0].Temporality, query.Aggregations[0].TimeAggregation, query.Aggregations[0].TableHints) if err != nil { return "", nil, err } baseSb.SelectMore(fmt.Sprintf("%s AS per_series_value", aggCol)) baseSb.From(fmt.Sprintf("%s.%s AS points", metertelemetryschema.DBName, tbl)) baseSb.Where( baseSb.In("metric_name", query.Aggregations[0].MetricName), baseSb.GTE("unix_milli", start), baseSb.LT("unix_milli", end), ) if query.Filter != nil && query.Filter.Expression != "" { filterWhere, err = querybuilder.PrepareWhereClause(query.Filter.Expression, querybuilder.FilterExprVisitorOpts{ Context: ctx, OrgID: orgID, Logger: b.logger, FieldMapper: b.fm, ConditionBuilder: b.cb, FieldKeys: keys, FullTextColumn: &telemetrytypes.TelemetryFieldKey{Name: "labels"}, Variables: variables, StartNs: start, EndNs: end, }) if err != nil { return "", nil, err } } if !filterWhere.IsEmpty() { baseSb.AddWhereClause(filterWhere.WhereClause) } if query.Aggregations[0].Temporality != metrictypes.Unknown { baseSb.Where(baseSb.ILike("temporality", query.Aggregations[0].Temporality.StringValue())) } baseSb.GroupBy("fingerprint", "ts") baseSb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...) baseSb.OrderBy("fingerprint", "ts") innerQuery, innerArgs := baseSb.BuildWithFlavor(sqlbuilder.ClickHouse) switch query.Aggregations[0].TimeAggregation { case metrictypes.TimeAggregationRate: wrapped := sqlbuilder.NewSelectBuilder() wrapped.Select("ts") for _, g := range query.GroupBy { wrapped.SelectMore(fmt.Sprintf("`%s`", g.Name)) } wrapped.SelectMore(fmt.Sprintf("%s AS per_series_value", metricsstatementbuilder.RateTmpl)) wrapped.From(fmt.Sprintf("(%s) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)", innerQuery)) q, args := wrapped.BuildWithFlavor(sqlbuilder.ClickHouse, innerArgs...) return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil case metrictypes.TimeAggregationIncrease: wrapped := sqlbuilder.NewSelectBuilder() wrapped.Select("ts") for _, g := range query.GroupBy { wrapped.SelectMore(fmt.Sprintf("`%s`", g.Name)) } wrapped.SelectMore(fmt.Sprintf("%s AS per_series_value", metricsstatementbuilder.IncreaseTmpl)) wrapped.From(fmt.Sprintf("(%s) WINDOW rate_window AS (PARTITION BY fingerprint ORDER BY fingerprint, ts)", innerQuery)) q, args := wrapped.BuildWithFlavor(sqlbuilder.ClickHouse, innerArgs...) return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", q), args, nil default: return fmt.Sprintf("__temporal_aggregation_cte AS (%s)", innerQuery), innerArgs, nil } } func (b *meterQueryStatementBuilder) buildSpatialAggregationCTE( _ context.Context, _ uint64, _ uint64, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation], _ map[string][]*telemetrytypes.TelemetryFieldKey, ) (string, []any, error) { if query.Aggregations[0].SpaceAggregation.IsZero() { return "", nil, errors.Newf( errors.TypeInvalidInput, errors.CodeInvalidInput, "invalid space aggregation, should be one of the following: [`sum`, `avg`, `min`, `max`, `count`]", ) } sb := sqlbuilder.NewSelectBuilder() sb.Select("ts") for _, g := range query.GroupBy { sb.SelectMore(fmt.Sprintf("`%s`", g.Name)) } sb.SelectMore(fmt.Sprintf("%s(per_series_value) AS value", query.Aggregations[0].SpaceAggregation.StringValue())) sb.From("__temporal_aggregation_cte") sb.Where(sb.EQ("isNaN(per_series_value)", 0)) if query.Aggregations[0].ValueFilter != nil { sb.Where(sb.EQ("per_series_value", query.Aggregations[0].ValueFilter.Value)) } sb.GroupBy("ts") sb.GroupBy(querybuilder.GroupByKeys(query.GroupBy)...) q, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse) return fmt.Sprintf("__spatial_aggregation_cte AS (%s)", q), args, nil }