mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-06 21:20:42 +01:00
Compare commits
3 Commits
feat/semco
...
test/semco
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
487d00b6c9 | ||
|
|
59329d3aaa | ||
|
|
7c58b6af49 |
4
Makefile
4
Makefile
@@ -220,6 +220,10 @@ 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"
|
||||
|
||||
@@ -6,6 +6,8 @@ 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{
|
||||
@@ -29,13 +31,61 @@ 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 := fmt.Sprintf("simpleJSONExtractString(labels, '%s')", key)
|
||||
searchKey := resourceValueExpression(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)
|
||||
|
||||
@@ -43,9 +93,9 @@ func buildResourceFilter(logsOp string, key string, op v3.FilterOperator, value
|
||||
|
||||
switch op {
|
||||
case v3.FilterOperatorExists:
|
||||
return fmt.Sprintf("simpleJSONHas(labels, '%s')", key)
|
||||
return resourcePresenceExpression(key, true)
|
||||
case v3.FilterOperatorNotExists:
|
||||
return fmt.Sprintf("not simpleJSONHas(labels, '%s')", key)
|
||||
return resourcePresenceExpression(key, false)
|
||||
case v3.FilterOperatorRegex, v3.FilterOperatorNotRegex:
|
||||
return fmt.Sprintf(logsOp, searchKey, chFmtVal)
|
||||
case v3.FilterOperatorContains, v3.FilterOperatorNotContains:
|
||||
@@ -110,6 +160,38 @@ 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)
|
||||
@@ -206,14 +288,31 @@ func buildResourceFiltersFromGroupBy(groupBy []v3.AttributeKey) []string {
|
||||
if attr.Type != v3.AttributeKeyTypeResource {
|
||||
continue
|
||||
}
|
||||
conditions = append(conditions, fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", attr.Key, attr.Key))
|
||||
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 ")))
|
||||
}
|
||||
return conditions
|
||||
}
|
||||
|
||||
func buildResourceFiltersFromAggregateAttribute(aggregateAttribute v3.AttributeKey) string {
|
||||
if aggregateAttribute.Key != "" && aggregateAttribute.Type == v3.AttributeKeyTypeResource {
|
||||
return fmt.Sprintf("(simpleJSONHas(labels, '%s') AND labels like '%%%s%%')", aggregateAttribute.Key, aggregateAttribute.Key)
|
||||
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 ""
|
||||
|
||||
@@ -5,6 +5,8 @@ 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) {
|
||||
@@ -552,3 +554,38 @@ 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,16 +6,32 @@ 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 = map[string]struct{}{
|
||||
"deployment_environment": {},
|
||||
"k8s_cluster_name": {},
|
||||
"k8s_namespace_name": {},
|
||||
}
|
||||
columns = serviceMapColumns()
|
||||
)
|
||||
|
||||
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{}
|
||||
@@ -24,39 +40,40 @@ func BuildServiceMapQuery(tags []model.TagQuery) (string, []interface{}) {
|
||||
operator := tag.GetOperator()
|
||||
value := tag.GetValues()
|
||||
|
||||
if _, ok := columns[key]; !ok {
|
||||
column, ok := columns[key]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
switch operator {
|
||||
case model.InOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s IN @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s IN @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.NotInOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT IN @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT IN @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.EqualOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s = @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s = @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.NotEqualOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s != @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s != @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, value))
|
||||
case model.ContainsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%%%s%%", value)))
|
||||
case model.NotContainsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%%%s%%", value)))
|
||||
case model.StartsWithOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s LIKE @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%s%%", value)))
|
||||
case model.NotStartsWithOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", key, key)
|
||||
filterQuery += fmt.Sprintf(" AND %s NOT LIKE @%s", column, key)
|
||||
namedArgs = append(namedArgs, clickhouse.Named(key, fmt.Sprintf("%s%%", value)))
|
||||
case model.ExistsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NOT NULL", key)
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NOT NULL", column)
|
||||
case model.NotExistsOperator:
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NULL", key)
|
||||
filterQuery += fmt.Sprintf(" AND %s IS NULL", column)
|
||||
}
|
||||
}
|
||||
return filterQuery, namedArgs
|
||||
|
||||
35
pkg/query-service/app/services/map_test.go
Normal file
35
pkg/query-service/app/services/map_test.go
Normal file
@@ -0,0 +1,35 @@
|
||||
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)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
|
||||
"github.com/SigNoz/signoz/pkg/errors"
|
||||
grammar "github.com/SigNoz/signoz/pkg/parser/filterquery/grammar"
|
||||
"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"
|
||||
@@ -982,25 +983,78 @@ 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 {
|
||||
fieldKeysForName := []*telemetrytypes.TelemetryFieldKey{}
|
||||
selector := telemetrytypes.FieldKeySelector{
|
||||
Name: field.Name,
|
||||
Signal: field.Signal,
|
||||
FieldContext: field.FieldContext,
|
||||
}
|
||||
members := semconv.Members(semconv.KindAttribute, selector)
|
||||
isFamily := len(members) > 1
|
||||
fieldKeysForName := make([]*telemetrytypes.TelemetryFieldKey, 0)
|
||||
indexByIdentity := make(map[string]int)
|
||||
|
||||
// match by name; keep items whose context and data type match (unspecified matches any)
|
||||
for _, item := range fieldKeys[field.Name] {
|
||||
if (field.FieldContext == telemetrytypes.FieldContextUnspecified || field.FieldContext == item.FieldContext) &&
|
||||
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || field.FieldDataType == item.FieldDataType) {
|
||||
fieldKeysForName = append(fieldKeysForName, item)
|
||||
appendMatches := func(lookupName string, memberName string, contextAlreadyMatched bool) {
|
||||
for _, item := range fieldKeys[lookupName] {
|
||||
if !contextAlreadyMatched && field.FieldContext != telemetrytypes.FieldContextUnspecified && field.FieldContext != item.FieldContext {
|
||||
continue
|
||||
}
|
||||
if field.FieldDataType != telemetrytypes.FieldDataTypeUnspecified && field.FieldDataType != item.FieldDataType {
|
||||
continue
|
||||
}
|
||||
|
||||
// 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.
|
||||
if memberName != field.Name {
|
||||
itemSelector := telemetrytypes.FieldKeySelector{
|
||||
Name: field.Name,
|
||||
Signal: item.Signal,
|
||||
FieldContext: item.FieldContext,
|
||||
}
|
||||
if !slices.Contains(semconv.Members(semconv.KindAttribute, itemSelector), memberName) {
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
physicalMembers := item.SemconvMembers
|
||||
if len(physicalMembers) == 0 {
|
||||
physicalMembers = []string{memberName}
|
||||
}
|
||||
identity := item.Signal.StringValue() + ";" + item.FieldContext.StringValue() + ";" + item.FieldDataType.StringValue()
|
||||
if isFamily {
|
||||
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)
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
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 isFamily {
|
||||
resolved.Name = field.Name
|
||||
resolved.SemconvMembers = slices.Clone(physicalMembers)
|
||||
}
|
||||
fieldKeysForName = append(fieldKeysForName, &resolved)
|
||||
}
|
||||
}
|
||||
|
||||
// A context may have been split off a name that legitimately contained it (e.g.
|
||||
// `attribute.key`); also look up the context-prefixed name so both readings resolve.
|
||||
// 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)
|
||||
}
|
||||
|
||||
// A context may have been split off a name that legitimately contained it
|
||||
// (e.g. `attribute.key`); preserve that historical alternate reading for
|
||||
// every family member.
|
||||
if field.FieldContext != telemetrytypes.FieldContextUnspecified {
|
||||
contextPrefixedFieldName := fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), field.Name)
|
||||
for _, item := range fieldKeys[contextPrefixedFieldName] {
|
||||
// Context already matched via the lookup key; only data type needs checking.
|
||||
if field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || item.FieldDataType == field.FieldDataType {
|
||||
fieldKeysForName = append(fieldKeysForName, item)
|
||||
}
|
||||
for _, member := range members {
|
||||
appendMatches(fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), member), member, true)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -14,6 +14,7 @@ import (
|
||||
"github.com/antlr4-go/antlr/v4"
|
||||
sqlbuilder "github.com/huandu/go-sqlbuilder"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestPrepareWhereClause_EmptyVariableList ensures PrepareWhereClause errors when a variable has an empty list value.
|
||||
@@ -685,6 +686,53 @@ func TestVisitKey(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestMatchingFieldKeysResolvesSemconvFamily(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,
|
||||
}
|
||||
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
|
||||
current.Name: {current},
|
||||
old.Name: {old},
|
||||
}
|
||||
|
||||
for _, requestedName := range []string{current.Name, old.Name} {
|
||||
requested := telemetrytypes.NewTelemetryFieldKey(
|
||||
requestedName,
|
||||
telemetrytypes.FieldContextResource,
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
matches := MatchingFieldKeys(requested, fieldKeys)
|
||||
require.Len(t, matches, 1, "family lookup must return one field before its metadata is inspected")
|
||||
assert.Equal(t, requestedName, matches[0].Name)
|
||||
assert.Equal(t, "current metadata", matches[0].Description)
|
||||
assert.Equal(t, []string{current.Name, old.Name}, matches[0].SemconvMembers)
|
||||
}
|
||||
|
||||
// A current-name query still resolves when metadata has seen only the old
|
||||
// spelling. The returned name remains the request identity.
|
||||
requested := telemetrytypes.NewTelemetryFieldKey(
|
||||
current.Name,
|
||||
telemetrytypes.FieldContextResource,
|
||||
telemetrytypes.FieldDataTypeString,
|
||||
)
|
||||
matches := MatchingFieldKeys(requested, map[string][]*telemetrytypes.TelemetryFieldKey{old.Name: {old}})
|
||||
require.Len(t, matches, 1, "family lookup must return one field before its metadata is inspected")
|
||||
assert.Equal(t, current.Name, matches[0].Name)
|
||||
assert.Equal(t, "old metadata", matches[0].Description)
|
||||
assert.Equal(t, []string{old.Name}, matches[0].SemconvMembers)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// TestVisitComparison
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -235,6 +235,7 @@ func NewSQLMigrationProviderFactories(
|
||||
sqlmigration.NewFillDashboardSpecCollectionsFactory(sqlstore, dashboardStore),
|
||||
sqlmigration.NewScrubEmailChannelTransportFactory(sqlstore),
|
||||
sqlmigration.NewAddDashboardTuplesFactory(sqlstore),
|
||||
sqlmigration.NewMigrateDeploymentEnvironmentQuickFilterFactory(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,127 @@
|
||||
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())
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
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)
|
||||
}
|
||||
@@ -44,6 +44,73 @@ 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 {
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
conditions = append(conditions, sb.Like(column, keyIndexFilter(memberKey(key, member))))
|
||||
}
|
||||
if len(conditions) == 1 {
|
||||
return conditions[0]
|
||||
}
|
||||
return sb.Or(conditions...)
|
||||
}
|
||||
|
||||
func valueIndexCondition(
|
||||
sb *sqlbuilder.SelectBuilder,
|
||||
column string,
|
||||
key *telemetrytypes.TelemetryFieldKey,
|
||||
members []string,
|
||||
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)
|
||||
switch values := patterns.(type) {
|
||||
case []string:
|
||||
for _, pattern := range values {
|
||||
conditions = append(conditions, sb.Like(column, pattern))
|
||||
}
|
||||
default:
|
||||
if caseInsensitive {
|
||||
conditions = append(conditions, sb.ILike(column, values))
|
||||
} else {
|
||||
conditions = append(conditions, sb.Like(column, values))
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(conditions) == 1 {
|
||||
return conditions[0]
|
||||
}
|
||||
return sb.Or(conditions...)
|
||||
}
|
||||
|
||||
func memberPresenceCondition(sb *sqlbuilder.SelectBuilder, column string, members []string, exists bool) string {
|
||||
conditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
field := fmt.Sprintf("simpleJSONHas(%s, '%s')", column, member)
|
||||
if exists {
|
||||
conditions = append(conditions, sb.E(field, true))
|
||||
} else {
|
||||
conditions = append(conditions, sb.NE(field, true))
|
||||
}
|
||||
}
|
||||
if exists {
|
||||
if len(conditions) == 1 {
|
||||
return conditions[0]
|
||||
}
|
||||
return sb.Or(conditions...)
|
||||
}
|
||||
return sb.And(conditions...)
|
||||
}
|
||||
|
||||
// SkipResourceFilter is not applicable here: the fingerprint table only stores resource attributes.
|
||||
func (b *defaultConditionBuilder) ConditionFor(
|
||||
ctx context.Context,
|
||||
@@ -115,8 +182,10 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
// as we have not changed the resource column in the resource fingerprint table.
|
||||
column := columns[0]
|
||||
|
||||
keyIdxFilter := sb.Like(column.Name, keyIndexFilter(key))
|
||||
valueForIndexFilter := valueForIndexFilter(op, key, value)
|
||||
members := resourceSemconvMembers(key)
|
||||
isFamily := len(members) > 1
|
||||
keyIdxFilter := keyIndexCondition(sb, column.Name, key, members)
|
||||
singleValueIndexFilter := valueForIndexFilter(op, memberKey(key, members[0]), value)
|
||||
|
||||
fieldName, err := b.fm.FieldFor(ctx, valuer.UUID{}, startNs, endNs, key)
|
||||
if err != nil {
|
||||
@@ -128,12 +197,15 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
return sb.And(
|
||||
sb.E(fieldName, formattedValue),
|
||||
keyIdxFilter,
|
||||
sb.Like(column.Name, valueForIndexFilter),
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, false),
|
||||
), nil
|
||||
case qbtypes.FilterOperatorNotEqual:
|
||||
if isFamily {
|
||||
return sb.NE(fieldName, formattedValue), nil
|
||||
}
|
||||
return sb.And(
|
||||
sb.NE(fieldName, formattedValue),
|
||||
sb.NotLike(column.Name, valueForIndexFilter),
|
||||
sb.NotLike(column.Name, singleValueIndexFilter),
|
||||
), nil
|
||||
case qbtypes.FilterOperatorGreaterThan:
|
||||
return sb.And(sb.GT(fieldName, formattedValue), keyIdxFilter), nil
|
||||
@@ -148,7 +220,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
return sb.And(
|
||||
sb.ILike(fieldName, formattedValue),
|
||||
keyIdxFilter,
|
||||
sb.ILike(column.Name, valueForIndexFilter),
|
||||
valueIndexCondition(sb, column.Name, key, 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
|
||||
@@ -185,13 +257,11 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
inConditions = append(inConditions, sb.E(fieldName, querybuilder.FormatValueForContains(v)))
|
||||
}
|
||||
mainCondition := sb.Or(inConditions...)
|
||||
valConditions := make([]string, 0, len(values))
|
||||
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
|
||||
for _, v := range valuesForIndexFilter {
|
||||
valConditions = append(valConditions, sb.Like(column.Name, v))
|
||||
}
|
||||
}
|
||||
mainCondition = sb.And(mainCondition, keyIdxFilter, sb.Or(valConditions...))
|
||||
mainCondition = sb.And(
|
||||
mainCondition,
|
||||
keyIdxFilter,
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, false),
|
||||
)
|
||||
|
||||
return mainCondition, nil
|
||||
case qbtypes.FilterOperatorNotIn:
|
||||
@@ -204,8 +274,11 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
notInConditions = append(notInConditions, sb.NE(fieldName, querybuilder.FormatValueForContains(v)))
|
||||
}
|
||||
mainCondition := sb.And(notInConditions...)
|
||||
if isFamily {
|
||||
return mainCondition, nil
|
||||
}
|
||||
valConditions := make([]string, 0, len(values))
|
||||
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
|
||||
if valuesForIndexFilter, ok := singleValueIndexFilter.([]string); ok {
|
||||
for _, v := range valuesForIndexFilter {
|
||||
valConditions = append(valConditions, sb.NotLike(column.Name, v))
|
||||
}
|
||||
@@ -215,13 +288,11 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
|
||||
case qbtypes.FilterOperatorExists:
|
||||
return sb.And(
|
||||
sb.E(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
|
||||
memberPresenceCondition(sb, column.Name, members, true),
|
||||
keyIdxFilter,
|
||||
), nil
|
||||
case qbtypes.FilterOperatorNotExists:
|
||||
return sb.And(
|
||||
sb.NE(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
|
||||
), nil
|
||||
return memberPresenceCondition(sb, column.Name, members, false), nil
|
||||
|
||||
case qbtypes.FilterOperatorRegexp:
|
||||
return sb.And(
|
||||
@@ -237,7 +308,7 @@ func (b *defaultConditionBuilder) conditionForKey(
|
||||
return sb.And(
|
||||
sb.ILike(fieldName, fmt.Sprintf(`%%%s%%`, formattedValue)),
|
||||
keyIdxFilter,
|
||||
sb.ILike(column.Name, valueForIndexFilter),
|
||||
valueIndexCondition(sb, column.Name, key, members, op, value, true),
|
||||
), nil
|
||||
case qbtypes.FilterOperatorNotContains:
|
||||
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
|
||||
|
||||
@@ -198,6 +198,75 @@ func TestConditionBuilder(t *testing.T) {
|
||||
expected: "match(simpleJSONExtractString(labels, 'k8s.namespace.name'), ?) AND labels LIKE ?",
|
||||
expectedArgs: []any{"ban.*", "%k8s.namespace.name%"},
|
||||
},
|
||||
{
|
||||
name: "semantic convention family equality uses current-first fallback",
|
||||
key: &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
},
|
||||
op: qbtypes.FilterOperatorEqual,
|
||||
value: "production",
|
||||
expected: "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), '')) = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)",
|
||||
expectedArgs: []any{
|
||||
"production",
|
||||
"%deployment.environment.name%",
|
||||
"%deployment.environment%",
|
||||
`%deployment.environment.name":"production%`,
|
||||
`%deployment.environment":"production%`,
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "old semantic convention request uses only current metadata member",
|
||||
key: &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name"},
|
||||
},
|
||||
op: qbtypes.FilterOperatorEqual,
|
||||
value: "production",
|
||||
expected: "simpleJSONExtractString(labels, 'deployment.environment.name') = ? AND labels LIKE ? AND labels LIKE ?",
|
||||
expectedArgs: []any{"production", "%deployment.environment.name%", `%deployment.environment.name":"production%`},
|
||||
},
|
||||
{
|
||||
name: "semantic convention family negative filter does not reject fallback rows",
|
||||
key: &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
},
|
||||
op: qbtypes.FilterOperatorNotEqual,
|
||||
value: "staging",
|
||||
expected: "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), '')) <> ?",
|
||||
expectedArgs: []any{"staging"},
|
||||
},
|
||||
{
|
||||
name: "semantic convention family exists checks every member",
|
||||
key: &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
},
|
||||
op: qbtypes.FilterOperatorExists,
|
||||
expected: "(simpleJSONHas(labels, 'deployment.environment.name') = ? OR simpleJSONHas(labels, 'deployment.environment') = ?) AND (labels LIKE ? OR labels LIKE ?)",
|
||||
expectedArgs: []any{true, true, "%deployment.environment.name%", "%deployment.environment%"},
|
||||
},
|
||||
{
|
||||
name: "semantic convention family not exists checks every member",
|
||||
key: &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
},
|
||||
op: qbtypes.FilterOperatorNotExists,
|
||||
expected: "(simpleJSONHas(labels, 'deployment.environment.name') <> ? AND simpleJSONHas(labels, 'deployment.environment') <> ?)",
|
||||
expectedArgs: []any{true, true},
|
||||
},
|
||||
}
|
||||
|
||||
fm := NewFieldMapper()
|
||||
|
||||
@@ -3,8 +3,10 @@ package resourcefilter
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
|
||||
"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"
|
||||
@@ -32,6 +34,20 @@ func NewFieldMapper() *defaultFieldMapper {
|
||||
return &defaultFieldMapper{}
|
||||
}
|
||||
|
||||
func resourceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
|
||||
if key.FieldContext != telemetrytypes.FieldContextResource {
|
||||
return []string{key.Name}
|
||||
}
|
||||
if len(key.SemconvMembers) > 0 {
|
||||
return key.SemconvMembers
|
||||
}
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
})
|
||||
}
|
||||
|
||||
func (m *defaultFieldMapper) getColumn(
|
||||
_ context.Context,
|
||||
_, _ uint64,
|
||||
@@ -66,7 +82,15 @@ func (m *defaultFieldMapper) FieldFor(
|
||||
return "", err
|
||||
}
|
||||
if key.FieldContext == telemetrytypes.FieldContextResource {
|
||||
return fmt.Sprintf("simpleJSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
|
||||
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 columns[0].Name, nil
|
||||
}
|
||||
|
||||
@@ -1,35 +0,0 @@
|
||||
package telemetrymetadata
|
||||
|
||||
import "github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
|
||||
type BackwardCompatibleKeyMap map[string]string
|
||||
|
||||
var (
|
||||
TracesBackwardCompatKeys = BackwardCompatibleKeyMap{
|
||||
"net.peer.name": "server.address",
|
||||
"server.address": "net.peer.name",
|
||||
"http.url": "url.full",
|
||||
"url.full": "http.url",
|
||||
}
|
||||
|
||||
// LogsBackwardCompatKeys contains bidirectional mappings for logs.
|
||||
// Currently empty, can be extended in the future.
|
||||
LogsBackwardCompatKeys = BackwardCompatibleKeyMap{}
|
||||
|
||||
// MetricsBackwardCompatKeys contains bidirectional mappings for metrics.
|
||||
// Currently empty, can be extended in the future.
|
||||
MetricsBackwardCompatKeys = BackwardCompatibleKeyMap{}
|
||||
)
|
||||
|
||||
func GetBackwardCompatKeysForSignal(signal telemetrytypes.Signal) BackwardCompatibleKeyMap {
|
||||
switch signal {
|
||||
case telemetrytypes.SignalTraces:
|
||||
return TracesBackwardCompatKeys
|
||||
case telemetrytypes.SignalLogs:
|
||||
return LogsBackwardCompatKeys
|
||||
case telemetrytypes.SignalMetrics:
|
||||
return MetricsBackwardCompatKeys
|
||||
default:
|
||||
return BackwardCompatibleKeyMap{}
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"slices"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -14,6 +15,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/factory"
|
||||
"github.com/SigNoz/signoz/pkg/flagger"
|
||||
"github.com/SigNoz/signoz/pkg/querybuilder"
|
||||
"github.com/SigNoz/signoz/pkg/semconv"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/audittelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
|
||||
"github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema"
|
||||
@@ -151,6 +153,102 @@ func (t *telemetryMetaStore) tracesTblStatementToFieldKeys(ctx context.Context)
|
||||
return materialisedKeys, nil
|
||||
}
|
||||
|
||||
func traceSemconvMembers(name string, fieldContext telemetrytypes.FieldContext) []string {
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: fieldContext,
|
||||
})
|
||||
}
|
||||
|
||||
func traceSemconvDuplicateFactor() int {
|
||||
factor := 1
|
||||
for _, family := range semconv.All() {
|
||||
if family.Kind != semconv.KindAttribute {
|
||||
continue
|
||||
}
|
||||
if _, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: family.Current,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
}); ok {
|
||||
factor = max(factor, len(family.Old)+1)
|
||||
}
|
||||
}
|
||||
return factor
|
||||
}
|
||||
|
||||
// canonicalizeTraceSemconvKeys presents one current-name key for each family.
|
||||
// If metadata contains both spellings, metadata attached to the current name
|
||||
// wins; otherwise the old entry is copied under the current response name.
|
||||
func canonicalizeTraceSemconvKeys(keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
|
||||
result := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
|
||||
indexByIdentity := make(map[string]int)
|
||||
currentSourceByIdentity := make(map[string]bool)
|
||||
|
||||
for _, key := range keys {
|
||||
if key.Signal != telemetrytypes.SignalTraces {
|
||||
result = append(result, key)
|
||||
continue
|
||||
}
|
||||
|
||||
family, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: key.FieldContext,
|
||||
})
|
||||
if !ok {
|
||||
result = append(result, key)
|
||||
continue
|
||||
}
|
||||
|
||||
resolved := *key
|
||||
resolved.Name = family.Current
|
||||
resolved.SemconvMembers = []string{key.Name}
|
||||
identity := resolved.Name + ";" + resolved.Signal.StringValue() + ";" + resolved.FieldContext.StringValue() + ";" + resolved.FieldDataType.StringValue()
|
||||
fromCurrent := key.Name == family.Current
|
||||
|
||||
if index, found := indexByIdentity[identity]; found {
|
||||
physicalMembers := result[index].SemconvMembers
|
||||
if fromCurrent && !currentSourceByIdentity[identity] {
|
||||
result[index] = &resolved
|
||||
currentSourceByIdentity[identity] = true
|
||||
}
|
||||
for _, member := range physicalMembers {
|
||||
if !slices.Contains(result[index].SemconvMembers, member) {
|
||||
result[index].SemconvMembers = append(result[index].SemconvMembers, member)
|
||||
}
|
||||
}
|
||||
if !slices.Contains(result[index].SemconvMembers, key.Name) {
|
||||
result[index].SemconvMembers = append(result[index].SemconvMembers, key.Name)
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
indexByIdentity[identity] = len(result)
|
||||
currentSourceByIdentity[identity] = fromCurrent
|
||||
result = append(result, &resolved)
|
||||
}
|
||||
|
||||
for _, key := range result {
|
||||
if len(key.SemconvMembers) < 2 {
|
||||
continue
|
||||
}
|
||||
present := make(map[string]bool, len(key.SemconvMembers))
|
||||
for _, member := range key.SemconvMembers {
|
||||
present[member] = true
|
||||
}
|
||||
ordered := make([]string, 0, len(key.SemconvMembers))
|
||||
for _, member := range traceSemconvMembers(key.Name, key.FieldContext) {
|
||||
if present[member] {
|
||||
ordered = append(ordered, member)
|
||||
}
|
||||
}
|
||||
key.SemconvMembers = ordered
|
||||
}
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
// getTracesKeys returns the keys from the spans that match the field selection criteria.
|
||||
func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelectors []*telemetrytypes.FieldKeySelector) ([]*telemetrytypes.TelemetryFieldKey, bool, error) {
|
||||
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
|
||||
@@ -202,10 +300,23 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
|
||||
|
||||
// key part of the selector
|
||||
fieldKeyConds := []string{}
|
||||
members := traceSemconvMembers(fieldKeySelector.Name, fieldKeySelector.FieldContext)
|
||||
if fieldKeySelector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeExact {
|
||||
fieldKeyConds = append(fieldKeyConds, sb.E("tagKey", fieldKeySelector.Name))
|
||||
if len(members) == 1 {
|
||||
fieldKeyConds = append(fieldKeyConds, sb.E("tagKey", members[0]))
|
||||
} else {
|
||||
memberValues := make([]any, 0, len(members))
|
||||
for _, member := range members {
|
||||
memberValues = append(memberValues, member)
|
||||
}
|
||||
fieldKeyConds = append(fieldKeyConds, sb.In("tagKey", memberValues...))
|
||||
}
|
||||
} else {
|
||||
fieldKeyConds = append(fieldKeyConds, sb.ILike("tagKey", "%"+escapeForLike(fieldKeySelector.Name)+"%"))
|
||||
memberConditions := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
memberConditions = append(memberConditions, sb.ILike("tagKey", "%"+escapeForLike(member)+"%"))
|
||||
}
|
||||
fieldKeyConds = append(fieldKeyConds, sb.Or(memberConditions...))
|
||||
}
|
||||
|
||||
searchTexts = append(searchTexts, fieldKeySelector.Name)
|
||||
@@ -238,8 +349,10 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
|
||||
mainSb.From(mainSb.BuilderAs(sb, "sub_query"))
|
||||
mainSb.GroupBy("tag_key", "tag_type", "tag_data_type")
|
||||
mainSb.OrderBy("priority")
|
||||
// query one extra to check if we hit the limit
|
||||
mainSb.Limit(limit + 1)
|
||||
// Family members collapse after the database query. In the worst case each
|
||||
// logical key occupies one row per family member, so fetch enough physical
|
||||
// rows to return the requested logical page, plus one to detect truncation.
|
||||
mainSb.Limit(limit*traceSemconvDuplicateFactor() + 1)
|
||||
|
||||
query, args := mainSb.BuildWithFlavor(sqlbuilder.ClickHouse)
|
||||
|
||||
@@ -249,14 +362,7 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
|
||||
}
|
||||
defer rows.Close()
|
||||
keys := []*telemetrytypes.TelemetryFieldKey{}
|
||||
rowCount := 0
|
||||
for rows.Next() {
|
||||
rowCount++
|
||||
// reached the limit, we know there are more results
|
||||
if rowCount > limit {
|
||||
break
|
||||
}
|
||||
|
||||
var name string
|
||||
var fieldContext telemetrytypes.FieldContext
|
||||
var fieldDataType telemetrytypes.FieldDataType
|
||||
@@ -285,8 +391,11 @@ func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelector
|
||||
return nil, false, errors.Wrap(rows.Err(), errors.TypeInternal, errors.CodeInternal, ErrFailedToGetTracesKeys.Error())
|
||||
}
|
||||
|
||||
// hit the limit? (only counting DB results)
|
||||
complete := rowCount <= limit
|
||||
keys = canonicalizeTraceSemconvKeys(keys)
|
||||
complete := len(keys) <= limit
|
||||
if !complete {
|
||||
keys = keys[:limit]
|
||||
}
|
||||
|
||||
staticKeys := []string{"isRoot", "isEntryPoint"}
|
||||
staticKeys = append(staticKeys, maps.Keys(tracestelemetryschema.IntrinsicFields)...)
|
||||
@@ -1108,40 +1217,6 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
|
||||
|
||||
}
|
||||
|
||||
// applyBackwardCompatibleKeys adds backward compatible key aliases to the map.
|
||||
func applyBackwardCompatibleKeys(mapOfKeys map[string][]*telemetrytypes.TelemetryFieldKey) {
|
||||
// Get backward compatible keys for all signals
|
||||
backwardCompatKeysBySignal := map[telemetrytypes.Signal]BackwardCompatibleKeyMap{
|
||||
telemetrytypes.SignalTraces: GetBackwardCompatKeysForSignal(telemetrytypes.SignalTraces),
|
||||
telemetrytypes.SignalLogs: GetBackwardCompatKeysForSignal(telemetrytypes.SignalLogs),
|
||||
telemetrytypes.SignalMetrics: GetBackwardCompatKeysForSignal(telemetrytypes.SignalMetrics),
|
||||
}
|
||||
|
||||
// Iterate over existing keys and add aliases if they exist in backward compat mapping
|
||||
for srcKey, srcKeys := range mapOfKeys {
|
||||
for _, srcKeyEntry := range srcKeys {
|
||||
backwardCompatKeys := backwardCompatKeysBySignal[srcKeyEntry.Signal]
|
||||
if backwardCompatKeys == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
if aliasKey, ok := backwardCompatKeys[srcKey]; ok {
|
||||
if _, aliasExists := mapOfKeys[aliasKey]; !aliasExists {
|
||||
aliasKeyEntry := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: aliasKey,
|
||||
Signal: srcKeyEntry.Signal,
|
||||
FieldContext: srcKeyEntry.FieldContext,
|
||||
FieldDataType: srcKeyEntry.FieldDataType,
|
||||
}
|
||||
mapOfKeys[aliasKey] = []*telemetrytypes.TelemetryFieldKey{aliasKeyEntry}
|
||||
}
|
||||
// Found the alias for this signal, no need to check other entries
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func enrichWithIntrinsicMetricKeys(keys map[string][]*telemetrytypes.TelemetryFieldKey, selectors []*telemetrytypes.FieldKeySelector) map[string][]*telemetrytypes.TelemetryFieldKey {
|
||||
if len(selectors) == 0 {
|
||||
return keys
|
||||
@@ -1272,7 +1347,6 @@ func (t *telemetryMetaStore) GetKeys(ctx context.Context, orgID valuer.UUID, fie
|
||||
mapOfKeys[key.Name] = append(mapOfKeys[key.Name], key)
|
||||
}
|
||||
|
||||
applyBackwardCompatibleKeys(mapOfKeys)
|
||||
mapOfKeys = enrichWithIntrinsicMetricKeys(mapOfKeys, selectors)
|
||||
if t.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableAIObservability, featuretypes.NewFlaggerEvaluationContext(orgID)) {
|
||||
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, selectors)
|
||||
@@ -1353,7 +1427,6 @@ func (t *telemetryMetaStore) GetKeysMulti(ctx context.Context, orgID valuer.UUID
|
||||
mapOfKeys[key.Name] = append(mapOfKeys[key.Name], key)
|
||||
}
|
||||
|
||||
applyBackwardCompatibleKeys(mapOfKeys)
|
||||
mapOfKeys = enrichWithIntrinsicMetricKeys(mapOfKeys, fieldKeySelectors)
|
||||
if t.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableAIObservability, featuretypes.NewFlaggerEvaluationContext(orgID)) {
|
||||
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, fieldKeySelectors)
|
||||
@@ -1367,7 +1440,18 @@ func (t *telemetryMetaStore) GetKey(ctx context.Context, orgID valuer.UUID, fiel
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return keys[fieldKeySelector.Name], nil
|
||||
members := semconv.Members(semconv.KindAttribute, *fieldKeySelector)
|
||||
resolved := make([]*telemetrytypes.TelemetryFieldKey, 0)
|
||||
seen := make(map[*telemetrytypes.TelemetryFieldKey]bool)
|
||||
for _, member := range members {
|
||||
for _, key := range keys[member] {
|
||||
if !seen[key] {
|
||||
resolved = append(resolved, key)
|
||||
seen[key] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
return resolved, nil
|
||||
}
|
||||
|
||||
func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *telemetrytypes.FieldValueSelector) ([]string, bool, error) {
|
||||
@@ -1542,7 +1626,16 @@ func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueS
|
||||
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(t.tracesDBName + "." + t.tracesFieldsTblName)
|
||||
|
||||
if fieldValueSelector.Name != "" {
|
||||
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
|
||||
members := traceSemconvMembers(fieldValueSelector.Name, fieldValueSelector.FieldContext)
|
||||
if len(members) == 1 {
|
||||
sb.Where(sb.E("tag_key", members[0]))
|
||||
} else {
|
||||
memberValues := make([]any, 0, len(members))
|
||||
for _, member := range members {
|
||||
memberValues = append(memberValues, member)
|
||||
}
|
||||
sb.Where(sb.In("tag_key", memberValues...))
|
||||
}
|
||||
}
|
||||
|
||||
// now look at the field context
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore"
|
||||
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
|
||||
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
|
||||
"github.com/SigNoz/signoz/pkg/valuer"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
@@ -83,3 +84,83 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
|
||||
t.Errorf("there were unfulfilled expectations: %s", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCanonicalizeTraceSemconvKeys(t *testing.T) {
|
||||
oldResource := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Description: "old resource metadata",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
currentResource := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
Description: "current resource metadata",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
oldAttribute := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
oldLogResource := &telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
Signal: telemetrytypes.SignalLogs,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
|
||||
result := canonicalizeTraceSemconvKeys([]*telemetrytypes.TelemetryFieldKey{
|
||||
oldResource,
|
||||
currentResource,
|
||||
oldAttribute,
|
||||
oldLogResource,
|
||||
})
|
||||
|
||||
require.Len(t, result, 3)
|
||||
assert.Equal(t, "deployment.environment.name", result[0].Name)
|
||||
assert.Equal(t, "current resource metadata", result[0].Description)
|
||||
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, result[0].SemconvMembers)
|
||||
assert.Equal(t, "deployment.environment.name", result[1].Name)
|
||||
assert.Equal(t, telemetrytypes.FieldContextAttribute, result[1].FieldContext)
|
||||
assert.Equal(t, []string{"deployment.environment"}, result[1].SemconvMembers)
|
||||
assert.Equal(t, "deployment.environment", result[2].Name, "phase 1 must not rewrite raw log metadata")
|
||||
}
|
||||
|
||||
func TestGetSpanFieldValuesMergesSemconvFamily(t *testing.T) {
|
||||
mockTelemetryStore := telemetrystoretest.New(telemetrystore.Config{}, ®exMatcher{})
|
||||
mock := mockTelemetryStore.Mock()
|
||||
|
||||
metadata := NewTelemetryMetaStore(
|
||||
instrumentationtest.New().ToProviderSettings(),
|
||||
mockTelemetryStore,
|
||||
flaggertest.New(t),
|
||||
)
|
||||
|
||||
mock.ExpectQuery(`SELECT DISTINCT string_value, number_value FROM signoz_traces\.distributed_tag_attributes_v2 WHERE tag_key IN \(\?, \?\) AND tag_type = \? AND tag_data_type = \? LIMIT \?`).
|
||||
WithArgs("deployment.environment.name", "deployment.environment", "resource", "string", 51).
|
||||
WillReturnRows(cmock.NewRows([]cmock.ColumnType{
|
||||
{Name: "string_value", Type: "String"},
|
||||
{Name: "number_value", Type: "Float64"},
|
||||
}, [][]any{
|
||||
{"production", float64(0)},
|
||||
{"staging", float64(0)},
|
||||
{"production", float64(0)},
|
||||
}))
|
||||
|
||||
values, complete, err := metadata.GetAllValues(context.Background(), valuer.UUID{}, &telemetrytypes.FieldValueSelector{
|
||||
FieldKeySelector: &telemetrytypes.FieldKeySelector{
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Name: "deployment.environment",
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
assert.True(t, complete)
|
||||
assert.Equal(t, []string{"production", "staging"}, values.StringValues)
|
||||
assert.NoError(t, mock.ExpectationsWereMet(), "all expected metadata queries should be executed")
|
||||
}
|
||||
|
||||
@@ -154,6 +154,15 @@ 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) {
|
||||
if fm, ok := c.fm.(*fieldMapper); ok {
|
||||
return fm.existsExpressionFor(ctx, orgID, startNs, endNs, key, operator == qbtypes.FilterOperatorExists)
|
||||
}
|
||||
}
|
||||
columns, err := c.fm.ColumnFor(ctx, orgID, startNs, endNs, key)
|
||||
if err != nil {
|
||||
return "", err
|
||||
|
||||
@@ -210,6 +210,32 @@ func TestConditionFor(t *testing.T) {
|
||||
expectedSQL: "NOT mapContains(attributes_string, 'user.id')",
|
||||
expectedError: nil,
|
||||
},
|
||||
{
|
||||
name: "Equal operator - semantic convention family",
|
||||
key: telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment.name",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
},
|
||||
operator: qbtypes.FilterOperatorEqual,
|
||||
value: "production",
|
||||
expectedSQL: "(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'))))",
|
||||
expectedArgs: []any{"production"},
|
||||
expectedError: nil,
|
||||
},
|
||||
{
|
||||
name: "Not Exists operator - semantic convention family",
|
||||
key: telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name", "deployment.environment"},
|
||||
},
|
||||
operator: qbtypes.FilterOperatorNotExists,
|
||||
expectedSQL: "NOT (((mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))))",
|
||||
expectedError: nil,
|
||||
},
|
||||
{
|
||||
name: "Exists operator - json field",
|
||||
key: telemetrytypes.TelemetryFieldKey{
|
||||
|
||||
@@ -8,6 +8,7 @@ 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"
|
||||
@@ -167,6 +168,32 @@ func NewFieldMapper() *fieldMapper {
|
||||
return &fieldMapper{}
|
||||
}
|
||||
|
||||
func traceSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
|
||||
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
|
||||
return []string{key.Name}
|
||||
}
|
||||
if len(key.SemconvMembers) > 0 {
|
||||
return key.SemconvMembers
|
||||
}
|
||||
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: key.FieldContext,
|
||||
})
|
||||
}
|
||||
|
||||
func isTraceSemconvFamily(key *telemetrytypes.TelemetryFieldKey) bool {
|
||||
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
|
||||
return false
|
||||
}
|
||||
_, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
|
||||
Name: key.Name,
|
||||
Signal: telemetrytypes.SignalTraces,
|
||||
FieldContext: key.FieldContext,
|
||||
})
|
||||
return ok
|
||||
}
|
||||
|
||||
func (m *fieldMapper) getColumn(
|
||||
_ context.Context,
|
||||
_, _ uint64,
|
||||
@@ -291,10 +318,25 @@ 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)
|
||||
}
|
||||
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
|
||||
// once clickHouse dependency is updated, we need to check if we can remove it.
|
||||
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
|
||||
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
|
||||
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))
|
||||
}
|
||||
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]))
|
||||
}
|
||||
case schema.ColumnTypeEnumString,
|
||||
schema.ColumnTypeEnumUInt64,
|
||||
schema.ColumnTypeEnumUInt32,
|
||||
@@ -319,13 +361,35 @@ func (m *fieldMapper) resolveColumnExprs(
|
||||
|
||||
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
|
||||
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumFloat64, schema.ColumnTypeEnumBool:
|
||||
// a key could have been materialized, if so return the materialized column name
|
||||
if key.Materialized {
|
||||
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
|
||||
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
|
||||
members := traceSemconvMembers(key)
|
||||
if len(members) > 1 {
|
||||
guards := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
guards = append(guards, fmt.Sprintf("mapContains(%s, '%s')", columnName, member))
|
||||
}
|
||||
if valueType.GetType() == schema.ColumnTypeEnumString {
|
||||
values := make([]string, 0, len(members))
|
||||
for _, member := range members {
|
||||
values = append(values, fmt.Sprintf("NULLIF(%s['%s'], '')", columnName, member))
|
||||
}
|
||||
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+")")
|
||||
} else {
|
||||
branches := make([]string, 0, len(members)*2)
|
||||
for i, member := range members {
|
||||
branches = append(branches, guards[i], fmt.Sprintf("%s['%s']", columnName, member))
|
||||
}
|
||||
exprs = append(exprs, "multiIf("+strings.Join(branches, ", ")+", NULL)")
|
||||
}
|
||||
existExprs = append(existExprs, "("+strings.Join(guards, " OR ")+")")
|
||||
} else 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))
|
||||
} else {
|
||||
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, key.Name))
|
||||
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, key.Name))
|
||||
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, members[0]))
|
||||
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, members[0]))
|
||||
}
|
||||
default:
|
||||
return nil, nil, nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "value type %s is not supported for map column type %s", valueType, column.Type)
|
||||
@@ -529,6 +593,25 @@ 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
|
||||
|
||||
@@ -80,7 +80,7 @@ func TestGetFieldKeyName(t *testing.T) {
|
||||
Materialized: true,
|
||||
Evolutions: mockEvolution,
|
||||
},
|
||||
expectedResult: "multiIf(resource.`deployment.environment` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
|
||||
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 mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(resources_string['deployment.environment'], '')), NULL)",
|
||||
expectedError: nil,
|
||||
},
|
||||
{
|
||||
@@ -120,6 +120,51 @@ func TestGetFieldKeyName(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestFieldForSemconvFamily(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())
|
||||
|
||||
for _, requestedName := range []string{"deployment.environment.name", "deployment.environment"} {
|
||||
attributeKey := telemetrytypes.TelemetryFieldKey{
|
||||
Name: requestedName,
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
}
|
||||
attributeExpression, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, &attributeKey)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t,
|
||||
"COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''))",
|
||||
attributeExpression,
|
||||
)
|
||||
|
||||
resourceKey := telemetrytypes.TelemetryFieldKey{
|
||||
Name: requestedName,
|
||||
FieldContext: telemetrytypes.FieldContextResource,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
Materialized: true,
|
||||
Evolutions: MockEvolutionData(time.Date(2024, 6, 2, 0, 0, 0, 0, time.UTC)),
|
||||
}
|
||||
resourceExpression, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, &resourceKey)
|
||||
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, '')), (mapContains(resources_string, 'deployment.environment.name') OR mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(resources_string['deployment.environment'], '')), NULL)",
|
||||
resourceExpression,
|
||||
)
|
||||
}
|
||||
|
||||
oldRequestWithCurrentOnlyMetadata := telemetrytypes.TelemetryFieldKey{
|
||||
Name: "deployment.environment",
|
||||
FieldContext: telemetrytypes.FieldContextAttribute,
|
||||
FieldDataType: telemetrytypes.FieldDataTypeString,
|
||||
SemconvMembers: []string{"deployment.environment.name"},
|
||||
}
|
||||
expression, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, &oldRequestWithCurrentOnlyMetadata)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "attributes_string['deployment.environment.name']", expression)
|
||||
}
|
||||
|
||||
func TestFieldForResourceWithEvolution(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
releaseTime := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)
|
||||
@@ -176,7 +221,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: "resource.`deployment.environment`::String",
|
||||
expectedResult: "COALESCE(NULLIF(resource.`deployment.environment.name`::String, ''), NULLIF(resource.`deployment.environment`::String, ''))",
|
||||
},
|
||||
{
|
||||
name: "Window straddles release - materialized resource",
|
||||
@@ -189,7 +234,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` IS NOT NULL, resource.`deployment.environment`::String, `resource_string_deployment$$environment_exists`, `resource_string_deployment$$environment`, NULL)",
|
||||
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 mapContains(resources_string, 'deployment.environment')), COALESCE(NULLIF(resources_string['deployment.environment.name'], ''), NULLIF(resources_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", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment.name", "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", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment.name", "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", "dataType": "string", "type": "resource"},
|
||||
{"key": "deployment.environment.name", "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"},
|
||||
|
||||
37
pkg/types/quickfiltertypes/filter_test.go
Normal file
37
pkg/types/quickfiltertypes/filter_test.go
Normal file
@@ -0,0 +1,37 @@
|
||||
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())
|
||||
}
|
||||
}
|
||||
@@ -47,7 +47,8 @@ type TelemetryFieldKey struct {
|
||||
Indexes []TelemetryFieldKeySkipIndex `json:"-"`
|
||||
Materialized bool `json:"-"` // refers to promoted in case of body.... fields
|
||||
|
||||
Evolutions []*EvolutionEntry `json:"-"`
|
||||
Evolutions []*EvolutionEntry `json:"-"`
|
||||
SemconvMembers []string `json:"-"`
|
||||
}
|
||||
|
||||
func (f *TelemetryFieldKey) KeyNameContainsArray() bool {
|
||||
@@ -128,6 +129,7 @@ func (f *TelemetryFieldKey) OverrideMetadataFrom(src *TelemetryFieldKey) {
|
||||
f.Materialized = src.Materialized
|
||||
f.JSONPlan = src.JSONPlan
|
||||
f.Evolutions = src.Evolutions
|
||||
f.SemconvMembers = src.SemconvMembers
|
||||
}
|
||||
|
||||
func (f *TelemetryFieldKey) Equal(key *TelemetryFieldKey) bool {
|
||||
|
||||
239
tests/integration/tests/queriertraces/13_semconv_evolution.py
Normal file
239
tests/integration/tests/queriertraces/13_semconv_evolution.py
Normal file
@@ -0,0 +1,239 @@
|
||||
"""Phase 1 end-to-end checks for semantic-convention name evolution.
|
||||
|
||||
The fixture models a fleet split across SDK generations and deliberately includes
|
||||
a dual-emitting conflict. Both request spellings must address one logical field,
|
||||
with the current spelling winning when a row contains both.
|
||||
"""
|
||||
|
||||
from collections.abc import Callable, Generator
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from http import HTTPStatus
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
import requests
|
||||
|
||||
from fixtures import types
|
||||
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD
|
||||
from fixtures.querier import Aggregation, BuilderQuery, OrderBy, RequestType, TelemetryFieldKey, make_query_request
|
||||
from fixtures.traces import TraceIdGenerator, Traces, TracesKind, TracesStatusCode
|
||||
|
||||
CURRENT = "deployment.environment.name"
|
||||
OLD = "deployment.environment"
|
||||
PREFIX = "semconv-phase1"
|
||||
|
||||
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 _span(timestamp: datetime, suffix: str, environment: dict[str, str]) -> Traces:
|
||||
service = f"{PREFIX}-{suffix}"
|
||||
return 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),
|
||||
)
|
||||
|
||||
|
||||
@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)
|
||||
insert_traces(
|
||||
[
|
||||
_span(now - timedelta(seconds=5), "old", {OLD: "production"}),
|
||||
_span(now - timedelta(seconds=4), "current", {CURRENT: "production"}),
|
||||
_span(now - timedelta(seconds=3), "both", {OLD: "production", CURRENT: "production"}),
|
||||
_span(now - timedelta(seconds=2), "conflict", {OLD: "staging", CURRENT: "production"}),
|
||||
_span(now - timedelta(seconds=1), "staging", {OLD: "staging"}),
|
||||
_span(now, "missing", {}),
|
||||
]
|
||||
)
|
||||
|
||||
# Service-map rows are derived by the collector in production. Seed the
|
||||
# derived table directly here so the backend alias allowlist is tested in
|
||||
# isolation; 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
|
||||
'{PREFIX}-map-{suffix}', '{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}' "
|
||||
f"DELETE WHERE startsWith(src, '{PREFIX}-map-') SETTINGS mutations_sync = 1"
|
||||
)
|
||||
|
||||
|
||||
def _result(response: requests.Response) -> dict[str, Any]:
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
results = response.json()["data"]["data"]["results"]
|
||||
assert len(results) == 1
|
||||
return results[0]
|
||||
|
||||
|
||||
def _raw_names(
|
||||
signoz: types.SigNoz,
|
||||
token: str,
|
||||
now: datetime,
|
||||
expression: str,
|
||||
) -> set[str]:
|
||||
response = make_query_request(
|
||||
signoz,
|
||||
token,
|
||||
start_ms=int((now - timedelta(minutes=2)).timestamp() * 1000),
|
||||
end_ms=int((now + timedelta(minutes=1)).timestamp() * 1000),
|
||||
request_type=RequestType.RAW,
|
||||
queries=[
|
||||
BuilderQuery(
|
||||
signal="traces",
|
||||
name="A",
|
||||
limit=100,
|
||||
filter_expression=expression,
|
||||
select_fields=[TelemetryFieldKey("span.name")],
|
||||
order=[OrderBy(TelemetryFieldKey("timestamp"), "asc")],
|
||||
).to_dict()
|
||||
],
|
||||
)
|
||||
return {row["data"]["name"] for row in (_result(response).get("rows") or [])}
|
||||
|
||||
|
||||
def _metadata_values(signoz: types.SigNoz, token: str, name: str, context: str) -> set[str]:
|
||||
response = requests.get(
|
||||
signoz.self.host_configs["8080"].get("/api/v1/fields/values"),
|
||||
timeout=5,
|
||||
headers={"authorization": f"Bearer {token}"},
|
||||
params={
|
||||
"signal": "traces",
|
||||
"name": name,
|
||||
"fieldContext": context,
|
||||
"fieldDataType": "string",
|
||||
},
|
||||
)
|
||||
assert response.status_code == HTTPStatus.OK, response.text
|
||||
return set(response.json()["data"]["values"].get("stringValues") or [])
|
||||
|
||||
|
||||
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)
|
||||
|
||||
# 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}"
|
||||
assert _raw_names(signoz, token, now, f"{field} = 'production'") == PRODUCTION_SPANS
|
||||
assert _raw_names(signoz, token, now, f"{field} = 'staging'") == STAGING_SPANS
|
||||
assert _raw_names(signoz, token, now, f"{field} != 'production'") == STAGING_SPANS
|
||||
assert _raw_names(signoz, token, now, f"{field} EXISTS") == PRODUCTION_SPANS | STAGING_SPANS
|
||||
assert _raw_names(signoz, token, now, f"{field} NOT EXISTS") == MISSING_SPANS
|
||||
|
||||
grouped = make_query_request(
|
||||
signoz,
|
||||
token,
|
||||
start_ms=int((now - timedelta(minutes=2)).timestamp() * 1000),
|
||||
end_ms=int((now + timedelta(minutes=1)).timestamp() * 1000),
|
||||
request_type=RequestType.SCALAR,
|
||||
queries=[
|
||||
BuilderQuery(
|
||||
signal="traces",
|
||||
name="A",
|
||||
filter_expression=f"{field} EXISTS",
|
||||
aggregations=[Aggregation("count()")],
|
||||
group_by=[TelemetryFieldKey(requested, "string", context)],
|
||||
order=[OrderBy(TelemetryFieldKey(requested, "string", context), "asc")],
|
||||
).to_dict()
|
||||
],
|
||||
)
|
||||
result = _result(grouped)
|
||||
assert result["columns"][0]["name"] == requested, "response identity must match the request spelling"
|
||||
assert result["data"] == [["production", 4], ["staging", 1]]
|
||||
|
||||
assert _metadata_values(signoz, token, CURRENT, context) == {"production", "staging"}
|
||||
assert _metadata_values(signoz, token, OLD, context) == {"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 not 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