mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-06 21:20:42 +01:00
Compare commits
1 Commits
test/semco
...
feat/semco
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6b4bb82efb |
4
Makefile
4
Makefile
@@ -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"
|
||||
|
||||
@@ -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: {},
|
||||
},
|
||||
|
||||
@@ -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())
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -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,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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{}, ®exMatcher{})
|
||||
mock := mockTelemetryStore.Mock()
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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: {¤t},
|
||||
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")
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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]))
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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]
|
||||
|
||||
157
tests/integration/tests/queriersemconv/02_cross_signal.py
Normal file
157
tests/integration/tests/queriersemconv/02_cross_signal.py
Normal 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)
|
||||
Reference in New Issue
Block a user