mirror of
https://github.com/SigNoz/signoz.git
synced 2026-08-24 21:50:32 +01:00
Compare commits
1 Commits
feat/trace
...
feat/semco
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b902642d59 |
2
.github/workflows/goci.yaml
vendored
2
.github/workflows/goci.yaml
vendored
@@ -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' ||
|
||||
|
||||
4
Makefile
4
Makefile
@@ -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"
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 != "" {
|
||||
|
||||
@@ -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())
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
125
pkg/querybuilder/family_condition.go
Normal file
125
pkg/querybuilder/family_condition.go
Normal 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
|
||||
}
|
||||
@@ -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] {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 != "" {
|
||||
|
||||
@@ -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"},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
89
pkg/statementbuilder/logsstatementbuilder/family_test.go
Normal file
89
pkg/statementbuilder/logsstatementbuilder/family_test.go
Normal 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)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
},
|
||||
)
|
||||
|
||||
@@ -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)
|
||||
|
||||
102
pkg/statementbuilder/metricsstatementbuilder/family_test.go
Normal file
102
pkg/statementbuilder/metricsstatementbuilder/family_test.go
Normal 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)
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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),
|
||||
)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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 != "" {
|
||||
|
||||
43
pkg/telemetrymetadata/family_values_test.go
Normal file
43
pkg/telemetrymetadata/family_values_test.go
Normal 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")
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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 != "" {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
|
||||
@@ -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]
|
||||
|
||||
Reference in New Issue
Block a user