Compare commits

...

2 Commits

Author SHA1 Message Date
srikanthccv
17b73fa357 proto: represent resolution output as LogicalField instead of key annotations
A LogicalField is one queryable field: the requested spelling plus the
physical member keys that store it, ordered current-first. The resolver
returns []*LogicalField — the slice expresses ambiguity (union across),
the group expresses a semantic-convention family (merge within). Members
alias the metadata map entries untouched; every member carries its own
physical facts, so no key ever needs sibling bookkeeping.

What this deletes:
- TelemetryFieldKey.SemconvMembers and SemconvMaterializedColumns — the
  group, previously flattened onto a representative key
- the deep-copy machinery in resolution (members are never mutated)
- resourcefilter's memberKey() fabrication and per-mapper member
  derivation with its static fallbacks
- the family branches inside resolveColumnExprs — FieldFor is a strict
  per-physical-key primitive again, byte-identical to main

What replaces them:
- MatchingLogicalFields / ResolveLogicalFields in querybuilder, with
  members sorted by family rank (fixes the context-prefixed-pass
  precedence inversion, pinned by a test)
- FieldForLogical / ExistsForLogical on the traces and resourcefilter
  mappers: the family expression is a composition of the members' own
  FieldFor outputs, so materialization and evolution state ride the
  member keys with no extra plumbing
- an upgrade pass in ColumnExpressionFor that swaps legacy-flow
  candidates for their family without changing candidate order or any
  non-family behavior
- non-family signals adapt with a two-line SingleKeys() flatten; their
  logical fields are single-member by construction

Generated SQL is unchanged for every case the phase-1 tests pin: the
keyless tails, member-OR guards and hints, pruning to present members,
and per-member materialized columns all carry over, asserted by the
rewritten unit tests.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-09 23:49:12 +05:30
srikanthccv
0c3223b23a feat: resolve semantic convention names in trace queries 2026-08-08 17:20:33 +05:30
27 changed files with 1481 additions and 160 deletions

View File

@@ -40,7 +40,8 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
logicalFields, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
keys := querybuilder.SingleKeys(logicalFields)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -0,0 +1,17 @@
package querybuilder
import "strings"
// ClickHouseStringLiteral quotes a value for a ClickHouse string literal.
func ClickHouseStringLiteral(value string) string {
escaped := strings.ReplaceAll(value, `\`, `\\`)
escaped = strings.ReplaceAll(escaped, `'`, `\'`)
return "'" + escaped + "'"
}
// ClickHouseIdentifier quotes a value for a ClickHouse identifier.
func ClickHouseIdentifier(value string) string {
escaped := strings.ReplaceAll(value, `\`, `\\`)
escaped = strings.ReplaceAll(escaped, "`", "\\`")
return "`" + escaped + "`"
}

View File

@@ -0,0 +1,17 @@
package querybuilder
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestClickHouseQuoting(t *testing.T) {
t.Run("string literal", func(t *testing.T) {
assert.Equal(t, `'name\'\\); SELECT 1 --'`, ClickHouseStringLiteral(`name'\); SELECT 1 --`))
})
t.Run("identifier", func(t *testing.T) {
assert.Equal(t, "`name\\`\\\\); SELECT 1 --`", ClickHouseIdentifier("name`\\); SELECT 1 --"))
})
}

View File

@@ -43,7 +43,7 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
if len(evolutionsEntries) > 0 && evolutionsEntries[0] != nil {
columnName = evolutionsEntries[0].ColumnName
}
rawPath := fmt.Sprintf("%s.`%s`", columnName, key.Name)
rawPath := fmt.Sprintf("%s.%s", columnName, ClickHouseIdentifier(key.Name))
if exists {
return rawPath + " IS NOT NULL", nil
}
@@ -88,7 +88,7 @@ func ExistsExpression(columns []*schema.Column, key *telemetrytypes.TelemetryFie
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumBool, schema.ColumnTypeEnumFloat64:
leftOperand := fmt.Sprintf("mapContains(%s, '%s')", column.Name, key.Name)
leftOperand := fmt.Sprintf("mapContains(%s, %s)", column.Name, ClickHouseStringLiteral(key.Name))
if key.Materialized {
leftOperand = telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key)
}

View File

@@ -0,0 +1,41 @@
package querybuilder
import (
"testing"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// ExistsExpression is a per-physical-key primitive: family composition happens
// at the logical-field layer, so this only ever sees one spelling.
func TestExistsExpressionIsPerKey(t *testing.T) {
columns := []*schema.Column{{
Name: "attributes_string",
Type: schema.MapColumnType{
KeyType: schema.LowCardinalityColumnType{ElementType: schema.ColumnTypeString},
ValueType: schema.ColumnTypeString,
},
}}
plain := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := ExistsExpression(columns, plain, 0, 0, "unused", true)
require.NoError(t, err)
assert.Equal(t, "mapContains(attributes_string, 'deployment.environment')", expression)
materialized := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
}
expression, err = ExistsExpression(columns, materialized, 0, 0, "unused", false)
require.NoError(t, err)
assert.Equal(t, "NOT `attribute_string_deployment$$environment_exists`", expression)
}

View File

@@ -21,24 +21,25 @@ const (
hasTokenFunctionDocURL = "https://signoz.io/docs/userguide/functions-reference/#hastoken-function"
)
// ResolveKeys picks which matching field keys a filter term builds conditions for.
// With 0 or 1 match it returns the input unchanged and no warning. When a name is
// ambiguous it returns a warning; a resource+attribute mix defaults to the resource
// keys (the common intent), noted in the warning.
func ResolveKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeysForName []*telemetrytypes.TelemetryFieldKey) ([]*telemetrytypes.TelemetryFieldKey, string) {
if len(fieldKeysForName) <= 1 {
return fieldKeysForName, ""
// ResolveLogicalFields picks which logical fields a filter term builds conditions
// for. With 0 or 1 field it returns the input unchanged and no warning. When a
// name is ambiguous (several logical fields — a family is one field and never
// ambiguous with itself) it returns a warning; a resource+attribute mix defaults
// to the resource fields (the common intent), noted in the warning.
func ResolveLogicalFields(field *telemetrytypes.TelemetryFieldKey, logicalFields []*telemetrytypes.LogicalField) ([]*telemetrytypes.LogicalField, string) {
if len(logicalFields) <= 1 {
return logicalFields, ""
}
warning := fmt.Sprintf(
"Key `%s` is ambiguous, found %d different combinations of field context / data type: %v.",
field.Name,
len(fieldKeysForName),
fieldKeysForName,
len(logicalFields),
logicalFields,
)
hasResource, hasAttribute := false, false
for _, item := range fieldKeysForName {
for _, item := range logicalFields {
switch item.FieldContext {
case telemetrytypes.FieldContextResource:
hasResource = true
@@ -49,18 +50,28 @@ func ResolveKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeysForName []*te
// when there is both resource and attribute context, default to resource only
if hasResource && hasAttribute {
filteredKeys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fieldKeysForName))
for _, item := range fieldKeysForName {
filtered := make([]*telemetrytypes.LogicalField, 0, len(logicalFields))
for _, item := range logicalFields {
if item.FieldContext == telemetrytypes.FieldContextResource {
filteredKeys = append(filteredKeys, item)
filtered = append(filtered, item)
}
}
fieldKeysForName = filteredKeys
logicalFields = filtered
warning += " " + "Using `resource` context by default. To query attributes explicitly, " +
fmt.Sprintf("use the fully qualified name (e.g., 'attribute.%s')", field.Name)
}
return fieldKeysForName, warning
return logicalFields, warning
}
// WrapAsLogicalFields wraps physical keys (candidate or synthesized) as
// single-member logical fields addressed by the requested spelling.
func WrapAsLogicalFields(requestedName string, keys []*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
fields := make([]*telemetrytypes.LogicalField, 0, len(keys))
for _, key := range keys {
fields = append(fields, telemetrytypes.SingleLogicalField(requestedName, key))
}
return fields
}
// NewKeyNotFoundError builds the error a condition builder returns when a filter term
@@ -175,3 +186,15 @@ func NewFunctionUnsupportedError(operator qbtypes.FilterOperator) error {
return nil
}
}
// SingleKeys flattens single-member logical fields to their member keys. It is
// the adapter for signals whose fields are never families (everything except
// traces today); a family in the input would be silently narrowed, so callers
// must be gated signals.
func SingleKeys(fields []*telemetrytypes.LogicalField) []*telemetrytypes.TelemetryFieldKey {
keys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(fields))
for _, logical := range fields {
keys = append(keys, logical.Single())
}
return keys
}

View File

@@ -0,0 +1,53 @@
package querybuilder_test
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/telemetryschema/tracestelemetryschema"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// A promoted historical member keeps its materialized column inside the family
// expression: the member key carries its own Materialized state, so the
// logical-field merge needs no sibling bookkeeping.
func TestTraceFamilyUsesMaterializedHistoricalMember(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
historical := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextAttribute,
telemetrytypes.FieldDataTypeString,
)
matches := querybuilder.MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
historical.Name: {historical},
})
require.Len(t, matches, 1, "family metadata should resolve to one logical field")
expression, err := tracestelemetryschema.NewFieldMapper().FieldForLogical(context.Background(), valuer.UUID{}, 0, 0, matches[0])
require.NoError(t, err, "resolved trace family should map to a value expression")
assert.Equal(
t,
"COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(`attribute_string_deployment$$environment`, ''), '')",
expression,
"family expression should retain the promoted historical member",
)
}

View File

@@ -10,6 +10,7 @@ import (
"github.com/SigNoz/signoz/pkg/errors"
grammar "github.com/SigNoz/signoz/pkg/parser/filterquery/grammar"
"github.com/SigNoz/signoz/pkg/semconv"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -360,7 +361,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
return ErrorConditionLiteral
}
}
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(searchText))
if !ok {
return ErrorConditionLiteral
}
@@ -379,7 +380,7 @@ func (v *filterExpressionVisitor) VisitPrimary(ctx *grammar.PrimaryContext) any
// VisitComparison handles all comparison operators.
func (v *filterExpressionVisitor) VisitComparison(ctx *grammar.ComparisonContext) any {
key := v.Visit(ctx.Key()).(*telemetrytypes.TelemetryFieldKey)
matching := MatchingFieldKeys(key, v.fieldKeys)
matching := MatchingLogicalFields(key, v.fieldKeys)
// Handle EXISTS specially
if ctx.EXISTS() != nil {
@@ -675,7 +676,7 @@ func (v *filterExpressionVisitor) VisitFullText(ctx *grammar.FullTextContext) an
v.errors = append(v.errors, "full text search is not supported")
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.TelemetryFieldKey{v.fullTextColumn}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
conds, ok := v.buildConditions(v.fullTextColumn, []*telemetrytypes.LogicalField{telemetrytypes.SingleLogicalField(v.fullTextColumn.Name, v.fullTextColumn)}, qbtypes.FilterOperatorRegexp, FormatFullTextSearch(text))
if !ok {
return ErrorConditionLiteral
}
@@ -730,7 +731,7 @@ func (v *filterExpressionVisitor) VisitFunctionCall(ctx *grammar.FunctionCallCon
return ErrorConditionLiteral
}
conds, ok := v.buildConditions(key, MatchingFieldKeys(key, v.fieldKeys), operator, value)
conds, ok := v.buildConditions(key, MatchingLogicalFields(key, v.fieldKeys), operator, value)
if !ok {
return ErrorConditionLiteral
}
@@ -922,7 +923,7 @@ func (v *filterExpressionVisitor) VisitKey(ctx *grammar.KeyContext) any {
// buildConditions invokes the condition builder for a filter term, folding its
// warnings/errors into visitor state; returns false if an error was recorded.
func (v *filterExpressionVisitor) buildConditions(key *telemetrytypes.TelemetryFieldKey, matching []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) ([]string, bool) {
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)
if err != nil {
_, _, _, _, errURL, _ := errors.Unwrapb(err)
@@ -979,30 +980,123 @@ func assignIfEmpty(s *string, value string) {
}
}
// MatchingFieldKeys returns the field keys from the map that match the given key,
// honoring any context/data type the user specified.
func MatchingFieldKeys(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
fieldKeysForName := []*telemetrytypes.TelemetryFieldKey{}
// familyMemberNames returns the physical spellings to look up for the
// referenced key: the semantic-convention family members (current-first) when
// 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(field *telemetrytypes.TelemetryFieldKey) []string {
if field.Signal != telemetrytypes.SignalUnspecified && field.Signal != telemetrytypes.SignalTraces {
return []string{field.Name}
}
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: field.FieldContext,
})
}
// match by name; keep items whose context and data type match (unspecified matches any)
for _, item := range fieldKeys[field.Name] {
if (field.FieldContext == telemetrytypes.FieldContextUnspecified || field.FieldContext == item.FieldContext) &&
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || field.FieldDataType == item.FieldDataType) {
fieldKeysForName = append(fieldKeysForName, item)
}
// 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 a single 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 therefore the
// length of the returned slice, and a family is never ambiguous with itself.
// Members alias the metadata map entries; nothing is copied or mutated.
func MatchingLogicalFields(field *telemetrytypes.TelemetryFieldKey, fieldKeys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
members := familyMemberNames(field)
memberRank := make(map[string]int, len(members))
for i, member := range members {
memberRank[member] = i
}
// A context may have been split off a name that legitimately contained it (e.g.
// `attribute.key`); also look up the context-prefixed name so both readings resolve.
if field.FieldContext != telemetrytypes.FieldContextUnspecified {
contextPrefixedFieldName := fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), field.Name)
for _, item := range fieldKeys[contextPrefixedFieldName] {
// Context already matched via the lookup key; only data type needs checking.
if field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || item.FieldDataType == field.FieldDataType {
fieldKeysForName = append(fieldKeysForName, item)
fields := make([]*telemetrytypes.LogicalField, 0)
indexByIdentity := make(map[string]int)
// rank of the family member each physical key matched under; the stored
// name of a context-prefixed match differs from the member name.
ranks := make(map[*telemetrytypes.TelemetryFieldKey]int)
appendMatches := func(lookupName string, memberName string, contextAlreadyMatched bool) {
for _, item := range fieldKeys[lookupName] {
if !contextAlreadyMatched && field.FieldContext != telemetrytypes.FieldContextUnspecified && field.FieldContext != item.FieldContext {
continue
}
if field.FieldDataType != telemetrytypes.FieldDataTypeUnspecified && field.FieldDataType != item.FieldDataType {
continue
}
// A wildcard lookup may have found a same-named field in a scope where
// this family does not apply. Keep exact names, but reject cross-member
// matches outside the generated family scope.
traceFamilyMatch := len(members) > 1 && item.Signal == telemetrytypes.SignalTraces
if memberName != field.Name {
if !traceFamilyMatch {
continue
}
itemSelector := telemetrytypes.FieldKeySelector{
Name: field.Name,
Signal: telemetrytypes.SignalTraces,
FieldContext: item.FieldContext,
}
if !slices.Contains(semconv.Members(semconv.KindAttribute, itemSelector), memberName) {
continue
}
}
if !traceFamilyMatch {
fields = append(fields, telemetrytypes.SingleLogicalField(field.Name, item))
continue
}
identity := item.Signal.StringValue() + ";" + item.FieldContext.StringValue() + ";" + item.FieldDataType.StringValue()
index, found := indexByIdentity[identity]
if !found {
index = len(fields)
indexByIdentity[identity] = index
fields = append(fields, &telemetrytypes.LogicalField{
Name: field.Name,
Signal: item.Signal,
FieldContext: item.FieldContext,
FieldDataType: item.FieldDataType,
})
}
logical := fields[index]
duplicate := false
for _, existing := range logical.Members {
if existing.Name == item.Name {
duplicate = true
break
}
}
if !duplicate {
ranks[item] = memberRank[memberName]
logical.Members = append(logical.Members, item)
}
}
}
return fieldKeysForName
for _, member := range members {
appendMatches(member, member, false)
}
// A context may have been split off a name that legitimately contained it
// (e.g. `attribute.key`); preserve that historical alternate reading for
// every family member.
if field.FieldContext != telemetrytypes.FieldContextUnspecified {
for _, member := range members {
appendMatches(fmt.Sprintf("%s.%s", field.FieldContext.StringValue(), member), member, true)
}
}
// Precedence is a property of the family, not of arrival order: members
// sort current-first no matter which lookup pass found them.
for _, logical := range fields {
slices.SortStableFunc(logical.Members, func(a, b *telemetrytypes.TelemetryFieldKey) int {
return ranks[a] - ranks[b]
})
}
return fields
}

View File

@@ -14,6 +14,7 @@ import (
"github.com/antlr4-go/antlr/v4"
sqlbuilder "github.com/huandu/go-sqlbuilder"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// TestPrepareWhereClause_EmptyVariableList ensures PrepareWhereClause errors when a variable has an empty list value.
@@ -589,8 +590,8 @@ func TestVisitKey(t *testing.T) {
// VisitKey only parses; the condition builder matches, resolves ambiguity
// and decides not-found handling. Replay that here against the generic
// builder behavior (error unless the key is ignored).
matching := MatchingFieldKeys(key, tt.fieldKeys)
keys, warning := ResolveKeys(key, matching)
matching := MatchingLogicalFields(key, tt.fieldKeys)
keys, warning := ResolveLogicalFields(key, matching)
var gotErrors []string
var gotMainErrURL, gotMainWrnURL string
@@ -612,15 +613,19 @@ func TestVisitKey(t *testing.T) {
t.Errorf("expected %d keys, got %d", len(tt.expectedKeys), len(keys))
}
// Check each expected key matches name, field context, and data type
// Check each expected key matches a member's stored name plus the
// logical field's context and data type (the logical Name is the
// requested spelling, members keep the stored spellings).
for _, expectedKey := range tt.expectedKeys {
found := false
for _, key := range keys {
if key.Name == expectedKey.Name &&
key.FieldContext == expectedKey.FieldContext &&
key.FieldDataType == expectedKey.FieldDataType {
found = true
break
for _, logical := range keys {
for _, member := range logical.Members {
if member.Name == expectedKey.Name &&
logical.FieldContext == expectedKey.FieldContext &&
logical.FieldDataType == expectedKey.FieldDataType {
found = true
break
}
}
}
if !found {
@@ -685,6 +690,156 @@ func TestVisitKey(t *testing.T) {
}
}
func memberNames(logical *telemetrytypes.LogicalField) []string {
names := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
names = append(names, member.Name)
}
return names
}
func TestMatchingLogicalFieldsResolvesCurrentTraceNameFromOldMetadata(t *testing.T) {
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Description: "old metadata",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
"deployment.environment.name",
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{old.Name: {old}})
require.Len(t, matches, 1, "a family is one logical field, not an ambiguity")
assert.Equal(t, "deployment.environment.name", matches[0].Name, "the requested spelling is the response identity")
assert.Equal(t, []string{"deployment.environment"}, memberNames(matches[0]))
assert.Same(t, old, matches[0].Members[0], "members alias metadata entries; nothing is copied")
}
func TestMatchingLogicalFieldsGroupsFamilyMembersCurrentFirst(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
old.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "a family is one logical field, not an ambiguity")
assert.Equal(t, old.Name, matches[0].Name, "the requested spelling is the response identity")
assert.True(t, matches[0].IsFamily())
assert.Equal(t, []string{current.Name, old.Name}, memberNames(matches[0]), "members order current-first")
}
// Precedence is a property of the family, not of arrival order: a member that
// only exists under its context-prefixed stored spelling still sorts by its
// family rank.
func TestMatchingLogicalFieldsOrdersMembersByFamilyRank(t *testing.T) {
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
prefixedCurrent := &telemetrytypes.TelemetryFieldKey{
Name: "resource.deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
"deployment.environment.name",
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
old.Name: {old},
prefixedCurrent.Name: {prefixedCurrent},
})
require.Len(t, matches, 1)
assert.Equal(t, []string{prefixedCurrent.Name, old.Name}, memberNames(matches[0]),
"the current-spelling member must coalesce before the old one")
}
func TestMatchingLogicalFieldsKeepsLogSemconvNamesLiteral(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "log lookup must keep the requested spelling literal")
assert.False(t, matches[0].IsFamily())
assert.Equal(t, current.Name, matches[0].Single().Name)
}
func TestMatchingLogicalFieldsKeepsMetricSemconvNamesLiteral(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalMetrics,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
requested := telemetrytypes.NewTelemetryFieldKey(
current.Name,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
matches := MatchingLogicalFields(requested, map[string][]*telemetrytypes.TelemetryFieldKey{
current.Name: {current},
old.Name: {old},
})
require.Len(t, matches, 1, "metric lookup must keep the requested spelling literal")
assert.False(t, matches[0].IsFamily())
assert.Equal(t, current.Name, matches[0].Single().Name)
}
// ---------------------------------------------------------------------------
// TestVisitComparison
// ---------------------------------------------------------------------------
@@ -766,7 +921,7 @@ func (b *resourceConditionBuilder) ConditionFor(
return nil, nil, nil
}
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
keys, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -774,11 +929,11 @@ func (b *resourceConditionBuilder) ConditionFor(
var conds []string
for _, k := range keys {
// only resource keys contribute; others (and unknown keys) are ignored
// only resource fields contribute; others (and unknown keys) are ignored
if k.FieldContext != telemetrytypes.FieldContextResource {
continue
}
conds = append(conds, fmt.Sprintf("%s_cond", k.Name))
conds = append(conds, fmt.Sprintf("%s_cond", k.Single().Name))
}
return conds, warnings, nil
}
@@ -808,7 +963,7 @@ func (b *conditionBuilder) ConditionFor(
return []string{fmt.Sprintf("%s_cond", key.Name)}, nil, nil
}
keys, warning := ResolveKeys(key, MatchingFieldKeys(key, fieldKeys))
keys, warning := ResolveLogicalFields(key, MatchingLogicalFields(key, fieldKeys))
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -820,7 +975,7 @@ func (b *conditionBuilder) ConditionFor(
// A resource sub-query already covers the term; drop resource keys from the main query.
if options.SkipResourceFilter {
filtered := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys))
filtered := make([]*telemetrytypes.LogicalField, 0, len(keys))
for _, k := range keys {
if k.FieldContext != telemetrytypes.FieldContextResource {
filtered = append(filtered, k)

View File

@@ -13,12 +13,14 @@ import (
)
type defaultConditionBuilder struct {
fm qbtypes.FieldMapper
// The builder composes family expressions, so it needs this package's
// mapper, not the narrower qbtypes.FieldMapper.
fm *defaultFieldMapper
}
var _ qbtypes.ConditionBuilder = (*defaultConditionBuilder)(nil)
func NewConditionBuilder(fm qbtypes.FieldMapper) *defaultConditionBuilder {
func NewConditionBuilder(fm *defaultFieldMapper) *defaultConditionBuilder {
return &defaultConditionBuilder{fm: fm}
}
@@ -44,6 +46,66 @@ func keyIndexFilter(key *telemetrytypes.TelemetryFieldKey) any {
return fmt.Sprintf(`%%%s%%`, key.Name)
}
func keyIndexCondition(sb *sqlbuilder.SelectBuilder, column string, members []*telemetrytypes.TelemetryFieldKey) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
conditions = append(conditions, sb.Like(column, keyIndexFilter(member)))
}
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
func valueIndexCondition(
sb *sqlbuilder.SelectBuilder,
column string,
members []*telemetrytypes.TelemetryFieldKey,
op qbtypes.FilterOperator,
value any,
caseInsensitive bool,
) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
patterns := valueForIndexFilter(op, member, value)
switch values := patterns.(type) {
case []string:
for _, pattern := range values {
conditions = append(conditions, sb.Like(column, pattern))
}
default:
if caseInsensitive {
conditions = append(conditions, sb.ILike(column, values))
} else {
conditions = append(conditions, sb.Like(column, values))
}
}
}
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
func memberPresenceCondition(sb *sqlbuilder.SelectBuilder, column string, members []*telemetrytypes.TelemetryFieldKey, exists bool) string {
conditions := make([]string, 0, len(members))
for _, member := range members {
field := fmt.Sprintf("simpleJSONHas(%s, %s)", column, querybuilder.ClickHouseStringLiteral(member.Name))
if exists {
conditions = append(conditions, sb.E(field, true))
} else {
conditions = append(conditions, sb.NE(field, true))
}
}
if exists {
if len(conditions) == 1 {
return conditions[0]
}
return sb.Or(conditions...)
}
return sb.And(conditions...)
}
// SkipResourceFilter is not applicable here: the fingerprint table only stores resource attributes.
func (b *defaultConditionBuilder) ConditionFor(
ctx context.Context,
@@ -57,7 +119,7 @@ func (b *defaultConditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(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).
@@ -65,21 +127,21 @@ func (b *defaultConditionBuilder) ConditionFor(
return nil, nil, nil
}
keys, warning := querybuilder.ResolveKeys(key, matches)
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
}
conds := make([]string, 0, len(keys))
for _, k := range keys {
// the resource fingerprint table only stores resource attributes; keys from
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
// the resource fingerprint table only stores resource attributes; fields from
// any other context contribute no condition and are omitted. An empty result
// (including an unknown key) lets the caller skip this filter entirely.
if k.FieldContext != telemetrytypes.FieldContextResource {
if logical.FieldContext != telemetrytypes.FieldContextResource {
continue
}
cond, err := b.conditionForKey(ctx, startNs, endNs, k, op, value, sb)
cond, err := b.conditionForLogicalField(ctx, startNs, endNs, logical, op, value, sb)
if err != nil {
return nil, nil, err
}
@@ -88,11 +150,11 @@ func (b *defaultConditionBuilder) ConditionFor(
return conds, warnings, nil
}
func (b *defaultConditionBuilder) conditionForKey(
func (b *defaultConditionBuilder) conditionForLogicalField(
ctx context.Context,
startNs uint64,
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
logical *telemetrytypes.LogicalField,
op qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -102,7 +164,7 @@ func (b *defaultConditionBuilder) conditionForKey(
// as we store resource values as string
formattedValue := querybuilder.FormatValueForContains(value)
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, key)
columns, err := b.fm.ColumnFor(ctx, valuer.UUID{}, startNs, endNs, logical.Single())
if err != nil {
return "", err
}
@@ -115,10 +177,12 @@ func (b *defaultConditionBuilder) conditionForKey(
// as we have not changed the resource column in the resource fingerprint table.
column := columns[0]
keyIdxFilter := sb.Like(column.Name, keyIndexFilter(key))
valueForIndexFilter := valueForIndexFilter(op, key, value)
members := logical.Members
isFamily := logical.IsFamily()
keyIdxFilter := keyIndexCondition(sb, column.Name, members)
singleValueIndexFilter := valueForIndexFilter(op, members[0], value)
fieldName, err := b.fm.FieldFor(ctx, valuer.UUID{}, startNs, endNs, key)
fieldName, err := b.fm.FieldForLogical(ctx, valuer.UUID{}, startNs, endNs, logical)
if err != nil {
return "", err
}
@@ -128,12 +192,15 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.E(fieldName, formattedValue),
keyIdxFilter,
sb.Like(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, members, op, value, false),
), nil
case qbtypes.FilterOperatorNotEqual:
if isFamily {
return sb.NE(fieldName, formattedValue), nil
}
return sb.And(
sb.NE(fieldName, formattedValue),
sb.NotLike(column.Name, valueForIndexFilter),
sb.NotLike(column.Name, singleValueIndexFilter),
), nil
case qbtypes.FilterOperatorGreaterThan:
return sb.And(sb.GT(fieldName, formattedValue), keyIdxFilter), nil
@@ -148,7 +215,7 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.ILike(fieldName, formattedValue),
keyIdxFilter,
sb.ILike(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, members, op, value, true),
), nil
case qbtypes.FilterOperatorNotLike, qbtypes.FilterOperatorNotILike:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else
@@ -185,13 +252,11 @@ func (b *defaultConditionBuilder) conditionForKey(
inConditions = append(inConditions, sb.E(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.Or(inConditions...)
valConditions := make([]string, 0, len(values))
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.Like(column.Name, v))
}
}
mainCondition = sb.And(mainCondition, keyIdxFilter, sb.Or(valConditions...))
mainCondition = sb.And(
mainCondition,
keyIdxFilter,
valueIndexCondition(sb, column.Name, members, op, value, false),
)
return mainCondition, nil
case qbtypes.FilterOperatorNotIn:
@@ -204,8 +269,11 @@ func (b *defaultConditionBuilder) conditionForKey(
notInConditions = append(notInConditions, sb.NE(fieldName, querybuilder.FormatValueForContains(v)))
}
mainCondition := sb.And(notInConditions...)
if isFamily {
return mainCondition, nil
}
valConditions := make([]string, 0, len(values))
if valuesForIndexFilter, ok := valueForIndexFilter.([]string); ok {
if valuesForIndexFilter, ok := singleValueIndexFilter.([]string); ok {
for _, v := range valuesForIndexFilter {
valConditions = append(valConditions, sb.NotLike(column.Name, v))
}
@@ -215,13 +283,11 @@ func (b *defaultConditionBuilder) conditionForKey(
case qbtypes.FilterOperatorExists:
return sb.And(
sb.E(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
memberPresenceCondition(sb, column.Name, members, true),
keyIdxFilter,
), nil
case qbtypes.FilterOperatorNotExists:
return sb.And(
sb.NE(fmt.Sprintf("simpleJSONHas(%s, '%s')", column.Name, key.Name), true),
), nil
return memberPresenceCondition(sb, column.Name, members, false), nil
case qbtypes.FilterOperatorRegexp:
return sb.And(
@@ -237,7 +303,7 @@ func (b *defaultConditionBuilder) conditionForKey(
return sb.And(
sb.ILike(fieldName, fmt.Sprintf(`%%%s%%`, formattedValue)),
keyIdxFilter,
sb.ILike(column.Name, valueForIndexFilter),
valueIndexCondition(sb, column.Name, members, op, value, true),
), nil
case qbtypes.FilterOperatorNotContains:
// no index filter: as cannot apply `not contains x%y` as y can be somewhere else

View File

@@ -220,3 +220,129 @@ func TestConditionBuilder(t *testing.T) {
})
}
}
// The family tests drive resolution through the metadata map, exactly as
// production does: two plain member keys in the map, one requested spelling.
// The keys carry no family bookkeeping — grouping is the resolver's job.
func familyConditionSQL(t *testing.T, requestedName string, memberNames []string, op qbtypes.FilterOperator, value any) (string, []any) {
t.Helper()
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{}
for _, name := range memberNames {
fieldKeys[name] = []*telemetrytypes.TelemetryFieldKey{{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}}
}
requested := telemetrytypes.NewTelemetryFieldKey(
requestedName,
telemetrytypes.FieldContextResource,
telemetrytypes.FieldDataTypeString,
)
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, requested,
fieldKeys,
qbtypes.ConditionBuilderOptions{}, op, value, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
return sb.BuildWithFlavor(sqlbuilder.ClickHouse)
}
var deploymentFamilyMembers = []string{"deployment.environment.name", "deployment.environment"}
func TestFamilyPositiveFilterExcludesKeylessRows(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorEqual, "production")
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') = ? AND (labels LIKE ? OR labels LIKE ?) AND (labels LIKE ? OR labels LIKE ?)")
assert.Equal(t, []any{
"production",
"%deployment.environment.name%",
"%deployment.environment%",
`%deployment.environment.name":"production%`,
`%deployment.environment":"production%`,
}, args)
}
func TestFamilyNotEqualIncludesKeylessRows(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotEqual, "staging")
assert.Contains(t, sql, "COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ?")
assert.Equal(t, []any{"staging"}, args)
}
func TestFamilyNotInIncludesKeylessRows(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotIn, []any{"staging", "dev"})
assert.Contains(t, sql, "(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ? AND COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '') <> ?)")
assert.Equal(t, []any{"staging", "dev"}, args)
}
func TestFamilyNotLikeIncludesKeylessRows(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotLike, "%stag%")
assert.Contains(t, sql, "LOWER(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')) NOT LIKE LOWER(?)")
assert.Equal(t, []any{"%stag%"}, args)
}
func TestFamilyNotContainsIncludesKeylessRows(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotContains, "stag")
assert.Contains(t, sql, "LOWER(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), '')) NOT LIKE LOWER(?)")
assert.Equal(t, []any{"%stag%"}, args)
}
func TestFamilyNotRegexpIncludesKeylessRows(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotRegexp, "stag.*")
assert.Contains(t, sql, "NOT match(COALESCE(NULLIF(simpleJSONExtractString(labels, 'deployment.environment.name'), ''), NULLIF(simpleJSONExtractString(labels, 'deployment.environment'), ''), ''), ?)")
assert.Equal(t, []any{"stag.*"}, args)
}
func TestFamilyExistsChecksEveryMember(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorExists, nil)
assert.Contains(t, sql, "(simpleJSONHas(labels, 'deployment.environment.name') = ? OR simpleJSONHas(labels, 'deployment.environment') = ?) AND (labels LIKE ? OR labels LIKE ?)")
assert.Equal(t, []any{true, true, "%deployment.environment.name%", "%deployment.environment%"}, args)
}
func TestFamilyNotExistsChecksEveryMember(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment.name", deploymentFamilyMembers, qbtypes.FilterOperatorNotExists, nil)
assert.Contains(t, sql, "simpleJSONHas(labels, 'deployment.environment.name') <> ? AND simpleJSONHas(labels, 'deployment.environment') <> ?")
assert.Equal(t, []any{true, true}, args)
}
// The old-name request with only the current spelling in metadata prunes to a
// single member: plain single-key SQL, no coalesce.
func TestFamilyPrunesToPresentMembers(t *testing.T) {
sql, args := familyConditionSQL(t, "deployment.environment", []string{"deployment.environment.name"}, qbtypes.FilterOperatorEqual, "production")
assert.Contains(t, sql, "simpleJSONExtractString(labels, 'deployment.environment.name') = ? AND labels LIKE ? AND labels LIKE ?")
assert.Equal(t, []any{"production", "%deployment.environment.name%", `%deployment.environment.name":"production%`}, args)
}
func TestLogSemconvNameStaysLiteral(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalLogs,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
sb := sqlbuilder.NewSelectBuilder()
conditions, _, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, key,
map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}},
qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, "production", sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "simpleJSONExtractString(labels, 'deployment.environment.name') = ? AND labels LIKE ? AND labels LIKE ?")
assert.NotContains(t, sql, "deployment.environment')")
assert.Equal(t, []any{"production", "%deployment.environment.name%", `%deployment.environment.name":"production%`}, args)
}

View File

@@ -3,8 +3,10 @@ package resourcefilter
import (
"context"
"fmt"
"strings"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/querybuilder"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -32,6 +34,31 @@ func NewFieldMapper() *defaultFieldMapper {
return &defaultFieldMapper{}
}
// FieldForLogical returns the value expression for a resolved logical field:
// the member's own expression for a single-member field, and a current-first
// merge for a family. Resource label values are strings, so the merge is a
// coalesce with a trailing '' that keeps single-key semantics for rows
// without any member (see AddDefaultExistsFilter).
func (m *defaultFieldMapper) FieldForLogical(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,
logical *telemetrytypes.LogicalField,
) (string, error) {
if !logical.IsFamily() {
return m.FieldFor(ctx, orgID, tsStart, tsEnd, logical.Single())
}
values := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
expr, err := m.FieldFor(ctx, orgID, tsStart, tsEnd, member)
if err != nil {
return "", err
}
values = append(values, fmt.Sprintf("NULLIF(%s, '')", expr))
}
return "COALESCE(" + strings.Join(values, ", ") + ", '')", nil
}
func (m *defaultFieldMapper) getColumn(
_ context.Context,
_, _ uint64,
@@ -66,7 +93,7 @@ func (m *defaultFieldMapper) FieldFor(
return "", err
}
if key.FieldContext == telemetrytypes.FieldContextResource {
return fmt.Sprintf("simpleJSONExtractString(%s, '%s')", columns[0].Name, key.Name), nil
return fmt.Sprintf("simpleJSONExtractString(%s, %s)", columns[0].Name, querybuilder.ClickHouseStringLiteral(key.Name)), nil
}
return columns[0].Name, nil
}

View File

@@ -0,0 +1,24 @@
package resourcefilter
import (
"context"
"testing"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestFieldForQuotesRequestKeyName(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "name'\\); SELECT 1 --",
FieldContext: telemetrytypes.FieldContextResource,
}
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
assert.Equal(t, "simpleJSONExtractString(labels, "+querybuilder.ClickHouseStringLiteral(key.Name)+")", expression)
}

View File

@@ -39,7 +39,8 @@ func (c *conditionBuilder) ConditionFor(
}
// an unknown key simply yields no condition rather than an error.
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
logicalFields, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
keys := querybuilder.SingleKeys(logicalFields)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -14,6 +14,7 @@ import (
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/semconv"
"github.com/SigNoz/signoz/pkg/telemetryschema/audittelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema"
"github.com/SigNoz/signoz/pkg/telemetryschema/metertelemetryschema"
@@ -151,6 +152,14 @@ func (t *telemetryMetaStore) tracesTblStatementToFieldKeys(ctx context.Context)
return materialisedKeys, nil
}
func traceSemconvMembers(name string, fieldContext telemetrytypes.FieldContext) []string {
return semconv.Members(semconv.KindAttribute, telemetrytypes.FieldKeySelector{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: fieldContext,
})
}
// getTracesKeys returns the keys from the spans that match the field selection criteria.
func (t *telemetryMetaStore) getTracesKeys(ctx context.Context, fieldKeySelectors []*telemetrytypes.FieldKeySelector) ([]*telemetrytypes.TelemetryFieldKey, bool, error) {
ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{
@@ -1320,6 +1329,17 @@ func (t *telemetryMetaStore) GetKeysMulti(ctx context.Context, orgID valuer.UUID
if err != nil {
return nil, false, err
}
// GetKeys backs key suggestions and remains literal. The internal multi-key
// lookup expands only trace selectors so query builders see stored family members.
expandedTraceSelectors := make([]*telemetrytypes.FieldKeySelector, 0, len(tracesSelectors))
for _, selector := range tracesSelectors {
for _, member := range traceSemconvMembers(selector.Name, selector.FieldContext) {
memberSelector := selector.Copy()
memberSelector.Name = member
expandedTraceSelectors = append(expandedTraceSelectors, memberSelector)
}
}
tracesSelectors = expandedTraceSelectors
tracesKeys, tracesComplete, err := t.getTracesKeys(ctx, tracesSelectors)
if err != nil {
return nil, false, err
@@ -1542,7 +1562,16 @@ func (t *telemetryMetaStore) getSpanFieldValues(ctx context.Context, fieldValueS
sb := sqlbuilder.Select("DISTINCT string_value, number_value").From(t.tracesDBName + "." + t.tracesFieldsTblName)
if fieldValueSelector.Name != "" {
sb.Where(sb.E("tag_key", fieldValueSelector.Name))
members := traceSemconvMembers(fieldValueSelector.Name, fieldValueSelector.FieldContext)
if len(members) == 1 {
sb.Where(sb.E("tag_key", members[0]))
} else {
memberValues := make([]any, 0, len(members))
for _, member := range members {
memberValues = append(memberValues, member)
}
sb.Where(sb.In("tag_key", memberValues...))
}
}
// now look at the field context

View File

@@ -12,6 +12,7 @@ import (
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -83,3 +84,38 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
t.Errorf("there were unfulfilled expectations: %s", err)
}
}
func TestGetAllValuesReturnsValuesFromEveryTraceSemconvFamilyMember(t *testing.T) {
mockTelemetryStore := telemetrystoretest.New(telemetrystore.Config{}, &regexMatcher{})
mock := mockTelemetryStore.Mock()
metadata := NewTelemetryMetaStore(
instrumentationtest.New().ToProviderSettings(),
mockTelemetryStore,
flaggertest.New(t),
)
mock.ExpectQuery(`SELECT DISTINCT string_value, number_value FROM signoz_traces\.distributed_tag_attributes_v2 WHERE tag_key IN \(\?, \?\) AND tag_type = \? AND tag_data_type = \? LIMIT \?`).
WithArgs("deployment.environment.name", "deployment.environment", "resource", "string", 51).
WillReturnRows(cmock.NewRows([]cmock.ColumnType{
{Name: "string_value", Type: "String"},
{Name: "number_value", Type: "Float64"},
}, [][]any{
{"production", float64(0)},
{"staging", float64(0)},
{"production", float64(0)},
}))
values, complete, err := metadata.GetAllValues(context.Background(), valuer.UUID{}, &telemetrytypes.FieldValueSelector{
FieldKeySelector: &telemetrytypes.FieldKeySelector{
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
Name: "deployment.environment",
},
})
require.NoError(t, err)
assert.True(t, complete)
assert.Equal(t, []string{"production", "staging"}, values.StringValues)
assert.NoError(t, mock.ExpectationsWereMet(), "all expected metadata queries should be executed")
}

View File

@@ -139,7 +139,8 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
keys, warning := querybuilder.ResolveKeys(key, querybuilder.MatchingFieldKeys(key, fieldKeys))
logicalFields, warning := querybuilder.ResolveLogicalFields(key, querybuilder.MatchingLogicalFields(key, fieldKeys))
keys := querybuilder.SingleKeys(logicalFields)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -452,7 +452,7 @@ func (c *conditionBuilder) ConditionFor(
value any,
sb *sqlbuilder.SelectBuilder,
) ([]string, []string, error) {
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
// search() resolves its own (optional) scope; handle it before key resolution.
@@ -460,7 +460,8 @@ func (c *conditionBuilder) ConditionFor(
return c.conditionForSearch(ctx, orgID, key, value, sb)
}
keys, warning := querybuilder.ResolveKeys(key, matches)
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
keys := querybuilder.SingleKeys(logicalFields)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)

View File

@@ -162,7 +162,7 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
keys := querybuilder.MatchingFieldKeys(key, fieldKeys)
keys := querybuilder.SingleKeys(querybuilder.MatchingLogicalFields(key, fieldKeys))
var warnings []string
if len(keys) == 0 {
if _, isColumn := timeSeriesV4Columns[key.Name]; isColumn {

View File

@@ -18,12 +18,14 @@ import (
)
type conditionBuilder struct {
fm qbtypes.FieldMapper
// The builder composes family expressions, so it needs this package's
// mapper, not the narrower qbtypes.FieldMapper.
fm *fieldMapper
}
var _ qbtypes.ConditionBuilder = (*conditionBuilder)(nil)
func NewConditionBuilder(fm qbtypes.FieldMapper) *conditionBuilder {
func NewConditionBuilder(fm *fieldMapper) *conditionBuilder {
return &conditionBuilder{fm: fm}
}
@@ -32,7 +34,7 @@ func (c *conditionBuilder) conditionFor(
orgID valuer.UUID,
startNs uint64,
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
logical *telemetrytypes.LogicalField,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
@@ -42,13 +44,13 @@ func (c *conditionBuilder) conditionFor(
value = querybuilder.FormatValueForContains(value)
}
fieldExpression, err := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
fieldExpression, err := c.fm.FieldForLogical(ctx, orgID, startNs, endNs, logical)
if err != nil {
return "", err
}
// TODO(srikanthccv): maybe extend this to every possible attribute
if key.Name == "duration_nano" || key.Name == "durationNano" { // QoL improvement
if logical.Name == "duration_nano" || logical.Name == "durationNano" { // QoL improvement
switch v := value.(type) {
case string:
if duration, err := time.ParseDuration(v); err == nil {
@@ -65,7 +67,7 @@ func (c *conditionBuilder) conditionFor(
}
}
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(key, value, fieldExpression, operator)
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), value, fieldExpression, operator)
// regular operators
switch operator {
@@ -154,11 +156,7 @@ func (c *conditionBuilder) conditionFor(
// in the query builder, `exists` and `not exists` are used for
// key membership checks, so depending on the column type, the condition changes
case qbtypes.FilterOperatorExists, qbtypes.FilterOperatorNotExists:
columns, err := c.fm.ColumnFor(ctx, orgID, startNs, endNs, key)
if err != nil {
return "", err
}
pred, err := querybuilder.ExistsExpression(columns, key, startNs, endNs, fieldExpression, operator == qbtypes.FilterOperatorExists)
pred, err := c.fm.ExistsForLogical(ctx, orgID, startNs, endNs, logical, operator == qbtypes.FilterOperatorExists)
if err != nil {
return "", err
}
@@ -210,10 +208,10 @@ func (c *conditionBuilder) ConditionFor(
return nil, nil, err
}
matches := querybuilder.MatchingFieldKeys(key, fieldKeys)
matches := querybuilder.MatchingLogicalFields(key, fieldKeys)
skipResourceFilter := options.SkipResourceFilter
keys, warning := querybuilder.ResolveKeys(key, matches)
logicalFields, warning := querybuilder.ResolveLogicalFields(key, matches)
var warnings []string
if warning != "" {
warnings = append(warnings, warning)
@@ -221,10 +219,10 @@ func (c *conditionBuilder) ConditionFor(
// A bare key that names a real column filters on the column too — first. When metadata
// only knows the name under other contexts, prepend the column and keep metadata matches
// only where their type is consistent with it (a corrupt entry can't degrade the column).
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(keys) > 0 {
if key.FieldContext == telemetrytypes.FieldContextUnspecified && len(logicalFields) > 0 {
hasColumn := false
for _, k := range keys {
if k.FieldContext == telemetrytypes.FieldContextSpan {
for _, logical := range logicalFields {
if logical.FieldContext == telemetrytypes.FieldContextSpan {
hasColumn = true
break
}
@@ -232,49 +230,49 @@ func (c *conditionBuilder) ConditionFor(
if !hasColumn {
probe := telemetrytypes.NewTelemetryFieldKey(key.Name, telemetrytypes.FieldContextSpan, key.FieldDataType)
if cols, colErr := c.fm.ColumnFor(ctx, orgID, startNs, endNs, probe); colErr == nil && len(cols) > 0 {
combined := make([]*telemetrytypes.TelemetryFieldKey, 0, len(keys)+1)
combined = append(combined, probe)
for _, k := range keys {
if columnMatchesDataType(cols[0], k.FieldDataType) {
combined = append(combined, k)
combined := make([]*telemetrytypes.LogicalField, 0, len(logicalFields)+1)
combined = append(combined, telemetrytypes.SingleLogicalField(key.Name, probe))
for _, logical := range logicalFields {
if columnMatchesDataType(cols[0], logical.FieldDataType) {
combined = append(combined, logical)
}
}
keys = combined
logicalFields = combined
}
}
}
synthesized := false
if len(keys) == 0 {
if len(logicalFields) == 0 {
// Not in metadata. CandidateKeys resolves it: fold contexts (span/trace) get the
// metadata map so it can honor a real column, correct to a stripped-name metadata
// match, or synthesize; strict contexts pass nil and keep their synthesize path.
keys = c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys))
if len(keys) == 0 {
logicalFields = querybuilder.WrapAsLogicalFields(key.Name, c.fm.CandidateKeys(ctx, orgID, key, value, candidateLookupKeys(key, fieldKeys)))
if len(logicalFields) == 0 {
return nil, warnings, querybuilder.NewKeyNotFoundError(key.Name)
}
synthesized = true
warnings = append(warnings, querybuilder.NewKeyNotFoundWarning(key.Name))
}
// When a resource sub-query already covers the term, drop resource keys from the main
// When a resource sub-query already covers the term, drop resource fields from the main
// query. Synthesized keys are exempt: the sub-query skips keys absent from metadata.
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 {
cond, err := c.conditionForKey(ctx, orgID, startNs, endNs, k, operator, value, sb)
conds := make([]string, 0, len(logicalFields))
for _, logical := range logicalFields {
cond, err := c.conditionForLogicalField(ctx, orgID, startNs, endNs, logical, operator, value, sb)
if err != nil {
return nil, nil, err
}
@@ -283,28 +281,28 @@ func (c *conditionBuilder) ConditionFor(
return conds, warnings, nil
}
func (c *conditionBuilder) conditionForKey(
func (c *conditionBuilder) conditionForLogicalField(
ctx context.Context,
orgID valuer.UUID,
startNs uint64,
endNs uint64,
key *telemetrytypes.TelemetryFieldKey,
logical *telemetrytypes.LogicalField,
operator qbtypes.FilterOperator,
value any,
sb *sqlbuilder.SelectBuilder,
) (string, error) {
if c.isSpanScopeField(key.Name) {
return c.buildSpanScopeCondition(key, operator, value, startNs)
if c.isSpanScopeField(logical.Name) {
return c.buildSpanScopeCondition(logical.Single(), operator, value, startNs)
}
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, operator, value, sb)
condition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, operator, value, sb)
if err != nil {
return "", err
}
if operator.AddDefaultExistsFilter() {
// skip adding exists filter for intrinsic fields
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, key)
field, _ := c.fm.FieldFor(ctx, orgID, startNs, endNs, logical.Single())
if slices.Contains(maps.Keys(IntrinsicFields), field) ||
slices.Contains(maps.Keys(IntrinsicFieldsDeprecated), field) ||
slices.Contains(maps.Keys(CalculatedFields), field) ||
@@ -312,7 +310,7 @@ func (c *conditionBuilder) conditionForKey(
return condition, nil
}
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, key, qbtypes.FilterOperatorExists, nil, sb)
existsCondition, err := c.conditionFor(ctx, orgID, startNs, endNs, logical, qbtypes.FilterOperatorExists, nil, sb)
if err != nil {
return "", err
}

View File

@@ -308,6 +308,93 @@ func TestConditionFor(t *testing.T) {
}
}
// The family tests drive resolution through the metadata map, exactly as
// production does: plain member keys in, the resolver groups them, and the
// builder composes one condition per logical field.
func traceFamilyConditionSQL(t *testing.T, requestedName string, members []*telemetrytypes.TelemetryFieldKey, op qbtypes.FilterOperator, value any) (string, []any, []string) {
t.Helper()
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{}
for _, member := range members {
fieldKeys[member.Name] = []*telemetrytypes.TelemetryFieldKey{member}
}
requested := telemetrytypes.NewTelemetryFieldKey(
requestedName,
telemetrytypes.FieldContextAttribute,
telemetrytypes.FieldDataTypeString,
)
sb := sqlbuilder.NewSelectBuilder()
conditions, warnings, err := NewConditionBuilder(NewFieldMapper()).ConditionFor(
context.Background(), valuer.UUID{}, 0, 0, requested,
fieldKeys,
qbtypes.ConditionBuilderOptions{}, op, value, sb,
)
require.NoError(t, err)
sb.Where(conditions...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return sql, args, warnings
}
func traceAttrMember(name string, materialized bool) *telemetrytypes.TelemetryFieldKey {
return &telemetrytypes.TelemetryFieldKey{
Name: name,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: materialized,
}
}
func TestConditionForSemconvFamilyPositiveFilterChecksPresence(t *testing.T) {
sql, args, warnings := traceFamilyConditionSQL(t,
"deployment.environment.name",
[]*telemetrytypes.TelemetryFieldKey{
traceAttrMember("deployment.environment.name", false),
traceAttrMember("deployment.environment", false),
},
qbtypes.FilterOperatorEqual, "production",
)
assert.Empty(t, warnings, "a family is one logical field, never an ambiguity warning")
assert.Contains(t, sql, "COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '') = ? AND (mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))")
assert.Equal(t, []any{"production"}, args)
}
func TestNewConditionBuilderTakesThisPackagesMapper(t *testing.T) {
// The builder composes family expressions, so it is deliberately coupled
// to this package's mapper rather than the narrower qbtypes.FieldMapper.
require.NotNil(t, NewConditionBuilder(NewFieldMapper()))
}
func TestConditionForSemconvFamilyPreservesMaterializedMemberExistsColumn(t *testing.T) {
sql, args, warnings := traceFamilyConditionSQL(t,
"deployment.environment.name",
[]*telemetrytypes.TelemetryFieldKey{
traceAttrMember("deployment.environment.name", false),
traceAttrMember("deployment.environment", true),
},
qbtypes.FilterOperatorEqual, "production",
)
assert.Empty(t, warnings)
assert.Contains(t, sql, "`attribute_string_deployment$$environment_exists`")
assert.NotContains(t, sql, "`attribute_string_deployment$environment_exists`")
assert.Equal(t, []any{"production"}, args)
}
func TestConditionForSemconvFamilyNotExistsChecksEveryMember(t *testing.T) {
sql, _, warnings := traceFamilyConditionSQL(t,
"deployment.environment",
[]*telemetrytypes.TelemetryFieldKey{
traceAttrMember("deployment.environment.name", false),
traceAttrMember("deployment.environment", false),
},
qbtypes.FilterOperatorNotExists, nil,
)
assert.Empty(t, warnings)
assert.Contains(t, sql, "NOT (mapContains(attributes_string, 'deployment.environment.name') OR mapContains(attributes_string, 'deployment.environment'))")
}
func TestConditionForResourceWithEvolution(t *testing.T) {
ctx := context.Background()
releaseTime := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)

View File

@@ -167,6 +167,145 @@ func NewFieldMapper() *fieldMapper {
return &fieldMapper{}
}
// FieldForLogical returns the value expression for a resolved logical field:
// the member's own expression for a single-member field, and a current-first
// merge across the members' expressions for a family. Each member expression
// comes from FieldFor and therefore honors that member's own materialization
// and evolution state — a family never needs sibling information on a key.
func (m *fieldMapper) FieldForLogical(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,
logical *telemetrytypes.LogicalField,
) (string, error) {
if !logical.IsFamily() {
return m.FieldFor(ctx, orgID, tsStart, tsEnd, logical.Single())
}
memberExprs := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
expr, err := m.FieldFor(ctx, orgID, tsStart, tsEnd, member)
if err != nil {
return "", err
}
memberExprs = append(memberExprs, expr)
}
if logical.FieldDataType == telemetrytypes.FieldDataTypeString {
// The trailing '' keeps single-key semantics for rows without any
// member: string maps read '' for an absent key, and negative
// operators must keep including such rows (see AddDefaultExistsFilter).
values := make([]string, 0, len(memberExprs))
for _, expr := range memberExprs {
values = append(values, fmt.Sprintf("NULLIF(%s, '')", expr))
}
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.
branches := make([]string, 0, len(logical.Members)*2)
for i, member := range logical.Members {
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, member, true)
if err != nil {
return "", err
}
branches = append(branches, guard, memberExprs[i])
}
return "multiIf(" + strings.Join(branches, ", ") + ", NULL)", nil
}
// ExistsForLogical renders the existence predicate for a resolved logical
// field: a member's own predicate for a single-member field, presence of any
// member for a family.
func (m *fieldMapper) ExistsForLogical(
ctx context.Context,
orgID valuer.UUID,
tsStart, tsEnd uint64,
logical *telemetrytypes.LogicalField,
exists bool,
) (string, error) {
if !logical.IsFamily() {
return m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, logical.Single(), exists)
}
guards := make([]string, 0, len(logical.Members))
for _, member := range logical.Members {
guard, err := m.existsExpressionFor(ctx, orgID, tsStart, tsEnd, member, true)
if err != nil {
return "", err
}
guards = append(guards, guard)
}
combined := "(" + strings.Join(guards, " OR ") + ")"
if exists {
return combined, nil
}
return "NOT " + combined, nil
}
// logicalForResolvedColumn upgrades a directly-resolvable key (the FieldFor
// probe succeeded) to its family when the metadata map proves membership;
// otherwise the key stays a single-member logical field.
func logicalForResolvedColumn(field *telemetrytypes.TelemetryFieldKey, keys map[string][]*telemetrytypes.TelemetryFieldKey) *telemetrytypes.LogicalField {
for _, logical := range querybuilder.MatchingLogicalFields(field, keys) {
if logical.IsFamily() &&
logical.FieldContext == field.FieldContext &&
(field.FieldDataType == telemetrytypes.FieldDataTypeUnspecified || logical.FieldDataType == field.FieldDataType) {
return logical
}
}
return telemetrytypes.SingleLogicalField(field.Name, field)
}
// upgradeToFamilies swaps single-member candidates for their family when the
// metadata map proves membership. Candidate order and every non-family
// candidate stay exactly as the legacy flow produced them; sibling candidates
// of an already-emitted family are dropped rather than duplicated.
func upgradeToFamilies(field *telemetrytypes.TelemetryFieldKey, candidates []*telemetrytypes.LogicalField, keys map[string][]*telemetrytypes.TelemetryFieldKey) []*telemetrytypes.LogicalField {
var families []*telemetrytypes.LogicalField
for _, logical := range querybuilder.MatchingLogicalFields(field, keys) {
if logical.IsFamily() {
families = append(families, logical)
}
}
if len(families) == 0 {
return candidates
}
out := make([]*telemetrytypes.LogicalField, 0, len(candidates))
emitted := make(map[*telemetrytypes.LogicalField]bool)
for _, candidate := range candidates {
var family *telemetrytypes.LogicalField
for _, fam := range families {
if fam.FieldContext != candidate.FieldContext || fam.FieldDataType != candidate.FieldDataType {
continue
}
memberOfFamily := candidate.Single().Name == field.Name
for _, member := range fam.Members {
if member.Name == candidate.Single().Name {
memberOfFamily = true
break
}
}
if memberOfFamily {
family = fam
break
}
}
if family == nil {
out = append(out, candidate)
continue
}
if emitted[family] {
continue
}
emitted[family] = true
out = append(out, family)
}
return out
}
func (m *fieldMapper) getColumn(
_ context.Context,
_, _ uint64,
@@ -292,9 +431,9 @@ func (m *fieldMapper) resolveColumnExprs(
return nil, nil, nil, errors.Newf(errors.TypeInvalidInput, errors.CodeInvalidInput, "only resource context fields are supported for json columns, got %s", key.FieldContext.String)
}
// have to add ::string as clickHouse throws an error :- data types Variant/Dynamic are not allowed in GROUP BY
// once clickHouse dependency is updated, we need to check if we can remove it.
exprs = append(exprs, fmt.Sprintf("%s.`%s`::String", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("%s.`%s` IS NOT NULL", columnName, key.Name))
// once ClickHouse is updated, check whether this cast can be removed.
exprs = append(exprs, fmt.Sprintf("%s.%s::String", columnName, querybuilder.ClickHouseIdentifier(key.Name)))
existExprs = append(existExprs, fmt.Sprintf("%s.%s IS NOT NULL", columnName, querybuilder.ClickHouseIdentifier(key.Name)))
case schema.ColumnTypeEnumString,
schema.ColumnTypeEnumUInt64,
schema.ColumnTypeEnumUInt32,
@@ -319,13 +458,13 @@ func (m *fieldMapper) resolveColumnExprs(
switch valueType := column.Type.(schema.MapColumnType).ValueType; valueType.GetType() {
case schema.ColumnTypeEnumString, schema.ColumnTypeEnumFloat64, schema.ColumnTypeEnumBool:
// a key could have been materialized, if so return the materialized column name
if key.Materialized {
// a key could have been materialized, if so return the materialized column name
exprs = append(exprs, telemetrytypes.FieldKeyToMaterializedColumnName(key))
existExprs = append(existExprs, telemetrytypes.FieldKeyToMaterializedColumnNameForExists(key))
} else {
exprs = append(exprs, fmt.Sprintf("%s['%s']", columnName, key.Name))
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, '%s')", columnName, key.Name))
exprs = append(exprs, fmt.Sprintf("%s[%s]", columnName, querybuilder.ClickHouseStringLiteral(key.Name)))
existExprs = append(existExprs, fmt.Sprintf("mapContains(%s, %s)", columnName, querybuilder.ClickHouseStringLiteral(key.Name)))
}
default:
return nil, nil, nil, errors.NewInvalidInputf(errors.CodeInvalidInput, "value type %s is not supported for map column type %s", valueType, column.Type)
@@ -348,18 +487,23 @@ func (m *fieldMapper) ColumnExpressionFor(
keys map[string][]*telemetrytypes.TelemetryFieldKey,
) (string, error) {
// Resolve the candidate column(s).
var candidates []*telemetrytypes.TelemetryFieldKey
// Resolve the candidate logical field(s).
var candidates []*telemetrytypes.LogicalField
switch _, err := m.FieldFor(ctx, orgID, startNs, endNs, field); {
case err == nil:
candidates = []*telemetrytypes.TelemetryFieldKey{field}
// A directly-resolvable key upgrades to its family when the metadata
// map proves membership; otherwise it stays single-member.
candidates = []*telemetrytypes.LogicalField{logicalForResolvedColumn(field, keys)}
case errors.Is(err, qbtypes.ErrColumnNotFound):
// column (when the bare name is one) plus metadata matches, else synthesized
// type-variant keys.
candidates = m.CandidateKeys(ctx, orgID, field, nil, keys)
if len(candidates) == 0 {
// The legacy candidate flow, unchanged: column (when the bare name is
// one) plus metadata matches, else synthesized type-variant keys. The
// family step below only swaps candidates for their family; it never
// changes candidate order or non-family behavior.
raw := m.CandidateKeys(ctx, orgID, field, nil, keys)
if len(raw) == 0 {
return "", errors.Wrapf(err, errors.TypeInvalidInput, errors.CodeInvalidInput, "field `%s` not found", field.Name).WithSuggestions(errors.NewSuggestionsOnLevenshteinDistance(field.Name, errors.NounKeys, maps.Keys(keys))...)
}
candidates = upgradeToFamilies(field, querybuilder.WrapAsLogicalFields(field.Name, raw), keys)
default:
return "", err
}
@@ -373,21 +517,21 @@ func (m *fieldMapper) ColumnExpressionFor(
dummyValue = 0.0
}
stmts := make([]string, 0, len(candidates)*2)
for _, key := range candidates {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
for _, logical := range candidates {
value, err := m.FieldForLogical(ctx, orgID, startNs, endNs, logical)
if err != nil {
return "", err
}
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
guard, err := m.ExistsForLogical(ctx, orgID, startNs, endNs, logical, true)
if err != nil {
return "", err
}
coerced := value
// a time column keeps its native type; coercing it would yield seconds
if temporal, err := m.columnIsTemporal(ctx, startNs, endNs, key); err != nil {
if temporal, err := m.logicalIsTemporal(ctx, startNs, endNs, logical); err != nil {
return "", err
} else if !temporal {
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(key, dummyValue, value, qbtypes.FilterOperatorUnknown)
coerced, _ = querybuilder.DataTypeCollisionHandledFieldName(logical.Single(), dummyValue, value, qbtypes.FilterOperatorUnknown)
}
stmts = append(stmts, guard, coerced)
}
@@ -395,13 +539,14 @@ func (m *fieldMapper) ColumnExpressionFor(
}
if len(candidates) == 1 {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, candidates[0])
logical := candidates[0]
value, err := m.FieldForLogical(ctx, orgID, startNs, endNs, logical)
if err != nil {
return "", err
}
exprs, existExprs, _, _ := m.resolveColumnExprs(ctx, startNs, endNs, candidates[0])
if len(exprs) == 1 && len(existExprs) == 1 {
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, candidates[0], true)
exprs, existExprs, _, _ := m.resolveColumnExprs(ctx, startNs, endNs, logical.Single())
if !logical.IsFamily() && len(exprs) == 1 && len(existExprs) == 1 {
guard, err := m.ExistsForLogical(ctx, orgID, startNs, endNs, logical, true)
if err != nil {
return "", err
}
@@ -413,12 +558,12 @@ func (m *fieldMapper) ColumnExpressionFor(
// Multiple candidates (collision / synth): multiIf picks the first that exists,
// stringified so branches share a type.
args := make([]string, 0, len(candidates))
for _, key := range candidates {
value, err := m.FieldFor(ctx, orgID, startNs, endNs, key)
for _, logical := range candidates {
value, err := m.FieldForLogical(ctx, orgID, startNs, endNs, logical)
if err != nil {
return "", err
}
guard, err := m.existsExpressionFor(ctx, orgID, startNs, endNs, key, true)
guard, err := m.ExistsForLogical(ctx, orgID, startNs, endNs, logical, true)
if err != nil {
return "", err
}
@@ -427,6 +572,15 @@ func (m *fieldMapper) ColumnExpressionFor(
return fmt.Sprintf("multiIf(%s, NULL)", strings.Join(args, ", ")), nil
}
// logicalIsTemporal reports whether the logical field resolves to a single time
// column. A family is attribute-backed and never temporal.
func (m *fieldMapper) logicalIsTemporal(ctx context.Context, startNs, endNs uint64, logical *telemetrytypes.LogicalField) (bool, error) {
if logical.IsFamily() {
return false, nil
}
return m.columnIsTemporal(ctx, startNs, endNs, logical.Single())
}
// columnIsTemporal reports whether key resolves to a single time column, after evolution
// selection. Multiple columns mean an attribute-map union, which is never temporal.
func (m *fieldMapper) columnIsTemporal(ctx context.Context, startNs, endNs uint64, key *telemetrytypes.TelemetryFieldKey) (bool, error) {

View File

@@ -5,6 +5,7 @@ import (
"testing"
"time"
"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"
@@ -12,6 +13,20 @@ import (
"github.com/stretchr/testify/require"
)
func TestFieldForQuotesRequestKeyNames(t *testing.T) {
key := &telemetrytypes.TelemetryFieldKey{
Name: "name'`\\); SELECT 1 --",
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, key)
require.NoError(t, err)
assert.Contains(t, expression, "resource."+querybuilder.ClickHouseIdentifier(key.Name))
assert.Contains(t, expression, "mapContains(resources_string, "+querybuilder.ClickHouseStringLiteral(key.Name)+")")
}
func TestGetFieldKeyName(t *testing.T) {
ctx := context.Background()
@@ -120,6 +135,118 @@ func TestGetFieldKeyName(t *testing.T) {
}
}
// FieldFor is a per-physical-key primitive: it never consults the family
// table. Family composition is FieldForLogical's job.
func TestFieldForIsPerKey(t *testing.T) {
key := telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
expression, err := NewFieldMapper().FieldFor(context.Background(), valuer.UUID{}, 0, 0, &key)
require.NoError(t, err)
assert.Equal(t, "attributes_string['deployment.environment.name']", expression)
}
func traceFamilyLogicalField(t *testing.T, requestedName string, members ...*telemetrytypes.TelemetryFieldKey) *telemetrytypes.LogicalField {
t.Helper()
fieldKeys := map[string][]*telemetrytypes.TelemetryFieldKey{}
for _, member := range members {
fieldKeys[member.Name] = []*telemetrytypes.TelemetryFieldKey{member}
}
requested := telemetrytypes.NewTelemetryFieldKey(requestedName, members[0].FieldContext, members[0].FieldDataType)
matches := querybuilder.MatchingLogicalFields(requested, fieldKeys)
require.Len(t, matches, 1)
return matches[0]
}
func TestFieldForLogicalMergesFamilyMembersCurrentFirst(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
for _, requestedName := range []string{current.Name, old.Name} {
logical := traceFamilyLogicalField(t, requestedName, current, old)
expression, err := NewFieldMapper().FieldForLogical(context.Background(), valuer.UUID{}, 0, 0, logical)
require.NoError(t, err)
assert.Equal(t,
"COALESCE(NULLIF(attributes_string['deployment.environment.name'], ''), NULLIF(attributes_string['deployment.environment'], ''), '')",
expression,
"both request spellings address one logical field",
)
}
}
// The old-name request with only the current spelling in metadata prunes to a
// single member: plain single-key SQL, no coalesce.
func TestFieldForLogicalPrunesToPresentMembers(t *testing.T) {
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
logical := traceFamilyLogicalField(t, "deployment.environment", current)
expression, err := NewFieldMapper().FieldForLogical(context.Background(), valuer.UUID{}, 0, 0, logical)
require.NoError(t, err)
assert.Equal(t, "attributes_string['deployment.environment.name']", expression)
}
// Every member brings its own storage state: the merge composes each member's
// FieldFor output, so a materialized-with-evolutions member keeps its column
// history inside the family expression with no sibling bookkeeping anywhere.
func TestFieldForLogicalComposesResourceMembersFromTheirOwnStorage(t *testing.T) {
ctx := context.Background()
fm := NewFieldMapper()
start := uint64(time.Date(2024, 6, 1, 0, 0, 0, 0, time.UTC).UnixNano())
end := uint64(time.Date(2024, 6, 5, 0, 0, 0, 0, time.UTC).UnixNano())
current := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
old := &telemetrytypes.TelemetryFieldKey{
Name: "deployment.environment",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
Materialized: true,
Evolutions: MockEvolutionData(time.Date(2024, 6, 2, 0, 0, 0, 0, time.UTC)),
}
currentExpr, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, current)
require.NoError(t, err)
oldExpr, err := fm.FieldFor(ctx, valuer.UUID{}, start, end, old)
require.NoError(t, err)
logical := traceFamilyLogicalField(t, current.Name, current, old)
expression, err := fm.FieldForLogical(ctx, valuer.UUID{}, start, end, logical)
require.NoError(t, err)
assert.Equal(t,
"COALESCE(NULLIF("+currentExpr+", ''), NULLIF("+oldExpr+", ''), '')",
expression,
"the family merge is exactly the members' own expressions, current-first, with the keyless tail",
)
assert.Contains(t, expression, "`resource_string_deployment$$environment`",
"the promoted member keeps its materialized column")
}
func TestFieldForResourceWithEvolution(t *testing.T) {
ctx := context.Background()
releaseTime := time.Date(2025, 1, 1, 0, 0, 0, 0, time.UTC)

View File

@@ -155,6 +155,8 @@ var operatorInverseMapping = map[FilterOperator]FilterOperator{
// doesn't have value "redis"
// Since we don't know the intent, we don't add the exists filter. They are expected
// to add exists filter themselves if exclusion is desired.
// Negative predicates therefore include rows where the key is absent; value
// expressions must preserve the storage column's absent-key default.
//
// For the positive predicates, the key existence is implied.
func (f FilterOperator) AddDefaultExistsFilter() bool {

View File

@@ -2,6 +2,7 @@ package telemetrytypes
import (
"fmt"
"slices"
"strings"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -50,6 +51,74 @@ type TelemetryFieldKey struct {
Evolutions []*EvolutionEntry `json:"-"`
}
// Copy returns an independent copy of f.
func (f *TelemetryFieldKey) Copy() *TelemetryFieldKey {
if f == nil {
return nil
}
copied := *f
copied.Indexes = slices.Clone(f.Indexes)
if f.Evolutions != nil {
copied.Evolutions = make([]*EvolutionEntry, len(f.Evolutions))
for index, evolution := range f.Evolutions {
if evolution != nil {
copiedEvolution := *evolution
copied.Evolutions[index] = &copiedEvolution
}
}
}
copied.JSONPlan = copyJSONAccessPlan(f.JSONPlan, f, &copied)
return &copied
}
func copyJSONAccessPlan(plan JSONAccessPlan, sourceKey, copiedKey *TelemetryFieldKey) JSONAccessPlan {
if plan == nil {
return nil
}
nodes := make(map[*JSONAccessNode]*JSONAccessNode)
var copyNode func(*JSONAccessNode) *JSONAccessNode
copyNode = func(node *JSONAccessNode) *JSONAccessNode {
if node == nil {
return nil
}
if copiedNode, ok := nodes[node]; ok {
return copiedNode
}
copiedNode := *node
nodes[node] = &copiedNode
copiedNode.Parent = copyNode(node.Parent)
if node.Branches != nil {
copiedNode.Branches = make(map[JSONAccessBranchType]*JSONAccessNode, len(node.Branches))
for branchType, branch := range node.Branches {
copiedNode.Branches[branchType] = copyNode(branch)
}
}
if node.TerminalConfig != nil {
copiedTerminal := *node.TerminalConfig
switch node.TerminalConfig.Key {
case nil:
case sourceKey:
copiedTerminal.Key = copiedKey
default:
copiedTerminal.Key = node.TerminalConfig.Key.Copy()
}
copiedNode.TerminalConfig = &copiedTerminal
}
return &copiedNode
}
copied := make(JSONAccessPlan, len(plan))
for index, node := range plan {
copied[index] = copyNode(node)
}
return copied
}
func (f *TelemetryFieldKey) KeyNameContainsArray() bool {
return strings.Contains(f.Name, ArraySep) || strings.Contains(f.Name, ArrayAnyIndex)
}
@@ -233,6 +302,15 @@ type MetricContext struct {
MetricNamespace string `json:"metricNamespace,omitempty"`
}
// Copy returns an independent copy of m.
func (m *MetricContext) Copy() *MetricContext {
if m == nil {
return nil
}
copied := *m
return &copied
}
type FieldKeySelector struct {
StartUnixMilli int64 `json:"startUnixMilli"`
EndUnixMilli int64 `json:"endUnixMilli"`
@@ -246,6 +324,16 @@ type FieldKeySelector struct {
MetricContext *MetricContext `json:"metricContext,omitempty"`
}
// Copy returns an independent copy of s.
func (s *FieldKeySelector) Copy() *FieldKeySelector {
if s == nil {
return nil
}
copied := *s
copied.MetricContext = s.MetricContext.Copy()
return &copied
}
type FieldValueSelector struct {
*FieldKeySelector
ExistingQuery string `json:"existingQuery"`

View File

@@ -0,0 +1,75 @@
package telemetrytypes
import (
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestTelemetryFieldKeyCopyOwnsMutableState(t *testing.T) {
original := &TelemetryFieldKey{
Name: "items.name",
FieldContext: FieldContextBody,
FieldDataType: FieldDataTypeString,
Indexes: []TelemetryFieldKeySkipIndex{
{Name: "items.name"},
},
Evolutions: []*EvolutionEntry{
{FieldName: "items.name"},
},
}
require.NoError(t, original.SetJSONAccessPlan(JSONColumnMetadata{BaseColumn: "body_v2"}, nil))
require.Len(t, original.JSONPlan, 1)
require.NotNil(t, original.JSONPlan[0].TerminalConfig)
copied := original.Copy()
require.NotNil(t, copied)
require.NotSame(t, original, copied)
require.Len(t, copied.JSONPlan, 1)
require.NotNil(t, copied.JSONPlan[0].TerminalConfig)
assert.NotSame(t, original.JSONPlan[0], copied.JSONPlan[0])
assert.NotSame(t, original.JSONPlan[0].Parent, copied.JSONPlan[0].Parent)
assert.Same(t, copied, copied.JSONPlan[0].TerminalConfig.Key)
assert.Equal(t, original.JSONPlan[0].Alias(), copied.JSONPlan[0].Alias())
copied.Name = "changed"
copied.Indexes[0].Name = "changed"
copied.Evolutions[0].FieldName = "changed"
copied.JSONPlan[0].Name = "changed"
copied.JSONPlan[0].Parent.Name = "changed"
assert.Equal(t, "items.name", original.Name)
assert.Equal(t, "items.name", original.Indexes[0].Name)
assert.Equal(t, "items.name", original.Evolutions[0].FieldName)
assert.Equal(t, "items.name", original.JSONPlan[0].Name)
assert.Equal(t, "body_v2", original.JSONPlan[0].Parent.Name)
}
func TestFieldKeySelectorCopyOwnsMetricContext(t *testing.T) {
original := &FieldKeySelector{
Name: "state",
MetricContext: &MetricContext{
MetricName: "system.cpu.time",
MetricNamespace: "system",
},
}
copied := original.Copy()
require.NotNil(t, copied)
require.NotNil(t, copied.MetricContext)
assert.NotSame(t, original.MetricContext, copied.MetricContext)
copied.Name = "changed"
copied.MetricContext.MetricName = "changed"
assert.Equal(t, "state", original.Name)
assert.Equal(t, "system.cpu.time", original.MetricContext.MetricName)
}
func TestNilFieldCopies(t *testing.T) {
assert.Nil(t, (*TelemetryFieldKey)(nil).Copy())
assert.Nil(t, (*FieldKeySelector)(nil).Copy())
assert.Nil(t, (*MetricContext)(nil).Copy())
}

View File

@@ -0,0 +1,78 @@
package telemetrytypes
// LogicalField is resolution output: one queryable field, addressed by the
// spelling the request used, backed by the physical member keys that store it.
//
// The resolver expresses ambiguity ("possibly different fields sharing a
// name") as a []*LogicalField — never inside one LogicalField. Within one
// LogicalField, members are alternate physical spellings of the same field
// (a semantic-convention family), ordered current-first; compilers merge
// them into one expression with current-wins precedence. Across the slice,
// compilers build one condition per LogicalField and combine per the
// operator, exactly as they previously combined ambiguous keys.
//
// Members always has at least one entry. A non-family field has exactly
// one. Members alias the metadata map entries and must not be mutated.
type LogicalField struct {
// Name is the requested spelling. It is the response identity: aliases,
// series labels, and warnings use it, so responses echo the request.
Name string
// The physical identity every member shares. Members with a different
// signal, field context, or data type belong to different logical
// fields by definition.
Signal Signal
FieldContext FieldContext
FieldDataType FieldDataType
// Members are the physical keys that store this field, ordered
// current-first. Each member carries its own physical facts
// (Materialized, Evolutions, JSONPlan, ...), so per-member accessors
// need no sibling information.
Members []*TelemetryFieldKey
}
// SingleLogicalField wraps one physical key as its own logical field.
func SingleLogicalField(name string, key *TelemetryFieldKey) *LogicalField {
return &LogicalField{
Name: name,
Signal: key.Signal,
FieldContext: key.FieldContext,
FieldDataType: key.FieldDataType,
Members: []*TelemetryFieldKey{key},
}
}
// Single returns the only member. It is the accessor for signals whose
// logical fields are always single-member (everything except traces today).
func (l *LogicalField) Single() *TelemetryFieldKey {
return l.Members[0]
}
// IsFamily reports whether the field has more than one physical member.
func (l *LogicalField) IsFamily() bool {
return len(l.Members) > 1
}
// String implements fmt.Stringer for warning messages.
func (l *LogicalField) String() string {
if len(l.Members) == 1 {
return l.Members[0].String()
}
names := make([]string, 0, len(l.Members))
for _, member := range l.Members {
names = append(names, member.Name)
}
return l.Name + "(" + l.FieldContext.StringValue() + ", " + l.FieldDataType.StringValue() + ", members: " + joinNames(names) + ")"
}
func joinNames(names []string) string {
out := ""
for i, name := range names {
if i > 0 {
out += ", "
}
out += name
}
return out
}