Compare commits

...

1 Commits

Author SHA1 Message Date
srikanthccv
6b4bb82efb feat: resolve semconv families across logs and metrics 2026-08-07 01:20:58 +05:30
31 changed files with 1037 additions and 137 deletions

View File

@@ -224,6 +224,10 @@ py-test: ## Runs integration tests
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-test-semconv-phase2
py-test-semconv-phase2: py-test-setup ## Rebuild the shared stack and run the Phase 1-2 cross-signal matrices
@cd tests && uv run pytest --basetemp=./tmp/ -vv --reuse --capture=no integration/tests/queriertraces/13_semconv_evolution.py integration/tests/queriersemconv/02_cross_signal.py
.PHONY: py-clean
py-clean: ## Clear all pycache and pytest cache from tests directory recursively
@echo ">> cleaning python cache files from tests directory"

View File

@@ -11,12 +11,21 @@ export type SemconvFamily = {
};
export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
{
current: 'container.cpu.usage',
old: ['container.cpu.utilization'],
kind: 'metric',
contexts: ['metric'],
signals: ['metrics'],
applyToMetrics: [],
valueMap: {},
},
{
current: 'db.system.name',
old: ['db.system'],
kind: 'attribute',
contexts: [],
signals: [],
contexts: ['attribute', 'resource'],
signals: ['logs', 'metrics', 'traces'],
applyToMetrics: [],
valueMap: {},
},
@@ -24,8 +33,26 @@ export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
current: 'deployment.environment.name',
old: ['deployment.environment'],
kind: 'attribute',
contexts: [],
signals: [],
contexts: ['attribute', 'resource'],
signals: ['logs', 'metrics', 'traces'],
applyToMetrics: [],
valueMap: {},
},
{
current: 'k8s.node.cpu.usage',
old: ['k8s.node.cpu.utilization'],
kind: 'metric',
contexts: ['metric'],
signals: ['metrics'],
applyToMetrics: [],
valueMap: {},
},
{
current: 'k8s.pod.cpu.usage',
old: ['k8s.pod.cpu.utilization'],
kind: 'metric',
contexts: ['metric'],
signals: ['metrics'],
applyToMetrics: [],
valueMap: {},
},

View File

@@ -52,10 +52,10 @@ import (
"github.com/SigNoz/signoz/pkg/query-service/constants"
chErrors "github.com/SigNoz/signoz/pkg/query-service/errors"
"github.com/SigNoz/signoz/pkg/query-service/metrics"
"github.com/SigNoz/signoz/pkg/query-service/model"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/query-service/utils"
"github.com/SigNoz/signoz/pkg/semconv"
)
const (
@@ -3190,11 +3190,6 @@ func (r *ClickHouseReader) GetMetricAttributeValues(ctx context.Context, orgID v
var rows driver.Rows
var attributeValues v3.FilterAttributeValueResponse
normalized := true
if constants.IsDotMetricsEnabled {
normalized = false
}
reductionEnabled := r.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableMetricsReduction, featuretypes.NewFlaggerEvaluationContext(orgID))
if reductionEnabled {
@@ -3205,8 +3200,7 @@ func (r *ClickHouseReader) GetMetricAttributeValues(ctx context.Context, orgID v
if req.Limit != 0 {
query = query + fmt.Sprintf(" LIMIT %d;", req.Limit)
}
names := []string{req.AggregateAttribute}
names = append(names, metrics.GetTransitionedMetric(req.AggregateAttribute, normalized))
names := semconv.MetricNames(req.AggregateAttribute)
rows, err = r.db.Query(ctx, query, req.FilterAttributeKey, names, req.FilterAttributeKey, fmt.Sprintf("%%%s%%", req.SearchText), common.PastDayRoundOff())

View File

@@ -1,27 +0,0 @@
package metrics
var MetricsUnderTransition = map[string]string{
"k8s_pod_cpu_utilization": "k8s_pod_cpu_usage",
"k8s_node_cpu_utilization": "k8s_node_cpu_usage",
"container_cpu_utilization": "container_cpu_usage",
}
var DotMetricsUnderTransition = map[string]string{
"k8s.pod.cpu.utilization": "k8s.pod.cpu.usage",
"k8s.node.cpu.utilization": "k8s.node.cpu.usage",
"container.cpu.utilization": "container.cpu.usage",
}
func GetTransitionedMetric(metric string, normalized bool) string {
if normalized {
if _, ok := MetricsUnderTransition[metric]; ok {
return MetricsUnderTransition[metric]
}
return metric
} else {
if _, ok := DotMetricsUnderTransition[metric]; ok {
return DotMetricsUnderTransition[metric]
}
return metric
}
}

View File

@@ -10,8 +10,8 @@ import (
"log/slog"
"github.com/SigNoz/signoz/pkg/query-service/constants"
"github.com/SigNoz/signoz/pkg/query-service/metrics"
v3 "github.com/SigNoz/signoz/pkg/query-service/model/v3"
"github.com/SigNoz/signoz/pkg/semconv"
)
// ValidateAndCastValue validates and casts the value of a key to the corresponding data type of the key
@@ -234,12 +234,12 @@ func ClickHouseFormattedValue(v interface{}) string {
func ClickHouseFormattedMetricNames(v interface{}) string {
if name, ok := v.(string); ok {
transitionedMetrics := metrics.GetTransitionedMetric(name, !constants.IsDotMetricsEnabled)
if transitionedMetrics != name {
return ClickHouseFormattedValue([]interface{}{transitionedMetrics})
} else {
return ClickHouseFormattedValue([]interface{}{name})
members := semconv.MetricNames(name)
values := make([]interface{}, 0, len(members))
for _, member := range members {
values = append(values, member)
}
return ClickHouseFormattedValue(values)
}
return ClickHouseFormattedValue(v)

View File

@@ -2,13 +2,29 @@ package querybuilder
import (
"fmt"
"strings"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
func physicalSemconvMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
if key.FieldContext != telemetrytypes.FieldContextResource && key.FieldContext != telemetrytypes.FieldContextAttribute {
return []string{key.Name}
}
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: key.Signal,
FieldContext: key.FieldContext,
})
}
// ExistsExpression renders the existence predicate for a key resolved to the given
// columns (negated when exists is false). Comparisons are against constants rendered
// as literals, so the expression carries no bind args and can guard column expressions
@@ -43,11 +59,26 @@ 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)
if exists {
return rawPath + " IS NOT NULL", nil
members := physicalSemconvMembers(key)
paths := make([]string, 0, len(members))
for _, member := range members {
paths = append(paths, fmt.Sprintf("%s.`%s`", columnName, member))
}
return rawPath + " IS NULL", nil
if len(paths) == 1 {
if exists {
return paths[0] + " IS NOT NULL", nil
}
return paths[0] + " IS NULL", nil
}
guards := make([]string, 0, len(paths))
for _, path := range paths {
guards = append(guards, path+" IS NOT NULL")
}
rawPath := "(" + strings.Join(guards, " OR ") + ")"
if exists {
return rawPath, nil
}
return "NOT " + rawPath, nil
case schema.ColumnTypeEnumString,
schema.ColumnTypeEnumFixedString:
if exists {
@@ -88,8 +119,16 @@ 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)
if key.Materialized {
members := physicalSemconvMembers(key)
operands := make([]string, 0, len(members))
for _, member := range members {
operands = append(operands, fmt.Sprintf("mapContains(%s, '%s')", column.Name, member))
}
leftOperand := strings.Join(operands, " OR ")
if len(operands) > 1 {
leftOperand = "(" + leftOperand + ")"
}
if key.Materialized && (len(members) == 1 || key.MaterializedSemconv) {
leftOperand = telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key)
}
if exists {

View File

@@ -39,6 +39,7 @@ type filterExpressionVisitor struct {
fullTextColumn *telemetrytypes.TelemetryFieldKey
skipResourceFilter bool
skipFullTextFilter bool
exactSemconv bool
variables map[string]qbtypes.VariableItem
keysWithWarnings map[string]bool
@@ -59,6 +60,7 @@ type FilterExprVisitorOpts struct {
FullTextColumn *telemetrytypes.TelemetryFieldKey
SkipResourceFilter bool
SkipFullTextFilter bool
ExactSemconv bool
Variables map[string]qbtypes.VariableItem
StartNs uint64
EndNs uint64
@@ -76,6 +78,7 @@ func newFilterExpressionVisitor(opts FilterExprVisitorOpts) *filterExpressionVis
fullTextColumn: opts.FullTextColumn,
skipResourceFilter: opts.SkipResourceFilter,
skipFullTextFilter: opts.SkipFullTextFilter,
exactSemconv: opts.ExactSemconv,
variables: opts.Variables,
keysWithWarnings: make(map[string]bool),
startNs: opts.StartNs,
@@ -380,7 +383,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 := v.matchingFieldKeys(key)
// Handle EXISTS specially
if ctx.EXISTS() != nil {
@@ -731,7 +734,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, v.matchingFieldKeys(key), operator, value)
if !ok {
return ErrorConditionLiteral
}
@@ -924,7 +927,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) {
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)
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter, ExactSemconv: v.exactSemconv}, op, value, v.builder)
if err != nil {
_, _, _, _, errURL, _ := errors.Unwrapb(err)
assignIfEmpty(&v.mainErrorURL, errURL)
@@ -983,12 +986,24 @@ 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 {
return matchingFieldKeys(field, fieldKeys, true)
}
// MatchingFieldKeysExact matches only the requested physical spelling.
func MatchingFieldKeysExact(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
return matchingFieldKeys(field, fieldKeys, false)
}
func matchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey, resolveSemconv bool) []*telemetrytypes.TelemetryFieldKey {
selector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: field.Signal,
FieldContext: field.FieldContext,
}
members := semconv.Members(semconv.KindAttribute, selector)
members := []string{field.Name}
if resolveSemconv {
members = semconv.AttributeMembers(selector)
}
isFamily := len(members) > 1
fieldKeysForName := make([]*telemetrytypes.TelemetryFieldKey, 0)
indexByIdentity := make(map[string]int)
@@ -1005,13 +1020,13 @@ 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.
if memberName != field.Name {
if resolveSemconv && memberName != field.Name {
itemSelector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
}
if !slices.Contains(semconv.Members(semconv.KindAttribute, itemSelector), memberName) {
if !slices.Contains(semconv.AttributeMembers(itemSelector), memberName) {
continue
}
}
@@ -1038,6 +1053,8 @@ func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[st
if isFamily {
resolved.Name = field.Name
resolved.SemconvMembers = slices.Clone(physicalMembers)
} else if !resolveSemconv {
resolved.SemconvMembers = []string{field.Name}
}
fieldKeysForName = append(fieldKeysForName, &resolved)
}
@@ -1060,3 +1077,25 @@ func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[st
return fieldKeysForName
}
func (v *filterExpressionVisitor) matchingFieldKeys(field *telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
if v.exactSemconv {
return MatchingFieldKeysExact(field, v.fieldKeys)
}
return MatchingFieldKeys(field, v.fieldKeys)
}
// ExactSemconvKeys returns copies pinned to their physical names, preventing a
// field mapper from expanding a synthesized or metadata-free key into a family.
func ExactSemconvKeys(keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
result := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
for _, key := range keys {
if key == nil {
continue
}
resolved := *key
resolved.SemconvMembers = []string{key.Name}
result = append(result, &resolved)
}
return result
}

View File

@@ -731,6 +731,11 @@ func TestMatchingFieldKeysResolvesSemconvFamily(t *testing.T) {
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)
exactMatches := MatchingFieldKeysExact(requested, fieldKeys)
require.Len(t, exactMatches, 1, "exact lookup must return one field before its metadata is inspected")
assert.Equal(t, current.Name, exactMatches[0].Name)
assert.Equal(t, []string{current.Name}, exactMatches[0].SemconvMembers)
}
// ---------------------------------------------------------------------------

View File

@@ -2,21 +2,47 @@
package semconv
import "github.com/SigNoz/signoz/pkg/types/telemetrytypes"
var families = []Family{
{
Current: "container.cpu.usage",
Old: []string{"container.cpu.utilization"},
Kind: KindMetric,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextMetric},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalMetrics},
ApplyToMetrics: nil,
},
{
Current: "db.system.name",
Old: []string{"db.system"},
Kind: KindAttribute,
Contexts: nil,
Signals: nil,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextAttribute, telemetrytypes.FieldContextResource},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalMetrics, telemetrytypes.SignalTraces},
ApplyToMetrics: nil,
},
{
Current: "deployment.environment.name",
Old: []string{"deployment.environment"},
Kind: KindAttribute,
Contexts: nil,
Signals: nil,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextAttribute, telemetrytypes.FieldContextResource},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalMetrics, telemetrytypes.SignalTraces},
ApplyToMetrics: nil,
},
{
Current: "k8s.node.cpu.usage",
Old: []string{"k8s.node.cpu.utilization"},
Kind: KindMetric,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextMetric},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalMetrics},
ApplyToMetrics: nil,
},
{
Current: "k8s.pod.cpu.usage",
Old: []string{"k8s.pod.cpu.utilization"},
Kind: KindMetric,
Contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextMetric},
Signals: []telemetrytypes.Signal{telemetrytypes.SignalMetrics},
ApplyToMetrics: nil,
},
}

View File

@@ -2,6 +2,7 @@ package semconv
import (
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -26,6 +27,13 @@ type Family struct {
ValueMap map[string]string
}
type metricSpelling uint8
const (
metricSpellingDotted metricSpelling = iota
metricSpellingNormalized
)
var (
KindAttribute = Kind{String: valuer.NewString("attribute")}
KindMetric = Kind{String: valuer.NewString("metric")}
@@ -60,6 +68,88 @@ func Members(kind Kind, selector telemetrytypes.FieldKeySelector) []string {
return familyMembers[idx]
}
// AttributeMembers returns the physical attribute spellings that may represent
// selector.Name. Metrics have used both dotted and normalized label layouts;
// resource labels have additionally used a resource_ prefix. Keeping that
// storage detail here prevents metrics readers from maintaining local
// transition tables.
func AttributeMembers(selector telemetrytypes.FieldKeySelector) []string {
if selector.Signal != telemetrytypes.SignalMetrics {
return Members(KindAttribute, selector)
}
lookupSelector := selector
lookupSelector.Name = strings.TrimPrefix(selector.Name, "resource_")
family, style, ok := lookupMetricSpelling(KindAttribute, lookupSelector)
if !ok {
return []string{selector.Name}
}
logicalMembers := familyMembers[family]
result := make([]string, 0, len(logicalMembers)*4)
for _, member := range logicalMembers {
dotted := member
normalized := normalizeMetricSpelling(member)
variants := []string{dotted, normalized}
if style == metricSpellingNormalized {
variants[0], variants[1] = variants[1], variants[0]
}
if selector.FieldContext == telemetrytypes.FieldContextResource ||
selector.FieldContext == telemetrytypes.FieldContextUnspecified ||
strings.HasPrefix(selector.Name, "resource_") {
for _, variant := range variants {
result = appendUniqueString(result, "resource_"+variant)
}
}
for _, variant := range variants {
result = appendUniqueString(result, variant)
}
}
return result
}
// MetricNames returns the current and historical storage names for a metric.
// The input's dotted or normalized style is preserved because both layouts are
// valid metric identities and must not be mixed in one query.
func MetricNames(name string) []string {
selector := telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
}
family, style, ok := lookupMetricSpelling(KindMetric, selector)
if !ok {
return []string{name}
}
logicalMembers := familyMembers[family]
result := make([]string, 0, len(logicalMembers))
for _, member := range logicalMembers {
if style == metricSpellingNormalized {
member = normalizeMetricSpelling(member)
}
result = appendUniqueString(result, member)
}
return result
}
// CurrentAttribute returns the canonical dotted name for an attribute
// spelling, or selector.Name if no enabled family matches.
func CurrentAttribute(selector telemetrytypes.FieldKeySelector) string {
if selector.Signal != telemetrytypes.SignalMetrics {
return Current(KindAttribute, selector)
}
lookupSelector := selector
lookupSelector.Name = strings.TrimPrefix(selector.Name, "resource_")
family, _, ok := lookupMetricSpelling(KindAttribute, lookupSelector)
if !ok {
return selector.Name
}
return families[family].Current
}
// Current returns the current name for selector.Name, or the input name when
// it does not belong to an enabled family.
func Current(kind Kind, selector telemetrytypes.FieldKeySelector) string {
@@ -103,6 +193,35 @@ func lookupIndex(kind Kind, selector telemetrytypes.FieldKeySelector) (int, bool
return 0, false
}
func lookupMetricSpelling(kind Kind, selector telemetrytypes.FieldKeySelector) (int, metricSpelling, bool) {
if idx, ok := lookupIndex(kind, selector); ok {
return idx, metricSpellingDotted, true
}
for idx, family := range families {
if !matchesSelector(family, kind, selector) {
continue
}
for _, member := range familyMembers[idx] {
if normalizeMetricSpelling(member) == selector.Name {
return idx, metricSpellingNormalized, true
}
}
}
return 0, metricSpellingDotted, false
}
func normalizeMetricSpelling(name string) string {
return strings.ReplaceAll(name, ".", "_")
}
func appendUniqueString(values []string, value string) []string {
if value == "" || slices.Contains(values, value) {
return values
}
return append(values, value)
}
func matchesSelector(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if family.Kind != kind {
return false

View File

@@ -57,3 +57,58 @@ func TestAllReturnsDefensiveCopies(t *testing.T) {
second := All()
assert.NotEqual(t, "mutated", second[0].Old[0], "All must not expose mutable generated data")
}
func TestAttributeMembersIncludesMetricResourceStorageSpellings(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, []string{
"resource_db.system.name", "resource_db_system_name", "db.system.name", "db_system_name",
"resource_db.system", "resource_db_system", "db.system", "db_system",
}, AttributeMembers(selector), "resource metric attributes should cover every historical storage layout")
}
func TestCurrentAttributeResolvesNormalizedMetricResourceSpelling(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "resource_db_system",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, "db.system.name", CurrentAttribute(selector), "normalized resource spelling should resolve to the dotted current name")
}
func TestAttributeMembersPreservesNormalizedMetricPointStyle(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "db_system",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
}
assert.Equal(t, []string{
"db_system_name", "db.system.name", "db_system", "db.system",
}, AttributeMembers(selector), "normalized point attribute should remain the preferred storage spelling")
}
func TestMetricNamesPreservesDottedStyle(t *testing.T) {
assert.Equal(t,
[]string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"},
MetricNames("k8s.pod.cpu.usage"),
"dotted metric input should produce dotted family names",
)
}
func TestMetricNamesPreservesNormalizedStyle(t *testing.T) {
assert.Equal(t,
[]string{"k8s_pod_cpu_usage", "k8s_pod_cpu_utilization"},
MetricNames("k8s_pod_cpu_utilization"),
"normalized metric input should produce normalized family names",
)
}
func TestMetricNamesReturnsUnknownNameUnchanged(t *testing.T) {
assert.Equal(t, []string{"custom_metric"}, MetricNames("custom_metric"), "unknown metrics should not be expanded")
}

View File

@@ -10,6 +10,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/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
"github.com/SigNoz/signoz/pkg/types/metrictypes"
@@ -31,6 +32,15 @@ const (
OthersMultiTemporality = `IF(LOWER(temporality) LIKE LOWER('delta'), %s, %s) AS per_series_value`
)
func metricNameValues(name string) []any {
members := semconv.MetricNames(name)
values := make([]any, 0, len(members))
for _, member := range members {
values = append(values, member)
}
return values
}
type StatementBuilder struct {
logger *slog.Logger
metadataStore telemetrytypes.MetadataStore
@@ -341,7 +351,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
sb.SelectMore(col)
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -385,7 +395,7 @@ func (b *StatementBuilder) buildReducedSpatialAggFastPath(
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", agg.MetricName),
sb.In("metric_name", metricNameValues(agg.MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -427,7 +437,7 @@ func (b *StatementBuilder) buildReducedTemporalAggregationCTE(
sb.From(fmt.Sprintf("%s.%s AS points FINAL", metricstelemetryschema.DBName, metricstelemetryschema.WhichReducedSamplesTableToUse(agg.Type)))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.reduced_fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", agg.MetricName),
sb.In("metric_name", metricNameValues(agg.MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -505,7 +515,7 @@ func (b *StatementBuilder) buildTemporalAggDeltaFastPath(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -560,7 +570,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -638,7 +648,7 @@ func (b *StatementBuilder) buildTemporalAggDelta(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -679,7 +689,7 @@ func (b *StatementBuilder) buildTemporalAggCumulativeOrUnspecified(
baseSb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
baseSb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
baseSb.Where(
baseSb.In("metric_name", query.Aggregations[0].MetricName),
baseSb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
baseSb.GTE("unix_milli", start),
baseSb.LT("unix_milli", end),
)
@@ -770,7 +780,7 @@ func (b *StatementBuilder) buildTemporalAggForMultipleTemporalities(
sb.From(fmt.Sprintf("%s.%s AS points", metricstelemetryschema.DBName, samplesTable))
sb.JoinWithOption(sqlbuilder.InnerJoin, timeSeriesCTE, "points.fingerprint = filtered_time_series.fingerprint")
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNameValues(query.Aggregations[0].MetricName)...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)

View File

@@ -13,9 +13,21 @@ import (
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestMetricNameValuesResolveFamily(t *testing.T) {
assert.Equal(t,
[]any{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"},
metricNameValues("k8s.pod.cpu.usage"),
)
assert.Equal(t,
[]any{"container_cpu_usage", "container_cpu_utilization"},
metricNameValues("container_cpu_utilization"),
)
}
func TestStatementBuilder(t *testing.T) {
cases := []struct {
name string

View File

@@ -153,15 +153,19 @@ 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{
func attributeSemconvMembers(name string, signal telemetrytypes.Signal, fieldContext telemetrytypes.FieldContext) []string {
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalTraces,
Signal: signal,
FieldContext: fieldContext,
})
}
func traceSemconvDuplicateFactor() int {
func traceSemconvMembers(name string, fieldContext telemetrytypes.FieldContext) []string {
return attributeSemconvMembers(name, telemetrytypes.SignalTraces, fieldContext)
}
func semconvDuplicateFactor(signal telemetrytypes.Signal) int {
factor := 1
for _, family := range semconv.All() {
if family.Kind != semconv.KindAttribute {
@@ -169,43 +173,66 @@ func traceSemconvDuplicateFactor() int {
}
if _, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: family.Current,
Signal: telemetrytypes.SignalTraces,
Signal: signal,
}); ok {
factor = max(factor, len(family.Old)+1)
factor = max(factor, len(attributeSemconvMembers(family.Current, signal, telemetrytypes.FieldContextUnspecified)))
}
}
return factor
}
func traceSemconvDuplicateFactor() int {
return semconvDuplicateFactor(telemetrytypes.SignalTraces)
}
func inStrings(sb *sqlbuilder.SelectBuilder, column string, values []string) string {
args := make([]any, 0, len(values))
for _, value := range values {
args = append(args, value)
}
return sb.In(column, args...)
}
// 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 {
return canonicalizeSemconvKeys(keys, telemetrytypes.SignalTraces)
}
// canonicalizeSemconvKeys presents one current-name key for each enabled
// attribute family and records only the physical spellings actually present in
// metadata. That lets field mappers retain their optimized single-key SQL for
// homogeneous ranges while mixed ranges coalesce current-first.
func canonicalizeSemconvKeys(keys []*telemetrytypes.TelemetryFieldKey, signal telemetrytypes.Signal) []*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 {
if key.Signal != signal {
result = append(result, key)
continue
}
family, ok := semconv.Lookup(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
selector := telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalTraces,
Signal: signal,
FieldContext: key.FieldContext,
})
if !ok {
}
members := semconv.AttributeMembers(selector)
current := semconv.CurrentAttribute(selector)
if len(members) == 1 && current == key.Name {
result = append(result, key)
continue
}
resolved := *key
resolved.Name = family.Current
resolved.Name = current
resolved.SemconvMembers = []string{key.Name}
identity := resolved.Name + ";" + resolved.Signal.StringValue() + ";" + resolved.FieldContext.StringValue() + ";" + resolved.FieldDataType.StringValue()
fromCurrent := key.Name == family.Current
physicalName := strings.TrimPrefix(key.Name, "resource_")
fromCurrent := physicalName == current || physicalName == strings.ReplaceAll(current, ".", "_")
if index, found := indexByIdentity[identity]; found {
physicalMembers := result[index].SemconvMembers
@@ -238,7 +265,7 @@ func canonicalizeTraceSemconvKeys(keys []*telemetrytypes.TelemetryFieldKey) []*t
present[member] = true
}
ordered := make([]string, 0, len(key.SemconvMembers))
for _, member := range traceSemconvMembers(key.Name, key.FieldContext) {
for _, member := range attributeSemconvMembers(key.Name, signal, key.FieldContext) {
if present[member] {
ordered = append(ordered, member)
}
@@ -569,10 +596,19 @@ func (t *telemetryMetaStore) getLogsKeys(ctx context.Context, orgID valuer.UUID,
continue
}
fieldKeyConds := []string{}
members := attributeSemconvMembers(sel.Name, telemetrytypes.SignalLogs, fieldContext)
if sel.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeExact {
fieldKeyConds = append(fieldKeyConds, sb.E("name", sel.Name))
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
fieldKeyConds = append(fieldKeyConds, sb.In("name", memberValues...))
} else {
fieldKeyConds = append(fieldKeyConds, sb.ILike("name", "%"+escapeForLike(sel.Name)+"%"))
memberConds := make([]string, 0, len(members))
for _, member := range members {
memberConds = append(memberConds, sb.ILike("name", "%"+escapeForLike(member)+"%"))
}
fieldKeyConds = append(fieldKeyConds, sb.Or(memberConds...))
}
if sel.FieldDataType != telemetrytypes.FieldDataTypeUnspecified {
fieldKeyConds = append(fieldKeyConds, sb.E("datatype", sel.FieldDataType.TagDataType()))
@@ -665,6 +701,7 @@ func (t *telemetryMetaStore) getLogsKeys(ctx context.Context, orgID valuer.UUID,
limit = 1000
}
dbLimit := limit * semconvDuplicateFactor(telemetrytypes.SignalLogs)
mainQuery := fmt.Sprintf(`
SELECT tag_key, tag_type, tag_data_type, max(priority) as priority
FROM (
@@ -673,7 +710,7 @@ func (t *telemetryMetaStore) getLogsKeys(ctx context.Context, orgID valuer.UUID,
GROUP BY tag_key, tag_type, tag_data_type
ORDER BY priority
LIMIT %d
`, strings.Join(queries, " UNION ALL "), limit+1)
`, strings.Join(queries, " UNION ALL "), dbLimit+1)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, mainQuery, allArgs...)
if err != nil {
@@ -694,7 +731,7 @@ func (t *telemetryMetaStore) getLogsKeys(ctx context.Context, orgID valuer.UUID,
for rows.Next() {
rowCount++
// reached the limit, we know there are more results
if rowCount > limit {
if rowCount > dbLimit {
break
}
@@ -742,7 +779,16 @@ func (t *telemetryMetaStore) getLogsKeys(ctx context.Context, orgID valuer.UUID,
}
// hit the limit? (only counting DB results)
complete := rowCount <= limit
complete := rowCount <= dbLimit
keys = canonicalizeSemconvKeys(keys, telemetrytypes.SignalLogs)
if len(keys) > limit {
keys = keys[:limit]
complete = false
}
mapOfKeys = make(map[string]*telemetrytypes.TelemetryFieldKey, len(keys))
for _, key := range keys {
mapOfKeys[key.Name+";"+key.FieldContext.StringValue()+";"+key.FieldDataType.StringValue()] = key
}
staticKeys := []string{}
staticKeys = append(staticKeys, maps.Keys(logstelemetryschema.IntrinsicFields)...)
@@ -1050,10 +1096,19 @@ func (t *telemetryMetaStore) getMetricsKeys(ctx context.Context, fieldKeySelecto
conds := []string{}
for _, fieldKeySelector := range fieldKeySelectors {
fieldConds := []string{}
members := attributeSemconvMembers(fieldKeySelector.Name, telemetrytypes.SignalMetrics, fieldKeySelector.FieldContext)
if fieldKeySelector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeExact {
fieldConds = append(fieldConds, sb.E("attr_name", fieldKeySelector.Name))
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
fieldConds = append(fieldConds, sb.In("attr_name", memberValues...))
} else {
fieldConds = append(fieldConds, sb.ILike("attr_name", "%"+escapeForLike(fieldKeySelector.Name)+"%"))
memberConds := make([]string, 0, len(members))
for _, member := range members {
memberConds = append(memberConds, sb.ILike("attr_name", "%"+escapeForLike(member)+"%"))
}
fieldConds = append(fieldConds, sb.Or(memberConds...))
}
fieldConds = append(fieldConds, sb.NotLike("attr_name", "\\_\\_%"))
@@ -1069,7 +1124,12 @@ func (t *telemetryMetaStore) getMetricsKeys(ctx context.Context, fieldKeySelecto
if fieldKeySelector.MetricContext != nil {
if fieldKeySelector.MetricContext.MetricName != "" {
fieldConds = append(fieldConds, sb.E("metric_name", fieldKeySelector.MetricContext.MetricName))
metricNames := semconv.MetricNames(fieldKeySelector.MetricContext.MetricName)
metricValues := make([]any, 0, len(metricNames))
for _, metricName := range metricNames {
metricValues = append(metricValues, metricName)
}
fieldConds = append(fieldConds, sb.In("metric_name", metricValues...))
}
if fieldKeySelector.MetricContext.MetricNamespace != "" {
fieldConds = append(fieldConds, sb.Like("metric_name", escapeForLike(fieldKeySelector.MetricContext.MetricNamespace)+"%"))
@@ -1091,7 +1151,8 @@ func (t *telemetryMetaStore) getMetricsKeys(ctx context.Context, fieldKeySelecto
mainSb.GroupBy("name", "field_context", "field_data_type")
mainSb.OrderBy("priority")
// query one extra to check if we hit the limit
mainSb.Limit(limit + 1)
dbLimit := limit * semconvDuplicateFactor(telemetrytypes.SignalMetrics)
mainSb.Limit(dbLimit + 1)
query, args := mainSb.BuildWithFlavor(sqlbuilder.ClickHouse)
@@ -1106,7 +1167,7 @@ func (t *telemetryMetaStore) getMetricsKeys(ctx context.Context, fieldKeySelecto
for rows.Next() {
rowCount++
// reached the limit, we know there are more results
if rowCount > limit {
if rowCount > dbLimit {
break
}
@@ -1131,7 +1192,12 @@ func (t *telemetryMetaStore) getMetricsKeys(ctx context.Context, fieldKeySelecto
}
// hit the limit?
complete := rowCount <= limit
complete := rowCount <= dbLimit
keys = canonicalizeSemconvKeys(keys, telemetrytypes.SignalMetrics)
if len(keys) > limit {
keys = keys[:limit]
complete = false
}
return keys, complete, nil
}
@@ -1153,16 +1219,30 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
var limit int
for _, fieldKeySelector := range fieldKeySelectors {
fieldConds := []string{}
members := attributeSemconvMembers(fieldKeySelector.Name, telemetrytypes.SignalMetrics, fieldKeySelector.FieldContext)
if fieldKeySelector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeExact {
fieldConds = append(fieldConds, sb.E("attr_name", fieldKeySelector.Name))
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
fieldConds = append(fieldConds, sb.In("attr_name", memberValues...))
} else {
fieldConds = append(fieldConds, sb.Like("attr_name", "%"+fieldKeySelector.Name+"%"))
memberConds := make([]string, 0, len(members))
for _, member := range members {
memberConds = append(memberConds, sb.Like("attr_name", "%"+member+"%"))
}
fieldConds = append(fieldConds, sb.Or(memberConds...))
}
fieldConds = append(fieldConds, sb.NotLike("attr_name", "\\_\\_%"))
if fieldKeySelector.MetricContext != nil {
if fieldKeySelector.MetricContext.MetricName != "" {
fieldConds = append(fieldConds, sb.E("metric_name", fieldKeySelector.MetricContext.MetricName))
metricNames := semconv.MetricNames(fieldKeySelector.MetricContext.MetricName)
metricValues := make([]any, 0, len(metricNames))
for _, metricName := range metricNames {
metricValues = append(metricValues, metricName)
}
fieldConds = append(fieldConds, sb.In("metric_name", metricValues...))
}
if fieldKeySelector.MetricContext.MetricNamespace != "" {
fieldConds = append(fieldConds, sb.Like("metric_name", escapeForLike(fieldKeySelector.MetricContext.MetricNamespace)+"%"))
@@ -1177,7 +1257,8 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
limit = 1000
}
sb.Limit(limit)
dbLimit := limit * semconvDuplicateFactor(telemetrytypes.SignalMetrics)
sb.Limit(dbLimit + 1)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
@@ -1191,7 +1272,7 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
for rows.Next() {
rowCount++
// reached the limit, we know there are more results
if rowCount > limit {
if rowCount > dbLimit {
break
}
@@ -1200,9 +1281,14 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
if err != nil {
return nil, false, errors.Wrap(err, errors.TypeInternal, errors.CodeInternal, ErrFailedToGetMeterKeys.Error())
}
fieldContext := telemetrytypes.FieldContextAttribute
if strings.HasPrefix(name, "resource_") {
fieldContext = telemetrytypes.FieldContextResource
}
keys = append(keys, &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalMetrics,
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: fieldContext,
})
}
@@ -1211,7 +1297,12 @@ func (t *telemetryMetaStore) getMeterSourceMetricKeys(ctx context.Context, field
}
// hit the limit?
complete := rowCount <= limit
complete := rowCount <= dbLimit
keys = canonicalizeSemconvKeys(keys, telemetrytypes.SignalMetrics)
if len(keys) > limit {
keys = keys[:limit]
complete = false
}
return keys, complete, nil
@@ -1570,7 +1661,6 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
if limit == 0 {
limit = 50
}
// query one extra to check if we hit the limit
sb.Limit(limit + 1)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
@@ -1725,7 +1815,12 @@ func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, fieldValueSe
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(t.logsDBName + "." + t.logsFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
members := attributeSemconvMembers(fieldValueSelector.Name, telemetrytypes.SignalLogs, fieldValueSelector.FieldContext)
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
sb.Where(sb.In("tag_key", memberValues...))
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1923,7 +2018,9 @@ func (t *telemetryMetaStore) getMetricFieldValues(ctx context.Context, orgID val
From(t.metricsDBName + "." + t.metricsFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("attr_name", fieldValueSelector.Name))
sb.Where(inStrings(sb, "attr_name", attributeSemconvMembers(
fieldValueSelector.Name, telemetrytypes.SignalMetrics, fieldValueSelector.FieldContext,
)))
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1935,7 +2032,7 @@ func (t *telemetryMetaStore) getMetricFieldValues(ctx context.Context, orgID val
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
sb.Where(inStrings(sb, "metric_name", semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)))
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2069,7 +2166,7 @@ func (t *telemetryMetaStore) getIntrinsicMetricFieldValuesForTable(ctx context.C
From(t.metricsDBName + "." + tableName)
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
sb.Where(inStrings(sb, "metric_name", semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)))
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2132,12 +2229,14 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
From(t.meterDBName + "." + t.meterFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("attr.1", fieldValueSelector.Name))
sb.Where(inStrings(sb, "attr.1", attributeSemconvMembers(
fieldValueSelector.Name, telemetrytypes.SignalMetrics, fieldValueSelector.FieldContext,
)))
}
sb.Where(sb.NotLike("attr.1", "\\_\\_%"))
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricName != "" {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
sb.Where(inStrings(sb, "metric_name", semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)))
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2156,8 +2255,10 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
if limit == 0 {
limit = 50
}
// query one extra to check if we hit the limit
sb.Limit(limit + 1)
// A value can be present under several physical spellings; over-fetch and
// de-duplicate after scanning so one family cannot consume the result limit.
dbLimit := limit * semconvDuplicateFactor(telemetrytypes.SignalMetrics)
sb.Limit(dbLimit + 1)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
@@ -2167,11 +2268,13 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
defer rows.Close()
values := &telemetrytypes.TelemetryFieldValues{}
seen := make(map[string]bool)
rowCount := 0
uniqueCount := 0
for rows.Next() {
rowCount++
// reached the limit, we know there are more results
if rowCount > limit {
if rowCount > dbLimit {
break
}
@@ -2180,12 +2283,16 @@ func (t *telemetryMetaStore) getMeterSourceMetricFieldValues(ctx context.Context
return nil, false, errors.Wrap(err, errors.TypeInternal, errors.CodeInternal, ErrFailedToGetMeterValues.Error())
}
if len(attribute) > 1 {
values.StringValues = append(values.StringValues, attribute[1])
if !seen[attribute[1]] && uniqueCount < limit {
values.StringValues = append(values.StringValues, attribute[1])
seen[attribute[1]] = true
uniqueCount++
}
}
}
// hit the limit?
complete := rowCount <= limit
complete := rowCount <= dbLimit && uniqueCount < limit
return values, complete, nil
}

View File

@@ -130,6 +130,44 @@ func TestCanonicalizeTraceSemconvKeys(t *testing.T) {
assert.Equal(t, "deployment.environment", result[2].Name, "phase 1 must not rewrite raw log metadata")
}
func TestCanonicalizeLogAndMetricSemconvKeys(t *testing.T) {
logKeys := canonicalizeSemconvKeys([]*telemetrytypes.TelemetryFieldKey{
{
Name: "db.system",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
{
Name: "db.system.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
}, telemetrytypes.SignalLogs)
require.Len(t, logKeys, 1)
assert.Equal(t, "db.system.name", logKeys[0].Name)
assert.Equal(t, []string{"db.system.name", "db.system"}, logKeys[0].SemconvMembers)
metricKeys := canonicalizeSemconvKeys([]*telemetrytypes.TelemetryFieldKey{
{
Name: "resource_db_system",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
}, telemetrytypes.SignalMetrics)
require.Len(t, metricKeys, 1)
assert.Equal(t, "db.system.name", metricKeys[0].Name)
assert.Equal(t, []string{"db.system.name", "resource_db_system"}, metricKeys[0].SemconvMembers)
}
func TestGetSpanFieldValuesMergesSemconvFamily(t *testing.T) {
mockTelemetryStore := telemetrystoretest.New(telemetrystore.Config{}, &regexMatcher{})
mock := mockTelemetryStore.Mock()

View File

@@ -89,10 +89,12 @@ func (v *TelemetryFieldVisitor) VisitColumnDef(expr *parser.ColumnDef) error {
// Create and store the TelemetryFieldKey
field := &telemetrytypes.TelemetryFieldKey{
Name: fieldName,
FieldContext: fieldContext,
FieldDataType: fieldDataType,
Materialized: true,
Name: fieldName,
FieldContext: fieldContext,
FieldDataType: fieldDataType,
Materialized: true,
MaterializedColumnName: columnName,
MaterializedSemconv: strings.Count(defaultExprStr, "['") > 1,
}
v.Fields = append(v.Fields, field)

View File

@@ -5,8 +5,25 @@ import (
"testing"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestExtractFieldKeyPreservesHistoricalMaterializedColumnName(t *testing.T) {
statement := `CREATE TABLE signoz_traces.signoz_index_v3
(
attributes_string Map(LowCardinality(String), String),
` + "`attribute_string_db$$system`" + ` LowCardinality(String)
DEFAULT if(mapContains(attributes_string, 'db.system.name'), attributes_string['db.system.name'], attributes_string['db.system'])
) ENGINE = MergeTree ORDER BY tuple()`
keys, err := ExtractFieldKeysFromTblStatement(statement)
require.NoError(t, err, "table statement should parse")
require.Len(t, keys, 1, "table statement should contain one materialized key")
assert.Equal(t, "db.system.name", keys[0].Name)
assert.Equal(t, "attribute_string_db$$system", keys[0].MaterializedColumnName)
assert.True(t, keys[0].MaterializedSemconv)
}
func TestExtractFieldKeysFromTblStatement(t *testing.T) {
var statement = `CREATE TABLE signoz_logs.logs_v2

View File

@@ -453,6 +453,9 @@ func (c *conditionBuilder) ConditionFor(
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
if options.ExactSemconv {
matches = querybuilder.MatchingFieldKeysExact(key, fieldKeys)
}
skipResourceFilter := options.SkipResourceFilter
// search() resolves its own (optional) scope; handle it before key resolution.
@@ -499,6 +502,9 @@ func (c *conditionBuilder) ConditionFor(
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
}
if options.ExactSemconv {
keys = querybuilder.ExactSemconvKeys(keys)
}
if skipResourceFilter && !synthesized {
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))

View File

@@ -10,6 +10,7 @@ import (
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -67,6 +68,20 @@ type fieldMapper struct {
fl flagger.Flagger
}
func logSemconvMembers(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.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalLogs,
FieldContext: key.FieldContext,
})
}
func NewFieldMapper(fl flagger.Flagger) qbtypes.FieldMapper {
return &fieldMapper{fl: fl}
}
@@ -141,8 +156,20 @@ func (m *fieldMapper) FieldFor(ctx context.Context, orgID valuer.UUID, tsStart,
case schema.ColumnTypeEnumJSON:
switch key.FieldContext {
case telemetrytypes.FieldContextResource:
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExpr = append(existExpr, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
members := logSemconvMembers(key)
if len(members) > 1 {
values := make([]string, 0, len(members))
guards := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(%s.`%s`::String, '')", columnName, member))
guards = append(guards, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, member))
}
exprs = append(exprs, "COALESCE("+strings.Join(values, ", ")+")")
existExpr = append(existExpr, "("+strings.Join(guards, " OR ")+")")
} else {
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, members[0]))
existExpr = append(existExpr, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, members[0]))
}
case telemetrytypes.FieldContextBody:
if key.Name == messageSubField {
exprs = append(exprs, messageSubColumn)
@@ -181,13 +208,32 @@ func (m *fieldMapper) FieldFor(ctx context.Context, orgID valuer.UUID, tsStart,
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumBool, schema.ColumnTypeEnumFloat64:
// a key could have been materialized, if so return the materialized column name
if key.Materialized {
members := logSemconvMembers(key)
if key.Materialized && (len(members) == 1 || key.MaterializedSemconv) {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExpr = append(existExpr, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
} else 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+1)
for i, member := range members {
branches = append(branches, guards[i], fmt.Sprintf("%s['%s']", columnName, member))
}
exprs = append(exprs, "multiIf("+strings.Join(branches, ", ")+", NULL)")
}
existExpr = append(existExpr, "("+strings.Join(guards, " OR ")+")")
} else {
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, key.Name))
existExpr = append(existExpr, fmt.Sprintf("mapContains(%s, '%s')", columnName, key.Name))
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, members[0]))
existExpr = append(existExpr, fmt.Sprintf("mapContains(%s, '%s')", columnName, members[0]))
}
default:
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "exists operator is not supported for map column type %s", valueType)

View File

@@ -2,6 +2,7 @@ package logstelemetryschema
import (
"context"
"strings"
"testing"
"time"
@@ -14,6 +15,32 @@ import (
"github.com/stretchr/testify/require"
)
func TestFieldForSemconvFamily(t *testing.T) {
ctx := context.Background()
fm := NewFieldMapper(flaggertest.New(t)).(*fieldMapper)
key := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := fm.FieldFor(ctx, valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
assert.Equal(t, "COALESCE(NULLIF(attributes_string['db.system.name'], ''), NULLIF(attributes_string['db.system'], ''))", expression)
assert.Less(t, strings.Index(expression, "db.system.name"), strings.Index(expression, "db.system']"), "current spelling must win")
exists, err := fm.existsExpressionFor(ctx, valuer.UUID{}, 0, 0, key, true)
require.NoError(t, err)
assert.Equal(t, "(mapContains(attributes_string, 'db.system.name') OR mapContains(attributes_string, 'db.system'))", exists)
exact := *key
exact.SemconvMembers = []string{"db.system"}
expression, err = fm.FieldFor(ctx, valuer.UUID{}, 0, 0, &exact)
require.NoError(t, err)
assert.Equal(t, "attributes_string['db.system']", expression)
}
func TestGetColumn(t *testing.T) {
ctx := context.Background()

View File

@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/querybuilder"
@@ -135,10 +136,22 @@ func (c *conditionBuilder) conditionFor(
return "true", nil
}
if operator == qbtypes.FilterOperatorExists {
return fmt.Sprintf("has(JSONExtractKeys(labels), '%s')", key.Name), nil
members := metricAttributeMembers(key)
guards := make([]string, 0, len(members))
for _, member := range members {
guards = append(guards, fmt.Sprintf("has(JSONExtractKeys(labels), '%s')", member))
}
return fmt.Sprintf("not has(JSONExtractKeys(labels), '%s')", key.Name), nil
guard := strings.Join(guards, " OR ")
if len(guards) > 1 {
guard = "(" + guard + ")"
}
if operator == qbtypes.FilterOperatorExists {
return guard, nil
}
if len(guards) == 1 {
return "not " + guard, nil
}
return "NOT " + guard, nil
}
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported operator: %v", operator)
}
@@ -151,7 +164,7 @@ func (c *conditionBuilder) ConditionFor(
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey,
_ qbtypes.ConditionBuilderOptions,
options qbtypes.ConditionBuilderOptions,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -162,7 +175,15 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
requestedKey := *key
if requestedKey.Signal == telemetrytypes.SignalUnspecified {
requestedKey.Signal = telemetrytypes.SignalMetrics
}
key = &requestedKey
keys := querybuilder.MatchingFieldKeys(key, fieldKeys)
if options.ExactSemconv {
keys = querybuilder.MatchingFieldKeysExact(key, fieldKeys)
}
var warnings []string
if len(keys) == 0 {
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {
@@ -180,6 +201,9 @@ func (c *conditionBuilder) ConditionFor(
}
}
}
if options.ExactSemconv {
keys = querybuilder.ExactSemconvKeys(keys)
}
conds := make([]string, 0, len(keys))
for _, k := range keys {

View File

@@ -390,3 +390,45 @@ func TestConditionForKeyNotInMetadata(t *testing.T) {
})
}
}
func TestConditionForSemconvMetricLabels(t *testing.T) {
ctx := context.Background()
fm := NewFieldMapper()
conditionBuilder := NewConditionBuilder(fm)
requested := telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
current := requested
legacyNormalized := requested
legacyNormalized.Name = "resource_db_system"
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {&current},
legacyNormalized.Name: {&legacyNormalized},
}
sb := sqlbuilder.NewSelectBuilder()
conditions, warnings, err := conditionBuilder.ConditionFor(
ctx, valuer.UUID{}, 0, 0, &requested, fieldKeys, qbtypes.ConditionBuilderOptions{},
qbtypes.FilterOperatorEqual, "postgresql", sb,
)
require.NoError(t, err)
assert.Empty(t, warnings)
sb.Where(conditions...)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, query, "COALESCE(NULLIF(JSONExtractString(labels, 'db.system.name'), ''), NULLIF(JSONExtractString(labels, 'resource_db_system'), '')) = ?")
assert.Equal(t, []any{"postgresql"}, args)
sb = sqlbuilder.NewSelectBuilder()
conditions, _, err = conditionBuilder.ConditionFor(
ctx, valuer.UUID{}, 0, 0, &requested, fieldKeys, qbtypes.ConditionBuilderOptions{ExactSemconv: true},
qbtypes.FilterOperatorExists, nil, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
query, _ = sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, query, "has(JSONExtractKeys(labels), 'db.system.name')")
assert.NotContains(t, query, "resource_db_system")
}

View File

@@ -4,8 +4,10 @@ import (
"context"
"fmt"
"slices"
"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"
@@ -38,6 +40,23 @@ var (
type fieldMapper struct{}
func metricAttributeMembers(key *telemetrytypes.TelemetryFieldKey) []string {
if key.FieldContext != telemetrytypes.FieldContextResource &&
key.FieldContext != telemetrytypes.FieldContextScope &&
key.FieldContext != telemetrytypes.FieldContextAttribute &&
key.FieldContext != telemetrytypes.FieldContextUnspecified {
return []string{key.Name}
}
if len(key.SemconvMembers) > 0 {
return key.SemconvMembers
}
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: key.Name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: key.FieldContext,
})
}
// CandidateKeys returns nil: metrics has no attribute-map fallback, so a context-missing
// key stays unresolved and the caller errors.
func (m *fieldMapper) CandidateKeys(_ context.Context, _ valuer.UUID, _ *telemetrytypes.TelemetryFieldKey, _ any, _ map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
@@ -80,19 +99,30 @@ func (m *fieldMapper) FieldFor(ctx context.Context, _ valuer.UUID, startNs, endN
switch key.FieldContext {
case telemetrytypes.FieldContextResource, telemetrytypes.FieldContextScope, telemetrytypes.FieldContextAttribute:
return fmt.Sprintf("JSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
return metricLabelExpression(columns[0].Name, metricAttributeMembers(key)), nil
case telemetrytypes.FieldContextMetric:
return columns[0].Name, nil
case telemetrytypes.FieldContextUnspecified:
if slices.Contains(IntrinsicFields, key.Name) {
return columns[0].Name, nil
}
return fmt.Sprintf("JSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
return metricLabelExpression(columns[0].Name, metricAttributeMembers(key)), nil
}
return columns[0].Name, nil
}
func metricLabelExpression(columnName string, members []string) string {
if len(members) == 1 {
return fmt.Sprintf("JSONExtractString(%s, '%s')", columnName, members[0])
}
values := make([]string, 0, len(members))
for _, member := range members {
values = append(values, fmt.Sprintf("NULLIF(JSONExtractString(%s, '%s'), '')", columnName, member))
}
return "COALESCE(" + strings.Join(values, ", ") + ")"
}
func (m *fieldMapper) ColumnFor(ctx context.Context, _ valuer.UUID, tsStart, tsEnd uint64, key *telemetrytypes.TelemetryFieldKey) ([]*schema.Column, error) {
return m.getColumn(ctx, tsStart, tsEnd, key)
}

View File

@@ -2,6 +2,7 @@ package metricstelemetryschema
import (
"context"
"strings"
"testing"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
@@ -12,6 +13,32 @@ import (
"github.com/stretchr/testify/require"
)
func TestMetricLabelSemconvSpellings(t *testing.T) {
fm := NewFieldMapper()
key := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := fm.FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
for _, member := range []string{
"resource_db.system.name", "resource_db_system_name", "db.system.name", "db_system_name",
"resource_db.system", "resource_db_system", "db.system", "db_system",
} {
assert.Contains(t, expression, "'"+member+"'")
}
assert.Less(t, strings.Index(expression, "resource_db.system.name"), strings.Index(expression, "resource_db.system'"))
exact := *key
exact.SemconvMembers = []string{"resource_db_system"}
expression, err = fm.FieldFor(context.Background(), valuer.UUID{}, 0, 0, &exact)
require.NoError(t, err)
assert.Equal(t, "JSONExtractString(labels, 'resource_db_system')", expression)
}
func TestGetColumn(t *testing.T) {
ctx := context.Background()

View File

@@ -220,6 +220,9 @@ func (c *conditionBuilder) ConditionFor(
}
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
if options.ExactSemconv {
matches = querybuilder.MatchingFieldKeysExact(key, fieldKeys)
}
skipResourceFilter := options.SkipResourceFilter
keys, warning := querybuilder.ResolveKeys(key, matches)
@@ -265,6 +268,9 @@ func (c *conditionBuilder) ConditionFor(
synthesized = true
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
if options.ExactSemconv {
keys = querybuilder.ExactSemconvKeys(keys)
}
// When a resource sub-query already covers the term, drop resource keys from the main
// query. Synthesized keys are exempt: the sub-query skips keys absent from metadata.

View File

@@ -362,7 +362,10 @@ 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 {
if key.Materialized && (len(members) == 1 || key.MaterializedSemconv) {
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
} else if len(members) > 1 {
guards := make([]string, 0, len(members))
for _, member := range members {
guards = append(guards, fmt.Sprintf("mapContains(%s, '%s')", columnName, member))
@@ -381,12 +384,6 @@ func (m *fieldMapper) resolveColumnExprs(
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, members[0]))
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, members[0]))

View File

@@ -5,6 +5,8 @@ import (
"testing"
"time"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"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"
@@ -348,3 +350,27 @@ func TestColumnExpressionForTimestampAttributeCollision(t *testing.T) {
assert.Contains(t, result, "attributes_number['timestamp']")
})
}
func TestDBSystemFamilyUsesSemconvAwareMaterializedColumn(t *testing.T) {
fm := NewFieldMapper()
key := &telemetrytypes.TelemetryFieldKey{
Name: "db.system.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
MaterializedColumnName: "attribute_string_db$$system",
MaterializedSemconv: true,
SemconvMembers: []string{"db.system.name", "db.system"},
}
expression, err := fm.FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
assert.Equal(t, "`attribute_string_db$$system`", expression)
exists, err := querybuilder.ExistsExpression(
[]*schema.Column{indexV3Columns["attributes_string"]}, key, 0, 0, expression, true,
)
require.NoError(t, err)
assert.Equal(t, "`attribute_string_db$$system_exists`", exists)
}

View File

@@ -45,6 +45,10 @@ type ConditionBuilder interface {
type ConditionBuilderOptions struct {
// SkipResourceFilter drops the resource context from the candidate set.
SkipResourceFilter bool
// ExactSemconv disables semantic-convention family expansion. It is an
// internal escape hatch for diagnostics and migrations that must address one
// physical spelling only; public query APIs continue to resolve families.
ExactSemconv bool
}
type AggExprRewriter interface {
// Rewrite rewrites the aggregation expression to be used in the query.

View File

@@ -46,6 +46,13 @@ type TelemetryFieldKey struct {
JSONPlan JSONAccessPlan `json:"-"`
Indexes []TelemetryFieldKeySkipIndex `json:"-"`
Materialized bool `json:"-"` // refers to promoted in case of body.... fields
// MaterializedColumnName preserves the physical column when its DEFAULT
// expression resolves a newer semantic-convention key than the historical
// column identifier (for example db.system.name in db$$system).
MaterializedColumnName string `json:"-"`
// MaterializedSemconv is true when the column's DEFAULT expression already
// coalesces every enabled family member and is therefore safe for a family query.
MaterializedSemconv bool `json:"-"`
Evolutions []*EvolutionEntry `json:"-"`
SemconvMembers []string `json:"-"`
@@ -127,6 +134,8 @@ func (f *TelemetryFieldKey) OverrideMetadataFrom(src *TelemetryFieldKey) {
f.FieldDataType = src.FieldDataType
f.Indexes = src.Indexes
f.Materialized = src.Materialized
f.MaterializedColumnName = src.MaterializedColumnName
f.MaterializedSemconv = src.MaterializedSemconv
f.JSONPlan = src.JSONPlan
f.Evolutions = src.Evolutions
f.SemconvMembers = src.SemconvMembers
@@ -204,6 +213,9 @@ func TelemetryFieldKeyToText(key *TelemetryFieldKey) string {
}
func FieldKeyToMaterializedColumnName(key *TelemetryFieldKey) string {
if key.MaterializedColumnName != "" {
return fmt.Sprintf("`%s`", key.MaterializedColumnName)
}
return fmt.Sprintf("`%s_%s_%s`",
key.FieldContext.String,
fieldDataTypes[key.FieldDataType.StringValue()].StringValue(),
@@ -212,6 +224,9 @@ func FieldKeyToMaterializedColumnName(key *TelemetryFieldKey) string {
}
func FieldKeyToMaterializedColumnNameForExists(key *TelemetryFieldKey) string {
if key.MaterializedColumnName != "" {
return fmt.Sprintf("`%s_exists`", key.MaterializedColumnName)
}
return fmt.Sprintf("`%s_%s_%s_exists`",
key.FieldContext.String,
fieldDataTypes[key.FieldDataType.StringValue()].StringValue(),

View File

@@ -7,5 +7,31 @@ default_enabled: false
families:
deployment.environment.name:
enabled: true
contexts: [resource, attribute]
signals: [traces, logs, metrics]
db.system.name:
enabled: true
contexts: [resource, attribute]
signals: [traces, logs, metrics]
# These metric renames predate the schema history vendored above. Keep them
# in the same generated registry so every v5 metric query uses one source of
# truth instead of the legacy hand-written transition table.
k8s.pod.cpu.usage:
enabled: true
kind: metric
old: [k8s.pod.cpu.utilization]
contexts: [metric]
signals: [metrics]
k8s.node.cpu.usage:
enabled: true
kind: metric
old: [k8s.node.cpu.utilization]
contexts: [metric]
signals: [metrics]
container.cpu.usage:
enabled: true
kind: metric
old: [container.cpu.utilization]
contexts: [metric]
signals: [metrics]

View File

@@ -0,0 +1,157 @@
"""Phase 2 semantic-convention checks across logs and metrics."""
from collections.abc import Callable
from datetime import UTC, 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.logs import Logs
from fixtures.metrics import Metrics
DB_CURRENT = "db.system.name"
DB_OLD = "db.system"
METRIC_CURRENT = "container.cpu.usage"
METRIC_OLD = "container.cpu.utilization"
PREFIX = "semconv-phase2"
def _raw_log_bodies(
signoz: types.SigNoz,
token: str,
now: datetime,
expression: str,
) -> set[str]:
response = querier.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=querier.RequestType.RAW,
queries=[querier.build_raw_query("A", "logs", limit=100, filter_expression=expression)],
)
assert response.status_code == HTTPStatus.OK, response.text
return {row["data"]["body"] for row in querier.get_rows(response)}
def _field_values(signoz: types.SigNoz, token: str, signal: 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": signal,
"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_logs_resolve_db_system_family(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_logs: Callable[[list[Logs]], None],
) -> None:
now = datetime.now(tz=UTC).replace(microsecond=0) - timedelta(minutes=1)
rows = [
("old", {DB_OLD: "postgresql"}),
("current", {DB_CURRENT: "postgresql"}),
("conflict", {DB_OLD: "mysql", DB_CURRENT: "postgresql"}),
("missing", {}),
]
insert_logs(
[
Logs(
timestamp=now + timedelta(seconds=index),
resources={"service.name": PREFIX, **attributes},
attributes=attributes,
body=f"{PREFIX}-{suffix}",
)
for index, (suffix, attributes) in enumerate(rows)
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
present = {f"{PREFIX}-old", f"{PREFIX}-current", f"{PREFIX}-conflict"}
for context in ("attribute", "resource"):
for requested in (DB_CURRENT, DB_OLD):
field = f"{context}.{requested}"
assert _raw_log_bodies(signoz, token, now, f'{field} = "postgresql"') == present
assert _raw_log_bodies(signoz, token, now, f"{field} EXISTS") == present
assert _raw_log_bodies(signoz, token, now, f"{field} NOT EXISTS") == {f"{PREFIX}-missing"}
assert _field_values(signoz, token, "logs", requested, context) == {"postgresql", "mysql"}
def test_metrics_resolve_label_and_metric_name_families(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_metrics: Callable[[list[Metrics]], None],
) -> None:
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
insert_metrics(
[
Metrics(
metric_name=METRIC_CURRENT,
labels={DB_CURRENT: "postgresql"},
timestamp=now - timedelta(seconds=3),
temporality="Unspecified",
type_="Gauge",
is_monotonic=False,
value=10,
),
Metrics(
metric_name=METRIC_OLD,
labels={"db_system": "mysql"},
timestamp=now - timedelta(seconds=2),
temporality="Unspecified",
type_="Gauge",
is_monotonic=False,
value=20,
),
Metrics(
metric_name=METRIC_OLD,
labels={DB_OLD: "mysql", DB_CURRENT: "postgresql", "series": "conflict"},
timestamp=now - timedelta(seconds=1),
temporality="Unspecified",
type_="Gauge",
is_monotonic=False,
value=30,
),
]
)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
for metric_name in (METRIC_CURRENT, METRIC_OLD):
for requested_label in (DB_CURRENT, DB_OLD):
response = querier.make_scalar_query_request(
signoz,
token,
now,
[
querier.build_scalar_query(
name="A",
signal="metrics",
aggregations=[
querier.build_metrics_aggregation(
metric_name,
"latest",
"sum",
"unspecified",
reduce_to="last",
)
],
group_by=[querier.build_group_by_field(requested_label, "string", "attribute")],
filter_expression=f"attribute.{requested_label} EXISTS",
)
],
)
assert response.status_code == HTTPStatus.OK, response.text
data = {row[0]: row[-1] for row in querier.get_scalar_table_data(response.json())}
assert data == {"postgresql": 40.0, "mysql": 20.0}, (metric_name, requested_label, data)