mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-09 22:50:38 +01:00
Compare commits
2 Commits
test/semco
...
proto/logi
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
17b73fa357 | ||
|
|
0c3223b23a |
4
Makefile
4
Makefile
@@ -220,10 +220,6 @@ py-test-teardown: ## Tear down the shared SigNoz backend
|
||||
py-test: ## Runs integration tests
|
||||
@cd tests && uv run pytest --basetemp=./tmp/ -vv --capture=no integration/tests/
|
||||
|
||||
.PHONY: py-test-semconv-phase1
|
||||
py-test-semconv-phase1: py-test-setup ## Rebuild the shared stack and run the semantic-convention Phase 1 matrix
|
||||
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --capture=no integration/tests/queriertraces/13_semconv_evolution.py
|
||||
|
||||
.PHONY: py-clean
|
||||
py-clean: ## Clear all pycache and pytest cache from tests directory recursively
|
||||
@echo ">> cleaning python cache files from tests directory"
|
||||
|
||||
@@ -40,7 +40,8 @@ func (c *conditionBuilder) ConditionFor(
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
|
||||
logicalFields, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
|
||||
keys := querybuilder.SingleKeys(logicalFields)
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
|
||||
@@ -6,8 +6,6 @@ import (
|
||||
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/query-service/utils"
|
||||
"github.com/SigNoz/signoz/pkg/semconv"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
var resourceLogOperators = map[v3.FilterOperator]string{
|
||||
@@ -31,61 +29,13 @@ var resourceLogOperators = map[v3.FilterOperator]string{
|
||||
v3.FilterOperatorNotILike: "NOT ILIKE",
|
||||
}
|
||||
|
||||
func resourceSemconvMembers(key string) []string {
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
})
|
||||
}
|
||||
|
||||
func resourceValueExpression(key string) string {
|
||||
members := resourceSemconvMembers(key)
|
||||
if len(members) == 1 {
|
||||
return fmt.Sprintf("simpleJSONExtractString(labels, '%s')", key)
|
||||
}
|
||||
|
||||
values := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
values = append(values, fmt.Sprintf("NULLIF(simpleJSONExtractString(labels, '%s'), '')", member))
|
||||
}
|
||||
return "COALESCE(" + strings.Join(values, ", ") + ")"
|
||||
}
|
||||
|
||||
func resourcePresenceExpression(key string, exists bool) string {
|
||||
members := resourceSemconvMembers(key)
|
||||
if len(members) == 1 {
|
||||
if exists {
|
||||
return fmt.Sprintf("simpleJSONHas(labels, '%s')", key)
|
||||
}
|
||||
return fmt.Sprintf("not simpleJSONHas(labels, '%s')", key)
|
||||
}
|
||||
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
if exists {
|
||||
conditions = append(conditions, fmt.Sprintf("simpleJSONHas(labels, '%s')", member))
|
||||
} else {
|
||||
conditions = append(conditions, fmt.Sprintf("not simpleJSONHas(labels, '%s')", member))
|
||||
}
|
||||
}
|
||||
separator := " OR "
|
||||
if !exists {
|
||||
separator = " AND "
|
||||
}
|
||||
return "(" + strings.Join(conditions, separator) + ")"
|
||||
}
|
||||
|
||||
// buildResourceFilter builds a clickhouse filter string for resource labels
|
||||
func buildResourceFilter(logsOp string, key string, op v3.FilterOperator, value interface{}) string {
|
||||
// for all operators except contains and like
|
||||
searchKey := resourceValueExpression(key)
|
||||
searchKey := fmt.Sprintf("simpleJSONExtractString(labels, '%s')", key)
|
||||
|
||||
// for contains and like it will be case insensitive
|
||||
lowerSearchKey := fmt.Sprintf("simpleJSONExtractString(lower(labels), '%s')", key)
|
||||
if len(resourceSemconvMembers(key)) > 1 {
|
||||
lowerSearchKey = "lower(" + searchKey + ")"
|
||||
}
|
||||
|
||||
chFmtVal := utils.ClickHouseFormattedValue(value)
|
||||
|
||||
@@ -93,9 +43,9 @@ func buildResourceFilter(logsOp string, key string, op v3.FilterOperator, value
|
||||
|
||||
switch op {
|
||||
case v3.FilterOperatorExists:
|
||||
return resourcePresenceExpression(key, true)
|
||||
return fmt.Sprintf("simpleJSONHas(labels, '%s')", key)
|
||||
case v3.FilterOperatorNotExists:
|
||||
return resourcePresenceExpression(key, false)
|
||||
return fmt.Sprintf("not simpleJSONHas(labels, '%s')", key)
|
||||
case v3.FilterOperatorRegex, v3.FilterOperatorNotRegex:
|
||||
return fmt.Sprintf(logsOp, searchKey, chFmtVal)
|
||||
case v3.FilterOperatorContains, v3.FilterOperatorNotContains:
|
||||
@@ -160,38 +110,6 @@ func buildIndexFilterForInOperator(key string, op v3.FilterOperator, value inter
|
||||
// we can use lower index for =, in etc but it's difficult to do it for !=, NIN etc
|
||||
// if as x != "ABC" we cannot predict something like "not lower(labels) like '%%x%%abc%%'". It has it be "not lower(labels) like '%%x%%ABC%%'"
|
||||
func buildResourceIndexFilter(key string, op v3.FilterOperator, value interface{}) string {
|
||||
return buildResourceIndexFilterForKey(key, op, value, true)
|
||||
}
|
||||
|
||||
func buildResourceIndexFilterForKey(key string, op v3.FilterOperator, value interface{}, resolveFamily bool) string {
|
||||
members := []string{key}
|
||||
if resolveFamily {
|
||||
members = resourceSemconvMembers(key)
|
||||
}
|
||||
if len(members) > 1 {
|
||||
switch op {
|
||||
case v3.FilterOperatorNotEqual,
|
||||
v3.FilterOperatorNotLike,
|
||||
v3.FilterOperatorNotILike,
|
||||
v3.FilterOperatorNotContains,
|
||||
v3.FilterOperatorNotExists,
|
||||
v3.FilterOperatorNotRegex,
|
||||
v3.FilterOperatorNotIn:
|
||||
return ""
|
||||
}
|
||||
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
if condition := buildResourceIndexFilterForKey(member, op, value, false); condition != "" {
|
||||
conditions = append(conditions, condition)
|
||||
}
|
||||
}
|
||||
if len(conditions) == 0 {
|
||||
return ""
|
||||
}
|
||||
return "(" + strings.Join(conditions, " OR ") + ")"
|
||||
}
|
||||
|
||||
// not using clickhouseFormattedValue as we don't wan't the quotes
|
||||
strVal := fmt.Sprintf("%s", value)
|
||||
fmtValEscapedForContains := utils.QuoteEscapedStringForContains(strVal, true)
|
||||
@@ -288,31 +206,14 @@ func buildResourceFiltersFromGroupBy(groupBy []v3.AttributeKey) []string {
|
||||
if attr.Type != v3.AttributeKeyTypeResource {
|
||||
continue
|
||||
}
|
||||
members := resourceSemconvMembers(attr.Key)
|
||||
if len(members) == 1 {
|
||||
conditions = append(conditions, fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", attr.Key, attr.Key))
|
||||
continue
|
||||
}
|
||||
indexConditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
indexConditions = append(indexConditions, fmt.Sprintf("labels like '%%%s%%'", member))
|
||||
}
|
||||
conditions = append(conditions, fmt.Sprintf("(%s AND (%s))", resourcePresenceExpression(attr.Key, true), strings.Join(indexConditions, " OR ")))
|
||||
conditions = append(conditions, fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", attr.Key, attr.Key))
|
||||
}
|
||||
return conditions
|
||||
}
|
||||
|
||||
func buildResourceFiltersFromAggregateAttribute(aggregateAttribute v3.AttributeKey) string {
|
||||
if aggregateAttribute.Key != "" && aggregateAttribute.Type == v3.AttributeKeyTypeResource {
|
||||
members := resourceSemconvMembers(aggregateAttribute.Key)
|
||||
if len(members) == 1 {
|
||||
return fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", aggregateAttribute.Key, aggregateAttribute.Key)
|
||||
}
|
||||
indexConditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
indexConditions = append(indexConditions, fmt.Sprintf("labels like '%%%s%%'", member))
|
||||
}
|
||||
return fmt.Sprintf("(%s AND (%s))", resourcePresenceExpression(aggregateAttribute.Key, true), strings.Join(indexConditions, " OR "))
|
||||
return fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", aggregateAttribute.Key, aggregateAttribute.Key)
|
||||
}
|
||||
|
||||
return ""
|
||||
|
||||
@@ -5,8 +5,6 @@ import (
|
||||
"testing"
|
||||
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func Test_buildResourceFilter(t *testing.T) {
|
||||
@@ -554,38 +552,3 @@ func Test_buildResourceSubQuery(t *testing.T) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestSemanticConventionResourceFamily(t *testing.T) {
|
||||
const resolvedValue = "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''))"
|
||||
|
||||
for _, requestedName := range []string{"deployment.environment.name", "deployment.environment"} {
|
||||
t.Run(requestedName, func(t *testing.T) {
|
||||
assert.Equal(t, resolvedValue+" = 'production'", buildResourceFilter("=", requestedName, v3.FilterOperatorEqual, "production"))
|
||||
assert.Equal(t, "(simpleJSONHas(labels, 'deployment.environment.name') OR simpleJSONHas(labels, 'deployment.environment'))", buildResourceFilter("", requestedName, v3.FilterOperatorExists, nil))
|
||||
assert.Equal(t, "(not simpleJSONHas(labels, 'deployment.environment.name') AND not simpleJSONHas(labels, 'deployment.environment'))", buildResourceFilter("", requestedName, v3.FilterOperatorNotExists, nil))
|
||||
assert.Equal(t, "(labels like '%deployment.environment.name\":\"production%' OR labels like '%deployment.environment\":\"production%')", buildResourceIndexFilter(requestedName, v3.FilterOperatorEqual, "production"))
|
||||
assert.Empty(t, buildResourceIndexFilter(requestedName, v3.FilterOperatorNotEqual, "production"), "negative family filter must not use a rejecting index hint")
|
||||
})
|
||||
}
|
||||
|
||||
filters, err := buildResourceFiltersFromFilterItems(&v3.FilterSet{Items: []v3.FilterItem{{
|
||||
Key: v3.AttributeKey{
|
||||
Key: "deployment.environment.name",
|
||||
DataType: v3.AttributeKeyDataTypeString,
|
||||
Type: v3.AttributeKeyTypeResource,
|
||||
},
|
||||
Operator: v3.FilterOperatorEqual,
|
||||
Value: "production",
|
||||
}}})
|
||||
require.NoError(t, err, "family filter items must build before their output is inspected")
|
||||
wantFilters := []string{
|
||||
resolvedValue + " = 'production'",
|
||||
"(labels like '%deployment.environment.name\":\"production%' OR labels like '%deployment.environment\":\"production%')",
|
||||
}
|
||||
assert.Equal(t, wantFilters, filters)
|
||||
|
||||
wantPresence := "((simpleJSONHas(labels, 'deployment.environment.name') OR simpleJSONHas(labels, 'deployment.environment')) AND (labels like '%deployment.environment.name%' OR labels like '%deployment.environment%'))"
|
||||
groupBy := buildResourceFiltersFromGroupBy([]v3.AttributeKey{{Key: "deployment.environment", Type: v3.AttributeKeyTypeResource}})
|
||||
assert.Equal(t, []string{wantPresence}, groupBy)
|
||||
assert.Equal(t, wantPresence, buildResourceFiltersFromAggregateAttribute(v3.AttributeKey{Key: "deployment.environment.name", Type: v3.AttributeKeyTypeResource}))
|
||||
}
|
||||
|
||||
@@ -6,32 +6,16 @@ import (
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2"
|
||||
"github.com/SigNoz/signoz/pkg/query-service/model"
|
||||
"github.com/SigNoz/signoz/pkg/semconv"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
)
|
||||
|
||||
var (
|
||||
columns = serviceMapColumns()
|
||||
columns = map[string]struct{}{
|
||||
"deployment_environment": {},
|
||||
"k8s_cluster_name": {},
|
||||
"k8s_namespace_name": {},
|
||||
}
|
||||
)
|
||||
|
||||
func serviceMapColumns() map[string]string {
|
||||
columns := map[string]string{
|
||||
"k8s_cluster_name": "k8s_cluster_name",
|
||||
"k8s_namespace_name": "k8s_namespace_name",
|
||||
}
|
||||
|
||||
// Dependency-graph rows keep their historical physical column name. Both
|
||||
// semantic-convention request spellings target that same derived column.
|
||||
for _, member := range semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
}) {
|
||||
columns[strings.ReplaceAll(member, ".", "_")] = "deployment_environment"
|
||||
}
|
||||
return columns
|
||||
}
|
||||
|
||||
func BuildServiceMapQuery(tags []model.TagQuery) (string, []interface{}) {
|
||||
var filterQuery string
|
||||
var namedArgs []interface{}
|
||||
@@ -40,40 +24,39 @@ func BuildServiceMapQuery(tags []model.TagQuery) (string, []interface{}) {
|
||||
operator := tag.GetOperator()
|
||||
value := tag.GetValues()
|
||||
|
||||
column, ok := columns[key]
|
||||
if !ok {
|
||||
if _, ok := columns[key]; !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
switch operator {
|
||||
case model.InOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s IN @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s IN @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.NotInOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT IN @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT IN @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.EqualOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s = @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s = @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.NotEqualOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s != @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s != @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.ContainsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%%%s%%", value)))
|
||||
case model.NotContainsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%%%s%%", value)))
|
||||
case model.StartsWithOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%s%%", value)))
|
||||
case model.NotStartsWithOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", column, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", key, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%s%%", value)))
|
||||
case model.ExistsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NOT NULL", column)
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NOT NULL", key)
|
||||
case model.NotExistsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NULL", column)
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NULL", key)
|
||||
}
|
||||
}
|
||||
return filterQuery, namedArgs
|
||||
|
||||
@@ -1,35 +0,0 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
|
||||
"github.com/SigNoz/signoz/pkg/query-service/model"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestBuildServiceMapQueryAcceptsEnvironmentFamily(t *testing.T) {
|
||||
for _, requestedName := range []string{"deployment.environment.name", "deployment.environment"} {
|
||||
t.Run(requestedName, func(t *testing.T) {
|
||||
tags := []model.TagQuery{model.NewTagQueryString(model.TagQueryParam{
|
||||
Key: requestedName,
|
||||
StringValues: []string{"production"},
|
||||
Operator: model.InOperator,
|
||||
TagType: model.ResourceAttributeTagType,
|
||||
})}
|
||||
|
||||
query, args := BuildServiceMapQuery(tags)
|
||||
argName := "deployment_environment"
|
||||
if requestedName == "deployment.environment.name" {
|
||||
argName = "deployment_environment_name"
|
||||
}
|
||||
assert.Equal(t, " AND deployment_environment IN @"+argName, query)
|
||||
require.Len(t, args, 1)
|
||||
named, ok := args[0].(driver.NamedValue)
|
||||
require.True(t, ok)
|
||||
assert.Equal(t, argName, named.Name)
|
||||
assert.Equal(t, []interface{}{"production"}, named.Value)
|
||||
})
|
||||
}
|
||||
}
|
||||
17
pkg/querybuilder/clickhouse_quote.go
Normal file
17
pkg/querybuilder/clickhouse_quote.go
Normal file
@@ -0,0 +1,17 @@
|
||||
package querybuilder
|
||||
|
||||
import "strings"
|
||||
|
||||
// ClickHouseStringLiteral quotes a value for a ClickHouse string literal.
|
||||
func ClickHouseStringLiteral(value string) string {
|
||||
escaped := strings.ReplaceAll(value, `\`, `\\`)
|
||||
escaped = strings.ReplaceAll(escaped, `'`, `\'`)
|
||||
return "'" + escaped + "'"
|
||||
}
|
||||
|
||||
// ClickHouseIdentifier quotes a value for a ClickHouse identifier.
|
||||
func ClickHouseIdentifier(value string) string {
|
||||
escaped := strings.ReplaceAll(value, `\`, `\\`)
|
||||
escaped = strings.ReplaceAll(escaped, "`", "\\`")
|
||||
return "`" + escaped + "`"
|
||||
}
|
||||
17
pkg/querybuilder/clickhouse_quote_test.go
Normal file
17
pkg/querybuilder/clickhouse_quote_test.go
Normal file
@@ -0,0 +1,17 @@
|
||||
package querybuilder
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestClickHouseQuoting(t *testing.T) {
|
||||
t.Run("string literal", func(t *testing.T) {
|
||||
assert.Equal(t, `'name\'\\); SELECT 1 --'`, ClickHouseStringLiteral(`name'\); SELECT 1 --`))
|
||||
})
|
||||
|
||||
t.Run("identifier", func(t *testing.T) {
|
||||
assert.Equal(t, "`name\\`\\\\); SELECT 1 --`", ClickHouseIdentifier("name`\\); SELECT 1 --"))
|
||||
})
|
||||
}
|
||||
@@ -43,7 +43,7 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
|
||||
if len(evolutionsEntries) > 0 && evolutionsEntries[0] != nil {
|
||||
columnName = evolutionsEntries[0].ColumnName
|
||||
}
|
||||
rawPath := fmt.Sprintf("%s.`%s`", columnName, key.Name)
|
||||
rawPath := fmt.Sprintf("%s.%s", columnName, ClickHouseIdentifier(key.Name))
|
||||
if exists {
|
||||
return rawPath + " IS NOT NULL", nil
|
||||
}
|
||||
@@ -88,7 +88,7 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
|
||||
|
||||
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
|
||||
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumBool, schema.ColumnTypeEnumFloat64:
|
||||
leftOperand := fmt.Sprintf("mapContains(%s, '%s')", column.Name, key.Name)
|
||||
leftOperand := fmt.Sprintf("mapContains(%s, %s)", column.Name, ClickHouseStringLiteral(key.Name))
|
||||
if key.Materialized {
|
||||
leftOperand = telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key)
|
||||
}
|
||||
|
||||
41
pkg/querybuilder/exists_expr_test.go
Normal file
41
pkg/querybuilder/exists_expr_test.go
Normal file
@@ -0,0 +1,41 @@
|
||||
package querybuilder
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// ExistsExpression is a per-physical-key primitive: family composition happens
|
||||
// at the logical-field layer, so this only ever sees one spelling.
|
||||
func TestExistsExpressionIsPerKey(t *testing.T) {
|
||||
columns := []*schema.Column{{
|
||||
Name: "attributes_string",
|
||||
Type: schema.MapColumnType{
|
||||
KeyType: schema.LowCardinalityColumnType{ElementType: schema.ColumnTypeString},
|
||||
ValueType: schema.ColumnTypeString,
|
||||
},
|
||||
}}
|
||||
|
||||
plain := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
expression, err := ExistsExpression(columns, plain, 0, 0, "unused", true)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "mapContains(attributes_string, 'deployment.environment')", expression)
|
||||
|
||||
materialized := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Materialized: true,
|
||||
}
|
||||
expression, err = ExistsExpression(columns, materialized, 0, 0, "unused", false)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "NOT `attribute_string_deployment$$environment_exists`", expression)
|
||||
}
|
||||
@@ -21,24 +21,25 @@ const (
|
||||
hasTokenFunctionDocURL = "https://signoz.io/docs/userguide/functions-reference/#hastoken-function"
|
||||
)
|
||||
|
||||
// ResolveKeys picks which matching field keys a filter term builds conditions for.
|
||||
// With 0 or 1 match it returns the input unchanged and no warning. When a name is
|
||||
// ambiguous it returns a warning; a resource+attribute mix defaults to the resource
|
||||
// keys (the common intent), noted in the warning.
|
||||
func ResolveKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeysForName []*telemetrytypes.TelemetryFieldKey) ([]*telemetrytypes.TelemetryFieldKey, string) {
|
||||
if len(fieldKeysForName) <= 1 {
|
||||
return fieldKeysForName, ""
|
||||
// ResolveLogicalFields picks which logical fields a filter term builds conditions
|
||||
// for. With 0 or 1 field it returns the input unchanged and no warning. When a
|
||||
// name is ambiguous (several logical fields — a family is one field and never
|
||||
// ambiguous with itself) it returns a warning; a resource+attribute mix defaults
|
||||
// to the resource fields (the common intent), noted in the warning.
|
||||
func ResolveLogicalFields(field *telemetrytypes.TelemetryFieldKey, logicalFields []*telemetrytypes.LogicalField) ([]*telemetrytypes.LogicalField, string) {
|
||||
if len(logicalFields) <= 1 {
|
||||
return logicalFields, ""
|
||||
}
|
||||
|
||||
warning := fmt.Sprintf(
|
||||
"Key `%s` is ambiguous, found %d different combinations of field context / data type: %v.",
|
||||
field.Name,
|
||||
len(fieldKeysForName),
|
||||
fieldKeysForName,
|
||||
len(logicalFields),
|
||||
logicalFields,
|
||||
)
|
||||
|
||||
hasResource, hasAttribute := false, false
|
||||
for _, item := range fieldKeysForName {
|
||||
for _, item := range logicalFields {
|
||||
switch item.FieldContext {
|
||||
case telemetrytypes.FieldContextResource:
|
||||
hasResource = true
|
||||
@@ -49,18 +50,28 @@ func ResolveKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeysForName []*te
|
||||
|
||||
// when there is both resource and attribute context, default to resource only
|
||||
if hasResource && hasAttribute {
|
||||
filteredKeys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fieldKeysForName))
|
||||
for _, item := range fieldKeysForName {
|
||||
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
|
||||
for _, item := range logicalFields {
|
||||
if item.FieldContext == telemetrytypes.FieldContextResource {
|
||||
filteredKeys = append(filteredKeys, item)
|
||||
filtered = append(filtered, item)
|
||||
}
|
||||
}
|
||||
fieldKeysForName = filteredKeys
|
||||
logicalFields = filtered
|
||||
warning += " " + "Using `resource` context by default. To query attributes explicitly, " +
|
||||
fmt.Sprintf("use the fully qualified name (e.g., 'attribute.%s')", field.Name)
|
||||
}
|
||||
|
||||
return fieldKeysForName, warning
|
||||
return logicalFields, warning
|
||||
}
|
||||
|
||||
// WrapAsLogicalFields wraps physical keys (candidate or synthesized) as
|
||||
// single-member logical fields addressed by the requested spelling.
|
||||
func WrapAsLogicalFields(requestedName string, keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
|
||||
fields := make([]*telemetrytypes.LogicalField, 0, len(keys))
|
||||
for _, key := range keys {
|
||||
fields = append(fields, telemetrytypes.SingleLogicalField(requestedName, key))
|
||||
}
|
||||
return fields
|
||||
}
|
||||
|
||||
// NewKeyNotFoundError builds the error a condition builder returns when a filter term
|
||||
@@ -175,3 +186,15 @@ func NewFunctionUnsupportedError(operator qbtypes.FilterOperator) error {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// SingleKeys flattens single-member logical fields to their member keys. It is
|
||||
// the adapter for signals whose fields are never families (everything except
|
||||
// traces today); a family in the input would be silently narrowed, so callers
|
||||
// must be gated signals.
|
||||
func SingleKeys(fields []*telemetrytypes.LogicalField) []*telemetrytypes.TelemetryFieldKey {
|
||||
keys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fields))
|
||||
for _, logical := range fields {
|
||||
keys = append(keys, logical.Single())
|
||||
}
|
||||
return keys
|
||||
}
|
||||
|
||||
@@ -12,6 +12,9 @@ import (
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// A promoted historical member keeps its materialized column inside the family
|
||||
// expression: the member key carries its own Materialized state, so the
|
||||
// logical-field merge needs no sibling bookkeeping.
|
||||
func TestTraceFamilyUsesMaterializedHistoricalMember(t *testing.T) {
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
@@ -32,12 +35,13 @@ func TestTraceFamilyUsesMaterializedHistoricalMember(t *testing.T) {
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
|
||||
matches := querybuilder.MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
matches := querybuilder.MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
current.Name: {current},
|
||||
historical.Name: {historical},
|
||||
})
|
||||
require.Len(t, matches, 1, "family metadata should resolve to one logical field")
|
||||
expression, err := tracestelemetryschema.NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, matches[0])
|
||||
|
||||
expression, err := tracestelemetryschema.NewFieldMapper().FieldForLogical(context.Background(), valuer.UUID{}, 0, 0, matches[0])
|
||||
require.NoError(t, err, "resolved trace family should map to a value expression")
|
||||
|
||||
assert.Equal(
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"maps"
|
||||
"slices"
|
||||
"strconv"
|
||||
"strings"
|
||||
@@ -362,7 +361,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
|
||||
return ErrorConditionLiteral
|
||||
}
|
||||
}
|
||||
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
|
||||
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
|
||||
if !ok {
|
||||
return ErrorConditionLiteral
|
||||
}
|
||||
@@ -381,7 +380,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
|
||||
// VisitComparison handles all comparison operators.
|
||||
func (v *filterExpressionVisitor) VisitComparison(ctx *grammar.ComparisonContext) any {
|
||||
key := v.Visit(ctx.Key()).(*telemetrytypes.TelemetryFieldKey)
|
||||
matching := MatchingFieldKeys(key, v.fieldKeys)
|
||||
matching := MatchingLogicalFields(key, v.fieldKeys)
|
||||
|
||||
// Handle EXISTS specially
|
||||
if ctx.EXISTS() != nil {
|
||||
@@ -677,7 +676,7 @@ func (v *filterExpressionVisitor) VisitFullText(ctx *grammar.FullTextContext) an
|
||||
v.errors = append(v.errors, "full text search is not supported")
|
||||
return ErrorConditionLiteral
|
||||
}
|
||||
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
|
||||
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
|
||||
if !ok {
|
||||
return ErrorConditionLiteral
|
||||
}
|
||||
@@ -732,7 +731,7 @@ func (v *filterExpressionVisitor) VisitFunctionCall(ctx *grammar.FunctionCallCon
|
||||
return ErrorConditionLiteral
|
||||
}
|
||||
|
||||
conds, ok := v.buildConditions(key, MatchingFieldKeys(key, v.fieldKeys), operator, value)
|
||||
conds, ok := v.buildConditions(key, MatchingLogicalFields(key, v.fieldKeys), operator, value)
|
||||
if !ok {
|
||||
return ErrorConditionLiteral
|
||||
}
|
||||
@@ -924,7 +923,7 @@ func (v *filterExpressionVisitor) VisitKey(ctx *grammar.KeyContext) any {
|
||||
|
||||
// buildConditions invokes the condition builder for a filter term, folding its
|
||||
// warnings/errors into visitor state; returns false if an error was recorded.
|
||||
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) ([]string, bool) {
|
||||
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.LogicalField, op qbtypes.FilterOperator, value any) ([]string, bool) {
|
||||
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter}, op, value, v.builder)
|
||||
if err != nil {
|
||||
_, _, _, _, errURL, _ := errors.Unwrapb(err)
|
||||
@@ -981,23 +980,43 @@ func assignIfEmpty(s *string, value string) {
|
||||
}
|
||||
}
|
||||
|
||||
// MatchingFieldKeys returns the field keys from the map that match the given key,
|
||||
// honoring any context/data type the user specified.
|
||||
func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
|
||||
selector := telemetrytypes.FieldKeySelector{
|
||||
// familyMemberNames returns the physical spellings to look up for the
|
||||
// referenced key: the semantic-convention family members (current-first) when
|
||||
// the key can resolve to traces, else just the requested name. Only trace
|
||||
// field mappers understand families today; logs and metrics keep the
|
||||
// requested spelling until theirs land.
|
||||
func familyMemberNames(field *telemetrytypes.TelemetryFieldKey) []string {
|
||||
if field.Signal != telemetrytypes.SignalUnspecified && field.Signal != telemetrytypes.SignalTraces {
|
||||
return []string{field.Name}
|
||||
}
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: field.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: field.FieldContext,
|
||||
})
|
||||
}
|
||||
|
||||
// MatchingLogicalFields resolves the referenced key against the metadata map
|
||||
// into logical fields, honoring any context/data type the user specified.
|
||||
//
|
||||
// Physical keys that are members of one semantic-convention family (traces
|
||||
// only today) group into a single logical field per (signal, context, data
|
||||
// type) identity, members ordered current-first. Every other matching key
|
||||
// becomes its own single-member logical field. Ambiguity is therefore the
|
||||
// length of the returned slice, and a family is never ambiguous with itself.
|
||||
// Members alias the metadata map entries; nothing is copied or mutated.
|
||||
func MatchingLogicalFields(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
|
||||
members := familyMemberNames(field)
|
||||
memberRank := make(map[string]int, len(members))
|
||||
for i, member := range members {
|
||||
memberRank[member] = i
|
||||
}
|
||||
members := []string{field.Name}
|
||||
// Only trace field mappers understand semantic-convention families today.
|
||||
// Logs and metrics must keep using the requested spelling until theirs land.
|
||||
if field.Signal == telemetrytypes.SignalUnspecified || field.Signal == telemetrytypes.SignalTraces {
|
||||
members = semconv.Members(semconv.KindAttribute, selector)
|
||||
}
|
||||
isFamily := len(members) > 1
|
||||
fieldKeysForName := make([]*telemetrytypes.TelemetryFieldKey, 0)
|
||||
|
||||
fields := make([]*telemetrytypes.LogicalField, 0)
|
||||
indexByIdentity := make(map[string]int)
|
||||
// rank of the family member each physical key matched under; the stored
|
||||
// name of a context-prefixed match differs from the member name.
|
||||
ranks := make(map[*telemetrytypes.TelemetryFieldKey]int)
|
||||
|
||||
appendMatches := func(lookupName string, memberName string, contextAlreadyMatched bool) {
|
||||
for _, item := range fieldKeys[lookupName] {
|
||||
@@ -1011,7 +1030,7 @@ func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[st
|
||||
// A wildcard lookup may have found a same-named field in a scope where
|
||||
// this family does not apply. Keep exact names, but reject cross-member
|
||||
// matches outside the generated family scope.
|
||||
traceFamilyMatch := isFamily && item.Signal == telemetrytypes.SignalTraces
|
||||
traceFamilyMatch := len(members) > 1 && item.Signal == telemetrytypes.SignalTraces
|
||||
if memberName != field.Name {
|
||||
if !traceFamilyMatch {
|
||||
continue
|
||||
@@ -1026,53 +1045,38 @@ func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[st
|
||||
}
|
||||
}
|
||||
|
||||
physicalMembers := item.SemconvMembers
|
||||
if len(physicalMembers) == 0 {
|
||||
physicalMembers = []string{memberName}
|
||||
}
|
||||
materializedColumns := maps.Clone(item.SemconvMaterializedColumns)
|
||||
if item.Materialized {
|
||||
if materializedColumns == nil {
|
||||
materializedColumns = make(map[string]string)
|
||||
}
|
||||
physicalKey := *item
|
||||
physicalKey.Name = memberName
|
||||
materializedColumns[memberName] = strings.Trim(telemetrytypes.FieldKeyToMaterializedColumnName(&physicalKey), "`")
|
||||
if !traceFamilyMatch {
|
||||
fields = append(fields, telemetrytypes.SingleLogicalField(field.Name, item))
|
||||
continue
|
||||
}
|
||||
|
||||
identity := item.Signal.StringValue() + ";" + item.FieldContext.StringValue() + ";" + item.FieldDataType.StringValue()
|
||||
if traceFamilyMatch {
|
||||
if index, found := indexByIdentity[identity]; found {
|
||||
for _, physicalMember := range physicalMembers {
|
||||
if !slices.Contains(fieldKeysForName[index].SemconvMembers, physicalMember) {
|
||||
fieldKeysForName[index].SemconvMembers = append(fieldKeysForName[index].SemconvMembers, physicalMember)
|
||||
}
|
||||
}
|
||||
if len(materializedColumns) > 0 {
|
||||
if fieldKeysForName[index].SemconvMaterializedColumns == nil {
|
||||
fieldKeysForName[index].SemconvMaterializedColumns = make(map[string]string)
|
||||
}
|
||||
maps.Copy(fieldKeysForName[index].SemconvMaterializedColumns, materializedColumns)
|
||||
}
|
||||
continue
|
||||
index, found := indexByIdentity[identity]
|
||||
if !found {
|
||||
index = len(fields)
|
||||
indexByIdentity[identity] = index
|
||||
fields = append(fields, &telemetrytypes.LogicalField{
|
||||
Name: field.Name,
|
||||
Signal: item.Signal,
|
||||
FieldContext: item.FieldContext,
|
||||
FieldDataType: item.FieldDataType,
|
||||
})
|
||||
}
|
||||
logical := fields[index]
|
||||
duplicate := false
|
||||
for _, existing := range logical.Members {
|
||||
if existing.Name == item.Name {
|
||||
duplicate = true
|
||||
break
|
||||
}
|
||||
indexByIdentity[identity] = len(fieldKeysForName)
|
||||
}
|
||||
resolved := *item
|
||||
// The requested spelling is the response identity. Field mappers use
|
||||
// it to resolve the available family members current-first.
|
||||
if traceFamilyMatch {
|
||||
resolved.Name = field.Name
|
||||
resolved.SemconvMembers = slices.Clone(physicalMembers)
|
||||
resolved.SemconvMaterializedColumns = materializedColumns
|
||||
// Materialization is member-specific after family keys are merged.
|
||||
resolved.Materialized = false
|
||||
if !duplicate {
|
||||
ranks[item] = memberRank[memberName]
|
||||
logical.Members = append(logical.Members, item)
|
||||
}
|
||||
fieldKeysForName = append(fieldKeysForName, &resolved)
|
||||
}
|
||||
}
|
||||
|
||||
// Members are current-first, so metadata from the current key wins when
|
||||
// both spellings describe the same signal/context/type.
|
||||
for _, member := range members {
|
||||
appendMatches(member, member, false)
|
||||
}
|
||||
@@ -1086,5 +1090,13 @@ func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[st
|
||||
}
|
||||
}
|
||||
|
||||
return fieldKeysForName
|
||||
// Precedence is a property of the family, not of arrival order: members
|
||||
// sort current-first no matter which lookup pass found them.
|
||||
for _, logical := range fields {
|
||||
slices.SortStableFunc(logical.Members, func(a, b *telemetrytypes.TelemetryFieldKey) int {
|
||||
return ranks[a] - ranks[b]
|
||||
})
|
||||
}
|
||||
|
||||
return fields
|
||||
}
|
||||
|
||||
@@ -590,8 +590,8 @@ func TestVisitKey(t *testing.T) {
|
||||
// VisitKey only parses; the condition builder matches, resolves ambiguity
|
||||
// and decides not-found handling. Replay that here against the generic
|
||||
// builder behavior (error unless the key is ignored).
|
||||
matching := MatchingFieldKeys(key, tt.fieldKeys)
|
||||
keys, warning := ResolveKeys(key, matching)
|
||||
matching := MatchingLogicalFields(key, tt.fieldKeys)
|
||||
keys, warning := ResolveLogicalFields(key, matching)
|
||||
|
||||
var gotErrors []string
|
||||
var gotMainErrURL, gotMainWrnURL string
|
||||
@@ -613,15 +613,19 @@ func TestVisitKey(t *testing.T) {
|
||||
t.Errorf("expected %d keys, got %d", len(tt.expectedKeys), len(keys))
|
||||
}
|
||||
|
||||
// Check each expected key matches name, field context, and data type
|
||||
// Check each expected key matches a member's stored name plus the
|
||||
// logical field's context and data type (the logical Name is the
|
||||
// requested spelling, members keep the stored spellings).
|
||||
for _, expectedKey := range tt.expectedKeys {
|
||||
found := false
|
||||
for _, key := range keys {
|
||||
if key.Name == expectedKey.Name &&
|
||||
key.FieldContext == expectedKey.FieldContext &&
|
||||
key.FieldDataType == expectedKey.FieldDataType {
|
||||
found = true
|
||||
break
|
||||
for _, logical := range keys {
|
||||
for _, member := range logical.Members {
|
||||
if member.Name == expectedKey.Name &&
|
||||
logical.FieldContext == expectedKey.FieldContext &&
|
||||
logical.FieldDataType == expectedKey.FieldDataType {
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
@@ -686,7 +690,15 @@ func TestVisitKey(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMatchingFieldKeysResolvesCurrentTraceNameFromOldMetadata(t *testing.T) {
|
||||
func memberNames(logical *telemetrytypes.LogicalField) []string {
|
||||
names := make([]string, 0, len(logical.Members))
|
||||
for _, member := range logical.Members {
|
||||
names = append(names, member.Name)
|
||||
}
|
||||
return names
|
||||
}
|
||||
|
||||
func TestMatchingLogicalFieldsResolvesCurrentTraceNameFromOldMetadata(t *testing.T) {
|
||||
old := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Description: "old metadata",
|
||||
@@ -700,25 +712,23 @@ func TestMatchingFieldKeysResolvesCurrentTraceNameFromOldMetadata(t *testing.T)
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
|
||||
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{old.Name: {old}})
|
||||
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{old.Name: {old}})
|
||||
|
||||
require.Len(t, matches, 1, "trace family lookup must resolve before inspecting metadata")
|
||||
assert.Equal(t, "deployment.environment.name", matches[0].Name)
|
||||
assert.Equal(t, "old metadata", matches[0].Description)
|
||||
assert.Equal(t, []string{"deployment.environment"}, matches[0].SemconvMembers)
|
||||
require.Len(t, matches, 1, "a family is one logical field, not an ambiguity")
|
||||
assert.Equal(t, "deployment.environment.name", matches[0].Name, "the requested spelling is the response identity")
|
||||
assert.Equal(t, []string{"deployment.environment"}, memberNames(matches[0]))
|
||||
assert.Same(t, old, matches[0].Members[0], "members alias metadata entries; nothing is copied")
|
||||
}
|
||||
|
||||
func TestMatchingFieldKeysUsesCurrentTraceMetadataForOldName(t *testing.T) {
|
||||
func TestMatchingLogicalFieldsGroupsFamilyMembersCurrentFirst(t *testing.T) {
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Description: "current metadata",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
old := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Description: "old metadata",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
@@ -729,18 +739,50 @@ func TestMatchingFieldKeysUsesCurrentTraceMetadataForOldName(t *testing.T) {
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
|
||||
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
current.Name: {current},
|
||||
old.Name: {old},
|
||||
})
|
||||
|
||||
require.Len(t, matches, 1, "trace family lookup must resolve before inspecting metadata")
|
||||
assert.Equal(t, old.Name, matches[0].Name)
|
||||
assert.Equal(t, "current metadata", matches[0].Description)
|
||||
assert.Equal(t, []string{current.Name, old.Name}, matches[0].SemconvMembers)
|
||||
require.Len(t, matches, 1, "a family is one logical field, not an ambiguity")
|
||||
assert.Equal(t, old.Name, matches[0].Name, "the requested spelling is the response identity")
|
||||
assert.True(t, matches[0].IsFamily())
|
||||
assert.Equal(t, []string{current.Name, old.Name}, memberNames(matches[0]), "members order current-first")
|
||||
}
|
||||
|
||||
func TestMatchingFieldKeysKeepsLogSemconvNamesLiteral(t *testing.T) {
|
||||
// Precedence is a property of the family, not of arrival order: a member that
|
||||
// only exists under its context-prefixed stored spelling still sorts by its
|
||||
// family rank.
|
||||
func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
|
||||
old := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
prefixedCurrent := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "resource.deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
requested := telemetrytypes.NewTelemetryFieldKey(
|
||||
"deployment.environment.name",
|
||||
telemetrytypes.FieldContextResource,
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
|
||||
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
old.Name: {old},
|
||||
prefixedCurrent.Name: {prefixedCurrent},
|
||||
})
|
||||
|
||||
require.Len(t, matches, 1)
|
||||
assert.Equal(t, []string{prefixedCurrent.Name, old.Name}, memberNames(matches[0]),
|
||||
"the current-spelling member must coalesce before the old one")
|
||||
}
|
||||
|
||||
func TestMatchingLogicalFieldsKeepsLogSemconvNamesLiteral(t *testing.T) {
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalLogs,
|
||||
@@ -759,17 +801,17 @@ func TestMatchingFieldKeysKeepsLogSemconvNamesLiteral(t *testing.T) {
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
|
||||
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
current.Name: {current},
|
||||
old.Name: {old},
|
||||
})
|
||||
|
||||
require.Len(t, matches, 1, "log lookup must keep the requested spelling literal")
|
||||
assert.Equal(t, current.Name, matches[0].Name)
|
||||
assert.Empty(t, matches[0].SemconvMembers)
|
||||
assert.False(t, matches[0].IsFamily())
|
||||
assert.Equal(t, current.Name, matches[0].Single().Name)
|
||||
}
|
||||
|
||||
func TestMatchingFieldKeysKeepsMetricSemconvNamesLiteral(t *testing.T) {
|
||||
func TestMatchingLogicalFieldsKeepsMetricSemconvNamesLiteral(t *testing.T) {
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalMetrics,
|
||||
@@ -788,14 +830,14 @@ func TestMatchingFieldKeysKeepsMetricSemconvNamesLiteral(t *testing.T) {
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
|
||||
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
current.Name: {current},
|
||||
old.Name: {old},
|
||||
})
|
||||
|
||||
require.Len(t, matches, 1, "metric lookup must keep the requested spelling literal")
|
||||
assert.Equal(t, current.Name, matches[0].Name)
|
||||
assert.Empty(t, matches[0].SemconvMembers)
|
||||
assert.False(t, matches[0].IsFamily())
|
||||
assert.Equal(t, current.Name, matches[0].Single().Name)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -879,7 +921,7 @@ func (b *resourceConditionBuilder) ConditionFor(
|
||||
return nil, nil, nil
|
||||
}
|
||||
|
||||
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
|
||||
keys, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
@@ -887,11 +929,11 @@ func (b *resourceConditionBuilder) ConditionFor(
|
||||
|
||||
var conds []string
|
||||
for _, k := range keys {
|
||||
// only resource keys contribute; others (and unknown keys) are ignored
|
||||
// only resource fields contribute; others (and unknown keys) are ignored
|
||||
if k.FieldContext != telemetrytypes.FieldContextResource {
|
||||
continue
|
||||
}
|
||||
conds = append(conds, fmt.Sprintf("%s_cond", k.Name))
|
||||
conds = append(conds, fmt.Sprintf("%s_cond", k.Single().Name))
|
||||
}
|
||||
return conds, warnings, nil
|
||||
}
|
||||
@@ -921,7 +963,7 @@ func (b *conditionBuilder) ConditionFor(
|
||||
return []string{fmt.Sprintf("%s_cond", key.Name)}, nil, nil
|
||||
}
|
||||
|
||||
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
|
||||
keys, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
@@ -933,7 +975,7 @@ func (b *conditionBuilder) ConditionFor(
|
||||
|
||||
// A resource sub-query already covers the term; drop resource keys from the main query.
|
||||
if options.SkipResourceFilter {
|
||||
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
|
||||
filtered := make([]*telemetrytypes.LogicalField, 0, len(keys))
|
||||
for _, k := range keys {
|
||||
if k.FieldContext != telemetrytypes.FieldContextResource {
|
||||
filtered = append(filtered, k)
|
||||
|
||||
@@ -237,7 +237,6 @@ func NewSQLMigrationProviderFactories(
|
||||
sqlmigration.NewAddDashboardTuplesFactory(sqlstore),
|
||||
sqlmigration.NewRestructureSavedViewSpecFactory(sqlstore, sqlschema),
|
||||
sqlmigration.NewAddSavedViewTuplesFactory(sqlstore),
|
||||
sqlmigration.NewMigrateDeploymentEnvironmentQuickFilterFactory(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -1,127 +0,0 @@
|
||||
package sqlmigration
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/semconv"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/uptrace/bun"
|
||||
"github.com/uptrace/bun/migrate"
|
||||
)
|
||||
|
||||
const deploymentEnvironmentCurrent = "deployment.environment.name"
|
||||
|
||||
type migrateDeploymentEnvironmentQuickFilter struct {
|
||||
logger *slog.Logger
|
||||
}
|
||||
|
||||
type semconvQuickFilterRow struct {
|
||||
bun.BaseModel `bun:"table:quick_filter"`
|
||||
|
||||
ID string `bun:"id"`
|
||||
Filter string `bun:"filter"`
|
||||
}
|
||||
|
||||
func NewMigrateDeploymentEnvironmentQuickFilterFactory() factory.ProviderFactory[SQLMigration, Config] {
|
||||
return factory.NewProviderFactory(
|
||||
factory.MustNewName("migrate_semconv_quick_filter"),
|
||||
func(_ context.Context, settings factory.ProviderSettings, _ Config) (SQLMigration, error) {
|
||||
return &migrateDeploymentEnvironmentQuickFilter{logger: settings.Logger}, nil
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
func (migration *migrateDeploymentEnvironmentQuickFilter) Register(migrations *migrate.Migrations) error {
|
||||
return migrations.Register(migration.Up, migration.Down)
|
||||
}
|
||||
|
||||
func deploymentEnvironmentOld() string {
|
||||
members := semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: deploymentEnvironmentCurrent,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
})
|
||||
if len(members) < 2 {
|
||||
return deploymentEnvironmentCurrent
|
||||
}
|
||||
return members[1]
|
||||
}
|
||||
|
||||
func rewriteQuickFilterSemconv(filterJSON, from, to string) (string, bool, error) {
|
||||
var filters []map[string]any
|
||||
if err := json.Unmarshal([]byte(filterJSON), &filters); err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
|
||||
changed := false
|
||||
for _, filter := range filters {
|
||||
if key, ok := filter["key"].(string); ok && key == from {
|
||||
filter["key"] = to
|
||||
changed = true
|
||||
}
|
||||
}
|
||||
if !changed {
|
||||
return filterJSON, false, nil
|
||||
}
|
||||
|
||||
rewritten, err := json.Marshal(filters)
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
return string(rewritten), true, nil
|
||||
}
|
||||
|
||||
func (migration *migrateDeploymentEnvironmentQuickFilter) migrate(ctx context.Context, db *bun.DB, from, to string) error {
|
||||
tx, err := db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer func() { _ = tx.Rollback() }()
|
||||
|
||||
rows := make([]*semconvQuickFilterRow, 0)
|
||||
if err := tx.NewSelect().
|
||||
Model(&rows).
|
||||
Where("signal IN (?)", bun.In([]string{"traces", "api_monitoring", "exceptions"})).
|
||||
Scan(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, row := range rows {
|
||||
rewritten, changed, err := rewriteQuickFilterSemconv(row.Filter, from, to)
|
||||
if err != nil {
|
||||
// Quick filters are user-editable. One malformed legacy row must not
|
||||
// prevent the application from starting or block every other org's
|
||||
// migration.
|
||||
if migration.logger != nil {
|
||||
migration.logger.WarnContext(ctx, "skipping quick filter with unreadable filter JSON",
|
||||
slog.String("quick_filter_id", row.ID), slog.Any("error", err))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if !changed {
|
||||
continue
|
||||
}
|
||||
if _, err := tx.NewUpdate().
|
||||
Model((*semconvQuickFilterRow)(nil)).
|
||||
Set("filter = ?", rewritten).
|
||||
Set("updated_at = ?", time.Now()).
|
||||
Where("id = ?", row.ID).
|
||||
Exec(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
func (migration *migrateDeploymentEnvironmentQuickFilter) Up(ctx context.Context, db *bun.DB) error {
|
||||
return migration.migrate(ctx, db, deploymentEnvironmentOld(), deploymentEnvironmentCurrent)
|
||||
}
|
||||
|
||||
func (migration *migrateDeploymentEnvironmentQuickFilter) Down(ctx context.Context, db *bun.DB) error {
|
||||
return migration.migrate(ctx, db, deploymentEnvironmentCurrent, deploymentEnvironmentOld())
|
||||
}
|
||||
@@ -1,38 +0,0 @@
|
||||
package sqlmigration
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestRewriteQuickFilterSemconv(t *testing.T) {
|
||||
oldName := deploymentEnvironmentOld()
|
||||
input := `[{"key":"service.name","dataType":"string","type":"resource"},{"key":"` + oldName + `","dataType":"string","type":"resource","custom":true}]`
|
||||
|
||||
rewritten, changed, err := rewriteQuickFilterSemconv(input, oldName, deploymentEnvironmentCurrent)
|
||||
require.NoError(t, err)
|
||||
assert.True(t, changed)
|
||||
|
||||
var filters []map[string]any
|
||||
require.NoError(t, json.Unmarshal([]byte(rewritten), &filters))
|
||||
assert.Equal(t, "service.name", filters[0]["key"])
|
||||
assert.Equal(t, deploymentEnvironmentCurrent, filters[1]["key"])
|
||||
assert.Equal(t, true, filters[1]["custom"], "unknown filter properties must be preserved")
|
||||
|
||||
restored, changed, err := rewriteQuickFilterSemconv(rewritten, deploymentEnvironmentCurrent, oldName)
|
||||
require.NoError(t, err)
|
||||
assert.True(t, changed)
|
||||
require.NoError(t, json.Unmarshal([]byte(restored), &filters))
|
||||
assert.Equal(t, oldName, filters[1]["key"])
|
||||
}
|
||||
|
||||
func TestRewriteQuickFilterSemconvNoop(t *testing.T) {
|
||||
input := `[{"key":"service.name","dataType":"string","type":"resource"}]`
|
||||
rewritten, changed, err := rewriteQuickFilterSemconv(input, deploymentEnvironmentOld(), deploymentEnvironmentCurrent)
|
||||
require.NoError(t, err)
|
||||
assert.False(t, changed)
|
||||
assert.Equal(t, input, rewritten)
|
||||
}
|
||||
@@ -13,12 +13,14 @@ import (
|
||||
)
|
||||
|
||||
type defaultConditionBuilder struct {
|
||||
fm qbtypes.FieldMapper
|
||||
// The builder composes family expressions, so it needs this package's
|
||||
// mapper, not the narrower qbtypes.FieldMapper.
|
||||
fm *defaultFieldMapper
|
||||
}
|
||||
|
||||
var _ qbtypes.ConditionBuilder = (*defaultConditionBuilder)(nil)
|
||||
|
||||
func NewConditionBuilder(fm qbtypes.FieldMapper) *defaultConditionBuilder {
|
||||
func NewConditionBuilder(fm *defaultFieldMapper) *defaultConditionBuilder {
|
||||
return &defaultConditionBuilder{fm: fm}
|
||||
}
|
||||
|
||||
@@ -44,16 +46,10 @@ func keyIndexFilter(key *telemetrytypes.TelemetryFieldKey) any {
|
||||
return fmt.Sprintf(`%%%s%%`, key.Name)
|
||||
}
|
||||
|
||||
func memberKey(key *telemetrytypes.TelemetryFieldKey, name string) *telemetrytypes.TelemetryFieldKey {
|
||||
member := *key
|
||||
member.Name = name
|
||||
return &member
|
||||
}
|
||||
|
||||
func keyIndexCondition(sb *sqlbuilder.SelectBuilder, column string, key *telemetrytypes.TelemetryFieldKey, members []string) string {
|
||||
func keyIndexCondition(sb *sqlbuilder.SelectBuilder, column string, members []*telemetrytypes.TelemetryFieldKey) string {
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
conditions = append(conditions, sb.Like(column, keyIndexFilter(memberKey(key, member))))
|
||||
conditions = append(conditions, sb.Like(column, keyIndexFilter(member)))
|
||||
}
|
||||
if len(conditions) == 1 {
|
||||
return conditions[0]
|
||||
@@ -64,15 +60,14 @@ func keyIndexCondition(sb *sqlbuilder.SelectBuilder, column string, key *telemet
|
||||
func valueIndexCondition(
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
column string,
|
||||
key *telemetrytypes.TelemetryFieldKey,
|
||||
members []string,
|
||||
members []*telemetrytypes.TelemetryFieldKey,
|
||||
op qbtypes.FilterOperator,
|
||||
value any,
|
||||
caseInsensitive bool,
|
||||
) string {
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
patterns := valueForIndexFilter(op, memberKey(key, member), value)
|
||||
patterns := valueForIndexFilter(op, member, value)
|
||||
switch values := patterns.(type) {
|
||||
case []string:
|
||||
for _, pattern := range values {
|
||||
@@ -92,10 +87,10 @@ func valueIndexCondition(
|
||||
return sb.Or(conditions...)
|
||||
}
|
||||
|
||||
func memberPresenceCondition(sb *sqlbuilder.SelectBuilder, column string, members []string, exists bool) string {
|
||||
func memberPresenceCondition(sb *sqlbuilder.SelectBuilder, column string, members []*telemetrytypes.TelemetryFieldKey, exists bool) string {
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
field := fmt.Sprintf("simpleJSONHas(%s, '%s')", column, member)
|
||||
field := fmt.Sprintf("simpleJSONHas(%s, %s)", column, querybuilder.ClickHouseStringLiteral(member.Name))
|
||||
if exists {
|
||||
conditions = append(conditions, sb.E(field, true))
|
||||
} else {
|
||||
@@ -124,7 +119,7 @@ func (b *defaultConditionBuilder) ConditionFor(
|
||||
value any,
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
) ([]string, []string, error) {
|
||||
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
|
||||
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
|
||||
|
||||
// has/hasAny/hasAll/hasToken are logs-body-only functions; they never apply to the
|
||||
// resource fingerprint table, so skip them (the main query still evaluates them).
|
||||
@@ -132,21 +127,21 @@ func (b *defaultConditionBuilder) ConditionFor(
|
||||
return nil, nil, nil
|
||||
}
|
||||
|
||||
keys, warning := querybuilder.ResolveKeys(key, matches)
|
||||
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
}
|
||||
|
||||
conds := make([]string, 0, len(keys))
|
||||
for _, k := range keys {
|
||||
// the resource fingerprint table only stores resource attributes; keys from
|
||||
conds := make([]string, 0, len(logicalFields))
|
||||
for _, logical := range logicalFields {
|
||||
// the resource fingerprint table only stores resource attributes; fields from
|
||||
// any other context contribute no condition and are omitted. An empty result
|
||||
// (including an unknown key) lets the caller skip this filter entirely.
|
||||
if k.FieldContext != telemetrytypes.FieldContextResource {
|
||||
if logical.FieldContext != telemetrytypes.FieldContextResource {
|
||||
continue
|
||||
}
|
||||
cond, err := b.conditionForKey(ctx, startNs, endNs, k, op, value, sb)
|
||||
cond, err := b.conditionForLogicalField(ctx, startNs, endNs, logical, op, value, sb)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
@@ -155,11 +150,11 @@ func (b *defaultConditionBuilder) ConditionFor(
|
||||
return conds, warnings, nil
|
||||
}
|
||||
|
||||
func (b *defaultConditionBuilder) conditionForKey(
|
||||
func (b *defaultConditionBuilder) conditionForLogicalField(
|
||||
ctx context.Context,
|
||||
startNs uint64,
|
||||
endNs uint64,
|
||||
key *telemetrytypes.TelemetryFieldKey,
|
||||
logical *telemetrytypes.LogicalField,
|
||||
op qbtypes.FilterOperator,
|
||||
value any,
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
@@ -169,7 +164,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
// as we store resource values as string
|
||||
formattedValue := querybuilder.FormatValueForContains(value)
|
||||
|
||||
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, key)
|
||||
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, logical.Single())
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -182,12 +177,12 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
// as we have not changed the resource column in the resource fingerprint table.
|
||||
column := columns[0]
|
||||
|
||||
members := resourceSemconvMembers(key)
|
||||
isFamily := len(members) > 1
|
||||
keyIdxFilter := keyIndexCondition(sb, column.Name, key, members)
|
||||
singleValueIndexFilter := valueForIndexFilter(op, memberKey(key, members[0]), value)
|
||||
members := logical.Members
|
||||
isFamily := logical.IsFamily()
|
||||
keyIdxFilter := keyIndexCondition(sb, column.Name, members)
|
||||
singleValueIndexFilter := valueForIndexFilter(op, members[0], value)
|
||||
|
||||
fieldName, err := b.fm.FieldFor(ctx, valuer.UUID{}, startNs, endNs, key)
|
||||
fieldName, err := b.fm.FieldForLogical(ctx, valuer.UUID{}, startNs, endNs, logical)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -197,7 +192,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
return sb.And(
|
||||
sb.E(fieldName, formattedValue),
|
||||
keyIdxFilter,
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, false),
|
||||
valueIndexCondition(sb, column.Name, members, op, value, false),
|
||||
), nil
|
||||
case qbtypes.FilterOperatorNotEqual:
|
||||
if isFamily {
|
||||
@@ -220,7 +215,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
return sb.And(
|
||||
sb.ILike(fieldName, formattedValue),
|
||||
keyIdxFilter,
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, true),
|
||||
valueIndexCondition(sb, column.Name, members, op, value, true),
|
||||
), nil
|
||||
case qbtypes.FilterOperatorNotLike, qbtypes.FilterOperatorNotILike:
|
||||
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
|
||||
@@ -260,7 +255,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
mainCondition = sb.And(
|
||||
mainCondition,
|
||||
keyIdxFilter,
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, false),
|
||||
valueIndexCondition(sb, column.Name, members, op, value, false),
|
||||
)
|
||||
|
||||
return mainCondition, nil
|
||||
@@ -308,7 +303,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
return sb.And(
|
||||
sb.ILike(fieldName, fmt.Sprintf(`%%%s%%`, formattedValue)),
|
||||
keyIdxFilter,
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, true),
|
||||
valueIndexCondition(sb, column.Name, members, op, value, true),
|
||||
), nil
|
||||
case qbtypes.FilterOperatorNotContains:
|
||||
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
|
||||
|
||||
@@ -221,24 +221,40 @@ func TestConditionBuilder(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestFamilyPositiveFilterExcludesKeylessRows(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
// The family tests drive resolution through the metadata map, exactly as
|
||||
// production does: two plain member keys in the map, one requested spelling.
|
||||
// The keys carry no family bookkeeping — grouping is the resolver's job.
|
||||
func familyConditionSQL(t *testing.T, requestedName string, memberNames []string, op qbtypes.FilterOperator, value any) (string, []any) {
|
||||
t.Helper()
|
||||
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{}
|
||||
for _, name := range memberNames {
|
||||
fieldKeys[name] = []*telemetrytypes.TelemetryFieldKey{{
|
||||
Name: name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}}
|
||||
}
|
||||
requested := telemetrytypes.NewTelemetryFieldKey(
|
||||
requestedName,
|
||||
telemetrytypes.FieldContextResource,
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
|
||||
context.Background(), valuer.UUID{}, 0, 0, requested,
|
||||
fieldKeys,
|
||||
qbtypes.ConditionBuilderOptions{}, op, value, sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
return sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
}
|
||||
|
||||
var deploymentFamilyMembers = []string{"deployment.environment.name", "deployment.environment"}
|
||||
|
||||
func TestFamilyPositiveFilterExcludesKeylessRows(t *testing.T) {
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorEqual, "production")
|
||||
|
||||
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)")
|
||||
assert.Equal(t, []any{
|
||||
@@ -251,166 +267,63 @@ func TestFamilyPositiveFilterExcludesKeylessRows(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestFamilyNotEqualIncludesKeylessRows(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotEqual, "staging", sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotEqual, "staging")
|
||||
|
||||
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ?")
|
||||
assert.Equal(t, []any{"staging"}, args)
|
||||
}
|
||||
|
||||
func TestFamilyNotInIncludesKeylessRows(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotIn, []any{"staging", "dev"}, sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotIn, []any{"staging", "dev"})
|
||||
|
||||
assert.Contains(t, sql, "(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ? AND COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ?)")
|
||||
assert.Equal(t, []any{"staging", "dev"}, args)
|
||||
}
|
||||
|
||||
func TestFamilyNotLikeIncludesKeylessRows(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotLike, "%stag%", sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotLike, "%stag%")
|
||||
|
||||
assert.Contains(t, sql, "LOWER(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')) NOT LIKE LOWER(?)")
|
||||
assert.Equal(t, []any{"%stag%"}, args)
|
||||
}
|
||||
|
||||
func TestFamilyNotContainsIncludesKeylessRows(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotContains, "stag", sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotContains, "stag")
|
||||
|
||||
assert.Contains(t, sql, "LOWER(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')) NOT LIKE LOWER(?)")
|
||||
assert.Equal(t, []any{"%stag%"}, args)
|
||||
}
|
||||
|
||||
func TestFamilyNotRegexpIncludesKeylessRows(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotRegexp, "stag.*", sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotRegexp, "stag.*")
|
||||
|
||||
assert.Contains(t, sql, "NOT match(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), ''), ?)")
|
||||
assert.Equal(t, []any{"stag.*"}, args)
|
||||
}
|
||||
|
||||
func TestFamilyExistsChecksEveryMember(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorExists, nil, sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorExists, nil)
|
||||
|
||||
assert.Contains(t, sql, "(simpleJSONHas(labels, 'deployment.environment.name') = ? OR simpleJSONHas(labels, 'deployment.environment') = ?) AND (labels LIKE ? OR labels LIKE ?)")
|
||||
assert.Equal(t, []any{true, true, "%deployment.environment.name%", "%deployment.environment%"}, args)
|
||||
}
|
||||
|
||||
func TestFamilyNotExistsChecksEveryMember(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotExists, nil, sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotExists, nil)
|
||||
|
||||
assert.Contains(t, sql, "simpleJSONHas(labels, 'deployment.environment.name') <> ? AND simpleJSONHas(labels, 'deployment.environment') <> ?")
|
||||
assert.Equal(t, []any{true, true}, args)
|
||||
}
|
||||
|
||||
// The old-name request with only the current spelling in metadata prunes to a
|
||||
// single member: plain single-key SQL, no coalesce.
|
||||
func TestFamilyPrunesToPresentMembers(t *testing.T) {
|
||||
sql, args := familyConditionSQL(t, "deployment.environment", []string{"deployment.environment.name"}, qbtypes.FilterOperatorEqual, "production")
|
||||
|
||||
assert.Contains(t, sql, "simpleJSONExtractString(labels, 'deployment.environment.name') = ? AND labels LIKE ? AND labels LIKE ?")
|
||||
assert.Equal(t, []any{"production", "%deployment.environment.name%", `%deployment.environment.name":"production%`}, args)
|
||||
}
|
||||
|
||||
func TestLogSemconvNameStaysLiteral(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
|
||||
@@ -6,7 +6,7 @@ import (
|
||||
"strings"
|
||||
|
||||
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
|
||||
"github.com/SigNoz/signoz/pkg/semconv"
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
@@ -34,18 +34,29 @@ func NewFieldMapper() *defaultFieldMapper {
|
||||
return &defaultFieldMapper{}
|
||||
}
|
||||
|
||||
func resourceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
|
||||
if key.Signal != telemetrytypes.SignalTraces || key.FieldContext != telemetrytypes.FieldContextResource {
|
||||
return []string{key.Name}
|
||||
// FieldForLogical returns the value expression for a resolved logical field:
|
||||
// the member's own expression for a single-member field, and a current-first
|
||||
// merge for a family. Resource label values are strings, so the merge is a
|
||||
// coalesce with a trailing '' that keeps single-key semantics for rows
|
||||
// without any member (see AddDefaultExistsFilter).
|
||||
func (m *defaultFieldMapper) FieldForLogical(
|
||||
ctx context.Context,
|
||||
orgID valuer.UUID,
|
||||
tsStart, tsEnd uint64,
|
||||
logical *telemetrytypes.LogicalField,
|
||||
) (string, error) {
|
||||
if !logical.IsFamily() {
|
||||
return m.FieldFor(ctx, orgID, tsStart, tsEnd, logical.Single())
|
||||
}
|
||||
if len(key.SemconvMembers) > 0 {
|
||||
return key.SemconvMembers
|
||||
values := make([]string, 0, len(logical.Members))
|
||||
for _, member := range logical.Members {
|
||||
expr, err := m.FieldFor(ctx, orgID, tsStart, tsEnd, member)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
values = append(values, fmt.Sprintf("NULLIF(%s, '')", expr))
|
||||
}
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
})
|
||||
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
|
||||
}
|
||||
|
||||
func (m *defaultFieldMapper) getColumn(
|
||||
@@ -82,15 +93,7 @@ func (m *defaultFieldMapper) FieldFor(
|
||||
return "", err
|
||||
}
|
||||
if key.FieldContext == telemetrytypes.FieldContextResource {
|
||||
members := resourceSemconvMembers(key)
|
||||
if len(members) > 1 {
|
||||
values := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
values = append(values, fmt.Sprintf("NULLIF(simpleJSONExtractString(%s, '%s'), '')", columns[0].Name, member))
|
||||
}
|
||||
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
|
||||
}
|
||||
return fmt.Sprintf("simpleJSONExtractString(%s, '%s')", columns[0].Name, members[0]), nil
|
||||
return fmt.Sprintf("simpleJSONExtractString(%s, %s)", columns[0].Name, querybuilder.ClickHouseStringLiteral(key.Name)), nil
|
||||
}
|
||||
return columns[0].Name, nil
|
||||
}
|
||||
|
||||
24
pkg/statementbuilder/resourcefilter/field_mapper_test.go
Normal file
24
pkg/statementbuilder/resourcefilter/field_mapper_test.go
Normal file
@@ -0,0 +1,24 @@
|
||||
package resourcefilter
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestFieldForQuotesRequestKeyName(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "name'\\); SELECT 1 --",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
}
|
||||
|
||||
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, "simpleJSONExtractString(labels, "+querybuilder.ClickHouseStringLiteral(key.Name)+")", expression)
|
||||
}
|
||||
@@ -39,7 +39,8 @@ func (c *conditionBuilder) ConditionFor(
|
||||
}
|
||||
|
||||
// an unknown key simply yields no condition rather than an error.
|
||||
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
|
||||
logicalFields, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
|
||||
keys := querybuilder.SingleKeys(logicalFields)
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
|
||||
@@ -1334,9 +1334,9 @@ func (t *telemetryMetaStore) GetKeysMulti(ctx context.Context, orgID valuer.UUID
|
||||
expandedTraceSelectors := make([]*telemetrytypes.FieldKeySelector, 0, len(tracesSelectors))
|
||||
for _, selector := range tracesSelectors {
|
||||
for _, member := range traceSemconvMembers(selector.Name, selector.FieldContext) {
|
||||
memberSelector := *selector
|
||||
memberSelector := selector.Copy()
|
||||
memberSelector.Name = member
|
||||
expandedTraceSelectors = append(expandedTraceSelectors, &memberSelector)
|
||||
expandedTraceSelectors = append(expandedTraceSelectors, memberSelector)
|
||||
}
|
||||
}
|
||||
tracesSelectors = expandedTraceSelectors
|
||||
|
||||
@@ -139,7 +139,8 @@ func (c *conditionBuilder) ConditionFor(
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
|
||||
logicalFields, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
|
||||
keys := querybuilder.SingleKeys(logicalFields)
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
|
||||
@@ -452,7 +452,7 @@ func (c *conditionBuilder) ConditionFor(
|
||||
value any,
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
) ([]string, []string, error) {
|
||||
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
|
||||
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
|
||||
skipResourceFilter := options.SkipResourceFilter
|
||||
|
||||
// search() resolves its own (optional) scope; handle it before key resolution.
|
||||
@@ -460,7 +460,8 @@ func (c *conditionBuilder) ConditionFor(
|
||||
return c.conditionForSearch(ctx, orgID, key, value, sb)
|
||||
}
|
||||
|
||||
keys, warning := querybuilder.ResolveKeys(key, matches)
|
||||
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
|
||||
keys := querybuilder.SingleKeys(logicalFields)
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
|
||||
@@ -162,7 +162,7 @@ func (c *conditionBuilder) ConditionFor(
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
keys := querybuilder.MatchingFieldKeys(key, fieldKeys)
|
||||
keys := querybuilder.SingleKeys(querybuilder.MatchingLogicalFields(key, fieldKeys))
|
||||
var warnings []string
|
||||
if len(keys) == 0 {
|
||||
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {
|
||||
|
||||
@@ -18,6 +18,8 @@ import (
|
||||
)
|
||||
|
||||
type conditionBuilder struct {
|
||||
// The builder composes family expressions, so it needs this package's
|
||||
// mapper, not the narrower qbtypes.FieldMapper.
|
||||
fm *fieldMapper
|
||||
}
|
||||
|
||||
@@ -32,7 +34,7 @@ func (c *conditionBuilder) conditionFor(
|
||||
orgID valuer.UUID,
|
||||
startNs uint64,
|
||||
endNs uint64,
|
||||
key *telemetrytypes.TelemetryFieldKey,
|
||||
logical *telemetrytypes.LogicalField,
|
||||
operator qbtypes.FilterOperator,
|
||||
value any,
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
@@ -42,13 +44,13 @@ func (c *conditionBuilder) conditionFor(
|
||||
value = querybuilder.FormatValueForContains(value)
|
||||
}
|
||||
|
||||
fieldExpression, err := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
|
||||
fieldExpression, err := c.fm.FieldForLogical(ctx, orgID, startNs, endNs, logical)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// TODO(srikanthccv): maybe extend this to every possible attribute
|
||||
if key.Name == "duration_nano" || key.Name == "durationNano" { // QoL improvement
|
||||
if logical.Name == "duration_nano" || logical.Name == "durationNano" { // QoL improvement
|
||||
switch v := value.(type) {
|
||||
case string:
|
||||
if duration, err := time.ParseDuration(v); err == nil {
|
||||
@@ -65,7 +67,7 @@ func (c *conditionBuilder) conditionFor(
|
||||
}
|
||||
}
|
||||
|
||||
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(key, value, fieldExpression, operator)
|
||||
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), value, fieldExpression, operator)
|
||||
|
||||
// regular operators
|
||||
switch operator {
|
||||
@@ -154,22 +156,7 @@ func (c *conditionBuilder) conditionFor(
|
||||
// in the query builder, `exists` and `not exists` are used for
|
||||
// key membership checks, so depending on the column type, the condition changes
|
||||
case qbtypes.FilterOperatorExists, qbtypes.FilterOperatorNotExists:
|
||||
// A semantic-convention family is represented by one current-first value
|
||||
// expression, but presence still has to inspect every physical member. In
|
||||
// particular, using ExistsExpression below with the requested key would add
|
||||
// a mapContains check for only that spelling and reject fallback-only rows.
|
||||
if isTraceSemconvFamily(key) {
|
||||
pred, err := c.fm.existsExpressionFor(ctx, orgID, startNs, endNs, key, operator == qbtypes.FilterOperatorExists)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return sqlbuilder.Escape(pred), nil
|
||||
}
|
||||
columns, err := c.fm.ColumnFor(ctx, orgID, startNs, endNs, key)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
pred, err := querybuilder.ExistsExpression(columns, key, startNs, endNs, fieldExpression, operator == qbtypes.FilterOperatorExists)
|
||||
pred, err := c.fm.ExistsForLogical(ctx, orgID, startNs, endNs, logical, operator == qbtypes.FilterOperatorExists)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -221,10 +208,10 @@ func (c *conditionBuilder) ConditionFor(
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
|
||||
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
|
||||
skipResourceFilter := options.SkipResourceFilter
|
||||
|
||||
keys, warning := querybuilder.ResolveKeys(key, matches)
|
||||
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
|
||||
var warnings []string
|
||||
if warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
@@ -232,10 +219,10 @@ func (c *conditionBuilder) ConditionFor(
|
||||
// A bare key that names a real column filters on the column too — first. When metadata
|
||||
// only knows the name under other contexts, prepend the column and keep metadata matches
|
||||
// only where their type is consistent with it (a corrupt entry can't degrade the column).
|
||||
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(keys) > 0 {
|
||||
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(logicalFields) > 0 {
|
||||
hasColumn := false
|
||||
for _, k := range keys {
|
||||
if k.FieldContext == telemetrytypes.FieldContextSpan {
|
||||
for _, logical := range logicalFields {
|
||||
if logical.FieldContext == telemetrytypes.FieldContextSpan {
|
||||
hasColumn = true
|
||||
break
|
||||
}
|
||||
@@ -243,49 +230,49 @@ func (c *conditionBuilder) ConditionFor(
|
||||
if !hasColumn {
|
||||
probe := telemetrytypes.NewTelemetryFieldKey(key.Name, telemetrytypes.FieldContextSpan, key.FieldDataType)
|
||||
if cols, colErr := c.fm.ColumnFor(ctx, orgID, startNs, endNs, probe); colErr == nil && len(cols) > 0 {
|
||||
combined := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys)+1)
|
||||
combined = append(combined, probe)
|
||||
for _, k := range keys {
|
||||
if columnMatchesDataType(cols[0], k.FieldDataType) {
|
||||
combined = append(combined, k)
|
||||
combined := make([]*telemetrytypes.LogicalField, 0, len(logicalFields)+1)
|
||||
combined = append(combined, telemetrytypes.SingleLogicalField(key.Name, probe))
|
||||
for _, logical := range logicalFields {
|
||||
if columnMatchesDataType(cols[0], logical.FieldDataType) {
|
||||
combined = append(combined, logical)
|
||||
}
|
||||
}
|
||||
keys = combined
|
||||
logicalFields = combined
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
synthesized := false
|
||||
if len(keys) == 0 {
|
||||
if len(logicalFields) == 0 {
|
||||
// Not in metadata. CandidateKeys resolves it: fold contexts (span/trace) get the
|
||||
// metadata map so it can honor a real column, correct to a stripped-name metadata
|
||||
// match, or synthesize; strict contexts pass nil and keep their synthesize path.
|
||||
keys = c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys))
|
||||
if len(keys) == 0 {
|
||||
logicalFields = querybuilder.WrapAsLogicalFields(key.Name, c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys)))
|
||||
if len(logicalFields) == 0 {
|
||||
return nil, warnings, querybuilder.NewKeyNotFoundError(key.Name)
|
||||
}
|
||||
synthesized = true
|
||||
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
|
||||
}
|
||||
|
||||
// When a resource sub-query already covers the term, drop resource keys from the main
|
||||
// When a resource sub-query already covers the term, drop resource fields from the main
|
||||
// query. Synthesized keys are exempt: the sub-query skips keys absent from metadata.
|
||||
if skipResourceFilter && !synthesized {
|
||||
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
|
||||
for _, k := range keys {
|
||||
if k.FieldContext != telemetrytypes.FieldContextResource {
|
||||
filtered = append(filtered, k)
|
||||
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
|
||||
for _, logical := range logicalFields {
|
||||
if logical.FieldContext != telemetrytypes.FieldContextResource {
|
||||
filtered = append(filtered, logical)
|
||||
}
|
||||
}
|
||||
if len(filtered) == 0 {
|
||||
return nil, warnings, nil
|
||||
}
|
||||
keys = filtered
|
||||
logicalFields = filtered
|
||||
}
|
||||
|
||||
conds := make([]string, 0, len(keys))
|
||||
for _, k := range keys {
|
||||
cond, err := c.conditionForKey(ctx, orgID, startNs, endNs, k, operator, value, sb)
|
||||
conds := make([]string, 0, len(logicalFields))
|
||||
for _, logical := range logicalFields {
|
||||
cond, err := c.conditionForLogicalField(ctx, orgID, startNs, endNs, logical, operator, value, sb)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
@@ -294,28 +281,28 @@ func (c *conditionBuilder) ConditionFor(
|
||||
return conds, warnings, nil
|
||||
}
|
||||
|
||||
func (c *conditionBuilder) conditionForKey(
|
||||
func (c *conditionBuilder) conditionForLogicalField(
|
||||
ctx context.Context,
|
||||
orgID valuer.UUID,
|
||||
startNs uint64,
|
||||
endNs uint64,
|
||||
key *telemetrytypes.TelemetryFieldKey,
|
||||
logical *telemetrytypes.LogicalField,
|
||||
operator qbtypes.FilterOperator,
|
||||
value any,
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
) (string, error) {
|
||||
if c.isSpanScopeField(key.Name) {
|
||||
return c.buildSpanScopeCondition(key, operator, value, startNs)
|
||||
if c.isSpanScopeField(logical.Name) {
|
||||
return c.buildSpanScopeCondition(logical.Single(), operator, value, startNs)
|
||||
}
|
||||
|
||||
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, operator, value, sb)
|
||||
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, operator, value, sb)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
if operator.AddDefaultExistsFilter() {
|
||||
// skip adding exists filter for intrinsic fields
|
||||
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
|
||||
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, logical.Single())
|
||||
if slices.Contains(maps.Keys(IntrinsicFields), field) ||
|
||||
slices.Contains(maps.Keys(IntrinsicFieldsDeprecated), field) ||
|
||||
slices.Contains(maps.Keys(CalculatedFields), field) ||
|
||||
@@ -323,7 +310,7 @@ func (c *conditionBuilder) conditionForKey(
|
||||
return condition, nil
|
||||
}
|
||||
|
||||
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, qbtypes.FilterOperatorExists, nil, sb)
|
||||
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, qbtypes.FilterOperatorExists, nil, sb)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
@@ -308,51 +308,72 @@ func TestConditionFor(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestConditionForSemconvFamilyPositiveFilterChecksPresence(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
// The family tests drive resolution through the metadata map, exactly as
|
||||
// production does: plain member keys in, the resolver groups them, and the
|
||||
// builder composes one condition per logical field.
|
||||
func traceFamilyConditionSQL(t *testing.T, requestedName string, members []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) (string, []any, []string) {
|
||||
t.Helper()
|
||||
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{}
|
||||
for _, member := range members {
|
||||
fieldKeys[member.Name] = []*telemetrytypes.TelemetryFieldKey{member}
|
||||
}
|
||||
requested := telemetrytypes.NewTelemetryFieldKey(
|
||||
requestedName,
|
||||
telemetrytypes.FieldContextAttribute,
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
|
||||
context.Background(), valuer.UUID{}, 0, 0, requested,
|
||||
fieldKeys,
|
||||
qbtypes.ConditionBuilderOptions{}, op, value, sb,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
return sql, args, warnings
|
||||
}
|
||||
|
||||
assert.Empty(t, warnings)
|
||||
assert.Contains(t, sql, "(COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '') = ? AND ((mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))))")
|
||||
func traceAttrMember(name string, materialized bool) *telemetrytypes.TelemetryFieldKey {
|
||||
return &telemetrytypes.TelemetryFieldKey{
|
||||
Name: name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Materialized: materialized,
|
||||
}
|
||||
}
|
||||
|
||||
func TestConditionForSemconvFamilyPositiveFilterChecksPresence(t *testing.T) {
|
||||
sql, args, warnings := traceFamilyConditionSQL(t,
|
||||
"deployment.environment.name",
|
||||
[]*telemetrytypes.TelemetryFieldKey{
|
||||
traceAttrMember("deployment.environment.name", false),
|
||||
traceAttrMember("deployment.environment", false),
|
||||
},
|
||||
qbtypes.FilterOperatorEqual, "production",
|
||||
)
|
||||
|
||||
assert.Empty(t, warnings, "a family is one logical field, never an ambiguity warning")
|
||||
assert.Contains(t, sql, "COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '') = ? AND (mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))")
|
||||
assert.Equal(t, []any{"production"}, args)
|
||||
}
|
||||
|
||||
func TestConditionForSemconvFamilyPreservesMaterializedMemberExistsColumn(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
SemconvMaterializedColumns: map[string]string{
|
||||
"deployment.environment": "attribute_string_deployment$$environment",
|
||||
},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
func TestNewConditionBuilderTakesThisPackagesMapper(t *testing.T) {
|
||||
// The builder composes family expressions, so it is deliberately coupled
|
||||
// to this package's mapper rather than the narrower qbtypes.FieldMapper.
|
||||
require.NotNil(t, NewConditionBuilder(NewFieldMapper()))
|
||||
}
|
||||
|
||||
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
|
||||
func TestConditionForSemconvFamilyPreservesMaterializedMemberExistsColumn(t *testing.T) {
|
||||
sql, args, warnings := traceFamilyConditionSQL(t,
|
||||
"deployment.environment.name",
|
||||
[]*telemetrytypes.TelemetryFieldKey{
|
||||
traceAttrMember("deployment.environment.name", false),
|
||||
traceAttrMember("deployment.environment", true),
|
||||
},
|
||||
qbtypes.FilterOperatorEqual, "production",
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
|
||||
assert.Empty(t, warnings)
|
||||
assert.Contains(t, sql, "`attribute_string_deployment$$environment_exists`")
|
||||
@@ -361,27 +382,17 @@ func TestConditionForSemconvFamilyPreservesMaterializedMemberExistsColumn(t *tes
|
||||
}
|
||||
|
||||
func TestConditionForSemconvFamilyNotExistsChecksEveryMember(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
}
|
||||
sb := sqlbuilder.NewSelectBuilder()
|
||||
|
||||
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
|
||||
context.Background(), valuer.UUID{}, 0, 0, key,
|
||||
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
|
||||
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorNotExists, nil, sb,
|
||||
sql, _, warnings := traceFamilyConditionSQL(t,
|
||||
"deployment.environment",
|
||||
[]*telemetrytypes.TelemetryFieldKey{
|
||||
traceAttrMember("deployment.environment.name", false),
|
||||
traceAttrMember("deployment.environment", false),
|
||||
},
|
||||
qbtypes.FilterOperatorNotExists, nil,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
sb.Where(conditions...)
|
||||
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
|
||||
assert.Empty(t, warnings)
|
||||
assert.Contains(t, sql, "NOT (((mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))))")
|
||||
assert.Empty(t, args)
|
||||
assert.Contains(t, sql, "NOT (mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))")
|
||||
}
|
||||
|
||||
func TestConditionForResourceWithEvolution(t *testing.T) {
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
"github.com/SigNoz/signoz/pkg/semconv"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
@@ -168,42 +167,143 @@ func NewFieldMapper() *fieldMapper {
|
||||
return &fieldMapper{}
|
||||
}
|
||||
|
||||
func traceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
|
||||
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
|
||||
return []string{key.Name}
|
||||
// FieldForLogical returns the value expression for a resolved logical field:
|
||||
// the member's own expression for a single-member field, and a current-first
|
||||
// merge across the members' expressions for a family. Each member expression
|
||||
// comes from FieldFor and therefore honors that member's own materialization
|
||||
// and evolution state — a family never needs sibling information on a key.
|
||||
func (m *fieldMapper) FieldForLogical(
|
||||
ctx context.Context,
|
||||
orgID valuer.UUID,
|
||||
tsStart, tsEnd uint64,
|
||||
logical *telemetrytypes.LogicalField,
|
||||
) (string, error) {
|
||||
if !logical.IsFamily() {
|
||||
return m.FieldFor(ctx, orgID, tsStart, tsEnd, logical.Single())
|
||||
}
|
||||
if len(key.SemconvMembers) > 0 {
|
||||
return key.SemconvMembers
|
||||
|
||||
memberExprs := make([]string, 0, len(logical.Members))
|
||||
for _, member := range logical.Members {
|
||||
expr, err := m.FieldFor(ctx, orgID, tsStart, tsEnd, member)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
memberExprs = append(memberExprs, expr)
|
||||
}
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: key.FieldContext,
|
||||
})
|
||||
|
||||
if logical.FieldDataType == telemetrytypes.FieldDataTypeString {
|
||||
// The trailing '' keeps single-key semantics for rows without any
|
||||
// member: string maps read '' for an absent key, and negative
|
||||
// operators must keep including such rows (see AddDefaultExistsFilter).
|
||||
values := make([]string, 0, len(memberExprs))
|
||||
for _, expr := range memberExprs {
|
||||
values = append(values, fmt.Sprintf("NULLIF(%s, '')", expr))
|
||||
}
|
||||
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
|
||||
}
|
||||
|
||||
// Numeric and boolean maps return zero for an absent key. If a family of
|
||||
// either type is enabled, this tail must become zero too.
|
||||
branches := make([]string, 0, len(logical.Members)*2)
|
||||
for i, member := range logical.Members {
|
||||
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, member, true)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
branches = append(branches, guard, memberExprs[i])
|
||||
}
|
||||
return "multiIf(" + strings.Join(branches, ", ") + ", NULL)", nil
|
||||
}
|
||||
|
||||
func isTraceSemconvFamily(key *telemetrytypes.TelemetryFieldKey) bool {
|
||||
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
|
||||
return false
|
||||
// ExistsForLogical renders the existence predicate for a resolved logical
|
||||
// field: a member's own predicate for a single-member field, presence of any
|
||||
// member for a family.
|
||||
func (m *fieldMapper) ExistsForLogical(
|
||||
ctx context.Context,
|
||||
orgID valuer.UUID,
|
||||
tsStart, tsEnd uint64,
|
||||
logical *telemetrytypes.LogicalField,
|
||||
exists bool,
|
||||
) (string, error) {
|
||||
if !logical.IsFamily() {
|
||||
return m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, logical.Single(), exists)
|
||||
}
|
||||
_, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: key.FieldContext,
|
||||
})
|
||||
return ok
|
||||
|
||||
guards := make([]string, 0, len(logical.Members))
|
||||
for _, member := range logical.Members {
|
||||
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, member, true)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
guards = append(guards, guard)
|
||||
}
|
||||
combined := "(" + strings.Join(guards, " OR ") + ")"
|
||||
if exists {
|
||||
return combined, nil
|
||||
}
|
||||
return "NOT " + combined, nil
|
||||
}
|
||||
|
||||
func traceSemconvMapMemberExpressions(columnName string, key *telemetrytypes.TelemetryFieldKey, member string) (string, string) {
|
||||
if materializedColumn, ok := key.SemconvMaterializedColumns[member]; ok {
|
||||
return fmt.Sprintf("`%s`", materializedColumn), fmt.Sprintf("`%s_exists`", materializedColumn)
|
||||
// logicalForResolvedColumn upgrades a directly-resolvable key (the FieldFor
|
||||
// probe succeeded) to its family when the metadata map proves membership;
|
||||
// otherwise the key stays a single-member logical field.
|
||||
func logicalForResolvedColumn(field *telemetrytypes.TelemetryFieldKey, keys map[string][]*telemetrytypes.TelemetryFieldKey) *telemetrytypes.LogicalField {
|
||||
for _, logical := range querybuilder.MatchingLogicalFields(field, keys) {
|
||||
if logical.IsFamily() &&
|
||||
logical.FieldContext == field.FieldContext &&
|
||||
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || logical.FieldDataType == field.FieldDataType) {
|
||||
return logical
|
||||
}
|
||||
}
|
||||
if key.Materialized && key.Name == member {
|
||||
physicalKey := *key
|
||||
physicalKey.Name = member
|
||||
return telemetrytypes.FieldKeyToMaterializedColumnName(&physicalKey), telemetrytypes.FieldKeyToMaterializedColumnNameForExists(&physicalKey)
|
||||
return telemetrytypes.SingleLogicalField(field.Name, field)
|
||||
}
|
||||
|
||||
// upgradeToFamilies swaps single-member candidates for their family when the
|
||||
// metadata map proves membership. Candidate order and every non-family
|
||||
// candidate stay exactly as the legacy flow produced them; sibling candidates
|
||||
// of an already-emitted family are dropped rather than duplicated.
|
||||
func upgradeToFamilies(field *telemetrytypes.TelemetryFieldKey, candidates []*telemetrytypes.LogicalField, keys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
|
||||
var families []*telemetrytypes.LogicalField
|
||||
for _, logical := range querybuilder.MatchingLogicalFields(field, keys) {
|
||||
if logical.IsFamily() {
|
||||
families = append(families, logical)
|
||||
}
|
||||
}
|
||||
return fmt.Sprintf("%s['%s']", columnName, member), fmt.Sprintf("mapContains(%s, '%s')", columnName, member)
|
||||
if len(families) == 0 {
|
||||
return candidates
|
||||
}
|
||||
|
||||
out := make([]*telemetrytypes.LogicalField, 0, len(candidates))
|
||||
emitted := make(map[*telemetrytypes.LogicalField]bool)
|
||||
for _, candidate := range candidates {
|
||||
var family *telemetrytypes.LogicalField
|
||||
for _, fam := range families {
|
||||
if fam.FieldContext != candidate.FieldContext || fam.FieldDataType != candidate.FieldDataType {
|
||||
continue
|
||||
}
|
||||
memberOfFamily := candidate.Single().Name == field.Name
|
||||
for _, member := range fam.Members {
|
||||
if member.Name == candidate.Single().Name {
|
||||
memberOfFamily = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if memberOfFamily {
|
||||
family = fam
|
||||
break
|
||||
}
|
||||
}
|
||||
if family == nil {
|
||||
out = append(out, candidate)
|
||||
continue
|
||||
}
|
||||
if emitted[family] {
|
||||
continue
|
||||
}
|
||||
emitted[family] = true
|
||||
out = append(out, family)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func (m *fieldMapper) getColumn(
|
||||
@@ -330,27 +430,10 @@ func (m *fieldMapper) resolveColumnExprs(
|
||||
if key.FieldContext != telemetrytypes.FieldContextResource {
|
||||
return nil, nil, nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "only resource context fields are supported for json columns, got %s", key.FieldContext.String)
|
||||
}
|
||||
members := traceSemconvMembers(key)
|
||||
if len(members) > 1 {
|
||||
values := make([]string, 0, len(members))
|
||||
guards := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
// The String cast is required because ClickHouse does not allow
|
||||
// Variant/Dynamic values in GROUP BY.
|
||||
value := fmt.Sprintf("%s.`%s`::String", columnName, member)
|
||||
values = append(values, fmt.Sprintf("NULLIF(%s, '')", value))
|
||||
guards = append(guards, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, member))
|
||||
}
|
||||
// Missing Dynamic paths are NULL, so this family expression must
|
||||
// retain the same NULL result as a single JSON-path lookup.
|
||||
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+")")
|
||||
existExprs = append(existExprs, "("+strings.Join(guards, " OR ")+")")
|
||||
} else {
|
||||
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
|
||||
// once ClickHouse is updated, check whether this cast can be removed.
|
||||
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, members[0]))
|
||||
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, members[0]))
|
||||
}
|
||||
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
|
||||
// once ClickHouse is updated, check whether this cast can be removed.
|
||||
exprs = append(exprs, fmt.Sprintf("%s.%s::String", columnName, querybuilder.ClickHouseIdentifier(key.Name)))
|
||||
existExprs = append(existExprs, fmt.Sprintf("%s.%s IS NOT NULL", columnName, querybuilder.ClickHouseIdentifier(key.Name)))
|
||||
case schema.ColumnTypeEnumString,
|
||||
schema.ColumnTypeEnumUInt64,
|
||||
schema.ColumnTypeEnumUInt32,
|
||||
@@ -375,40 +458,13 @@ func (m *fieldMapper) resolveColumnExprs(
|
||||
|
||||
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
|
||||
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumFloat64, schema.ColumnTypeEnumBool:
|
||||
members := traceSemconvMembers(key)
|
||||
if len(members) > 1 {
|
||||
guards := make([]string, 0, len(members))
|
||||
memberValues := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
valueExpression, existsExpression := traceSemconvMapMemberExpressions(columnName, key, member)
|
||||
memberValues = append(memberValues, valueExpression)
|
||||
guards = append(guards, existsExpression)
|
||||
}
|
||||
if valueType.GetType() == schema.ColumnTypeEnumString {
|
||||
values := make([]string, 0, len(members))
|
||||
for _, memberValue := range memberValues {
|
||||
values = append(values, fmt.Sprintf("NULLIF(%s, '')", memberValue))
|
||||
}
|
||||
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+", '')")
|
||||
} else {
|
||||
branches := make([]string, 0, len(members)*2)
|
||||
for i, memberValue := range memberValues {
|
||||
branches = append(branches, guards[i], memberValue)
|
||||
}
|
||||
// Numeric and boolean maps return zero for an absent key. If a
|
||||
// family of either type is enabled, this tail must become zero too.
|
||||
exprs = append(exprs, "multiIf("+strings.Join(branches, ", ")+", NULL)")
|
||||
}
|
||||
existExprs = append(existExprs, "("+strings.Join(guards, " OR ")+")")
|
||||
} else if key.Materialized {
|
||||
if key.Materialized {
|
||||
// a key could have been materialized, if so return the materialized column name
|
||||
physicalKey := *key
|
||||
physicalKey.Name = members[0]
|
||||
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(&physicalKey))
|
||||
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(&physicalKey))
|
||||
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
|
||||
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
|
||||
} else {
|
||||
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, members[0]))
|
||||
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, members[0]))
|
||||
exprs = append(exprs, fmt.Sprintf("%s[%s]", columnName, querybuilder.ClickHouseStringLiteral(key.Name)))
|
||||
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, %s)", columnName, querybuilder.ClickHouseStringLiteral(key.Name)))
|
||||
}
|
||||
default:
|
||||
return nil, nil, nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "value type %s is not supported for map column type %s", valueType, column.Type)
|
||||
@@ -431,18 +487,23 @@ func (m *fieldMapper) ColumnExpressionFor(
|
||||
keys map[string][]*telemetrytypes.TelemetryFieldKey,
|
||||
) (string, error) {
|
||||
|
||||
// Resolve the candidate column(s).
|
||||
var candidates []*telemetrytypes.TelemetryFieldKey
|
||||
// Resolve the candidate logical field(s).
|
||||
var candidates []*telemetrytypes.LogicalField
|
||||
switch _, err := m.FieldFor(ctx, orgID, startNs, endNs, field); {
|
||||
case err == nil:
|
||||
candidates = []*telemetrytypes.TelemetryFieldKey{field}
|
||||
// A directly-resolvable key upgrades to its family when the metadata
|
||||
// map proves membership; otherwise it stays single-member.
|
||||
candidates = []*telemetrytypes.LogicalField{logicalForResolvedColumn(field, keys)}
|
||||
case errors.Is(err, qbtypes.ErrColumnNotFound):
|
||||
// column (when the bare name is one) plus metadata matches, else synthesized
|
||||
// type-variant keys.
|
||||
candidates = m.CandidateKeys(ctx, orgID, field, nil, keys)
|
||||
if len(candidates) == 0 {
|
||||
// The legacy candidate flow, unchanged: column (when the bare name is
|
||||
// one) plus metadata matches, else synthesized type-variant keys. The
|
||||
// family step below only swaps candidates for their family; it never
|
||||
// changes candidate order or non-family behavior.
|
||||
raw := m.CandidateKeys(ctx, orgID, field, nil, keys)
|
||||
if len(raw) == 0 {
|
||||
return "", errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "field `%s` not found", field.Name).WithSuggestions(errors.NewSuggestionsOnLevenshteinDistance(field.Name, errors.NounKeys, maps.Keys(keys))...)
|
||||
}
|
||||
candidates = upgradeToFamilies(field, querybuilder.WrapAsLogicalFields(field.Name, raw), keys)
|
||||
default:
|
||||
return "", err
|
||||
}
|
||||
@@ -456,21 +517,21 @@ func (m *fieldMapper) ColumnExpressionFor(
|
||||
dummyValue = 0.0
|
||||
}
|
||||
stmts := make([]string, 0, len(candidates)*2)
|
||||
for _, key := range candidates {
|
||||
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
|
||||
for _, logical := range candidates {
|
||||
value, err := m.FieldForLogical(ctx, orgID, startNs, endNs, logical)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
|
||||
guard, err := m.ExistsForLogical(ctx, orgID, startNs, endNs, logical, true)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
coerced := value
|
||||
// a time column keeps its native type; coercing it would yield seconds
|
||||
if temporal, err := m.columnIsTemporal(ctx, startNs, endNs, key); err != nil {
|
||||
if temporal, err := m.logicalIsTemporal(ctx, startNs, endNs, logical); err != nil {
|
||||
return "", err
|
||||
} else if !temporal {
|
||||
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(key, dummyValue, value, qbtypes.FilterOperatorUnknown)
|
||||
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), dummyValue, value, qbtypes.FilterOperatorUnknown)
|
||||
}
|
||||
stmts = append(stmts, guard, coerced)
|
||||
}
|
||||
@@ -478,13 +539,14 @@ func (m *fieldMapper) ColumnExpressionFor(
|
||||
}
|
||||
|
||||
if len(candidates) == 1 {
|
||||
value, err := m.FieldFor(ctx, orgID, startNs, endNs, candidates[0])
|
||||
logical := candidates[0]
|
||||
value, err := m.FieldForLogical(ctx, orgID, startNs, endNs, logical)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
exprs, existExprs, _, _ := m.resolveColumnExprs(ctx, startNs, endNs, candidates[0])
|
||||
if len(exprs) == 1 && len(existExprs) == 1 {
|
||||
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, candidates[0], true)
|
||||
exprs, existExprs, _, _ := m.resolveColumnExprs(ctx, startNs, endNs, logical.Single())
|
||||
if !logical.IsFamily() && len(exprs) == 1 && len(existExprs) == 1 {
|
||||
guard, err := m.ExistsForLogical(ctx, orgID, startNs, endNs, logical, true)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -496,12 +558,12 @@ func (m *fieldMapper) ColumnExpressionFor(
|
||||
// Multiple candidates (collision / synth): multiIf picks the first that exists,
|
||||
// stringified so branches share a type.
|
||||
args := make([]string, 0, len(candidates))
|
||||
for _, key := range candidates {
|
||||
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
|
||||
for _, logical := range candidates {
|
||||
value, err := m.FieldForLogical(ctx, orgID, startNs, endNs, logical)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
|
||||
guard, err := m.ExistsForLogical(ctx, orgID, startNs, endNs, logical, true)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
@@ -510,6 +572,15 @@ func (m *fieldMapper) ColumnExpressionFor(
|
||||
return fmt.Sprintf("multiIf(%s, NULL)", strings.Join(args, ", ")), nil
|
||||
}
|
||||
|
||||
// logicalIsTemporal reports whether the logical field resolves to a single time
|
||||
// column. A family is attribute-backed and never temporal.
|
||||
func (m *fieldMapper) logicalIsTemporal(ctx context.Context, startNs, endNs uint64, logical *telemetrytypes.LogicalField) (bool, error) {
|
||||
if logical.IsFamily() {
|
||||
return false, nil
|
||||
}
|
||||
return m.columnIsTemporal(ctx, startNs, endNs, logical.Single())
|
||||
}
|
||||
|
||||
// columnIsTemporal reports whether key resolves to a single time column, after evolution
|
||||
// selection. Multiple columns mean an attribute-map union, which is never temporal.
|
||||
func (m *fieldMapper) columnIsTemporal(ctx context.Context, startNs, endNs uint64, key *telemetrytypes.TelemetryFieldKey) (bool, error) {
|
||||
@@ -612,25 +683,6 @@ func (m *fieldMapper) existsExpressionFor(
|
||||
key *telemetrytypes.TelemetryFieldKey,
|
||||
exists bool,
|
||||
) (string, error) {
|
||||
if isTraceSemconvFamily(key) {
|
||||
_, existExprs, _, err := m.resolveColumnExprs(ctx, tsStart, tsEnd, key)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if len(existExprs) == 0 {
|
||||
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "no existence expression found for field %s", key.Name)
|
||||
}
|
||||
parts := make([]string, 0, len(existExprs))
|
||||
for _, expression := range existExprs {
|
||||
parts = append(parts, "("+expression+")")
|
||||
}
|
||||
combined := strings.Join(parts, " OR ")
|
||||
if exists {
|
||||
return combined, nil
|
||||
}
|
||||
return "NOT (" + combined + ")", nil
|
||||
}
|
||||
|
||||
columns, err := m.getColumn(ctx, tsStart, tsEnd, key)
|
||||
if err != nil {
|
||||
return "", err
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
@@ -12,6 +13,20 @@ import (
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestFieldForQuotesRequestKeyNames(t *testing.T) {
|
||||
key := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "name'`\\); SELECT 1 --",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
|
||||
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Contains(t, expression, "resource."+querybuilder.ClickHouseIdentifier(key.Name))
|
||||
assert.Contains(t, expression, "mapContains(resources_string, "+querybuilder.ClickHouseStringLiteral(key.Name)+")")
|
||||
}
|
||||
|
||||
func TestGetFieldKeyName(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
|
||||
@@ -80,7 +95,7 @@ func TestGetFieldKeyName(t *testing.T) {
|
||||
Materialized: true,
|
||||
Evolutions: mockEvolution,
|
||||
},
|
||||
expectedResult: "multiIf((resource.`deployment.environment.name` IS NOT NULL OR resource.`deployment.environment` IS NOT NULL), COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, '')), (mapContains(resources_string, 'deployment.environment.name') OR `resource_string_deployment$$environment_exists`), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(`resource_string_deployment$$environment`, ''), ''), NULL)",
|
||||
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
|
||||
expectedError: nil,
|
||||
},
|
||||
{
|
||||
@@ -120,7 +135,9 @@ func TestGetFieldKeyName(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestFieldForResolvesCurrentTraceSemconvAttributeName(t *testing.T) {
|
||||
// FieldFor is a per-physical-key primitive: it never consults the family
|
||||
// table. Family composition is FieldForLogical's job.
|
||||
func TestFieldForIsPerKey(t *testing.T) {
|
||||
key := telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
@@ -129,54 +146,107 @@ func TestFieldForResolvesCurrentTraceSemconvAttributeName(t *testing.T) {
|
||||
|
||||
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '')", expression)
|
||||
}
|
||||
|
||||
func TestFieldForResolvesOldTraceSemconvAttributeName(t *testing.T) {
|
||||
key := telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
|
||||
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '')", expression)
|
||||
}
|
||||
|
||||
func TestFieldForPreservesResourceStorageDefaultsForSemconvFamily(t *testing.T) {
|
||||
key := telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Materialized: true,
|
||||
Evolutions: MockEvolutionData(time.Date(2024, 6, 2, 0, 0, 0, 0, time.UTC)),
|
||||
}
|
||||
start := uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano())
|
||||
end := uint64(time.Date(2024, 6, 5, 0, 0, 0, 0, time.UTC).UnixNano())
|
||||
|
||||
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, start, end, &key)
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "multiIf((resource.`deployment.environment.name` IS NOT NULL OR resource.`deployment.environment` IS NOT NULL), COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, '')), (`resource_string_deployment$$environment$$name_exists` OR mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(`resource_string_deployment$$environment$$name`, ''), NULLIF(resources_string['deployment.environment'], ''), ''), NULL)", expression)
|
||||
}
|
||||
|
||||
func TestFieldForUsesAvailableTraceSemconvMember(t *testing.T) {
|
||||
key := telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name"},
|
||||
}
|
||||
|
||||
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "attributes_string['deployment.environment.name']", expression)
|
||||
}
|
||||
|
||||
func traceFamilyLogicalField(t *testing.T, requestedName string, members ...*telemetrytypes.TelemetryFieldKey) *telemetrytypes.LogicalField {
|
||||
t.Helper()
|
||||
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{}
|
||||
for _, member := range members {
|
||||
fieldKeys[member.Name] = []*telemetrytypes.TelemetryFieldKey{member}
|
||||
}
|
||||
requested := telemetrytypes.NewTelemetryFieldKey(requestedName, members[0].FieldContext, members[0].FieldDataType)
|
||||
matches := querybuilder.MatchingLogicalFields(requested, fieldKeys)
|
||||
require.Len(t, matches, 1)
|
||||
return matches[0]
|
||||
}
|
||||
|
||||
func TestFieldForLogicalMergesFamilyMembersCurrentFirst(t *testing.T) {
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
old := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
|
||||
for _, requestedName := range []string{current.Name, old.Name} {
|
||||
logical := traceFamilyLogicalField(t, requestedName, current, old)
|
||||
expression, err := NewFieldMapper().FieldForLogical(context.Background(), valuer.UUID{}, 0, 0, logical)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t,
|
||||
"COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '')",
|
||||
expression,
|
||||
"both request spellings address one logical field",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// The old-name request with only the current spelling in metadata prunes to a
|
||||
// single member: plain single-key SQL, no coalesce.
|
||||
func TestFieldForLogicalPrunesToPresentMembers(t *testing.T) {
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
|
||||
logical := traceFamilyLogicalField(t, "deployment.environment", current)
|
||||
expression, err := NewFieldMapper().FieldForLogical(context.Background(), valuer.UUID{}, 0, 0, logical)
|
||||
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "attributes_string['deployment.environment.name']", expression)
|
||||
}
|
||||
|
||||
// Every member brings its own storage state: the merge composes each member's
|
||||
// FieldFor output, so a materialized-with-evolutions member keeps its column
|
||||
// history inside the family expression with no sibling bookkeeping anywhere.
|
||||
func TestFieldForLogicalComposesResourceMembersFromTheirOwnStorage(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
fm := NewFieldMapper()
|
||||
start := uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano())
|
||||
end := uint64(time.Date(2024, 6, 5, 0, 0, 0, 0, time.UTC).UnixNano())
|
||||
|
||||
current := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
old := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Materialized: true,
|
||||
Evolutions: MockEvolutionData(time.Date(2024, 6, 2, 0, 0, 0, 0, time.UTC)),
|
||||
}
|
||||
|
||||
currentExpr, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, current)
|
||||
require.NoError(t, err)
|
||||
oldExpr, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, old)
|
||||
require.NoError(t, err)
|
||||
|
||||
logical := traceFamilyLogicalField(t, current.Name, current, old)
|
||||
expression, err := fm.FieldForLogical(ctx, valuer.UUID{}, start, end, logical)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t,
|
||||
"COALESCE(NULLIF("+currentExpr+", ''), NULLIF("+oldExpr+", ''), '')",
|
||||
expression,
|
||||
"the family merge is exactly the members' own expressions, current-first, with the keyless tail",
|
||||
)
|
||||
assert.Contains(t, expression, "`resource_string_deployment$$environment`",
|
||||
"the promoted member keeps its materialized column")
|
||||
}
|
||||
|
||||
func TestFieldForResourceWithEvolution(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
releaseTime := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)
|
||||
@@ -233,7 +303,7 @@ func TestFieldForResourceWithEvolution(t *testing.T) {
|
||||
},
|
||||
tsStart: uint64(time.Date(2025, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
|
||||
tsEnd: uint64(time.Date(2025, 7, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
|
||||
expectedResult: "COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, ''))",
|
||||
expectedResult: "resource.`deployment.environment`::String",
|
||||
},
|
||||
{
|
||||
name: "Window straddles release - materialized resource",
|
||||
@@ -246,7 +316,7 @@ func TestFieldForResourceWithEvolution(t *testing.T) {
|
||||
},
|
||||
tsStart: uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
|
||||
tsEnd: uint64(time.Date(2025, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano()),
|
||||
expectedResult: "multiIf((resource.`deployment.environment.name` IS NOT NULL OR resource.`deployment.environment` IS NOT NULL), COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, '')), (mapContains(resources_string, 'deployment.environment.name') OR `resource_string_deployment$$environment_exists`), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(`resource_string_deployment$$environment`, ''), ''), NULL)",
|
||||
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -141,7 +141,7 @@ func NewSignalFilterFromStorableQuickFilter(storableQuickFilter *StorableQuickFi
|
||||
func NewDefaultQuickFilter(orgID valuer.UUID) ([]*StorableQuickFilter, error) {
|
||||
tracesFilters := []map[string]interface{}{
|
||||
{"key": "duration_nano", "dataType": "float64", "type": "tag"},
|
||||
{"key": "deployment.environment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "hasError", "dataType": "bool", "type": "tag"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "name", "dataType": "string", "type": "tag"},
|
||||
@@ -166,13 +166,13 @@ func NewDefaultQuickFilter(orgID valuer.UUID) ([]*StorableQuickFilter, error) {
|
||||
}
|
||||
|
||||
apiMonitoringFilters := []map[string]interface{}{
|
||||
{"key": "deployment.environment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "rpc.method", "dataType": "string", "type": "tag"},
|
||||
}
|
||||
|
||||
exceptionsFilters := []map[string]interface{}{
|
||||
{"key": "deployment.environment.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment", "dataType": "string", "type": "resource"},
|
||||
{"key": "service.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "host.name", "dataType": "string", "type": "resource"},
|
||||
{"key": "k8s.cluster.name", "dataType": "string", "type": "resource"},
|
||||
|
||||
@@ -1,37 +0,0 @@
|
||||
package quickfiltertypes
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestDefaultTraceQuickFiltersUseCurrentEnvironmentName(t *testing.T) {
|
||||
filters, err := NewDefaultQuickFilter(valuer.GenerateUUID())
|
||||
require.NoError(t, err)
|
||||
|
||||
traceSignals := map[string]bool{
|
||||
SignalTraces.StringValue(): true,
|
||||
SignalApiMonitoring.StringValue(): true,
|
||||
SignalExceptions.StringValue(): true,
|
||||
}
|
||||
for _, filter := range filters {
|
||||
if !traceSignals[filter.Signal.StringValue()] {
|
||||
continue
|
||||
}
|
||||
var keys []v3.AttributeKey
|
||||
require.NoError(t, json.Unmarshal([]byte(filter.Filter), &keys))
|
||||
found := false
|
||||
for _, key := range keys {
|
||||
if key.Key == "deployment.environment.name" {
|
||||
found = true
|
||||
}
|
||||
assert.NotEqual(t, "deployment.environment", key.Key)
|
||||
}
|
||||
assert.True(t, found, "missing environment quick filter for %s", filter.Signal.StringValue())
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,7 @@ package telemetrytypes
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"slices"
|
||||
"strings"
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
@@ -47,11 +48,75 @@ type TelemetryFieldKey struct {
|
||||
Indexes []TelemetryFieldKeySkipIndex `json:"-"`
|
||||
Materialized bool `json:"-"` // refers to promoted in case of body.... fields
|
||||
|
||||
Evolutions []*EvolutionEntry `json:"-"`
|
||||
SemconvMembers []string `json:"-"`
|
||||
// SemconvMaterializedColumns maps a physical family spelling to its
|
||||
// materialized column name. It is populated only on resolved query keys.
|
||||
SemconvMaterializedColumns map[string]string `json:"-"`
|
||||
Evolutions []*EvolutionEntry `json:"-"`
|
||||
}
|
||||
|
||||
// Copy returns an independent copy of f.
|
||||
func (f *TelemetryFieldKey) Copy() *TelemetryFieldKey {
|
||||
if f == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
copied := *f
|
||||
copied.Indexes = slices.Clone(f.Indexes)
|
||||
if f.Evolutions != nil {
|
||||
copied.Evolutions = make([]*EvolutionEntry, len(f.Evolutions))
|
||||
for index, evolution := range f.Evolutions {
|
||||
if evolution != nil {
|
||||
copiedEvolution := *evolution
|
||||
copied.Evolutions[index] = &copiedEvolution
|
||||
}
|
||||
}
|
||||
}
|
||||
copied.JSONPlan = copyJSONAccessPlan(f.JSONPlan, f, &copied)
|
||||
|
||||
return &copied
|
||||
}
|
||||
|
||||
func copyJSONAccessPlan(plan JSONAccessPlan, sourceKey, copiedKey *TelemetryFieldKey) JSONAccessPlan {
|
||||
if plan == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
nodes := make(map[*JSONAccessNode]*JSONAccessNode)
|
||||
var copyNode func(*JSONAccessNode) *JSONAccessNode
|
||||
copyNode = func(node *JSONAccessNode) *JSONAccessNode {
|
||||
if node == nil {
|
||||
return nil
|
||||
}
|
||||
if copiedNode, ok := nodes[node]; ok {
|
||||
return copiedNode
|
||||
}
|
||||
|
||||
copiedNode := *node
|
||||
nodes[node] = &copiedNode
|
||||
copiedNode.Parent = copyNode(node.Parent)
|
||||
if node.Branches != nil {
|
||||
copiedNode.Branches = make(map[JSONAccessBranchType]*JSONAccessNode, len(node.Branches))
|
||||
for branchType, branch := range node.Branches {
|
||||
copiedNode.Branches[branchType] = copyNode(branch)
|
||||
}
|
||||
}
|
||||
if node.TerminalConfig != nil {
|
||||
copiedTerminal := *node.TerminalConfig
|
||||
switch node.TerminalConfig.Key {
|
||||
case nil:
|
||||
case sourceKey:
|
||||
copiedTerminal.Key = copiedKey
|
||||
default:
|
||||
copiedTerminal.Key = node.TerminalConfig.Key.Copy()
|
||||
}
|
||||
copiedNode.TerminalConfig = &copiedTerminal
|
||||
}
|
||||
|
||||
return &copiedNode
|
||||
}
|
||||
|
||||
copied := make(JSONAccessPlan, len(plan))
|
||||
for index, node := range plan {
|
||||
copied[index] = copyNode(node)
|
||||
}
|
||||
return copied
|
||||
}
|
||||
|
||||
func (f *TelemetryFieldKey) KeyNameContainsArray() bool {
|
||||
@@ -132,8 +197,6 @@ func (f *TelemetryFieldKey) OverrideMetadataFrom(src *TelemetryFieldKey) {
|
||||
f.Materialized = src.Materialized
|
||||
f.JSONPlan = src.JSONPlan
|
||||
f.Evolutions = src.Evolutions
|
||||
f.SemconvMembers = src.SemconvMembers
|
||||
f.SemconvMaterializedColumns = src.SemconvMaterializedColumns
|
||||
}
|
||||
|
||||
func (f *TelemetryFieldKey) Equal(key *TelemetryFieldKey) bool {
|
||||
@@ -239,6 +302,15 @@ type MetricContext struct {
|
||||
MetricNamespace string `json:"metricNamespace,omitempty"`
|
||||
}
|
||||
|
||||
// Copy returns an independent copy of m.
|
||||
func (m *MetricContext) Copy() *MetricContext {
|
||||
if m == nil {
|
||||
return nil
|
||||
}
|
||||
copied := *m
|
||||
return &copied
|
||||
}
|
||||
|
||||
type FieldKeySelector struct {
|
||||
StartUnixMilli int64 `json:"startUnixMilli"`
|
||||
EndUnixMilli int64 `json:"endUnixMilli"`
|
||||
@@ -252,6 +324,16 @@ type FieldKeySelector struct {
|
||||
MetricContext *MetricContext `json:"metricContext,omitempty"`
|
||||
}
|
||||
|
||||
// Copy returns an independent copy of s.
|
||||
func (s *FieldKeySelector) Copy() *FieldKeySelector {
|
||||
if s == nil {
|
||||
return nil
|
||||
}
|
||||
copied := *s
|
||||
copied.MetricContext = s.MetricContext.Copy()
|
||||
return &copied
|
||||
}
|
||||
|
||||
type FieldValueSelector struct {
|
||||
*FieldKeySelector
|
||||
ExistingQuery string `json:"existingQuery"`
|
||||
|
||||
75
pkg/types/telemetrytypes/field_copy_test.go
Normal file
75
pkg/types/telemetrytypes/field_copy_test.go
Normal file
@@ -0,0 +1,75 @@
|
||||
package telemetrytypes
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestTelemetryFieldKeyCopyOwnsMutableState(t *testing.T) {
|
||||
original := &TelemetryFieldKey{
|
||||
Name: "items.name",
|
||||
FieldContext: FieldContextBody,
|
||||
FieldDataType: FieldDataTypeString,
|
||||
Indexes: []TelemetryFieldKeySkipIndex{
|
||||
{Name: "items.name"},
|
||||
},
|
||||
Evolutions: []*EvolutionEntry{
|
||||
{FieldName: "items.name"},
|
||||
},
|
||||
}
|
||||
require.NoError(t, original.SetJSONAccessPlan(JSONColumnMetadata{BaseColumn: "body_v2"}, nil))
|
||||
require.Len(t, original.JSONPlan, 1)
|
||||
require.NotNil(t, original.JSONPlan[0].TerminalConfig)
|
||||
|
||||
copied := original.Copy()
|
||||
require.NotNil(t, copied)
|
||||
require.NotSame(t, original, copied)
|
||||
require.Len(t, copied.JSONPlan, 1)
|
||||
require.NotNil(t, copied.JSONPlan[0].TerminalConfig)
|
||||
|
||||
assert.NotSame(t, original.JSONPlan[0], copied.JSONPlan[0])
|
||||
assert.NotSame(t, original.JSONPlan[0].Parent, copied.JSONPlan[0].Parent)
|
||||
assert.Same(t, copied, copied.JSONPlan[0].TerminalConfig.Key)
|
||||
assert.Equal(t, original.JSONPlan[0].Alias(), copied.JSONPlan[0].Alias())
|
||||
|
||||
copied.Name = "changed"
|
||||
copied.Indexes[0].Name = "changed"
|
||||
copied.Evolutions[0].FieldName = "changed"
|
||||
copied.JSONPlan[0].Name = "changed"
|
||||
copied.JSONPlan[0].Parent.Name = "changed"
|
||||
|
||||
assert.Equal(t, "items.name", original.Name)
|
||||
assert.Equal(t, "items.name", original.Indexes[0].Name)
|
||||
assert.Equal(t, "items.name", original.Evolutions[0].FieldName)
|
||||
assert.Equal(t, "items.name", original.JSONPlan[0].Name)
|
||||
assert.Equal(t, "body_v2", original.JSONPlan[0].Parent.Name)
|
||||
}
|
||||
|
||||
func TestFieldKeySelectorCopyOwnsMetricContext(t *testing.T) {
|
||||
original := &FieldKeySelector{
|
||||
Name: "state",
|
||||
MetricContext: &MetricContext{
|
||||
MetricName: "system.cpu.time",
|
||||
MetricNamespace: "system",
|
||||
},
|
||||
}
|
||||
|
||||
copied := original.Copy()
|
||||
require.NotNil(t, copied)
|
||||
require.NotNil(t, copied.MetricContext)
|
||||
assert.NotSame(t, original.MetricContext, copied.MetricContext)
|
||||
|
||||
copied.Name = "changed"
|
||||
copied.MetricContext.MetricName = "changed"
|
||||
|
||||
assert.Equal(t, "state", original.Name)
|
||||
assert.Equal(t, "system.cpu.time", original.MetricContext.MetricName)
|
||||
}
|
||||
|
||||
func TestNilFieldCopies(t *testing.T) {
|
||||
assert.Nil(t, (*TelemetryFieldKey)(nil).Copy())
|
||||
assert.Nil(t, (*FieldKeySelector)(nil).Copy())
|
||||
assert.Nil(t, (*MetricContext)(nil).Copy())
|
||||
}
|
||||
78
pkg/types/telemetrytypes/logical_field.go
Normal file
78
pkg/types/telemetrytypes/logical_field.go
Normal file
@@ -0,0 +1,78 @@
|
||||
package telemetrytypes
|
||||
|
||||
// LogicalField is resolution output: one queryable field, addressed by the
|
||||
// spelling the request used, backed by the physical member keys that store it.
|
||||
//
|
||||
// The resolver expresses ambiguity ("possibly different fields sharing a
|
||||
// name") as a []*LogicalField — never inside one LogicalField. Within one
|
||||
// LogicalField, members are alternate physical spellings of the same field
|
||||
// (a semantic-convention family), ordered current-first; compilers merge
|
||||
// them into one expression with current-wins precedence. Across the slice,
|
||||
// compilers build one condition per LogicalField and combine per the
|
||||
// operator, exactly as they previously combined ambiguous keys.
|
||||
//
|
||||
// Members always has at least one entry. A non-family field has exactly
|
||||
// one. Members alias the metadata map entries and must not be mutated.
|
||||
type LogicalField struct {
|
||||
// Name is the requested spelling. It is the response identity: aliases,
|
||||
// series labels, and warnings use it, so responses echo the request.
|
||||
Name string
|
||||
|
||||
// The physical identity every member shares. Members with a different
|
||||
// signal, field context, or data type belong to different logical
|
||||
// fields by definition.
|
||||
Signal Signal
|
||||
FieldContext FieldContext
|
||||
FieldDataType FieldDataType
|
||||
|
||||
// Members are the physical keys that store this field, ordered
|
||||
// current-first. Each member carries its own physical facts
|
||||
// (Materialized, Evolutions, JSONPlan, ...), so per-member accessors
|
||||
// need no sibling information.
|
||||
Members []*TelemetryFieldKey
|
||||
}
|
||||
|
||||
// SingleLogicalField wraps one physical key as its own logical field.
|
||||
func SingleLogicalField(name string, key *TelemetryFieldKey) *LogicalField {
|
||||
return &LogicalField{
|
||||
Name: name,
|
||||
Signal: key.Signal,
|
||||
FieldContext: key.FieldContext,
|
||||
FieldDataType: key.FieldDataType,
|
||||
Members: []*TelemetryFieldKey{key},
|
||||
}
|
||||
}
|
||||
|
||||
// Single returns the only member. It is the accessor for signals whose
|
||||
// logical fields are always single-member (everything except traces today).
|
||||
func (l *LogicalField) Single() *TelemetryFieldKey {
|
||||
return l.Members[0]
|
||||
}
|
||||
|
||||
// IsFamily reports whether the field has more than one physical member.
|
||||
func (l *LogicalField) IsFamily() bool {
|
||||
return len(l.Members) > 1
|
||||
}
|
||||
|
||||
// String implements fmt.Stringer for warning messages.
|
||||
func (l *LogicalField) String() string {
|
||||
if len(l.Members) == 1 {
|
||||
return l.Members[0].String()
|
||||
}
|
||||
names := make([]string, 0, len(l.Members))
|
||||
for _, member := range l.Members {
|
||||
names = append(names, member.Name)
|
||||
}
|
||||
return l.Name + "(" + l.FieldContext.StringValue() + ", " + l.FieldDataType.StringValue() + ", members: " + joinNames(names) + ")"
|
||||
}
|
||||
|
||||
func joinNames(names []string) string {
|
||||
out := ""
|
||||
for i, name := range names {
|
||||
if i > 0 {
|
||||
out += ", "
|
||||
}
|
||||
out += name
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -34,7 +34,6 @@ pytest_plugins = [
|
||||
"fixtures.role",
|
||||
"fixtures.savedview",
|
||||
"fixtures.seed_golden_dataset",
|
||||
"fixtures.semconv",
|
||||
]
|
||||
|
||||
|
||||
|
||||
66
tests/fixtures/semconv.py
vendored
66
tests/fixtures/semconv.py
vendored
@@ -1,66 +0,0 @@
|
||||
from collections.abc import Callable, Generator
|
||||
from datetime import UTC, datetime, timedelta
|
||||
|
||||
import pytest
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
|
||||
|
||||
SEMCONV_PHASE1_CURRENT = "deployment.environment.name"
|
||||
SEMCONV_PHASE1_OLD = "deployment.environment"
|
||||
SEMCONV_PHASE1_PREFIX = "semconv-phase1"
|
||||
|
||||
|
||||
@pytest.fixture(name="semconv_phase1_data")
|
||||
def semconv_phase1_data(
|
||||
insert_traces: Callable[[list[Traces]], None],
|
||||
clickhouse: types.TestContainerClickhouse,
|
||||
) -> Generator[datetime]:
|
||||
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=2)
|
||||
records = [
|
||||
(now - timedelta(seconds=5), "old", {SEMCONV_PHASE1_OLD: "production"}),
|
||||
(now - timedelta(seconds=4), "current", {SEMCONV_PHASE1_CURRENT: "production"}),
|
||||
(now - timedelta(seconds=3), "both", {SEMCONV_PHASE1_OLD: "production", SEMCONV_PHASE1_CURRENT: "production"}),
|
||||
(now - timedelta(seconds=2), "conflict", {SEMCONV_PHASE1_OLD: "staging", SEMCONV_PHASE1_CURRENT: "production"}),
|
||||
(now - timedelta(seconds=1), "staging", {SEMCONV_PHASE1_OLD: "staging"}),
|
||||
(now, "missing", {}),
|
||||
]
|
||||
traces = []
|
||||
for timestamp, suffix, environment in records:
|
||||
service = f"{SEMCONV_PHASE1_PREFIX}-{suffix}"
|
||||
traces.append(
|
||||
Traces(
|
||||
timestamp=timestamp,
|
||||
duration=timedelta(milliseconds=10),
|
||||
trace_id=TraceIdGenerator.trace_id(),
|
||||
span_id=TraceIdGenerator.span_id(),
|
||||
name=service,
|
||||
kind=TracesKind.SPAN_KIND_SERVER,
|
||||
status_code=TracesStatusCode.STATUS_CODE_OK,
|
||||
resources={"service.name": service, **environment},
|
||||
attributes=dict(environment),
|
||||
)
|
||||
)
|
||||
insert_traces(traces)
|
||||
|
||||
# Service-map rows are derived by the collector in production. Seed the
|
||||
# derived table directly so this test isolates the backend alias allowlist;
|
||||
# the collector repository owns its write-path integration test.
|
||||
for environment, suffix in (("production", "production"), ("staging", "staging")):
|
||||
clickhouse.conn.command(
|
||||
f"""
|
||||
INSERT INTO signoz_traces.distributed_dependency_graph_minutes_v2
|
||||
(src, dest, duration_quantiles_state, error_count, total_count, timestamp,
|
||||
deployment_environment, k8s_cluster_name, k8s_namespace_name)
|
||||
SELECT
|
||||
'{SEMCONV_PHASE1_PREFIX}-map-{suffix}', '{SEMCONV_PHASE1_PREFIX}-map-child',
|
||||
quantilesState(0.5, 0.75, 0.9, 0.95, 0.99)(toFloat64(1000000)),
|
||||
toUInt64(0), toUInt64(1), toDateTime({int(now.timestamp())}),
|
||||
'{environment}', '', ''
|
||||
"""
|
||||
)
|
||||
|
||||
yield now
|
||||
|
||||
cluster = clickhouse.env["SIGNOZ_TELEMETRYSTORE_CLICKHOUSE_CLUSTER"]
|
||||
clickhouse.conn.command(f"ALTER TABLE signoz_traces.dependency_graph_minutes_v2 ON CLUSTER '{cluster}' DELETE WHERE startsWith(src, '{SEMCONV_PHASE1_PREFIX}-map-') SETTINGS mutations_sync = 1")
|
||||
@@ -1,166 +0,0 @@
|
||||
"""Phase 1 end-to-end checks for semantic-convention name evolution."""
|
||||
|
||||
from collections.abc import Callable
|
||||
from datetime import datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
|
||||
import requests
|
||||
|
||||
from fixtures import querier, types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.semconv import SEMCONV_PHASE1_CURRENT as CURRENT
|
||||
from fixtures.semconv import SEMCONV_PHASE1_OLD as OLD
|
||||
from fixtures.semconv import SEMCONV_PHASE1_PREFIX as PREFIX
|
||||
|
||||
PRODUCTION_SPANS = {
|
||||
f"{PREFIX}-old",
|
||||
f"{PREFIX}-current",
|
||||
f"{PREFIX}-both",
|
||||
f"{PREFIX}-conflict",
|
||||
}
|
||||
STAGING_SPANS = {f"{PREFIX}-staging"}
|
||||
MISSING_SPANS = {f"{PREFIX}-missing"}
|
||||
|
||||
|
||||
def test_semconv_phase1_mixed_sdk_generations( # pylint: disable=too-many-statements
|
||||
signoz: types.SigNoz,
|
||||
create_user_admin: None, # pylint: disable=unused-argument
|
||||
get_token: Callable[[str, str], str],
|
||||
semconv_phase1_data: datetime,
|
||||
) -> None:
|
||||
now = semconv_phase1_data
|
||||
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
|
||||
start_ms = int((now - timedelta(minutes=2)).timestamp() * 1000)
|
||||
end_ms = int((now + timedelta(minutes=1)).timestamp() * 1000)
|
||||
|
||||
# Resource and span-attribute paths share the same matrix. Run every
|
||||
# operator with both the saved-query (old) and current request spellings.
|
||||
for context in ("resource", "attribute"):
|
||||
for requested in (CURRENT, OLD):
|
||||
field = f"{context}.{requested}"
|
||||
query_cases = {
|
||||
"production": (f"{field} = 'production'", PRODUCTION_SPANS),
|
||||
"staging": (f"{field} = 'staging'", STAGING_SPANS),
|
||||
# Negative operators intentionally include rows where no family
|
||||
# member exists; explicit EXISTS is the opt-in presence filter.
|
||||
"negative": (f"{field} != 'production'", STAGING_SPANS | MISSING_SPANS),
|
||||
"exists": (f"{field} EXISTS", PRODUCTION_SPANS | STAGING_SPANS),
|
||||
"not_exists": (f"{field} NOT EXISTS", MISSING_SPANS),
|
||||
}
|
||||
matrix_response = querier.make_query_request(
|
||||
signoz,
|
||||
token,
|
||||
start_ms=start_ms,
|
||||
end_ms=end_ms,
|
||||
request_type=querier.RequestType.RAW,
|
||||
queries=[
|
||||
querier.BuilderQuery(
|
||||
signal="traces",
|
||||
name=name,
|
||||
limit=100,
|
||||
filter_expression=expression,
|
||||
select_fields=[querier.TelemetryFieldKey("span.name")],
|
||||
order=[querier.OrderBy(querier.TelemetryFieldKey("timestamp"), "asc")],
|
||||
).to_dict()
|
||||
for name, (expression, _) in query_cases.items()
|
||||
],
|
||||
)
|
||||
assert matrix_response.status_code == HTTPStatus.OK, matrix_response.text
|
||||
matrix_results = matrix_response.json()["data"]["data"]["results"]
|
||||
for name, (_, expected_names) in query_cases.items():
|
||||
result = querier.find_named_result(matrix_results, name)
|
||||
assert result is not None, name
|
||||
assert {row["data"]["name"] for row in (result.get("rows") or [])} == expected_names
|
||||
|
||||
grouped_response = querier.make_query_request(
|
||||
signoz,
|
||||
token,
|
||||
start_ms=start_ms,
|
||||
end_ms=end_ms,
|
||||
request_type=querier.RequestType.SCALAR,
|
||||
queries=[
|
||||
querier.BuilderQuery(
|
||||
signal="traces",
|
||||
name="A",
|
||||
filter_expression=f"{field} EXISTS",
|
||||
aggregations=[querier.Aggregation("count()")],
|
||||
group_by=[querier.TelemetryFieldKey(requested, "string", context)],
|
||||
order=[querier.OrderBy(querier.TelemetryFieldKey(requested, "string", context), "asc")],
|
||||
).to_dict()
|
||||
],
|
||||
)
|
||||
assert grouped_response.status_code == HTTPStatus.OK, grouped_response.text
|
||||
grouped_results = grouped_response.json()["data"]["data"]["results"]
|
||||
assert len(grouped_results) == 1
|
||||
assert grouped_results[0]["columns"][0]["name"] == requested, "response identity must match the request spelling"
|
||||
assert grouped_results[0]["data"] == [["production", 4], ["staging", 1]]
|
||||
|
||||
for requested in (CURRENT, OLD):
|
||||
values_response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/fields/values"),
|
||||
timeout=5,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
params={
|
||||
"signal": "traces",
|
||||
"name": requested,
|
||||
"fieldContext": context,
|
||||
"fieldDataType": "string",
|
||||
},
|
||||
)
|
||||
assert values_response.status_code == HTTPStatus.OK, values_response.text
|
||||
assert set(values_response.json()["data"]["values"].get("stringValues") or []) == {"production", "staging"}
|
||||
|
||||
keys_response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/fields/keys"),
|
||||
timeout=5,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
params={"signal": "traces", "searchText": OLD},
|
||||
)
|
||||
assert keys_response.status_code == HTTPStatus.OK, keys_response.text
|
||||
keys = keys_response.json()["data"]["keys"]
|
||||
assert CURRENT in keys
|
||||
assert OLD in keys
|
||||
|
||||
start_ns = str(int((now - timedelta(minutes=2)).timestamp() * 1_000_000_000))
|
||||
end_ns = str(int((now + timedelta(minutes=1)).timestamp() * 1_000_000_000))
|
||||
for requested in (CURRENT, OLD):
|
||||
services_response = requests.post(
|
||||
signoz.self.host_configs["8080"].get("/api/v2/services"),
|
||||
timeout=30,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
json={
|
||||
"start": start_ns,
|
||||
"end": end_ns,
|
||||
"tags": [
|
||||
{
|
||||
"Key": requested,
|
||||
"Operator": "In",
|
||||
"StringValues": ["production"],
|
||||
"TagType": "ResourceAttribute",
|
||||
}
|
||||
],
|
||||
},
|
||||
)
|
||||
assert services_response.status_code == HTTPStatus.OK, services_response.text
|
||||
services = {item["serviceName"] for item in services_response.json()["data"]}
|
||||
assert services == PRODUCTION_SPANS
|
||||
|
||||
map_response = requests.post(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/dependency_graph"),
|
||||
timeout=30,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
json={
|
||||
"start": start_ns,
|
||||
"end": end_ns,
|
||||
"tags": [
|
||||
{
|
||||
"key": requested,
|
||||
"operator": "In",
|
||||
"stringValues": ["production"],
|
||||
"tagType": "ResourceAttribute",
|
||||
}
|
||||
],
|
||||
},
|
||||
)
|
||||
assert map_response.status_code == HTTPStatus.OK, map_response.text
|
||||
assert {edge["parent"] for edge in map_response.json()} == {f"{PREFIX}-map-production"}
|
||||
Reference in New Issue
Block a user