Compare commits

...

1 Commits

Author SHA1 Message Date
srikanthccv
b902642d59 feat: resolve semconv families across logs and metrics
Phase 2 extends family resolution to logs and metrics behind the
resolve_semconv_families flag, and reworks the generated registry to
per-member scope records.

- The registry stores one scope per member, from the member's own
  rename edges: a fanned-out old name keeps one membership per target,
  a scoped-and-unscoped merge stays unbounded, and a cross-context
  rename keeps each spelling where it existed. The
  missing-information policy is one rule on every axis: an unset
  selector axis constrains nothing, and a name that stays ambiguous
  stays literal. The family fields are unexported; All() iterates.
- AttributeMembers expands metric members into their stored label
  spellings (dotted, normalized, and resource_-prefixed) and
  MetricNames returns the storage names of a metric-name family in the
  requested style.
- The logs and metrics condition builders compile family logical
  fields through one shared compiler; the metrics group-by column
  merges the spellings; every metric_name filter unions the family
  storage names; the metrics builder threads the queried metric name
  into resolution and fans the metadata prefetch out per storage name.
- The numeric and boolean family tails read 0 and false, so keyless
  rows keep single-key semantics.
- Values suggestions union the family spellings for traces, logs, and
  metrics, and related values merge the member columns current-first.
- The legacy transition table is gone: the v3 and v4 readers resolve
  metric renames through the registry with unchanged semantics, and
  the overlay declares the three CPU metric-name families.
- The overlay gates deployment.environment.name to
  traces+logs+metrics and db.system.name to traces+logs: the db.system
  value domain also renamed, so metrics stays literal for it until a
  value-mapping reader exists.

Flag-off pins: literal statement-builder goldens for logs and metrics,
generator output text pins, and the full pkg and ee suites.

Assisted-by: Claude Fable 5
2026-08-20 11:37:44 +05:30
44 changed files with 2038 additions and 455 deletions

View File

@@ -67,7 +67,7 @@ jobs:
with:
go-version: "1.24"
- name: check-semconv-generated-files
run: go run ./scripts/semconv -check
run: make semconv-check
build:
if: |
github.event_name == 'merge_group' ||

View File

@@ -237,6 +237,10 @@ py-clean: ## Clear all pycache and pytest cache from tests directory recursively
semconv-generate: ## Regenerate semantic-convention families for Go and TypeScript
@go run ./scripts/semconv
.PHONY: semconv-check
semconv-check: ## Fail if the generated semantic-convention files are stale
@go run ./scripts/semconv -check
.PHONY: gen-mocks
gen-mocks:
@echo ">> Generating mocks"

View File

@@ -1,32 +1,72 @@
// Code generated by scripts/semconv. DO NOT EDIT.
export type SemconvFamily = {
readonly current: string;
readonly old: readonly string[];
readonly kind: 'attribute' | 'metric';
// An empty contexts/signals/applyToMetrics array places no constraint on
// that axis.
export type SemconvMember = {
readonly name: string;
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly applyToMetrics: readonly string[];
};
export type SemconvFamily = {
readonly current: string;
readonly kind: 'attribute' | 'metric';
readonly members: readonly SemconvMember[];
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly valueMap: Readonly<Record<string, string>>;
};
export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
{
current: 'db.system.name',
old: ['db.system'],
kind: 'attribute',
current: 'container.cpu.usage',
kind: 'metric',
members: [
{ name: 'container.cpu.utilization', contexts: [], signals: [], applyToMetrics: [] },
],
contexts: [],
signals: [],
applyToMetrics: [],
valueMap: {},
},
{
current: 'db.system.name',
kind: 'attribute',
members: [
{ name: 'db.system', contexts: [], signals: [], applyToMetrics: [] },
],
contexts: [],
signals: ['logs', 'traces'],
valueMap: {},
},
{
current: 'deployment.environment.name',
old: ['deployment.environment'],
kind: 'attribute',
members: [
{ name: 'deployment.environment', contexts: [], signals: [], applyToMetrics: [] },
],
contexts: [],
signals: ['logs', 'metrics', 'traces'],
valueMap: {},
},
{
current: 'k8s.node.cpu.usage',
kind: 'metric',
members: [
{ name: 'k8s.node.cpu.utilization', contexts: [], signals: [], applyToMetrics: [] },
],
contexts: [],
signals: [],
valueMap: {},
},
{
current: 'k8s.pod.cpu.usage',
kind: 'metric',
members: [
{ name: 'k8s.pod.cpu.utilization', contexts: [], signals: [], applyToMetrics: [] },
],
contexts: [],
signals: [],
applyToMetrics: [],
valueMap: {},
},
] as const;

View File

@@ -40,8 +40,8 @@ func NewModule(
providerSettings factory.ProviderSettings,
cfg inframonitoring.Config,
) inframonitoring.Module {
fieldMapper := metricstelemetryschema.NewFieldMapper()
condBuilder := metricstelemetryschema.NewConditionBuilder(fieldMapper)
fieldMapper := metricstelemetryschema.NewFieldMapper(fl)
condBuilder := metricstelemetryschema.NewConditionBuilder(fieldMapper, fl)
return &module{
telemetryStore: telemetryStore,
telemetryMetadataStore: telemetryMetadataStore,

View File

@@ -49,8 +49,8 @@ type module struct {
// NewModule constructs the metrics module with the provided dependencies.
func NewModule(ts telemetrystore.TelemetryStore, telemetryMetadataStore telemetrytypes.MetadataStore, cache cache.Cache, ruleStore ruletypes.RuleStore, dashboardModule dashboard.Module, fl flagger.Flagger, providerSettings factory.ProviderSettings, cfg metricsexplorer.Config) metricsexplorer.Module {
fieldMapper := metricstelemetryschema.NewFieldMapper()
condBuilder := metricstelemetryschema.NewConditionBuilder(fieldMapper)
fieldMapper := metricstelemetryschema.NewFieldMapper(fl)
condBuilder := metricstelemetryschema.NewConditionBuilder(fieldMapper, fl)
return &module{
telemetryStore: ts,
fieldMapper: fieldMapper,

View File

@@ -42,7 +42,7 @@ func (c *conditionBuilder) ConditionFor(
// Rule state history fields have no family support, so every logical field
// is single-member and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(ctx, orgID, nil, key, fieldKeys))
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(ctx, orgID, nil, telemetrytypes.SignalUnspecified, nil, key, fieldKeys))
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {

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 (
@@ -3201,8 +3201,13 @@ 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))
// Exact-name members keep the legacy union semantics: only the canonical
// dotted spellings widen, exactly like the transition table this replaced.
names := semconv.Members(semconv.KindMetric, telemetrytypes.FieldKeySelector{
Name: req.AggregateAttribute,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
})
rows, err = r.db.Query(ctx, query, req.FilterAttributeKey, names, req.FilterAttributeKey, fmt.Sprintf("%%%s%%", req.SearchText), common.PastDayRoundOff())

View File

@@ -1,14 +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",
}
func GetTransitionedMetric(metric string) string {
if transitionedMetric, ok := MetricsUnderTransition[metric]; ok {
return transitionedMetric
}
return metric
}

View File

@@ -10,8 +10,9 @@ 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"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// ValidateAndCastValue validates and casts the value of a key to the corresponding data type of the key
@@ -232,14 +233,17 @@ func ClickHouseFormattedValue(v interface{}) string {
}
}
// The exact-name lookup keeps the legacy substitution semantics: only the
// canonical dotted spellings redirect, exactly like the transition table this
// replaced. Style-aware matching (MetricNames) is for the flagged v5 path.
func ClickHouseFormattedMetricNames(v interface{}) string {
if name, ok := v.(string); ok {
transitionedMetrics := metrics.GetTransitionedMetric(name)
if transitionedMetrics != name {
return ClickHouseFormattedValue([]interface{}{transitionedMetrics})
} else {
return ClickHouseFormattedValue([]interface{}{name})
}
current := semconv.Current(semconv.KindMetric, telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextMetric,
})
return ClickHouseFormattedValue([]interface{}{current})
}
return ClickHouseFormattedValue(v)

View File

@@ -483,3 +483,22 @@ func TestGetEpochNanoSecs(t *testing.T) {
})
}
}
// Only the canonical dotted spelling of a metric-name family redirects on the
// legacy path; a normalized spelling keeps reading its own series.
func TestClickHouseFormattedMetricNames(t *testing.T) {
cases := []struct {
name string
expected string
}{
{name: "k8s.pod.cpu.utilization", expected: "['k8s.pod.cpu.usage']"},
{name: "k8s.pod.cpu.usage", expected: "['k8s.pod.cpu.usage']"},
{name: "k8s_pod_cpu_utilization", expected: "['k8s_pod_cpu_utilization']"},
{name: "http.server.duration", expected: "['http.server.duration']"},
}
for _, c := range cases {
if got := ClickHouseFormattedMetricNames(c.name); got != c.expected {
t.Errorf("ClickHouseFormattedMetricNames(%q) = %q, want %q", c.name, got, c.expected)
}
}
}

View File

@@ -0,0 +1,125 @@
package querybuilder
import (
"context"
"fmt"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
)
// LogicalFamilyCondition compiles one condition for a family logical field
// from the mapper's primitives. Family members are map-backed attribute or
// resource keys by construction, so the compiler has no column-specific
// branches; a signal keeps its own single-member paths and hands only
// families here.
func LogicalFamilyCondition(
ctx context.Context,
orgID valuer.UUID,
startNs, endNs uint64,
fm qbtypes.FieldMapper,
logical *telemetrytypes.LogicalField,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
) (string, error) {
if operator.IsStringSearchOperator() {
value = FormatValueForContains(value)
}
fieldExpression, err := LogicalValueExpr(ctx, orgID, startNs, endNs, fm, logical)
if err != nil {
return "", err
}
// Coercion switches only on the data type, which every member shares, so
// the first member stands in for the field.
fieldExpression, value = DataTypeCollisionHandledFieldName(logical.Single(), value, fieldExpression, operator)
switch operator {
case qbtypes.FilterOperatorEqual:
return sb.E(fieldExpression, value), nil
case qbtypes.FilterOperatorNotEqual:
return sb.NE(fieldExpression, value), nil
case qbtypes.FilterOperatorGreaterThan:
return sb.G(fieldExpression, value), nil
case qbtypes.FilterOperatorGreaterThanOrEq:
return sb.GE(fieldExpression, value), nil
case qbtypes.FilterOperatorLessThan:
return sb.LT(fieldExpression, value), nil
case qbtypes.FilterOperatorLessThanOrEq:
return sb.LE(fieldExpression, value), nil
case qbtypes.FilterOperatorLike:
return sb.Like(fieldExpression, value), nil
case qbtypes.FilterOperatorNotLike:
return sb.NotLike(fieldExpression, value), nil
case qbtypes.FilterOperatorILike:
return sb.ILike(fieldExpression, value), nil
case qbtypes.FilterOperatorNotILike:
return sb.NotILike(fieldExpression, value), nil
case qbtypes.FilterOperatorContains:
return sb.ILike(fieldExpression, fmt.Sprintf("%%%s%%", value)), nil
case qbtypes.FilterOperatorNotContains:
return sb.NotILike(fieldExpression, fmt.Sprintf("%%%s%%", value)), nil
case qbtypes.FilterOperatorRegexp:
return fmt.Sprintf(`match(%s, %s)`, sqlbuilder.Escape(fieldExpression), sb.Var(value)), nil
case qbtypes.FilterOperatorNotRegexp:
return fmt.Sprintf(`NOT match(%s, %s)`, sqlbuilder.Escape(fieldExpression), sb.Var(value)), nil
case qbtypes.FilterOperatorBetween:
values, ok := value.([]any)
if !ok || len(values) != 2 {
return "", qbtypes.ErrBetweenValues
}
return sb.Between(fieldExpression, values[0], values[1]), nil
case qbtypes.FilterOperatorNotBetween:
values, ok := value.([]any)
if !ok || len(values) != 2 {
return "", qbtypes.ErrBetweenValues
}
return sb.NotBetween(fieldExpression, values[0], values[1]), nil
// `=`+OR / `!=`+AND instead of IN / NOT IN, to make use of the index
case qbtypes.FilterOperatorIn:
values, ok := value.([]any)
if !ok {
return "", qbtypes.ErrInValues
}
conditions := make([]string, 0, len(values))
for _, item := range values {
cond, err := LogicalFamilyCondition(ctx, orgID, startNs, endNs, fm, logical, qbtypes.FilterOperatorEqual, item, sb)
if err != nil {
return "", err
}
conditions = append(conditions, cond)
}
return sb.Or(conditions...), nil
case qbtypes.FilterOperatorNotIn:
values, ok := value.([]any)
if !ok {
return "", qbtypes.ErrInValues
}
conditions := make([]string, 0, len(values))
for _, item := range values {
cond, err := LogicalFamilyCondition(ctx, orgID, startNs, endNs, fm, logical, qbtypes.FilterOperatorNotEqual, item, sb)
if err != nil {
return "", err
}
conditions = append(conditions, cond)
}
return sb.And(conditions...), nil
case qbtypes.FilterOperatorExists, qbtypes.FilterOperatorNotExists:
pred, err := LogicalExistsExpr(ctx, orgID, startNs, endNs, fm, logical, operator == qbtypes.FilterOperatorExists)
if err != nil {
return "", err
}
return sqlbuilder.Escape(pred), nil
}
return "", qbtypes.ErrUnsupportedOperator
}

View File

@@ -10,27 +10,26 @@ import (
"github.com/SigNoz/signoz/pkg/valuer"
)
// semconvFamiliesEnabled evaluates the resolve_semconv_families flag for the
// SemconvFamiliesEnabled evaluates the resolve_semconv_families flag for the
// org. A nil flagger means off, so a caller without family support stays
// literal by default.
func semconvFamiliesEnabled(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger) bool {
func SemconvFamiliesEnabled(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger) bool {
if fl == nil {
return false
}
return fl.BooleanOrEmpty(ctx, flagger.FeatureResolveSemconvFamilies, featuretypes.NewFlaggerEvaluationContext(orgID))
}
// ExpandKeySelectorsForFamilies adds selectors for the other members of each
// semantic-convention family that a selector names. The metadata fetched for
// a query then contains each spelling that MatchingLogicalFields can group.
// This function is the prefetch of the resolution layer: statement builders
// call it after they derive the selectors, and the metadata store stays
// family-blind (autocomplete responses keep the literal spelling that the
// user typed). It does nothing when the resolve_semconv_families flag is off
// for the org. Only trace selectors expand today, because that matches the
// family support. Fuzzy (search-style) selectors never expand.
// ExpandKeySelectorsForFamilies adds selectors for the other spellings of
// each semantic-convention family that a selector names. The metadata fetched
// for a query then contains each spelling that MatchingLogicalFields can
// group. This function is the prefetch of the resolution layer: statement
// builders call it after they derive the selectors, and the metadata store
// stays family-blind (autocomplete responses keep the literal spelling that
// the user typed). It does nothing when the resolve_semconv_families flag is
// off for the org. Fuzzy (search-style) selectors never expand.
func ExpandKeySelectorsForFamilies(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, selectors []*telemetrytypes.FieldKeySelector) []*telemetrytypes.FieldKeySelector {
if !semconvFamiliesEnabled(ctx, orgID, fl) {
if !SemconvFamiliesEnabled(ctx, orgID, fl) {
return selectors
}
@@ -41,14 +40,14 @@ func ExpandKeySelectorsForFamilies(ctx context.Context, orgID valuer.UUID, fl fl
}
for _, selector := range selectors {
if selector.Signal != telemetrytypes.SignalTraces ||
selector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeFuzzy {
if selector.SelectorMatchType == telemetrytypes.FieldSelectorMatchTypeFuzzy {
continue
}
members := semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: selector.Name,
Signal: selector.Signal,
FieldContext: selector.FieldContext,
members := semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: selector.Name,
Signal: selector.Signal,
FieldContext: selector.FieldContext,
MetricContext: selector.MetricContext,
})
for _, member := range members {
if seen[member] {

View File

@@ -52,8 +52,13 @@ func LogicalValueExpr(
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
}
// Numeric and boolean maps return zero for an absent key. If a family of
// either type is enabled, this tail must become zero too.
// Numeric and boolean maps read their zero value for an absent key, so the
// tail keeps single-key semantics for rows without any member — the same
// contract as the '' tail above.
tail := "0"
if logical.FieldDataType == telemetrytypes.FieldDataTypeBool {
tail = "false"
}
branches := make([]string, 0, len(logical.Members)*2)
for i, member := range logical.Members {
guard, err := fm.ExistsFor(ctx, orgID, tsStart, tsEnd, member, true)
@@ -62,7 +67,7 @@ func LogicalValueExpr(
}
branches = append(branches, guard, memberExprs[i])
}
return "multiIf(" + strings.Join(branches, ", ") + ", NULL)", nil
return "multiIf(" + strings.Join(branches, ", ") + ", " + tail + ")", nil
}
// LogicalExistsExpr returns the existence predicate for a resolved logical

View File

@@ -72,7 +72,23 @@ func TestLogicalValueExprNumericFamilyGuardsEveryMember(t *testing.T) {
}
expr, err := LogicalValueExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, logical)
require.NoError(t, err)
assert.Equal(t, "multiIf(has(current), value(current), has(old), value(old), NULL)", expr)
// The 0 tail mirrors the '' tail of the string branch: numeric maps read 0
// for an absent key, so keyless rows keep single-key semantics.
assert.Equal(t, "multiIf(has(current), value(current), has(old), value(old), 0)", expr)
}
func TestLogicalValueExprBoolFamilyReadsFalseForKeylessRows(t *testing.T) {
logical := &telemetrytypes.LogicalField{
Name: "flag",
FieldDataType: telemetrytypes.FieldDataTypeBool,
Members: []*telemetrytypes.TelemetryFieldKey{
{Name: "current", FieldDataType: telemetrytypes.FieldDataTypeBool},
{Name: "old", FieldDataType: telemetrytypes.FieldDataTypeBool},
},
}
expr, err := LogicalValueExpr(context.Background(), valuer.UUID{}, 0, 0, stubFieldMapper{}, logical)
require.NoError(t, err)
assert.Equal(t, "multiIf(has(current), value(current), has(old), value(old), false)", expr)
}
func TestLogicalExistsExprSingleMemberDelegatesToExistsFor(t *testing.T) {

View File

@@ -48,7 +48,7 @@ func TestFamiliesOffByDefault(t *testing.T) {
}},
}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, flaggertest.New(t), &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, flaggertest.New(t), telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 1)
assert.False(t, fields[0].IsFamily())
assert.Equal(t, []string{"deployment.environment.name"}, memberNames(fields[0]))
@@ -76,7 +76,7 @@ func TestMatchingLogicalFieldsGroupsFamilyMembers(t *testing.T) {
}
for _, requested := range []string{"deployment.environment.name", "deployment.environment"} {
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), &telemetrytypes.TelemetryFieldKey{Name: requested}, fieldKeys)
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{Name: requested}, fieldKeys)
require.Len(t, fields, 1, "a family is one logical field, requested via %s", requested)
logical := fields[0]
assert.Equal(t, requested, logical.Name, "response identity is the requested spelling")
@@ -106,7 +106,7 @@ func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
}},
}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), &telemetrytypes.TelemetryFieldKey{
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextResource,
}, fieldKeys)
@@ -115,9 +115,8 @@ func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
assert.Equal(t, []string{"resource.deployment.environment.name", "deployment.environment"}, memberNames(fields[0]))
}
// Non-trace signals have no family support: the requested spelling stays
// literal, and a family member name never pulls in its siblings.
func TestMatchingLogicalFieldsKeepsLogsLiteral(t *testing.T) {
// Log entries group into families exactly like trace entries.
func TestMatchingLogicalFieldsGroupsLogEntries(t *testing.T) {
logsKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
@@ -131,10 +130,60 @@ func TestMatchingLogicalFieldsKeepsLogsLiteral(t *testing.T) {
"deployment.environment": {logsKey("deployment.environment")},
}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalLogs, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 1)
assert.True(t, fields[0].IsFamily())
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, memberNames(fields[0]))
}
// A family gated away from a signal stays literal there: db.system carries
// signals [traces, logs], so a metrics lookup never pulls in siblings.
func TestMatchingLogicalFieldsHonorsTheFamilySignalGate(t *testing.T) {
metricsKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"db.system.name": {metricsKey("db.system.name")},
"db.system": {metricsKey("db.system")},
}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalMetrics, nil, &telemetrytypes.TelemetryFieldKey{Name: "db.system.name"}, fieldKeys)
require.Len(t, fields, 1)
assert.False(t, fields[0].IsFamily())
assert.Equal(t, []string{"deployment.environment.name"}, memberNames(fields[0]))
assert.Equal(t, []string{"db.system.name"}, memberNames(fields[0]))
}
// Metric entries group across the stored label spellings of the family, in
// member-major order: every spelling of the current name precedes the first
// spelling of the old one.
func TestMatchingLogicalFieldsGroupsMetricSpellings(t *testing.T) {
metricsKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {metricsKey("deployment.environment.name")},
"deployment_environment_name": {metricsKey("deployment_environment_name")},
"deployment.environment": {metricsKey("deployment.environment")},
"deployment_environment": {metricsKey("deployment_environment")},
}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalMetrics, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment"}, fieldKeys)
require.Len(t, fields, 1)
assert.True(t, fields[0].IsFamily())
assert.Equal(t, []string{
"deployment.environment.name", "deployment_environment_name",
"deployment.environment", "deployment_environment",
}, memberNames(fields[0]))
}
// A family and a genuine same-name collision stack cleanly: the family stays
@@ -165,7 +214,7 @@ func TestResolveLogicalFieldsKeepsFamilyThroughAmbiguity(t *testing.T) {
}
requested := &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), requested, fieldKeys)
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalUnspecified, nil, requested, fieldKeys)
require.Len(t, fields, 2, "resource family + attribute collision")
resolved, warning := ResolveLogicalFields(requested, fields)
@@ -193,7 +242,7 @@ func TestMatchingLogicalFieldsNeverMergesAcrossDataTypes(t *testing.T) {
}},
}
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
fields := MatchingLogicalFields(context.Background(), valuer.UUID{}, familiesOn(t), telemetrytypes.SignalUnspecified, nil, &telemetrytypes.TelemetryFieldKey{Name: "deployment.environment.name"}, fieldKeys)
require.Len(t, fields, 2)
for _, logical := range fields {
assert.False(t, logical.IsFamily())

View File

@@ -42,6 +42,8 @@ type filterExpressionVisitor struct {
skipResourceFilter bool
skipFullTextFilter bool
variables map[string]qbtypes.VariableItem
signal telemetrytypes.Signal
metricContext *telemetrytypes.MetricContext
keysWithWarnings map[string]bool
startNs uint64
@@ -67,6 +69,11 @@ type FilterExprVisitorOpts struct {
Variables map[string]qbtypes.VariableItem
StartNs uint64
EndNs uint64
// Signal is the signal the statement builder compiles for; family
// resolution uses it for keys that do not carry their own. MetricContext
// carries the queried metric name so metric-scoped families resolve.
Signal telemetrytypes.Signal
MetricContext *telemetrytypes.MetricContext
}
// newFilterExpressionVisitor creates a new filterExpressionVisitor.
@@ -83,6 +90,8 @@ func newFilterExpressionVisitor(opts FilterExprVisitorOpts) *filterExpressionVis
skipResourceFilter: opts.SkipResourceFilter,
skipFullTextFilter: opts.SkipFullTextFilter,
variables: opts.Variables,
signal: opts.Signal,
metricContext: opts.MetricContext,
keysWithWarnings: make(map[string]bool),
startNs: opts.StartNs,
endNs: opts.EndNs,
@@ -386,7 +395,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 := MatchingLogicalFields(v.context, v.orgID, v.fl, key, v.fieldKeys)
matching := MatchingLogicalFields(v.context, v.orgID, v.fl, v.signal, v.metricContext, key, v.fieldKeys)
// Handle EXISTS specially
if ctx.EXISTS() != nil {
@@ -737,7 +746,7 @@ func (v *filterExpressionVisitor) VisitFunctionCall(ctx *grammar.FunctionCallCon
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(key, MatchingLogicalFields(v.context, v.orgID, v.fl, key, v.fieldKeys), operator, value)
conds, ok := v.buildConditions(key, MatchingLogicalFields(v.context, v.orgID, v.fl, v.signal, v.metricContext, key, v.fieldKeys), operator, value)
if !ok {
return ErrorConditionLiteral
}
@@ -930,7 +939,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.LogicalField, op qbtypes.FilterOperator, value any) ([]string, bool) {
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter}, op, value, v.builder)
conds, warns, err := v.conditionBuilder.ConditionFor(v.context, v.orgID, v.startNs, v.endNs, key, v.fieldKeys, qbtypes.ConditionBuilderOptions{SkipResourceFilter: v.skipResourceFilter, MetricContext: v.metricContext}, op, value, v.builder)
if err != nil {
_, _, _, _, errURL, _ := errors.Unwrapb(err)
assignIfEmpty(&v.mainErrorURL, errURL)
@@ -987,44 +996,47 @@ func assignIfEmpty(s *string, value string) {
}
// familyMemberNames returns the physical spellings to look up for the
// referenced key: the semantic-convention family members (current-first) when
// the resolve_semconv_families flag is on for the org and the key can resolve
// to traces, else just the requested name. Only trace field mappers understand
// families today; logs and metrics keep the requested spelling until theirs
// land.
func familyMemberNames(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, field *telemetrytypes.TelemetryFieldKey) []string {
if !semconvFamiliesEnabled(ctx, orgID, fl) {
// referenced key: the semantic-convention family spellings (current-first)
// when the resolve_semconv_families flag is on for the org, else just the
// requested name. The key's own signal wins over the caller's; metric lookups
// additionally expand each member into its stored label spellings.
func familyMemberNames(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, signal telemetrytypes.Signal, metricCtx *telemetrytypes.MetricContext, field *telemetrytypes.TelemetryFieldKey) []string {
if !SemconvFamiliesEnabled(ctx, orgID, fl) {
return []string{field.Name}
}
if field.Signal != telemetrytypes.SignalUnspecified && field.Signal != telemetrytypes.SignalTraces {
return []string{field.Name}
if field.Signal != telemetrytypes.SignalUnspecified {
signal = field.Signal
}
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: field.FieldContext,
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: signal,
FieldContext: field.FieldContext,
MetricContext: metricCtx,
})
}
// MatchingLogicalFields resolves the referenced key against the metadata map
// into logical fields, honoring any context/data type the user specified.
//
// Physical keys that are members of one semantic-convention family (traces
// only today) group into one logical field per (signal, context, data type)
// identity, members ordered current-first. Every other matching key becomes
// its own single-member logical field. Ambiguity is the length of the
// returned slice: one family is one element and is never ambiguous with
// itself, but the slice can hold several logical fields — including several
// family fields, one per identity, when the family exists under more than
// one context or data type. Members alias the metadata map entries; nothing
// is copied or mutated.
// Physical keys that are members of one semantic-convention family group into
// one logical field per (signal, context, data type) identity, members
// ordered current-first. Every other matching key becomes its own
// single-member logical field. Ambiguity is the length of the returned slice:
// one family is one element and is never ambiguous with itself, but the slice
// can hold several logical fields — including several family fields, one per
// identity, when the family exists under more than one context or data type.
// Members alias the metadata map entries; nothing is copied or mutated.
//
// signal is the signal the caller compiles for and metricCtx the queried
// metric, both used only for family resolution; a key that carries its own
// signal wins over the caller's.
//
// Family grouping only happens when the resolve_semconv_families flag is on
// for the org. A nil flagger means off: every match then stays a
// single-member logical field.
func MatchingLogicalFields(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
members := familyMemberNames(ctx, orgID, fl, field)
matches := collectMemberMatches(field, members, fieldKeys)
func MatchingLogicalFields(ctx context.Context, orgID valuer.UUID, fl flagger.Flagger, signal telemetrytypes.Signal, metricCtx *telemetrytypes.MetricContext, field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
members := familyMemberNames(ctx, orgID, fl, signal, metricCtx, field)
matches := collectMemberMatches(field, members, metricCtx, fieldKeys)
return groupIntoLogicalFields(field.Name, len(members) > 1, matches)
}
@@ -1050,32 +1062,29 @@ func matchesRequestedIdentity(field, item *telemetrytypes.TelemetryFieldKey, con
}
// inFamilyScope reports whether a match found under a sibling member name is
// legitimate: the entry must be trace metadata, and the member must be in the
// family of the requested name for the entry's context. A member lookup can
// otherwise find a same-named field in a scope where the family does not
// apply.
func inFamilyScope(field, item *telemetrytypes.TelemetryFieldKey, memberName string) bool {
if item.Signal != telemetrytypes.SignalTraces {
return false
}
return slices.Contains(semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: item.FieldContext,
// legitimate: the member must be a family spelling of the requested name for
// the entry's own signal and context. A member lookup can otherwise find a
// same-named field in a scope where the family does not apply.
func inFamilyScope(field, item *telemetrytypes.TelemetryFieldKey, memberName string, metricCtx *telemetrytypes.MetricContext) bool {
return slices.Contains(semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
MetricContext: metricCtx,
}), memberName)
}
// collectMemberMatches finds the metadata entries for every member spelling:
// first under the member names, then under their context-prefixed spellings
// (a context can be a legitimate part of a stored name, e.g. `attribute.key`).
func collectMemberMatches(field *telemetrytypes.TelemetryFieldKey, members []string, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []memberMatch {
func collectMemberMatches(field *telemetrytypes.TelemetryFieldKey, members []string, metricCtx *telemetrytypes.MetricContext, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []memberMatch {
matches := make([]memberMatch, 0)
collect := func(lookupName string, rank int, memberName string, contextMatched bool) {
for _, item := range fieldKeys[lookupName] {
if !matchesRequestedIdentity(field, item, contextMatched) {
continue
}
if memberName != field.Name && !inFamilyScope(field, item, memberName) {
if memberName != field.Name && !inFamilyScope(field, item, memberName, metricCtx) {
continue
}
matches = append(matches, memberMatch{key: item, rank: rank})
@@ -1093,18 +1102,18 @@ func collectMemberMatches(field *telemetrytypes.TelemetryFieldKey, members []str
return matches
}
// groupIntoLogicalFields turns matches into logical fields. Trace entries in
// family mode group by their (signal, context, data type) identity; every
// other entry becomes its own single-member field. Members sort by family
// rank at the end: precedence is a property of the family, not of the order
// in which the lookups found the members.
// groupIntoLogicalFields turns matches into logical fields. Entries of a
// family-capable signal in family mode group by their (signal, context, data
// type) identity; every other entry becomes its own single-member field.
// Members sort by family rank at the end: precedence is a property of the
// family, not of the order in which the lookups found the members.
func groupIntoLogicalFields(requestedName string, familyMode bool, matches []memberMatch) []*telemetrytypes.LogicalField {
fields := make([]*telemetrytypes.LogicalField, 0, len(matches))
groups := make(map[string]*telemetrytypes.LogicalField)
ranks := make(map[*telemetrytypes.TelemetryFieldKey]int)
for _, match := range matches {
if !familyMode || match.key.Signal != telemetrytypes.SignalTraces {
if !familyMode || !familySignal(match.key.Signal) {
fields = append(fields, telemetrytypes.SingleLogicalField(requestedName, match.key))
continue
}
@@ -1144,3 +1153,11 @@ func groupHasMemberNamed(group *telemetrytypes.LogicalField, name string) bool {
}
return false
}
func familySignal(signal telemetrytypes.Signal) bool {
switch signal {
case telemetrytypes.SignalTraces, telemetrytypes.SignalLogs, telemetrytypes.SignalMetrics:
return true
}
return false
}

View File

@@ -590,7 +590,7 @@ func TestVisitKey(t *testing.T) {
// and decides not-found handling. Replay that here against the generic
// builder behavior (error unless the key is ignored). The test maps carry
// no signal, so every logical field is single-member and flattens losslessly.
matching := MatchingLogicalFields(context.Background(), valuer.UUID{}, nil, key, tt.fieldKeys)
matching := MatchingLogicalFields(context.Background(), valuer.UUID{}, nil, telemetrytypes.SignalUnspecified, nil, key, tt.fieldKeys)
resolved, warning := ResolveLogicalFields(key, matching)
keys := SingleKeys(resolved)
@@ -768,7 +768,7 @@ func (b *resourceConditionBuilder) ConditionFor(
return nil, nil, nil
}
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(context.Background(), valuer.UUID{}, nil, key, fieldKeys))
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(context.Background(), valuer.UUID{}, nil, telemetrytypes.SignalUnspecified, nil, key, fieldKeys))
keys := SingleKeys(resolved)
var warnings []string
if warning != "" {
@@ -811,7 +811,7 @@ func (b *conditionBuilder) ConditionFor(
return []string{fmt.Sprintf("%s_cond", key.Name)}, nil, nil
}
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(context.Background(), valuer.UUID{}, nil, key, fieldKeys))
resolved, warning := ResolveLogicalFields(key, MatchingLogicalFields(context.Background(), valuer.UUID{}, nil, telemetrytypes.SignalUnspecified, nil, key, fieldKeys))
keys := SingleKeys(resolved)
var warnings []string
if warning != "" {

View File

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

View File

@@ -1,7 +1,9 @@
package semconv
import (
"iter"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -14,16 +16,51 @@ type Kind struct {
valuer.String
}
// Family is one logical telemetry field. Old is ordered from the most recent
// predecessor to the oldest one and therefore also defines fallback order.
// Member is one historical spelling of a family, with the scope its rename
// edges declared. A nil axis places no constraint on that axis.
type Member struct {
name string
contexts []telemetrytypes.FieldContext
signals []telemetrytypes.Signal
applyToMetrics []string
}
// Name returns the member spelling.
func (m Member) Name() string {
return m.name
}
// Family is one logical telemetry field. Members are ordered from the most
// recent predecessor to the oldest one and therefore also define fallback
// order. The family-level contexts and signals come from the overlay and gate
// where the family may resolve at all; member scopes come from the schema
// edges and gate which members apply for a given selector.
type Family struct {
Current string
Old []string
Kind Kind
Contexts []telemetrytypes.FieldContext
Signals []telemetrytypes.Signal
ApplyToMetrics []string
ValueMap map[string]string
current string
kind Kind
members []Member
contexts []telemetrytypes.FieldContext
signals []telemetrytypes.Signal
valueMap map[string]string
}
// Current returns the current name of the family.
func (f Family) Current() string {
return f.current
}
// Kind returns the family kind.
func (f Family) Kind() Kind {
return f.kind
}
// Old returns the historical spellings in fallback order.
func (f Family) Old() []string {
names := make([]string, len(f.members))
for i, member := range f.members {
names[i] = member.name
}
return names
}
var (
@@ -38,8 +75,13 @@ func (Kind) Enum() []any {
return []any{KindAttribute, KindMetric}
}
// Lookup returns the enabled family containing selector.Name for kind. The
// returned family must not be modified.
// Lookup returns the family that resolves selector.Name for kind.
//
// The missing-information policy is the same on every axis (signal, field
// context, metric name): an axis the selector does not populate is a wildcard
// and constrains nothing. When the wildcards leave more than one family
// admitted, the name does not resolve — resolution never picks an arbitrary
// winner, it asks for more information by staying literal.
func Lookup(kind Kind, selector telemetrytypes.FieldKeySelector) (Family, bool) {
idx, ok := lookupIndex(kind, selector)
if !ok {
@@ -48,80 +90,284 @@ func Lookup(kind Kind, selector telemetrytypes.FieldKeySelector) (Family, bool)
return families[idx], true
}
// Members returns the current name first, followed by historical names in
// fallback order. A name outside an enabled family is returned unchanged. The
// Members returns the current name first, followed by the historical spellings
// admitted for the selector, in fallback order. A name outside an enabled
// family — or one the selector leaves ambiguous — is returned unchanged. The
// returned slice must not be modified.
func Members(kind Kind, selector telemetrytypes.FieldKeySelector) []string {
idx, ok := lookupIndex(kind, selector)
if !ok {
return []string{selector.Name}
}
return familyMembers[idx]
return admittedMembers(idx, selector)
}
func admittedMembers(idx int, selector telemetrytypes.FieldKeySelector) []string {
admitted := 0
for _, member := range families[idx].members {
if memberAdmits(member, selector) {
admitted++
}
}
if admitted == len(families[idx].members) {
return familyMembers[idx]
}
names := make([]string, 0, admitted+1)
names = append(names, families[idx].current)
for _, member := range families[idx].members {
if memberAdmits(member, selector) {
names = append(names, member.name)
}
}
return names
}
// AttributeMembers returns the physical attribute spellings that may represent
// selector.Name. Metrics have used both dotted and normalized label layouts,
// and resource labels have additionally used a resource_ prefix, so metric
// selectors expand every admitted member into those spellings; keeping that
// storage detail here prevents metrics readers from maintaining local
// transition tables. Every other signal gets the plain member list.
func AttributeMembers(selector telemetrytypes.FieldKeySelector) []string {
if selector.Signal != telemetrytypes.SignalMetrics {
return Members(KindAttribute, selector)
}
lookupSelector := selector
lookupSelector.Name = strings.TrimPrefix(selector.Name, "resource_")
idx, style, ok := lookupMetricSpelling(KindAttribute, lookupSelector)
if !ok {
return []string{selector.Name}
}
logicalMembers := admittedMembers(idx, lookupSelector)
result := make([]string, 0, len(logicalMembers)*4)
for _, member := range logicalMembers {
variants := []string{member, normalizeMetricSpelling(member)}
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,
}
idx, style, ok := lookupMetricSpelling(KindMetric, selector)
if !ok {
return []string{name}
}
logicalMembers := admittedMembers(idx, selector)
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_")
idx, _, ok := lookupMetricSpelling(KindAttribute, lookupSelector)
if !ok {
return selector.Name
}
return families[idx].current
}
// Current returns the current name for selector.Name, or the input name when
// it does not belong to an enabled family.
// it does not resolve to a family.
func Current(kind Kind, selector telemetrytypes.FieldKeySelector) string {
idx, ok := lookupIndex(kind, selector)
if !ok {
return selector.Name
}
return families[idx].Current
return families[idx].current
}
// All returns every enabled family. The returned slice and families must not be
// modified.
func All() []Family {
return families
// All iterates over every enabled family.
func All() iter.Seq[Family] {
return func(yield func(Family) bool) {
for _, family := range families {
if !yield(family) {
return
}
}
}
}
func buildIndexes() (map[string][]int, [][]string) {
index := make(map[string][]int)
members := make([][]string, len(families))
for i, family := range families {
members[i] = make([]string, 0, len(family.Old)+1)
members[i] = append(members[i], family.Current)
members[i] = append(members[i], family.Old...)
index[family.Current] = append(index[family.Current], i)
for _, old := range family.Old {
index[old] = append(index[old], i)
members[i] = make([]string, 0, len(family.members)+1)
members[i] = append(members[i], family.current)
index[family.current] = append(index[family.current], i)
for _, member := range family.members {
members[i] = append(members[i], member.name)
if !slices.Contains(index[member.name], i) {
index[member.name] = append(index[member.name], i)
}
}
}
return index, members
}
func lookupIndex(kind Kind, selector telemetrytypes.FieldKeySelector) (int, bool) {
found, foundIdx := 0, 0
for _, idx := range memberToFamilies[selector.Name] {
if matchesSelector(families[idx], kind, selector) {
return idx, true
if familyAdmits(families[idx], kind, selector) {
found++
foundIdx = idx
}
}
return 0, false
if found != 1 {
return 0, false
}
return foundIdx, true
}
func matchesSelector(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if family.Kind != kind {
type metricSpelling int
const (
metricSpellingDotted metricSpelling = iota
metricSpellingNormalized
)
// lookupMetricSpelling resolves a metric spelling to its family: first as the
// canonical dotted name, then by comparing the normalized form of every
// family spelling. The ambiguity rule of lookupIndex applies to both styles.
func lookupMetricSpelling(kind Kind, selector telemetrytypes.FieldKeySelector) (int, metricSpelling, bool) {
if idx, ok := lookupIndex(kind, selector); ok {
return idx, metricSpellingDotted, true
}
found, foundIdx := 0, 0
for idx := range families {
if normalizedFamilyAdmits(families[idx], kind, selector) {
found++
foundIdx = idx
}
}
if found != 1 {
return 0, metricSpellingDotted, false
}
return foundIdx, metricSpellingNormalized, true
}
func normalizedFamilyAdmits(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if normalizeMetricSpelling(family.current) == selector.Name ||
slices.ContainsFunc(family.members, func(member Member) bool {
return normalizeMetricSpelling(member.name) == selector.Name
}) {
dotted := selector
dotted.Name = denormalizeAgainst(family, selector.Name)
return familyAdmits(family, kind, dotted)
}
return false
}
func denormalizeAgainst(family Family, name string) string {
if normalizeMetricSpelling(family.current) == name {
return family.current
}
for _, member := range family.members {
if normalizeMetricSpelling(member.name) == name {
return member.name
}
}
return name
}
func normalizeMetricSpelling(name string) string {
return strings.ReplaceAll(name, ".", "_")
}
func appendUniqueString(values []string, value string) []string {
if slices.Contains(values, value) {
return values
}
return append(values, value)
}
// familyAdmits reports whether the family resolves selector.Name: the
// family-level gate must admit the selector, and the name must be the current
// name or an admitted member.
func familyAdmits(family Family, kind Kind, selector telemetrytypes.FieldKeySelector) bool {
if family.kind != kind {
return false
}
if selector.Signal != telemetrytypes.SignalUnspecified && len(family.Signals) > 0 {
if !slices.Contains(family.Signals, selector.Signal) {
return false
if !axisAdmits(family.signals, selector.Signal, telemetrytypes.SignalUnspecified) {
return false
}
if !axisAdmits(family.contexts, selector.FieldContext, telemetrytypes.FieldContextUnspecified) {
return false
}
if selector.Name == family.current {
for _, member := range family.members {
if memberAdmits(member, selector) {
return true
}
}
return false
}
for _, member := range family.members {
if member.name == selector.Name && memberAdmits(member, selector) {
return true
}
}
return false
}
if selector.FieldContext != telemetrytypes.FieldContextUnspecified && len(family.Contexts) > 0 {
if !slices.Contains(family.Contexts, selector.FieldContext) {
return false
}
// memberAdmits reports whether the member applies for the selector under the
// wildcard policy: a selector axis without a value never constrains, and a
// member axis without a value admits every selector value.
func memberAdmits(member Member, selector telemetrytypes.FieldKeySelector) bool {
if !axisAdmits(member.signals, selector.Signal, telemetrytypes.SignalUnspecified) {
return false
}
if selector.Signal == telemetrytypes.SignalMetrics && len(family.ApplyToMetrics) > 0 {
if selector.MetricContext == nil {
return false
}
return slices.Contains(family.ApplyToMetrics, selector.MetricContext.MetricName)
if !axisAdmits(member.contexts, selector.FieldContext, telemetrytypes.FieldContextUnspecified) {
return false
}
if len(member.applyToMetrics) > 0 &&
selector.MetricContext != nil && selector.MetricContext.MetricName != "" &&
!slices.Contains(member.applyToMetrics, selector.MetricContext.MetricName) {
return false
}
return true
}
func axisAdmits[T comparable](scope []T, value T, unspecified T) bool {
if len(scope) == 0 || value == unspecified {
return true
}
return slices.Contains(scope, value)
}

View File

@@ -78,3 +78,177 @@ func TestMembersReturnsInputWhenKindDoesNotMatch(t *testing.T) {
"an attribute family must not match a metric-name lookup",
)
}
func TestFamilySignalsGateResolution(t *testing.T) {
metrics := telemetrytypes.FieldKeySelector{Name: "db.system", Signal: telemetrytypes.SignalMetrics}
logs := telemetrytypes.FieldKeySelector{Name: "db.system", Signal: telemetrytypes.SignalLogs}
assert.Equal(t, []string{"db.system"}, Members(KindAttribute, metrics),
"a family gated to traces and logs must stay literal for metrics")
assert.Equal(t, []string{"db.system.name", "db.system"}, Members(KindAttribute, logs),
"the gate admits the signals it lists")
}
func TestMetricNameFamilyResolves(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{Name: "k8s.pod.cpu.utilization", Signal: telemetrytypes.SignalMetrics}
assert.Equal(t, []string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"}, Members(KindMetric, selector))
assert.Equal(t, "k8s.pod.cpu.usage", Current(KindMetric, selector))
assert.Equal(t, []string{"k8s.pod.cpu.utilization"}, Members(KindAttribute, selector),
"a metric-name family must not match an attribute lookup")
}
func TestMembersReturnsSharedSliceForUnscopedFamily(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{Name: "deployment.environment", Signal: telemetrytypes.SignalTraces}
first := Members(KindAttribute, selector)
second := Members(KindAttribute, selector)
assert.Equal(t, &first[0], &second[0],
"a family whose members all admit must return the precomputed slice, not a copy")
}
func TestAttributeMembersExpandsMetricSpellings(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "deployment.environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, []string{
"resource_deployment.environment.name", "resource_deployment_environment_name",
"deployment.environment.name", "deployment_environment_name",
"resource_deployment.environment", "resource_deployment_environment",
"deployment.environment", "deployment_environment",
}, AttributeMembers(selector), "metric members expand into dotted, normalized, and resource_-prefixed spellings")
}
func TestAttributeMembersPreservesRequestedSpellingStyle(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "resource_deployment_environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
}
assert.Equal(t, []string{
"resource_deployment_environment_name", "resource_deployment.environment.name",
"deployment_environment_name", "deployment.environment.name",
"resource_deployment_environment", "resource_deployment.environment",
"deployment_environment", "deployment.environment",
}, AttributeMembers(selector), "a normalized request lists normalized spellings first and keeps resource_ variants")
}
func TestAttributeMembersRespectsTheFamilyGateForMetrics(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "db.system",
Signal: telemetrytypes.SignalMetrics,
}
assert.Equal(t, []string{"db.system"}, AttributeMembers(selector),
"a family gated away from metrics must not expand metric spellings")
}
func TestAttributeMembersIsPlainForOtherSignals(t *testing.T) {
selector := telemetrytypes.FieldKeySelector{
Name: "deployment.environment",
Signal: telemetrytypes.SignalLogs,
}
assert.Equal(t, []string{"deployment.environment.name", "deployment.environment"}, AttributeMembers(selector))
}
func TestMetricNamesPreservesStyle(t *testing.T) {
assert.Equal(t, []string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"}, MetricNames("k8s.pod.cpu.utilization"))
assert.Equal(t, []string{"k8s_pod_cpu_usage", "k8s_pod_cpu_utilization"}, MetricNames("k8s_pod_cpu_utilization"))
assert.Equal(t, []string{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization"}, MetricNames("k8s.pod.cpu.usage"))
assert.Equal(t, []string{"http.server.duration"}, MetricNames("http.server.duration"))
}
func TestCurrentAttributeCanonicalizesMetricSpellings(t *testing.T) {
assert.Equal(t, "deployment.environment.name", CurrentAttribute(telemetrytypes.FieldKeySelector{
Name: "resource_deployment_environment",
Signal: telemetrytypes.SignalMetrics,
}))
assert.Equal(t, "unrelated_label", CurrentAttribute(telemetrytypes.FieldKeySelector{
Name: "unrelated_label",
Signal: telemetrytypes.SignalMetrics,
}))
}
func TestAllIteratesEnabledFamilies(t *testing.T) {
currents := []string{}
for family := range All() {
currents = append(currents, family.Current())
}
assert.Contains(t, currents, "deployment.environment.name")
assert.Contains(t, currents, "k8s.pod.cpu.usage")
}
// swapFamilies replaces the generated table for one test so scoped-member and
// fan-out behavior can be pinned without enabling such families for real.
func swapFamilies(t *testing.T, replacement []Family) {
t.Helper()
prevFamilies, prevIndex, prevMembers := families, memberToFamilies, familyMembers
families = replacement
memberToFamilies, familyMembers = buildIndexes()
t.Cleanup(func() {
families, memberToFamilies, familyMembers = prevFamilies, prevIndex, prevMembers
})
}
func TestFanOutResolvesOnlyWithEnoughInformation(t *testing.T) {
swapFamilies(t, []Family{
{
current: "cpu.mode",
kind: KindAttribute,
members: []Member{{name: "state", applyToMetrics: []string{"system.cpu.time"}}},
},
{
current: "db.client.connection.state",
kind: KindAttribute,
members: []Member{{name: "state", applyToMetrics: []string{"db.client.connections.usage"}}},
},
})
ambiguous := telemetrytypes.FieldKeySelector{Name: "state", Signal: telemetrytypes.SignalMetrics}
assert.Equal(t, []string{"state"}, Members(KindAttribute, ambiguous),
"without a metric name, a fanned-out member admits several families and must stay literal")
pinned := ambiguous
pinned.MetricContext = &telemetrytypes.MetricContext{MetricName: "system.cpu.time"}
assert.Equal(t, []string{"cpu.mode", "state"}, Members(KindAttribute, pinned),
"the metric name disambiguates the fan-out")
outside := ambiguous
outside.MetricContext = &telemetrytypes.MetricContext{MetricName: "http.server.duration"}
assert.Equal(t, []string{"state"}, Members(KindAttribute, outside),
"a metric outside every apply_to_metrics list resolves no family")
}
func TestMemberScopesFilterMembers(t *testing.T) {
swapFamilies(t, []Family{{
current: "user_agent.original",
kind: KindAttribute,
members: []Member{
{name: "http.user_agent", contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextAttribute}, signals: []telemetrytypes.Signal{telemetrytypes.SignalTraces}},
{name: "browser.user_agent", contexts: []telemetrytypes.FieldContext{telemetrytypes.FieldContextResource}},
},
}})
resource := telemetrytypes.FieldKeySelector{
Name: "user_agent.original",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
}
assert.Equal(t, []string{"user_agent.original", "browser.user_agent"}, Members(KindAttribute, resource),
"a strict resource lookup must not include the span-only member")
attribute := resource
attribute.FieldContext = telemetrytypes.FieldContextAttribute
assert.Equal(t, []string{"user_agent.original", "http.user_agent"}, Members(KindAttribute, attribute),
"a strict attribute lookup must not include the resource-only member")
strictResourceOldSpan := resource
strictResourceOldSpan.Name = "http.user_agent"
assert.Equal(t, []string{"http.user_agent"}, Members(KindAttribute, strictResourceOldSpan),
"an old spelling outside its own scope stays literal")
}

View File

@@ -0,0 +1,89 @@
package logsstatementbuilder
import (
"context"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/statementbuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"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/require"
)
// A filter on either spelling of an enabled family compiles to one merged
// condition over the log attribute maps; the flag default keeps it literal.
func TestStatementBuilderResolvesLogFamilies(t *testing.T) {
releaseTime := time.Date(2024, 1, 15, 10, 0, 0, 0, time.UTC)
releaseTimeNano := uint64(releaseTime.UnixNano())
logsKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
cases := []struct {
name string
familyOn bool
expected string
}{
{
name: "families on",
familyOn: true,
expected: "SELECT count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '') = ? AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC",
},
{
name: "families off",
familyOn: false,
expected: "SELECT count() AS __result_0 FROM signoz_logs.distributed_logs_v2 WHERE (attributes_string['deployment.environment'] = ? AND mapContains(attributes_string, 'deployment.environment')) AND timestamp >= ? AND ts_bucket_start >= ? AND timestamp < ? AND ts_bucket_start <= ? ORDER BY __result_0 DESC",
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): c.familyOn,
})
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys := logstelemetryschema.BuildCompleteFieldKeyMap(releaseTime)
keys["deployment.environment.name"] = []*telemetrytypes.TelemetryFieldKey{logsKey("deployment.environment.name")}
keys["deployment.environment"] = []*telemetrytypes.TelemetryFieldKey{logsKey("deployment.environment")}
mockMetadataStore.KeysMap = keys
fm := logstelemetryschema.NewFieldMapper(fl)
cb := logstelemetryschema.NewConditionBuilder(fm, fl)
aggExprRewriter := querybuilder.NewAggExprRewriter(instrumentationtest.New().ToProviderSettings(), nil, fm, cb, fl)
statementBuilder := NewLogQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore, fm, cb, aggExprRewriter,
logstelemetryschema.DefaultFullTextColumn, fl, nil,
statementbuilder.Config{SkipResourceFingerprint: statementbuilder.SkipResourceFingerprint{Enabled: false, Threshold: 100000}},
)
query := qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]{
Signal: telemetrytypes.SignalLogs,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.LogAggregation{{Expression: "count()"}},
Filter: &qbtypes.Filter{
Expression: "attribute.deployment.environment = 'production'",
},
}
q, err := statementBuilder.Build(context.Background(), valuer.UUID{},
releaseTimeNano+uint64(24*time.Hour.Nanoseconds()),
releaseTimeNano+uint64(48*time.Hour.Nanoseconds()),
qbtypes.RequestTypeScalar, query, nil)
require.NoError(t, err)
require.Equal(t, c.expected, q.Query)
})
}
}

View File

@@ -128,6 +128,7 @@ func (b *logQueryStatementBuilder) Build(
bodyJSONEnabled := b.fl.BooleanOrEmpty(ctx, flagger.FeatureUseJSONBody, featuretypes.NewFlaggerEvaluationContext(orgID))
keySelectors, warnings := getKeySelectors(query, bodyJSONEnabled)
keySelectors = querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, b.fl, keySelectors)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, err
@@ -741,6 +742,8 @@ func (b *logQueryStatementBuilder) addFilterCondition(
Variables: variables,
StartNs: start,
EndNs: end,
Flagger: b.fl,
Signal: telemetrytypes.SignalLogs,
})
if err != nil {

View File

@@ -44,8 +44,8 @@ func NewFactory(
if err != nil {
return nil, err
}
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(nil)
cb := metricstelemetryschema.NewConditionBuilder(fm, nil)
return NewMeterQueryStatementBuilder(settings, metadataStore, fm, cb, metricsStatementBuilder), nil
},
)

View File

@@ -162,8 +162,8 @@ func TestStatementBuilder(t *testing.T) {
},
}
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(nil)
cb := metricstelemetryschema.NewConditionBuilder(fm, nil)
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys, err := telemetrytypestest.LoadFieldKeysFromJSON("testdata/keys_map.json")
if err != nil {
@@ -202,8 +202,8 @@ func TestStatementBuilder(t *testing.T) {
}
func TestGroupByAliasAvoidsColumnCollision(t *testing.T) {
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(nil)
cb := metricstelemetryschema.NewConditionBuilder(fm, nil)
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys, err := telemetrytypestest.LoadFieldKeysFromJSON("testdata/keys_map.json")
require.NoError(t, err)

View File

@@ -0,0 +1,102 @@
package metricsstatementbuilder
import (
"context"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/telemetryschema/metricstelemetryschema"
"github.com/SigNoz/signoz/pkg/types/metrictypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/require"
)
func familyKeysMap() map[string][]*telemetrytypes.TelemetryFieldKey {
metricsKey := func(name string) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
}
return map[string][]*telemetrytypes.TelemetryFieldKey{
"deployment.environment.name": {metricsKey("deployment.environment.name")},
"deployment.environment": {metricsKey("deployment.environment")},
"deployment_environment": {metricsKey("deployment_environment")},
}
}
func familyStatementBuilder(t *testing.T, fl flagger.Flagger) *StatementBuilder {
t.Helper()
fm := metricstelemetryschema.NewFieldMapper(fl)
cb := metricstelemetryschema.NewConditionBuilder(fm, fl)
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
mockMetadataStore.KeysMap = familyKeysMap()
return NewMetricQueryStatementBuilder(
instrumentationtest.New().ToProviderSettings(),
mockMetadataStore,
fm,
cb,
fl,
)
}
func familyQuery() qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation] {
return qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]{
Signal: telemetrytypes.SignalMetrics,
StepInterval: qbtypes.Step{Duration: 30 * time.Second},
Aggregations: []qbtypes.MetricAggregation{
{
MetricName: "k8s.pod.cpu.utilization",
Type: metrictypes.GaugeType,
Temporality: metrictypes.Unspecified,
TimeAggregation: metrictypes.TimeAggregationAvg,
SpaceAggregation: metrictypes.SpaceAggregationAvg,
},
},
Filter: &qbtypes.Filter{
Expression: "deployment.environment = 'production'",
},
GroupBy: []qbtypes.GroupByKey{
{
TelemetryFieldKey: telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
},
},
},
}
}
// The flag merges the label spellings of the family and unions the storage
// names of the metric-name family, in the filter, the group-by column, and
// every metric_name filter.
func TestStatementBuilderResolvesFamilies(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})
statementBuilder := familyStatementBuilder(t, fl)
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, qbtypes.RequestTypeTimeSeries, familyQuery(), nil)
require.NoError(t, err)
require.Equal(t, "WITH __temporal_aggregation_cte AS (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, COALESCE(NULLIF(JSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(JSONExtractString(labels, 'deployment.environment'), ''), NULLIF(JSONExtractString(labels, 'deployment_environment'), ''), '') AS `__GROUP_BY_KEY_0_deployment.environment.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?, ?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND COALESCE(NULLIF(JSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(JSONExtractString(labels, 'deployment.environment'), ''), NULLIF(JSONExtractString(labels, 'deployment_environment'), ''), '') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_deployment.environment.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?, ?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY fingerprint, ts), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_deployment.environment.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_deployment.environment.name`, ts", q.Query)
require.Equal(t, []any{"k8s.pod.cpu.usage", "k8s.pod.cpu.utilization", uint64(1747936800000), uint64(1747983420000), "unspecified", "production", "k8s.pod.cpu.usage", "k8s.pod.cpu.utilization", uint64(1747947390000), uint64(1747983420000), 0}, q.Args)
}
// With the flag at its default, both the labels and the metric name stay
// literal.
func TestStatementBuilderKeepsFamiliesLiteralByDefault(t *testing.T) {
fl := flaggertest.WithBooleanFlags(t, map[string]bool{})
statementBuilder := familyStatementBuilder(t, fl)
q, err := statementBuilder.Build(context.Background(), valuer.UUID{}, 1747947419000, 1747983448000, qbtypes.RequestTypeTimeSeries, familyQuery(), nil)
require.NoError(t, err)
require.Equal(t, "WITH __temporal_aggregation_cte AS (SELECT fingerprint, toStartOfInterval(toDateTime(intDiv(unix_milli, 1000)), toIntervalSecond(30)) AS ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(value) AS per_series_value FROM signoz_metrics.distributed_samples_v4 AS points INNER JOIN (SELECT fingerprint, JSONExtractString(labels, 'deployment.environment.name') AS `__GROUP_BY_KEY_0_deployment.environment.name` FROM signoz_metrics.time_series_v4_6hrs WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli <= ? AND LOWER(temporality) LIKE LOWER(?) AND JSONExtractString(labels, 'deployment.environment') = ? GROUP BY fingerprint, `__GROUP_BY_KEY_0_deployment.environment.name`) AS filtered_time_series ON points.fingerprint = filtered_time_series.fingerprint WHERE metric_name IN (?) AND unix_milli >= ? AND unix_milli < ? GROUP BY fingerprint, ts, `__GROUP_BY_KEY_0_deployment.environment.name` ORDER BY fingerprint, ts), __spatial_aggregation_cte AS (SELECT ts, `__GROUP_BY_KEY_0_deployment.environment.name`, avg(per_series_value) AS value FROM __temporal_aggregation_cte WHERE isNaN(per_series_value) = ? GROUP BY ts, `__GROUP_BY_KEY_0_deployment.environment.name`) SELECT * FROM __spatial_aggregation_cte ORDER BY `__GROUP_BY_KEY_0_deployment.environment.name`, ts", q.Query)
require.Equal(t, []any{"k8s.pod.cpu.utilization", uint64(1747936800000), uint64(1747983420000), "unspecified", "production", "k8s.pod.cpu.utilization", uint64(1747947390000), uint64(1747983420000), 0}, q.Args)
}

View File

@@ -144,8 +144,8 @@ func TestReducedStatementBuilder(t *testing.T) {
},
}
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(nil)
cb := metricstelemetryschema.NewConditionBuilder(fm, nil)
fl, err := flagger.New(context.Background(), instrumentationtest.New().ToProviderSettings(), flagger.Config{}, flagger.MustNewRegistry())
require.NoError(t, err)
sb := NewMetricQueryStatementBuilder(instrumentationtest.New().ToProviderSettings(), telemetrytypestest.NewMockMetadataStore(), fm, cb, fl)

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"
@@ -51,8 +52,8 @@ func NewFactory(
return factory.NewProviderFactory(
factory.MustNewName("metrics"),
func(_ context.Context, settings factory.ProviderSettings, _ statementbuilder.Config) (*StatementBuilder, error) {
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(fl)
cb := metricstelemetryschema.NewConditionBuilder(fm, fl)
return NewMetricQueryStatementBuilder(settings, metadataStore, fm, cb, fl), nil
},
)
@@ -117,7 +118,9 @@ func (b *StatementBuilder) Build(
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
keySelectors := GetKeySelectors(query)
keySelectors := querybuilder.ExpandKeySelectorsForFamilies(ctx, orgID, b.flagger, GetKeySelectors(query))
metricNames := b.familyMetricNames(ctx, orgID, query.Aggregations[0].MetricName)
keySelectors = expandSelectorsForMetricNames(keySelectors, metricNames)
keys, _, err := b.metadataStore.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, err
@@ -125,7 +128,7 @@ func (b *StatementBuilder) Build(
start, end = querybuilder.AdjustedMetricTimeRange(start, end, uint64(query.StepInterval.Seconds()), query)
return b.buildPipelineStatement(ctx, orgID, start, end, query, keys, variables)
return b.buildPipelineStatement(ctx, orgID, start, end, query, keys, metricNames, variables)
}
func (b *StatementBuilder) buildPipelineStatement(
@@ -134,6 +137,7 @@ func (b *StatementBuilder) buildPipelineStatement(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
keys map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []any,
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
var (
@@ -165,13 +169,13 @@ func (b *StatementBuilder) buildPipelineStatement(
var filterWarnings []string
var err error
if timeSeriesCTE, timeSeriesCTEArgs, filterWarnings, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, variables, tsTable); err != nil {
if timeSeriesCTE, timeSeriesCTEArgs, filterWarnings, err = b.buildTimeSeriesCTE(ctx, orgID, tsStart, tsEnd, cteQuery, keys, metricNames, variables, tsTable); err != nil {
return nil, err
}
if qbtypes.CanShortCircuitDelta(agg) {
// spatial_aggregation_cte directly for certain delta queries
if frag, args, err := b.buildTemporalAggDeltaFastPath(start, end, cteQuery, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
if frag, args, err := b.buildTemporalAggDeltaFastPath(start, end, cteQuery, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
return nil, err
} else if frag != "" {
cteFragments = append(cteFragments, frag)
@@ -179,7 +183,7 @@ func (b *StatementBuilder) buildPipelineStatement(
}
} else {
// temporal_aggregation_cte
if frag, args, err := b.buildTemporalAggregationCTE(ctx, start, end, cteQuery, keys, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
if frag, args, err := b.buildTemporalAggregationCTE(ctx, start, end, cteQuery, keys, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs); err != nil {
return nil, err
} else if frag != "" {
cteFragments = append(cteFragments, frag)
@@ -200,16 +204,16 @@ func (b *StatementBuilder) buildPipelineStatement(
var tsArgs []any
// time series rows are written on hour boundaries
tsStart := start - (start % metricstelemetryschema.OneHourInMilliseconds)
if tsCTE, tsArgs, err = b.buildReducedTimeSeriesCTE(ctx, orgID, tsStart, end, cteQuery, keys, variables); err != nil {
if tsCTE, tsArgs, err = b.buildReducedTimeSeriesCTE(ctx, orgID, tsStart, end, cteQuery, keys, metricNames, variables); err != nil {
return nil, err
}
if qbtypes.CanShortCircuitReduced(agg) {
// spatial_aggregation_cte directly, no per-series level
if spatialFrag, spatialArgs, ok := b.buildReducedSpatialAggFastPath(start, end, cteQuery, tsCTE, tsArgs); ok {
if spatialFrag, spatialArgs, ok := b.buildReducedSpatialAggFastPath(start, end, cteQuery, metricNames, tsCTE, tsArgs); ok {
reducedFragments = []string{spatialFrag}
reducedArgs = [][]any{spatialArgs}
}
} else if temporalFrag, temporalArgs, ok := b.buildReducedTemporalAggregationCTE(start, end, cteQuery, tsCTE, tsArgs); ok {
} else if temporalFrag, temporalArgs, ok := b.buildReducedTemporalAggregationCTE(start, end, cteQuery, metricNames, tsCTE, tsArgs); ok {
spatialFrag, spatialArgs := b.buildReducedSpatialAggregationCTE(cteQuery)
reducedFragments = []string{temporalFrag, spatialFrag}
reducedArgs = [][]any{temporalArgs, spatialArgs}
@@ -231,6 +235,46 @@ func (b *StatementBuilder) buildPipelineStatement(
return unionStatements(mainStmt, reducedStmt, query)
}
// familyMetricNames returns the storage names the query must read: the
// requested name plus the other spellings of its metric-name family when the
// resolve_semconv_families flag is on for the org.
func (b *StatementBuilder) familyMetricNames(ctx context.Context, orgID valuer.UUID, metricName string) []any {
if !querybuilder.SemconvFamiliesEnabled(ctx, orgID, b.flagger) {
return []any{metricName}
}
members := semconv.MetricNames(metricName)
names := make([]any, len(members))
for i, member := range members {
names[i] = member
}
return names
}
// expandSelectorsForMetricNames duplicates the selectors for each storage
// name of a metric-name family: label-key metadata is filtered by exact
// metric_name, so the old-named series must contribute their keys too.
func expandSelectorsForMetricNames(selectors []*telemetrytypes.FieldKeySelector, metricNames []any) []*telemetrytypes.FieldKeySelector {
if len(metricNames) <= 1 {
return selectors
}
out := selectors
for _, selector := range selectors {
if selector.MetricContext == nil {
continue
}
for _, name := range metricNames {
metricName, ok := name.(string)
if !ok || metricName == selector.MetricContext.MetricName {
continue
}
expanded := *selector
expanded.MetricContext = &telemetrytypes.MetricContext{MetricName: metricName}
out = append(out, &expanded)
}
}
return out
}
func unionStatements(main, reduced *qbtypes.Statement, query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation]) (*qbtypes.Statement, error) {
orderBy := "ts"
for i, g := range query.GroupBy {
@@ -251,6 +295,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
keys map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []any,
variables map[string]qbtypes.VariableItem,
) (string, []any, error) {
sb := sqlbuilder.NewSelectBuilder()
@@ -269,6 +314,9 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
Variables: variables,
StartNs: start,
EndNs: end,
Flagger: b.flagger,
Signal: telemetrytypes.SignalMetrics,
MetricContext: &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName},
})
if err != nil {
return "", nil, err
@@ -285,7 +333,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
sb.SelectMore(fmt.Sprintf("%s AS `%s`", sqlbuilder.Escape(col), GroupByColumnAlias(i, g.Name)))
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNames...),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -309,6 +357,7 @@ func (b *StatementBuilder) buildReducedTimeSeriesCTE(
func (b *StatementBuilder) buildReducedSpatialAggFastPath(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []any,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
) (string, []any, bool) {
@@ -329,7 +378,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", metricNames...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -343,6 +392,7 @@ func (b *StatementBuilder) buildReducedSpatialAggFastPath(
func (b *StatementBuilder) buildReducedTemporalAggregationCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []any,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
) (string, []any, bool) {
@@ -371,7 +421,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", metricNames...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -412,6 +462,7 @@ func (b *StatementBuilder) buildReducedSpatialAggregationCTE(
func (b *StatementBuilder) buildTemporalAggDeltaFastPath(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []any,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -449,7 +500,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", metricNames...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -466,6 +517,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
keys map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []any,
variables map[string]qbtypes.VariableItem,
tsTable string,
) (string, []any, []string, error) {
@@ -486,6 +538,9 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
Variables: variables,
StartNs: start,
EndNs: end,
Flagger: b.flagger,
Signal: telemetrytypes.SignalMetrics,
MetricContext: &telemetrytypes.MetricContext{MetricName: query.Aggregations[0].MetricName},
})
if err != nil {
return "", nil, nil, err
@@ -504,7 +559,7 @@ func (b *StatementBuilder) buildTimeSeriesCTE(
}
sb.Where(
sb.In("metric_name", query.Aggregations[0].MetricName),
sb.In("metric_name", metricNames...),
sb.GTE("unix_milli", start),
sb.LTE("unix_milli", end),
)
@@ -535,22 +590,24 @@ func (b *StatementBuilder) buildTemporalAggregationCTE(
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
_ map[string][]*telemetrytypes.TelemetryFieldKey,
metricNames []any,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
) (string, []any, error) {
if query.Aggregations[0].Temporality == metrictypes.Delta {
return b.buildTemporalAggDelta(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
return b.buildTemporalAggDelta(ctx, start, end, query, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
} else if query.Aggregations[0].Temporality != metrictypes.Multiple {
return b.buildTemporalAggCumulativeOrUnspecified(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
return b.buildTemporalAggCumulativeOrUnspecified(ctx, start, end, query, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
}
return b.buildTemporalAggForMultipleTemporalities(ctx, start, end, query, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
return b.buildTemporalAggForMultipleTemporalities(ctx, start, end, query, metricNames, samplesTable, timeSeriesCTE, timeSeriesCTEArgs)
}
func (b *StatementBuilder) buildTemporalAggDelta(
_ context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []any,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -582,7 +639,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", metricNames...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)
@@ -598,6 +655,7 @@ func (b *StatementBuilder) buildTemporalAggCumulativeOrUnspecified(
_ context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []any,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -623,7 +681,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", metricNames...),
baseSb.GTE("unix_milli", start),
baseSb.LT("unix_milli", end),
)
@@ -664,6 +722,7 @@ func (b *StatementBuilder) buildTemporalAggForMultipleTemporalities(
_ context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.MetricAggregation],
metricNames []any,
samplesTable string,
timeSeriesCTE string,
timeSeriesCTEArgs []any,
@@ -714,7 +773,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", metricNames...),
sb.GTE("unix_milli", start),
sb.LT("unix_milli", end),
)

View File

@@ -356,8 +356,8 @@ func TestStatementBuilder(t *testing.T) {
},
}
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(nil)
cb := metricstelemetryschema.NewConditionBuilder(fm, nil)
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys, err := telemetrytypestest.LoadFieldKeysFromJSON("testdata/keys_map.json")
if err != nil {
@@ -404,8 +404,8 @@ func TestStatementBuilder(t *testing.T) {
}
func TestGroupByAliasAvoidsColumnCollision(t *testing.T) {
fm := metricstelemetryschema.NewFieldMapper()
cb := metricstelemetryschema.NewConditionBuilder(fm)
fm := metricstelemetryschema.NewFieldMapper(nil)
cb := metricstelemetryschema.NewConditionBuilder(fm, nil)
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
keys, err := telemetrytypestest.LoadFieldKeysFromJSON("testdata/keys_map.json")
require.NoError(t, err)

View File

@@ -125,7 +125,7 @@ func (b *defaultConditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingLogicalFields(ctx, orgID, b.fl, key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(ctx, orgID, b.fl, telemetrytypes.SignalUnspecified, nil, key, fieldKeys)
// has/hasAny/hasAll/hasToken are logs-body-only functions; they never apply to the
// resource fingerprint table, so skip them (the main query still evaluates them).

View File

@@ -41,7 +41,7 @@ func (c *conditionBuilder) ConditionFor(
// an unknown key simply yields no condition rather than an error. Metadata
// fields have no family support, so every logical field is single-member
// and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(ctx, orgID, nil, key, fieldKeys))
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(ctx, orgID, nil, telemetrytypes.SignalUnspecified, nil, key, fieldKeys))
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {

View File

@@ -0,0 +1,43 @@
package telemetrymetadata
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
)
// The flagger provider registration is process-global and keyed by provider
// name, so each flagger must be used before the next one is created.
func TestFamilyValueNames(t *testing.T) {
selector := &telemetrytypes.FieldValueSelector{
FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "deployment.environment"},
}
off := &telemetryMetaStore{fl: flaggertest.WithBooleanFlags(t, map[string]bool{})}
assert.Equal(t,
[]string{"deployment.environment"},
off.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalLogs, selector),
"the flag default keeps values literal")
on := &telemetryMetaStore{fl: flaggertest.WithBooleanFlags(t, map[string]bool{
flagger.FeatureResolveSemconvFamilies.String(): true,
})}
assert.Equal(t,
[]string{"deployment.environment.name", "deployment.environment"},
on.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalLogs, selector),
"values for one spelling must cover the whole family")
assert.Equal(t,
[]string{
"resource_deployment.environment.name", "resource_deployment_environment_name",
"deployment.environment.name", "deployment_environment_name",
"resource_deployment.environment", "resource_deployment_environment",
"deployment.environment", "deployment_environment",
},
on.familyValueNames(context.Background(), valuer.UUID{}, telemetrytypes.SignalMetrics, selector),
"metric values must cover the stored label spellings")
}

View File

@@ -14,6 +14,7 @@ import (
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/telemetryschema/audittelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema"
@@ -75,6 +76,30 @@ func escapeForLike(s string) string {
return strings.ReplaceAll(strings.ReplaceAll(s, `_`, `\_`), `%`, `\%`)
}
// familyValueNames returns the spellings whose stored values merge into the
// requested name's suggestions when the resolve_semconv_families flag is on
// for the org; otherwise just the requested name. A filter on any family
// spelling matches every member, so the suggested values must cover them all.
func (t *telemetryMetaStore) familyValueNames(ctx context.Context, orgID valuer.UUID, signal telemetrytypes.Signal, fieldValueSelector *telemetrytypes.FieldValueSelector) []string {
if !querybuilder.SemconvFamiliesEnabled(ctx, orgID, t.fl) {
return []string{fieldValueSelector.Name}
}
return semconv.AttributeMembers(telemetrytypes.FieldKeySelector{
Name: fieldValueSelector.Name,
Signal: signal,
FieldContext: fieldValueSelector.FieldContext,
MetricContext: fieldValueSelector.MetricContext,
})
}
func anySlice(values []string) []any {
out := make([]any, len(values))
for i, value := range values {
out[i] = value
}
return out
}
func NewTelemetryMetaStore(
settings factory.ProviderSettings,
telemetrystore telemetrystore.TelemetryStore,
@@ -1362,23 +1387,39 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
FieldDataType: fieldValueSelector.FieldDataType,
}
selectColumn, err := t.fm.FieldFor(ctx, orgID, 0, 0, key)
if err != nil {
// we don't have a explicit column to select from the related metadata table
// so we will select either from resource_attributes or attributes table
// in that order
resourceColumn, _ := t.fm.FieldFor(ctx, orgID, 0, 0, &telemetrytypes.TelemetryFieldKey{
Name: key.Name,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
attributeColumn, _ := t.fm.FieldFor(ctx, orgID, 0, 0, &telemetrytypes.TelemetryFieldKey{
Name: key.Name,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
selectColumn = fmt.Sprintf("if(notEmpty(%s), %s, %s)", resourceColumn, resourceColumn, attributeColumn)
// One column per family spelling, merged current-first so suggestions
// cover rows that carry only an old spelling.
names := t.familyValueNames(ctx, orgID, fieldValueSelector.Signal, fieldValueSelector)
memberColumns := make([]string, 0, len(names))
for _, name := range names {
memberKey := &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: fieldValueSelector.Signal,
FieldContext: fieldValueSelector.FieldContext,
FieldDataType: fieldValueSelector.FieldDataType,
}
memberColumn, err := t.fm.FieldFor(ctx, orgID, 0, 0, memberKey)
if err != nil {
// we don't have a explicit column to select from the related metadata table
// so we will select either from resource_attributes or attributes table
// in that order
resourceColumn, _ := t.fm.FieldFor(ctx, orgID, 0, 0, &telemetrytypes.TelemetryFieldKey{
Name: name,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
attributeColumn, _ := t.fm.FieldFor(ctx, orgID, 0, 0, &telemetrytypes.TelemetryFieldKey{
Name: name,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
memberColumn = fmt.Sprintf("if(notEmpty(%s), %s, %s)", resourceColumn, resourceColumn, attributeColumn)
}
memberColumns = append(memberColumns, memberColumn)
}
selectColumn := memberColumns[len(memberColumns)-1]
for i := len(memberColumns) - 2; i >= 0; i-- {
selectColumn = fmt.Sprintf("if(notEmpty(%s), %s, %s)", memberColumns[i], memberColumns[i], selectColumn)
}
sb := sqlbuilder.Select("DISTINCT " + selectColumn).From(t.relatedMetadataDBName + "." + t.relatedMetadataTblName)
@@ -1500,7 +1541,7 @@ func (t *telemetryMetaStore) GetRelatedValues(ctx context.Context, orgID valuer.
return t.getRelatedValues(ctx, orgID, fieldValueSelector)
}
func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalTraces.StringValue(),
instrumentationtypes.CodeNamespace: "metadata",
@@ -1515,7 +1556,11 @@ func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueS
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(t.tracesDBName + "." + t.tracesFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
if names := t.familyValueNames(ctx, orgID, telemetrytypes.SignalTraces, fieldValueSelector); len(names) > 1 {
sb.Where(sb.In("tag_key", anySlice(names)...))
} else {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
}
}
// now look at the field context
@@ -1590,7 +1635,7 @@ func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueS
return values, complete, nil
}
func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
func (t *telemetryMetaStore) getLogFieldValues(ctx context.Context, orgID valuer.UUID, fieldValueSelector *telemetrytypes.FieldValueSelector) (*telemetrytypes.TelemetryFieldValues, bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
instrumentationtypes.TelemetrySignal: telemetrytypes.SignalLogs.StringValue(),
instrumentationtypes.CodeNamespace: "metadata",
@@ -1605,7 +1650,11 @@ 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))
if names := t.familyValueNames(ctx, orgID, telemetrytypes.SignalLogs, fieldValueSelector); len(names) > 1 {
sb.Where(sb.In("tag_key", anySlice(names)...))
} else {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
}
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1803,7 +1852,11 @@ 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))
if names := t.familyValueNames(ctx, orgID, telemetrytypes.SignalMetrics, fieldValueSelector); len(names) > 1 {
sb.Where(sb.In("attr_name", anySlice(names)...))
} else {
sb.Where(sb.E("attr_name", fieldValueSelector.Name))
}
}
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextUnspecified {
@@ -1815,7 +1868,15 @@ 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))
metricNames := []string{fieldValueSelector.MetricContext.MetricName}
if querybuilder.SemconvFamiliesEnabled(ctx, orgID, t.fl) {
metricNames = semconv.MetricNames(fieldValueSelector.MetricContext.MetricName)
}
if len(metricNames) > 1 {
sb.Where(sb.In("metric_name", anySlice(metricNames)...))
} else {
sb.Where(sb.E("metric_name", fieldValueSelector.MetricContext.MetricName))
}
}
if fieldValueSelector.MetricContext != nil && fieldValueSelector.MetricContext.MetricNamespace != "" {
sb.Where(sb.Like("metric_name", escapeForLike(fieldValueSelector.MetricContext.MetricNamespace)+"%"))
@@ -2125,12 +2186,12 @@ func (t *telemetryMetaStore) GetAllValues(ctx context.Context, orgID valuer.UUID
switch fieldValueSelector.Signal {
case telemetrytypes.SignalTraces:
values, complete, err = t.getSpanFieldValues(ctx, fieldValueSelector)
values, complete, err = t.getSpanFieldValues(ctx, orgID, fieldValueSelector)
case telemetrytypes.SignalLogs:
if fieldValueSelector.Source == telemetrytypes.SourceAudit {
values, complete, err = t.getAuditFieldValues(ctx, fieldValueSelector)
} else {
values, complete, err = t.getLogFieldValues(ctx, fieldValueSelector)
values, complete, err = t.getLogFieldValues(ctx, orgID, fieldValueSelector)
}
case telemetrytypes.SignalMetrics:
if fieldValueSelector.Source == telemetrytypes.SourceMeter {
@@ -2143,13 +2204,13 @@ func (t *telemetryMetaStore) GetAllValues(ctx context.Context, orgID valuer.UUID
mapOfRelatedValues := make(map[any]bool)
allUnspecifiedValues := &telemetrytypes.TelemetryFieldValues{}
tracesValues, tracesComplete, err := t.getSpanFieldValues(ctx, fieldValueSelector)
tracesValues, tracesComplete, err := t.getSpanFieldValues(ctx, orgID, fieldValueSelector)
if err == nil {
populateComplete := populateAllUnspecifiedValues(allUnspecifiedValues, mapOfValues, mapOfRelatedValues, tracesValues, limit)
complete = complete && tracesComplete && populateComplete
}
logsValues, logsComplete, err := t.getLogFieldValues(ctx, fieldValueSelector)
logsValues, logsComplete, err := t.getLogFieldValues(ctx, orgID, fieldValueSelector)
if err == nil {
populateComplete := populateAllUnspecifiedValues(allUnspecifiedValues, mapOfValues, mapOfRelatedValues, logsValues, limit)
complete = complete && logsComplete && populateComplete

View File

@@ -149,7 +149,7 @@ func (c *conditionBuilder) ConditionFor(
// Audit fields have no family support, so every logical field is
// single-member and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(ctx, orgID, nil, key, fieldKeys))
resolved, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(ctx, orgID, nil, telemetrytypes.SignalUnspecified, nil, key, fieldKeys))
keys := querybuilder.SingleKeys(resolved)
var warnings []string
if warning != "" {

View File

@@ -460,7 +460,7 @@ func (c *conditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingLogicalFields(ctx, orgID, nil, key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(ctx, orgID, c.fl, telemetrytypes.SignalLogs, nil, key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
// search() resolves its own (optional) scope; handle it before key resolution.
@@ -468,18 +468,16 @@ func (c *conditionBuilder) ConditionFor(
return c.conditionForSearch(ctx, orgID, key, value, sb)
}
// Logs fields have no family support yet, so every logical field is
// single-member and flattens losslessly to its physical key.
resolved, warning := querybuilder.ResolveLogicalFields(key, matches)
keys := querybuilder.SingleKeys(resolved)
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
}
synthesized := false
if len(keys) == 0 {
if len(logicalFields) == 0 {
_, isIntrinsicColumn := logsV2Columns[key.Name]
var keys []*telemetrytypes.TelemetryFieldKey
switch {
case key.FieldContext == telemetrytypes.FieldContextBody && key.Name == "":
return nil, warnings, errors.NewInvalidInputf(errors.CodeInvalidInput, "missing key for body json search - expected key of the form `body.key` (ex: `body.status`)")
@@ -509,23 +507,38 @@ func (c *conditionBuilder) ConditionFor(
synthesized = true
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
logicalFields = querybuilder.WrapAsLogicalFields(key.Name, keys)
}
if skipResourceFilter && !synthesized {
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
for _, k := range keys {
if k.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, k)
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
for _, logical := range logicalFields {
if logical.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, logical)
}
}
if len(filtered) == 0 {
return nil, warnings, nil
}
keys = filtered
logicalFields = filtered
}
conds := make([]string, 0, len(keys))
for _, k := range keys {
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
if logical.IsFamily() {
// has/hasAny/hasAll/hasToken are body-JSON only; family members are
// attribute and resource keys, so keep the descriptive error.
if err := querybuilder.NewFunctionUnsupportedError(operator); err != nil {
return nil, nil, err
}
cond, err := querybuilder.LogicalFamilyCondition(ctx, orgID, startNs, endNs, c.fm, logical, operator, value, sb)
if err != nil {
return nil, nil, err
}
conds = append(conds, cond)
continue
}
k := logical.Single()
cond, err := c.conditionForKey(ctx, orgID, startNs, endNs, k, operator, value, sb)
if err != nil {
return nil, nil, err

View File

@@ -7,6 +7,7 @@ import (
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -17,10 +18,13 @@ import (
type conditionBuilder struct {
fm qbtypes.FieldMapper
// fl evaluates the resolve_semconv_families flag during resolution.
// A nil flagger keeps resolution literal.
fl flagger.Flagger
}
func NewConditionBuilder(fm qbtypes.FieldMapper) *conditionBuilder {
return &conditionBuilder{fm: fm}
func NewConditionBuilder(fm qbtypes.FieldMapper, fl flagger.Flagger) *conditionBuilder {
return &conditionBuilder{fm: fm, fl: fl}
}
// Labels read back as String from the `labels` JSON whatever type the metadata claims, so the
@@ -170,7 +174,7 @@ func (c *conditionBuilder) conditionFor(
return "", errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported operator: %v", operator)
}
// Metrics has no resource sub-query, so options are unused.
// Metrics has no resource sub-query; options carry only the metric context.
func (c *conditionBuilder) ConditionFor(
ctx context.Context,
orgID valuer.UUID,
@@ -178,7 +182,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,
@@ -189,11 +193,10 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
// Metric labels have no family support, so every logical field is
// single-member and flattens losslessly to its physical key.
keys := querybuilder.SingleKeys(querybuilder.MatchingLogicalFields(ctx, orgID, nil, key, fieldKeys))
logicalFields := querybuilder.MatchingLogicalFields(ctx, orgID, c.fl, telemetrytypes.SignalMetrics, options.MetricContext, key, fieldKeys)
var warnings []string
if len(keys) == 0 {
if len(logicalFields) == 0 {
var keys []*telemetrytypes.TelemetryFieldKey
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {
keys = []*telemetrytypes.TelemetryFieldKey{key}
} else {
@@ -208,11 +211,20 @@ func (c *conditionBuilder) ConditionFor(
key.FieldContext.StringValue()+"."+key.Name, telemetrytypes.FieldContextAttribute, key.FieldDataType))
}
}
logicalFields = querybuilder.WrapAsLogicalFields(key.Name, keys)
}
conds := make([]string, 0, len(keys))
for _, k := range keys {
cond, err := c.conditionFor(ctx, orgID, startNs, endNs, k, operator, value, sb)
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
if logical.IsFamily() {
cond, err := querybuilder.LogicalFamilyCondition(ctx, orgID, startNs, endNs, c.fm, logical, operator, value, sb)
if err != nil {
return nil, nil, err
}
conds = append(conds, cond)
continue
}
cond, err := c.conditionFor(ctx, orgID, startNs, endNs, logical.Single(), operator, value, sb)
if err != nil {
return nil, nil, err
}

View File

@@ -343,8 +343,8 @@ func TestConditionFor(t *testing.T) {
},
}
fm := NewFieldMapper()
conditionBuilder := NewConditionBuilder(fm)
fm := NewFieldMapper(nil)
conditionBuilder := NewConditionBuilder(fm, nil)
for _, tc := range testCases {
sb := sqlbuilder.NewSelectBuilder()
@@ -396,8 +396,8 @@ func TestConditionForMultipleKeys(t *testing.T) {
},
}
fm := NewFieldMapper()
conditionBuilder := NewConditionBuilder(fm)
fm := NewFieldMapper(nil)
conditionBuilder := NewConditionBuilder(fm, nil)
for _, tc := range testCases {
sb := sqlbuilder.NewSelectBuilder()
@@ -483,8 +483,8 @@ func TestConditionForKeyNotInMetadata(t *testing.T) {
},
}
fm := NewFieldMapper()
conditionBuilder := NewConditionBuilder(fm)
fm := NewFieldMapper(nil)
conditionBuilder := NewConditionBuilder(fm, nil)
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {

View File

@@ -6,6 +6,8 @@ import (
"slices"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/flagger"
"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"
@@ -35,7 +37,11 @@ var (
}
)
type fieldMapper struct{}
type fieldMapper struct {
// fl evaluates the resolve_semconv_families flag during resolution.
// A nil flagger keeps resolution literal.
fl flagger.Flagger
}
// CandidateKeys returns nil: metrics has no attribute-map fallback, so a context-missing
// key stays unresolved and the caller errors.
@@ -43,8 +49,8 @@ func (m *fieldMapper) CandidateKeys(_ context.Context, _ valuer.UUID, _ *telemet
return nil
}
func NewFieldMapper() qbtypes.FieldMapper {
return &fieldMapper{}
func NewFieldMapper(fl flagger.Flagger) qbtypes.FieldMapper {
return &fieldMapper{fl: fl}
}
func (m *fieldMapper) getColumn(_ context.Context, _, _ uint64, key *telemetrytypes.TelemetryFieldKey) ([]*schema.Column, error) {
@@ -108,13 +114,21 @@ func (m *fieldMapper) ExistsFor(_ context.Context, _ valuer.UUID, _, _ uint64, k
return fmt.Sprintf("not has(JSONExtractKeys(labels), '%s')", key.Name), nil
}
// ColumnExpressionFor merges the stored spellings of a semantic-convention
// family when the field resolves to exactly one family; everything else keeps
// the literal label read. An ambiguous name (several logical fields) stays
// literal, because a group-by column holds one expression.
func (m *fieldMapper) ColumnExpressionFor(
ctx context.Context,
orgID valuer.UUID,
startNs, endNs uint64,
field *telemetrytypes.TelemetryFieldKey,
_ telemetrytypes.FieldDataType,
_ map[string][]*telemetrytypes.TelemetryFieldKey,
keys map[string][]*telemetrytypes.TelemetryFieldKey,
) (string, error) {
logicalFields := querybuilder.MatchingLogicalFields(ctx, orgID, m.fl, telemetrytypes.SignalMetrics, nil, field, keys)
if len(logicalFields) == 1 && logicalFields[0].IsFamily() {
return querybuilder.LogicalValueExpr(ctx, orgID, startNs, endNs, m, logicalFields[0])
}
return m.FieldFor(ctx, orgID, startNs, endNs, field)
}

View File

@@ -120,7 +120,7 @@ func TestGetColumn(t *testing.T) {
},
}
fm := NewFieldMapper()
fm := NewFieldMapper(nil)
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
@@ -204,7 +204,7 @@ func TestGetFieldKeyName(t *testing.T) {
},
}
fm := NewFieldMapper()
fm := NewFieldMapper(nil)
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {

View File

@@ -220,7 +220,7 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
matches := querybuilder.MatchingLogicalFields(ctx, orgID, c.fl, key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(ctx, orgID, c.fl, telemetrytypes.SignalTraces, nil, key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)

View File

@@ -345,7 +345,7 @@ func (m *fieldMapper) resolveColumnExprs(
// probe succeeded) to its family when the metadata map proves membership;
// otherwise the key stays a single-member logical field.
func (m *fieldMapper) logicalForResolvedColumn(ctx context.Context, orgID valuer.UUID, field *telemetrytypes.TelemetryFieldKey, keys map[string][]*telemetrytypes.TelemetryFieldKey) *telemetrytypes.LogicalField {
for _, logical := range querybuilder.MatchingLogicalFields(ctx, orgID, m.fl, field, keys) {
for _, logical := range querybuilder.MatchingLogicalFields(ctx, orgID, m.fl, telemetrytypes.SignalTraces, nil, field, keys) {
if logical.IsFamily() &&
logical.FieldContext == field.FieldContext &&
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || logical.FieldDataType == field.FieldDataType) {
@@ -361,7 +361,7 @@ func (m *fieldMapper) logicalForResolvedColumn(ctx context.Context, orgID valuer
// of an already-emitted family are dropped rather than duplicated.
func (m *fieldMapper) upgradeToFamilies(ctx context.Context, orgID valuer.UUID, field *telemetrytypes.TelemetryFieldKey, candidates []*telemetrytypes.LogicalField, keys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
var families []*telemetrytypes.LogicalField
for _, logical := range querybuilder.MatchingLogicalFields(ctx, orgID, m.fl, field, keys) {
for _, logical := range querybuilder.MatchingLogicalFields(ctx, orgID, m.fl, telemetrytypes.SignalTraces, nil, field, keys) {
if logical.IsFamily() {
families = append(families, logical)
}

View File

@@ -57,6 +57,9 @@ type ConditionBuilder interface {
type ConditionBuilderOptions struct {
// SkipResourceFilter drops the resource context from the candidate set.
SkipResourceFilter bool
// MetricContext carries the queried metric name so metric-scoped
// semantic-convention families can resolve.
MetricContext *telemetrytypes.MetricContext
}
type AggExprRewriter interface {
// Rewrite rewrites the aggregation expression to be used in the query.

View File

@@ -64,6 +64,9 @@ type overlayFile struct {
Families map[string]overlayFamily `yaml:"families"`
}
// Contexts and Signals set the family-level gate: the axes a family may
// resolve on at all. The Add* fields widen the schema-derived member scopes,
// for renames SigNoz applies beyond where the schema published them.
type overlayFamily struct {
Enabled *bool `yaml:"enabled"`
Kind string `yaml:"kind"`
@@ -79,27 +82,37 @@ type overlayFamily struct {
ValueMap map[string]string `yaml:"value_map"`
}
type edge struct {
old string
current string
kind string
// scope is the constraint a rename edge carries. A nil axis places no
// constraint on that axis.
type scope struct {
contexts []string
signals []string
allContexts bool
allSignals bool
applyToMetrics []string
}
type edge struct {
old string
current string
kind string
scope scope
}
type graphKey struct{ kind, name string }
type generatedFamily struct {
Current string
Old []string
Kind string
type generatedMember struct {
Name string
Contexts []string
Signals []string
ApplyToMetrics []string
ValueMap map[string]string
}
type generatedFamily struct {
Current string
Kind string
Members []generatedMember
Contexts []string
Signals []string
ValueMap map[string]string
}
func main() {
@@ -232,25 +245,25 @@ func collectEdges(schemas []schemaFile) ([]edge, error) {
{name: "metrics", section: version.Metrics},
}
for _, scoped := range sections {
contexts, signals, allContexts, allSignals, err := scopeForSection(scoped.name)
sectionScope, err := scopeForSection(scoped.name)
if err != nil {
return nil, err
}
for _, change := range scoped.section.Changes {
if change.RenameAttributes != nil {
edgeScope := sectionScope
edgeScope.applyToMetrics = change.RenameAttributes.ApplyToMetrics
for _, old := range sortedMapKeys(change.RenameAttributes.AttributeMap) {
versionEdges = append(versionEdges, edge{
old: old, current: change.RenameAttributes.AttributeMap[old], kind: kindAttribute,
contexts: contexts, signals: signals,
allContexts: allContexts, allSignals: allSignals,
applyToMetrics: change.RenameAttributes.ApplyToMetrics,
scope: edgeScope,
})
}
}
for _, old := range sortedMapKeys(change.RenameMetrics) {
versionEdges = append(versionEdges, edge{
old: old, current: change.RenameMetrics[old], kind: kindMetric,
contexts: []string{"metric"}, signals: []string{"metrics"},
scope: scope{contexts: []string{"metric"}, signals: []string{"metrics"}},
})
}
}
@@ -311,99 +324,165 @@ func compareVersionParts(left, right [3]int) int {
return 0
}
func scopeForSection(section string) (contexts, signals []string, allContexts, allSignals bool, err error) {
func scopeForSection(section string) (scope, error) {
switch section {
case "all":
return nil, nil, true, true, nil
return scope{}, nil
case "resources":
return []string{"resource"}, nil, false, true, nil
return scope{contexts: []string{"resource"}}, nil
case "spans":
return []string{"attribute"}, []string{"traces"}, false, false, nil
return scope{contexts: []string{"attribute"}, signals: []string{"traces"}}, nil
case "logs":
return []string{"attribute"}, []string{"logs"}, false, false, nil
return scope{contexts: []string{"attribute"}, signals: []string{"logs"}}, nil
case "metrics":
return []string{"attribute"}, []string{"metrics"}, false, false, nil
return scope{contexts: []string{"attribute"}, signals: []string{"metrics"}}, nil
default:
return nil, nil, false, false, fmt.Errorf("unsupported schema section %q", section)
return scope{}, fmt.Errorf("unsupported schema section %q", section)
}
}
// unionScope widens per axis: no constraint on either side widens to no
// constraint.
func unionScope(left, right scope) scope {
return scope{
contexts: unionAxis(left.contexts, right.contexts),
signals: unionAxis(left.signals, right.signals),
applyToMetrics: unionAxis(left.applyToMetrics, right.applyToMetrics),
}
}
func unionAxis(left, right []string) []string {
if left == nil || right == nil {
return nil
}
merged := appendUnique(append([]string(nil), left...), right...)
sort.Strings(merged)
return merged
}
// pathResult is one resolution path from a name to a family root: the root
// name and the hop count. Hops carry only reachability — a member's scope
// comes from its own rename edges, because the section a later rename is
// filed under says nothing about where the older spelling existed (the
// vendored schema files chained renames under different sections).
type pathResult struct {
root string
distance int
}
func rootsFor(next map[graphKey][]edge, kind, name string, distance int, seen map[string]bool) ([]pathResult, error) {
if seen[name] {
return nil, fmt.Errorf("rename cycle for %s %q", kind, name)
}
outgoing := next[graphKey{kind: kind, name: name}]
if len(outgoing) == 0 {
return []pathResult{{root: name, distance: distance}}, nil
}
seen[name] = true
defer delete(seen, name)
var results []pathResult
for _, hop := range outgoing {
hopResults, err := rootsFor(next, kind, hop.current, distance+1, seen)
if err != nil {
return nil, err
}
results = append(results, hopResults...)
}
return results, nil
}
func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily, error) {
edges, err := collectEdges(schemas)
if err != nil {
return nil, err
}
next := make(map[graphKey]string)
// One old name can fan out into several families when its rename edges are
// scoped differently, so the graph keeps every successor.
next := make(map[graphKey][]edge)
for _, item := range edges {
key := graphKey{kind: item.kind, name: item.old}
if existing, ok := next[key]; ok && existing == item.current {
// Repeated entries are common in chained schema histories. Treat an
// identical edge as a no-op so it cannot sever a later edge in the
// same chain (A -> B, B -> C, then a repeated A -> B).
merged := false
for i, existing := range next[key] {
if existing.current == item.current {
// Repeated entries are common in chained schema histories. Merge
// the scopes instead of appending so a repeat cannot sever a
// later edge in the same chain (A -> B, B -> C, then a repeated
// A -> B).
next[key][i].scope = unionScope(existing.scope, item.scope)
merged = true
break
}
}
if merged {
continue
}
// Schema history occasionally repeats an old name with a newer direct
// destination or rolls a rename back. Edges are collected
// oldest-to-newest, so the latest published current name must be a root.
delete(next, graphKey{kind: item.kind, name: item.current})
next[key] = item.current
next[key] = append(next[key], item)
}
type memberState struct {
sc scope
distance int
}
type familyState struct {
family generatedFamily
distance map[string]int
allContexts bool
allSignals bool
family generatedFamily
members map[string]*memberState
}
states := map[graphKey]*familyState{}
for _, item := range edges {
root, distance, err := rootFor(next, item.kind, item.old)
if err != nil {
return nil, err
}
key := graphKey{kind: item.kind, name: root}
state := states[key]
if state == nil {
state = &familyState{
family: generatedFamily{Current: root, Kind: item.kind},
distance: map[string]int{},
for key, outgoing := range next {
for _, item := range outgoing {
results, err := rootsFor(next, key.kind, item.current, 1, map[string]bool{key.name: true})
if err != nil {
return nil, err
}
for _, result := range results {
rootKey := graphKey{kind: key.kind, name: result.root}
state := states[rootKey]
if state == nil {
state = &familyState{
family: generatedFamily{Current: result.root, Kind: key.kind},
members: map[string]*memberState{},
}
states[rootKey] = state
}
member := state.members[key.name]
if member == nil {
state.members[key.name] = &memberState{sc: item.scope, distance: result.distance}
continue
}
member.sc = unionScope(member.sc, item.scope)
if result.distance < member.distance {
member.distance = result.distance
}
}
states[key] = state
}
if prior, ok := state.distance[item.old]; !ok || distance < prior {
state.distance[item.old] = distance
}
state.allContexts = state.allContexts || item.allContexts
state.allSignals = state.allSignals || item.allSignals
state.family.Contexts = appendUnique(state.family.Contexts, item.contexts...)
state.family.Signals = appendUnique(state.family.Signals, item.signals...)
state.family.ApplyToMetrics = appendUnique(state.family.ApplyToMetrics, item.applyToMetrics...)
}
for _, state := range states {
for old := range state.distance {
if old != state.family.Current {
state.family.Old = append(state.family.Old, old)
}
names := make([]string, 0, len(state.members))
for name := range state.members {
names = append(names, name)
}
sort.Slice(state.family.Old, func(i, j int) bool {
left, right := state.family.Old[i], state.family.Old[j]
if state.distance[left] != state.distance[right] {
return state.distance[left] < state.distance[right]
sort.Slice(names, func(i, j int) bool {
left, right := state.members[names[i]], state.members[names[j]]
if left.distance != right.distance {
return left.distance < right.distance
}
return left < right
return names[i] < names[j]
})
if state.allContexts {
state.family.Contexts = nil
} else {
sort.Strings(state.family.Contexts)
for _, name := range names {
member := state.members[name]
state.family.Members = append(state.family.Members, generatedMember{
Name: name,
Contexts: sortedCopy(member.sc.contexts),
Signals: sortedCopy(member.sc.signals),
ApplyToMetrics: sortedCopy(member.sc.applyToMetrics),
})
}
if state.allSignals {
state.family.Signals = nil
} else {
sort.Strings(state.family.Signals)
}
sort.Strings(state.family.ApplyToMetrics)
}
for _, current := range sortedMapKeys(overlay.Families) {
@@ -424,10 +503,7 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
kind,
)
}
state = &familyState{
family: generatedFamily{Current: current, Kind: kind, Old: append([]string(nil), policy.Old...)},
distance: map[string]int{},
}
state = &familyState{family: generatedFamily{Current: current, Kind: kind}}
states[key] = state
}
applyOverlay(&state.family, policy)
@@ -446,7 +522,7 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
if !enabled {
continue
}
if len(state.family.Old) == 0 {
if len(state.family.Members) == 0 {
return nil, fmt.Errorf(
"enabled family %q with kind %q has no old members",
state.family.Current,
@@ -455,7 +531,6 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
}
sort.Strings(state.family.Contexts)
sort.Strings(state.family.Signals)
sort.Strings(state.family.ApplyToMetrics)
result = append(result, state.family)
}
@@ -468,21 +543,13 @@ func buildFamilies(schemas []schemaFile, overlay overlayFile) ([]generatedFamily
return result, nil
}
func rootFor(next map[graphKey]string, kind, name string) (string, int, error) {
seen := map[string]bool{}
distance := 0
for {
if seen[name] {
return "", 0, fmt.Errorf("rename cycle for %s %q", kind, name)
}
seen[name] = true
current, ok := next[graphKey{kind: kind, name: name}]
if !ok {
return name, distance, nil
}
name = current
distance++
func sortedCopy(values []string) []string {
if values == nil {
return nil
}
out := append([]string(nil), values...)
sort.Strings(out)
return out
}
func normalizedOverlayKind(current string, policy overlayFamily) (string, error) {
@@ -501,15 +568,29 @@ func applyOverlay(family *generatedFamily, policy overlayFamily) {
family.Kind = policy.Kind
}
if policy.Old != nil {
family.Old = append([]string(nil), policy.Old...)
family.Members = nil
for _, old := range policy.Old {
family.Members = append(family.Members, generatedMember{Name: old})
}
}
for _, old := range policy.AddOld {
if familyHasMember(family, old) {
continue
}
family.Members = append(family.Members, generatedMember{Name: old})
}
family.Old = appendUnique(family.Old, policy.AddOld...)
if len(policy.ExcludeOld) > 0 {
excluded := make(map[string]bool, len(policy.ExcludeOld))
for _, old := range policy.ExcludeOld {
excluded[old] = true
}
family.Old = deleteMatching(family.Old, excluded)
kept := family.Members[:0]
for _, member := range family.Members {
if !excluded[member.Name] {
kept = append(kept, member)
}
}
family.Members = kept
}
if policy.Contexts != nil {
family.Contexts = append([]string(nil), policy.Contexts...)
@@ -517,12 +598,28 @@ func applyOverlay(family *generatedFamily, policy overlayFamily) {
if policy.Signals != nil {
family.Signals = append([]string(nil), policy.Signals...)
}
family.Contexts = appendUnique(family.Contexts, policy.AddContexts...)
family.Signals = appendUnique(family.Signals, policy.AddSignals...)
if policy.ApplyToMetrics != nil {
family.ApplyToMetrics = append([]string(nil), policy.ApplyToMetrics...)
// Add* fields only widen: a nil gate already admits everything, so they
// extend a gate only when the overlay set one.
if family.Contexts != nil {
family.Contexts = appendUnique(family.Contexts, policy.AddContexts...)
}
if family.Signals != nil {
family.Signals = appendUnique(family.Signals, policy.AddSignals...)
}
for i := range family.Members {
if len(policy.AddContexts) > 0 {
family.Members[i].Contexts = unionAxis(family.Members[i].Contexts, sortedCopy(policy.AddContexts))
}
if len(policy.AddSignals) > 0 {
family.Members[i].Signals = unionAxis(family.Members[i].Signals, sortedCopy(policy.AddSignals))
}
if policy.ApplyToMetrics != nil {
family.Members[i].ApplyToMetrics = sortedCopy(policy.ApplyToMetrics)
}
if len(policy.AddApplyToMetrics) > 0 && family.Members[i].ApplyToMetrics != nil {
family.Members[i].ApplyToMetrics = unionAxis(family.Members[i].ApplyToMetrics, sortedCopy(policy.AddApplyToMetrics))
}
}
family.ApplyToMetrics = appendUnique(family.ApplyToMetrics, policy.AddApplyToMetrics...)
if policy.ValueMap != nil {
family.ValueMap = make(map[string]string, len(policy.ValueMap))
for old, current := range policy.ValueMap {
@@ -531,6 +628,15 @@ func applyOverlay(family *generatedFamily, policy overlayFamily) {
}
}
func familyHasMember(family *generatedFamily, name string) bool {
for _, member := range family.Members {
if member.Name == name {
return true
}
}
return false
}
func appendUnique(values []string, additions ...string) []string {
seen := make(map[string]bool, len(values)+len(additions))
for _, value := range values {
@@ -546,16 +652,6 @@ func appendUnique(values []string, additions ...string) []string {
return values
}
func deleteMatching(values []string, excluded map[string]bool) []string {
result := values[:0]
for _, value := range values {
if !excluded[value] {
result = append(result, value)
}
}
return result
}
func renderGo(families []generatedFamily) ([]byte, error) {
var out bytes.Buffer
out.WriteString("// Code generated by scripts/semconv. DO NOT EDIT.\n\n")
@@ -564,7 +660,11 @@ func renderGo(families []generatedFamily) ([]byte, error) {
for _, family := range families {
if len(family.Contexts) > 0 || len(family.Signals) > 0 {
needsTelemetryTypes = true
break
}
for _, member := range family.Members {
if len(member.Contexts) > 0 || len(member.Signals) > 0 {
needsTelemetryTypes = true
}
}
}
if needsTelemetryTypes {
@@ -572,27 +672,36 @@ func renderGo(families []generatedFamily) ([]byte, error) {
}
out.WriteString("var families = []Family{\n")
for _, family := range families {
contexts, err := goFieldContextSlice(family.Contexts)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
signals, err := goSignalSlice(family.Signals)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
out.WriteString("\t{\n")
fmt.Fprintf(&out, "\t\tCurrent: %s,\n", strconv.Quote(family.Current))
fmt.Fprintf(&out, "\t\tOld: %s,\n", goStringSlice(family.Old))
fmt.Fprintf(&out, "\t\tcurrent: %s,\n", strconv.Quote(family.Current))
if family.Kind == kindMetric {
out.WriteString("\t\tKind: KindMetric,\n")
out.WriteString("\t\tkind: KindMetric,\n")
} else {
out.WriteString("\t\tKind: KindAttribute,\n")
out.WriteString("\t\tkind: KindAttribute,\n")
}
out.WriteString("\t\tmembers: []Member{\n")
for _, member := range family.Members {
if err := writeGoMember(&out, family.Current, member); err != nil {
return nil, err
}
}
out.WriteString("\t\t},\n")
if len(family.Contexts) > 0 {
contexts, err := goFieldContextSlice(family.Contexts)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
fmt.Fprintf(&out, "\t\tcontexts: %s,\n", contexts)
}
if len(family.Signals) > 0 {
signals, err := goSignalSlice(family.Signals)
if err != nil {
return nil, fmt.Errorf("render family %q: %w", family.Current, err)
}
fmt.Fprintf(&out, "\t\tsignals: %s,\n", signals)
}
fmt.Fprintf(&out, "\t\tContexts: %s,\n", contexts)
fmt.Fprintf(&out, "\t\tSignals: %s,\n", signals)
fmt.Fprintf(&out, "\t\tApplyToMetrics: %s,\n", goStringSlice(family.ApplyToMetrics))
if len(family.ValueMap) > 0 {
out.WriteString("\t\tValueMap: map[string]string{\n")
out.WriteString("\t\tvalueMap: map[string]string{\n")
keys := sortedMapKeys(family.ValueMap)
for _, key := range keys {
fmt.Fprintf(&out, "\t\t\t%s: %s,\n", strconv.Quote(key), strconv.Quote(family.ValueMap[key]))
@@ -605,6 +714,29 @@ func renderGo(families []generatedFamily) ([]byte, error) {
return format.Source(out.Bytes())
}
func writeGoMember(out *bytes.Buffer, current string, member generatedMember) error {
parts := []string{fmt.Sprintf("name: %s", strconv.Quote(member.Name))}
if len(member.Contexts) > 0 {
contexts, err := goFieldContextSlice(member.Contexts)
if err != nil {
return fmt.Errorf("render family %q member %q: %w", current, member.Name, err)
}
parts = append(parts, "contexts: "+contexts)
}
if len(member.Signals) > 0 {
signals, err := goSignalSlice(member.Signals)
if err != nil {
return fmt.Errorf("render family %q member %q: %w", current, member.Name, err)
}
parts = append(parts, "signals: "+signals)
}
if len(member.ApplyToMetrics) > 0 {
parts = append(parts, "applyToMetrics: "+goStringSlice(member.ApplyToMetrics))
}
fmt.Fprintf(out, "\t\t\t{%s},\n", strings.Join(parts, ", "))
return nil
}
func goStringSlice(values []string) string {
if len(values) == 0 {
return "nil"
@@ -659,21 +791,38 @@ func goSignalSlice(values []string) (string, error) {
func renderTypeScript(families []generatedFamily) []byte {
var out bytes.Buffer
out.WriteString("// Code generated by scripts/semconv. DO NOT EDIT.\n\n")
out.WriteString("// An empty contexts/signals/applyToMetrics array places no constraint on\n")
out.WriteString("// that axis.\n")
out.WriteString("export type SemconvMember = {\n")
out.WriteString("\treadonly name: string;\n")
out.WriteString("\treadonly contexts: readonly string[];\n")
out.WriteString("\treadonly signals: readonly string[];\n")
out.WriteString("\treadonly applyToMetrics: readonly string[];\n};\n\n")
out.WriteString("export type SemconvFamily = {\n")
out.WriteString("\treadonly current: string;\n\treadonly old: readonly string[];\n")
out.WriteString("\treadonly current: string;\n")
out.WriteString("\treadonly kind: 'attribute' | 'metric';\n")
out.WriteString("\treadonly members: readonly SemconvMember[];\n")
out.WriteString("\treadonly contexts: readonly string[];\n\treadonly signals: readonly string[];\n")
out.WriteString("\treadonly applyToMetrics: readonly string[];\n")
out.WriteString("\treadonly valueMap: Readonly<Record<string, string>>;\n};\n\n")
out.WriteString("export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [\n")
for _, family := range families {
out.WriteString("\t{\n")
fmt.Fprintf(&out, "\t\tcurrent: %s,\n", tsString(family.Current))
fmt.Fprintf(&out, "\t\told: %s,\n", tsStringSlice(family.Old))
fmt.Fprintf(&out, "\t\tkind: %s,\n", tsString(family.Kind))
out.WriteString("\t\tmembers: [\n")
for _, member := range family.Members {
fmt.Fprintf(
&out,
"\t\t\t{ name: %s, contexts: %s, signals: %s, applyToMetrics: %s },\n",
tsString(member.Name),
tsStringSlice(member.Contexts),
tsStringSlice(member.Signals),
tsStringSlice(member.ApplyToMetrics),
)
}
out.WriteString("\t\t],\n")
fmt.Fprintf(&out, "\t\tcontexts: %s,\n", tsStringSlice(family.Contexts))
fmt.Fprintf(&out, "\t\tsignals: %s,\n", tsStringSlice(family.Signals))
fmt.Fprintf(&out, "\t\tapplyToMetrics: %s,\n", tsStringSlice(family.ApplyToMetrics))
out.WriteString("\t\tvalueMap: {")
keys := sortedMapKeys(family.ValueMap)
for i, key := range keys {

View File

@@ -37,6 +37,11 @@ versions:
assert.ErrorContains(t, err, `schema version "latest"`, "malformed versions must not be silently reordered")
}
func TestParseSchemaVersionRejectsNonNumericComponent(t *testing.T) {
_, err := parseSchemaVersion("1.2.x")
assert.ErrorContains(t, err, `invalid numeric component "x"`, "non-numeric version components must fail generation")
}
func TestBuildFamiliesResolvesRenameChain(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
@@ -71,11 +76,13 @@ versions:
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "c",
Old: []string{"b", "x", "a"},
Kind: kindAttribute,
Contexts: []string{"attribute"},
Signals: []string{"traces"},
Current: "c",
Kind: kindAttribute,
Members: []generatedMember{
{Name: "b", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
{Name: "x", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
{Name: "a", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
},
}}, families, "predecessors should be ordered by distance and then name")
}
@@ -120,27 +127,184 @@ versions:
require.NoError(t, err)
assert.Equal(t, []generatedFamily{
{
Current: "all.current", Old: []string{"all.old"}, Kind: kindAttribute,
Contexts: nil, Signals: nil,
Current: "all.current", Kind: kindAttribute,
Members: []generatedMember{{Name: "all.old"}},
},
{
Current: "cpu.mode", Old: []string{"state"}, Kind: kindAttribute,
Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"system.cpu.time"},
Current: "cpu.mode", Kind: kindAttribute,
Members: []generatedMember{{
Name: "state", Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"system.cpu.time"},
}},
},
{
Current: "current.metric", Old: []string{"old.metric"}, Kind: kindMetric,
Contexts: []string{"metric"}, Signals: []string{"metrics"},
Current: "current.metric", Kind: kindMetric,
Members: []generatedMember{{Name: "old.metric", Contexts: []string{"metric"}, Signals: []string{"metrics"}}},
},
{
Current: "log.current", Old: []string{"log.old"}, Kind: kindAttribute,
Contexts: []string{"attribute"}, Signals: []string{"logs"},
Current: "log.current", Kind: kindAttribute,
Members: []generatedMember{{Name: "log.old", Contexts: []string{"attribute"}, Signals: []string{"logs"}}},
},
{
Current: "resource.current", Old: []string{"resource.old"}, Kind: kindAttribute,
Contexts: []string{"resource"},
Current: "resource.current", Kind: kindAttribute,
Members: []generatedMember{{Name: "resource.old", Contexts: []string{"resource"}}},
},
}, families, "schema sections should produce their documented signal and context scopes")
}, families, "schema sections should produce their documented per-member signal and context scopes")
}
func TestBuildFamiliesKeepsFanOutSeparate(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
state: cpu.mode
apply_to_metrics: [system.cpu.time]
1.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
state: db.client.connection.state
apply_to_metrics: [db.client.connections.usage]
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{
{
Current: "cpu.mode", Kind: kindAttribute,
Members: []generatedMember{{
Name: "state", Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"system.cpu.time"},
}},
},
{
Current: "db.client.connection.state", Kind: kindAttribute,
Members: []generatedMember{{
Name: "state", Contexts: []string{"attribute"}, Signals: []string{"metrics"},
ApplyToMetrics: []string{"db.client.connections.usage"},
}},
},
}, families, "an old name with differently scoped rename targets must keep one membership per target")
}
func TestBuildFamiliesKeepsUnscopedRenameUnscoped(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
direction: network.io.direction
apply_to_metrics: [system.disk.io, system.disk.merged]
1.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
system.network.io.direction: network.io.direction
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "network.io.direction", Kind: kindAttribute,
Members: []generatedMember{
{Name: "direction", Contexts: []string{"attribute"}, Signals: []string{"metrics"}, ApplyToMetrics: []string{"system.disk.io", "system.disk.merged"}},
{Name: "system.network.io.direction", Contexts: []string{"attribute"}, Signals: []string{"metrics"}},
},
}}, families, "a rename without apply_to_metrics stays unbounded; a scoped sibling must not bound it")
}
func TestBuildFamiliesKeepsMemberScopesThroughChains(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
metrics:
changes:
- rename_attributes:
attribute_map:
messaging.client_id: messaging.client.id
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
messaging.kafka.client_id: messaging.client_id
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "messaging.client.id", Kind: kindAttribute,
Members: []generatedMember{
{Name: "messaging.client_id", Contexts: []string{"attribute"}, Signals: []string{"metrics"}},
{Name: "messaging.kafka.client_id", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
},
}}, families, "a member keeps the scope of its own rename edge; later hops only carry it to the root")
}
func TestBuildFamiliesRecordsCrossContextMembersSeparately(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
http.user_agent: user_agent.original
resources:
changes:
- rename_attributes:
attribute_map:
browser.user_agent: user_agent.original
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "user_agent.original", Kind: kindAttribute,
Members: []generatedMember{
{Name: "browser.user_agent", Contexts: []string{"resource"}},
{Name: "http.user_agent", Contexts: []string{"attribute"}, Signals: []string{"traces"}},
},
}}, families, "members renamed from different contexts must keep their own context scopes")
}
func TestBuildFamiliesMergesRepeatedEdgeScopes(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
2.0.0:
logs:
changes:
- rename_attributes:
attribute_map:
old: current
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
old: current
`), &schema), "test schema must decode")
families, err := buildFamilies([]schemaFile{schema}, overlayFile{DefaultEnabled: true})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{
{Name: "old", Contexts: []string{"attribute"}, Signals: []string{"logs", "traces"}},
},
}}, families, "the same rename filed under several sections widens the member scope")
}
func TestOverlayAddsFamilyWithoutSchemaHistory(t *testing.T) {
@@ -157,8 +321,8 @@ func TestOverlayAddsFamilyWithoutSchemaHistory(t *testing.T) {
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "added.current",
Old: []string{"added.old"},
Kind: kindAttribute,
Members: []generatedMember{{Name: "added.old"}},
Contexts: []string{"resource"},
Signals: []string{"traces"},
}}, families, "an explicit overlay family should not require schema history")
@@ -179,26 +343,76 @@ versions:
enabled := true
families, err := buildFamilies([]schemaFile{schema}, overlayFile{Families: map[string]overlayFamily{
"current": {
Enabled: &enabled,
AddOld: []string{"older"},
ExcludeOld: []string{"old"},
AddContexts: []string{"resource"},
AddSignals: []string{"logs"},
ValueMap: map[string]string{"legacy": "current"},
Enabled: &enabled,
AddOld: []string{"older"},
ExcludeOld: []string{"old"},
AddSignals: []string{"logs"},
ValueMap: map[string]string{"legacy": "current"},
},
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current",
Old: []string{"older"},
Kind: kindAttribute,
Contexts: []string{"attribute", "resource"},
Signals: []string{"logs", "traces"},
Members: []generatedMember{{Name: "older"}},
ValueMap: map[string]string{"legacy": "current"},
}}, families, "overlay additions and exclusions should be applied to the generated family")
}
func TestOverlayAddSignalsWidensMemberScopes(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
spans:
changes:
- rename_attributes:
attribute_map:
old: current
`), &schema), "test schema must decode")
enabled := true
families, err := buildFamilies([]schemaFile{schema}, overlayFile{Families: map[string]overlayFamily{
"current": {Enabled: &enabled, AddSignals: []string{"logs"}},
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current",
Kind: kindAttribute,
Members: []generatedMember{
{Name: "old", Contexts: []string{"attribute"}, Signals: []string{"logs", "traces"}},
},
}}, families, "add_signals widens the schema-derived member scopes and never narrows the family gate")
}
func TestOverlaySignalsSetTheFamilyGate(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
versions:
1.0.0:
all:
changes:
- rename_attributes:
attribute_map:
old: current
`), &schema), "test schema must decode")
enabled := true
families, err := buildFamilies([]schemaFile{schema}, overlayFile{Families: map[string]overlayFamily{
"current": {Enabled: &enabled, Signals: []string{"logs", "traces"}},
}})
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "current",
Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
Signals: []string{"logs", "traces"},
}}, families, "the overlay signals list gates the family without touching member scopes")
}
func TestOverlayDisablesFamilyWhenDefaultIsEnabled(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
@@ -225,7 +439,8 @@ versions:
func TestRenderGoIsDeterministic(t *testing.T) {
families := []generatedFamily{{
Current: "current", Old: []string{"old"}, Kind: kindAttribute,
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
ValueMap: map[string]string{"b": "2", "a": "1"},
}}
@@ -238,8 +453,8 @@ func TestRenderGoIsDeterministic(t *testing.T) {
func TestRenderGoUsesCanonicalTelemetryTypes(t *testing.T) {
families := []generatedFamily{{
Current: "current", Old: []string{"old"}, Kind: kindAttribute,
Contexts: []string{"resource"}, Signals: []string{"traces"},
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old", Contexts: []string{"resource"}, Signals: []string{"traces"}}},
}}
output, err := renderGo(families)
@@ -248,15 +463,107 @@ func TestRenderGoUsesCanonicalTelemetryTypes(t *testing.T) {
assert.Contains(t, string(output), "telemetrytypes.SignalTraces", "generated signals should use telemetrytypes")
}
func TestRenderGoRejectsUnknownContext(t *testing.T) {
families := []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old", Contexts: []string{"bogus"}}},
}}
_, err := renderGo(families)
assert.ErrorContains(t, err, `unsupported field context "bogus"`, "a bad overlay context must fail generation, not compile")
}
func TestRenderGoRejectsUnknownSignal(t *testing.T) {
families := []generatedFamily{{
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
Signals: []string{"bogus"},
}}
_, err := renderGo(families)
assert.ErrorContains(t, err, `unsupported signal "bogus"`, "a bad overlay signal must fail generation, not compile")
}
func TestRenderGoPinsOutput(t *testing.T) {
families := []generatedFamily{{
Current: "deployment.environment.name", Kind: kindAttribute,
Members: []generatedMember{{Name: "deployment.environment"}},
Signals: []string{"logs", "traces"},
}}
output, err := renderGo(families)
require.NoError(t, err)
assert.Equal(t, `// Code generated by scripts/semconv. DO NOT EDIT.
package semconv
import "github.com/SigNoz/signoz/pkg/types/telemetrytypes"
var families = []Family{
{
current: "deployment.environment.name",
kind: KindAttribute,
members: []Member{
{name: "deployment.environment"},
},
signals: []telemetrytypes.Signal{telemetrytypes.SignalLogs, telemetrytypes.SignalTraces},
},
}
`, string(output), "the emitted Go text is a contract; regeneration must be reviewable")
}
func TestRenderTypeScriptIsDeterministic(t *testing.T) {
families := []generatedFamily{{
Current: "current", Old: []string{"old"}, Kind: kindAttribute,
Current: "current", Kind: kindAttribute,
Members: []generatedMember{{Name: "old"}},
ValueMap: map[string]string{"b": "2", "a": "1"},
}}
assert.Equal(t, renderTypeScript(families), renderTypeScript(families), "TypeScript generation must not depend on map iteration order")
}
func TestRenderTypeScriptPinsOutput(t *testing.T) {
families := []generatedFamily{{
Current: "deployment.environment.name", Kind: kindAttribute,
Members: []generatedMember{{Name: "deployment.environment"}},
Signals: []string{"logs", "traces"},
}}
assert.Equal(t, `// Code generated by scripts/semconv. DO NOT EDIT.
// An empty contexts/signals/applyToMetrics array places no constraint on
// that axis.
export type SemconvMember = {
readonly name: string;
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly applyToMetrics: readonly string[];
};
export type SemconvFamily = {
readonly current: string;
readonly kind: 'attribute' | 'metric';
readonly members: readonly SemconvMember[];
readonly contexts: readonly string[];
readonly signals: readonly string[];
readonly valueMap: Readonly<Record<string, string>>;
};
export const SEMCONV_FAMILIES: readonly SemconvFamily[] = [
{
current: 'deployment.environment.name',
kind: 'attribute',
members: [
{ name: 'deployment.environment', contexts: [], signals: [], applyToMetrics: [] },
],
contexts: [],
signals: ['logs', 'traces'],
valueMap: {},
},
] as const;
`, string(renderTypeScript(families)), "the emitted TypeScript text is a contract; regeneration must be reviewable")
}
func TestBuildFamiliesHandlesRenameRollback(t *testing.T) {
var schema schemaFile
require.NoError(t, decodeKnownFields([]byte(`
@@ -280,11 +587,9 @@ versions:
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "original",
Old: []string{"temporary"},
Kind: kindMetric,
Contexts: []string{"metric"},
Signals: []string{"metrics"},
Current: "original",
Kind: kindMetric,
Members: []generatedMember{{Name: "temporary", Contexts: []string{"metric"}, Signals: []string{"metrics"}}},
}}, families, "the latest rollback destination should remain the family root")
}
@@ -356,11 +661,9 @@ versions:
require.NoError(t, err)
assert.Equal(t, []generatedFamily{{
Current: "shared.current",
Old: []string{"attribute.old"},
Kind: kindAttribute,
Contexts: []string{"attribute"},
Signals: []string{"traces"},
Current: "shared.current",
Kind: kindAttribute,
Members: []generatedMember{{Name: "attribute.old", Contexts: []string{"attribute"}, Signals: []string{"traces"}}},
}}, families, "a kind-less overlay policy should affect only the attribute family")
}

View File

@@ -2,10 +2,30 @@
#
# Families are keyed by their current OpenTelemetry name. Schema-derived
# families are disabled by default so rollout remains explicit and reversible.
# The signals list is the per-family rollout gate.
default_enabled: false
families:
deployment.environment.name:
enabled: true
signals: [traces, logs, metrics]
db.system.name:
enabled: true
# The db.system value domain also renamed, and no value mapping is read
# yet, so metrics resolution stays literal for this family.
signals: [traces, logs]
# These metric renames predate the vendored schema history, so the overlay
# declares the old names itself.
k8s.pod.cpu.usage:
enabled: true
kind: metric
old: [k8s.pod.cpu.utilization]
k8s.node.cpu.usage:
enabled: true
kind: metric
old: [k8s.node.cpu.utilization]
container.cpu.usage:
enabled: true
kind: metric
old: [container.cpu.utilization]