Compare commits

...

6 Commits

Author SHA1 Message Date
srikanthccv
e0de0fcf6f fix(telemetrymetadata): make the existence check exact and fail on its bound
- the check is a dedicated row lookup on the signal's tag table that
  matches the value as a string or as the number it parses to; the values
  reader dropped numeric zeros, so `shard = 0` read as absent
- the check runs without timeout_overflow_mode=break, so a bound hit is
  an error and proves nothing, instead of a partial empty result
- the fields window start is floored to six hours for every signal again,
  matching the collector's metrics bucket default

Assisted-by: Claude Fable 5.1
Claude-Session: https://claude.ai/code/session_01NNQrizqt5XQ1N2NNAcuAi8
2026-09-07 21:02:07 +05:30
srikanthccv
c456041089 fix(telemetrymetadata): review fixes for the related-values lookup
- the key-lookup cache is gone: it never evicted and keyed on whole
  selectors, timestamps included
- a resource key seen in several signals scans all of them again; the
  smallest signal need not hold the same values
- the search text is matched in the maps the key is read from, so typing
  into a span or log field no longer filters attribute maps that do not
  hold it
- the org's related-values slot and the query bounds are taken before the
  existence checks, not only before the final scan
- max_threads is passed through the settings map; the typed context key
  the settings hook never read is left alone
- the end bound is kept as requested instead of rounded up into the next
  bucket
- a clamped window marks the result incomplete, also when an existence
  check returns early
- the equality-term extractor reads literals the way the where-clause
  visitor does: quotes and escapes removed, numbers in decimal form
- the execution-time bound is checked from the start of the query
  (timeout_before_checking_execution_speed 0) rather than estimated

Assisted-by: Claude Fable 5.1
Claude-Session: https://claude.ai/code/session_01NNQrizqt5XQ1N2NNAcuAi8
2026-09-07 20:39:19 +05:30
srikanthccv
96128e9f87 feat(telemetrymetadata): bound, scope and cache the related-values lookup
- new telemetrymetadata config: keys_cache_ttl and related_values
  (max_execution_time, max_window, max_concurrency, max_threads,
  read_buffer_size, max_existence_checks)
- the related-values query runs with max_execution_time and
  timeout_overflow_mode=break, returns what it found by then, and marks
  the response incomplete; its window is clamped to max_window, floored
  and ceiled to the metadata bucket
- the requested key is resolved through the keys lookup and read from the
  map it resolves to, instead of decoding every map
- a request without a signal scans the one signal the key was seen in, or
  the smallest signal for a resource key seen in several
- equality terms of the existing query are checked against the all-values
  set first; a value absent from the window returns an empty related set
  without a scan
- related-values queries are capped per org; max_threads and the read
  buffer size can be set for them
- key lookups (table statements, tag keys, column evolution) are cached
  per org for keys_cache_ttl, and keys named in the existing query are
  looked up by exact name
- query settings can be added per request through the context

Assisted-by: Claude Fable 5.1
Claude-Session: https://claude.ai/code/session_01NNQrizqt5XQ1N2NNAcuAi8
2026-09-07 19:51:02 +05:30
srikanthccv
16a0c2375d perf(telemetrymetadata): simplify the related-values predicate shapes
- positive operators emit `mapContains(m, k) AND cond` instead of the
  if(...) form, which ClickHouse can analyse for skip indexes
- IN and NOT IN operands are emitted as IN lists instead of OR and AND
  chains of comparisons
- numeric operands are compared in their decimal string form; the
  metadata maps hold strings only, so the toFloat64OrNull cast never
  matched and forced a full map decode

Assisted-by: Claude Fable 5.1
Claude-Session: https://claude.ai/code/session_01NNQrizqt5XQ1N2NNAcuAi8
2026-09-07 19:51:02 +05:30
srikanthccv
14df7ff810 feat(telemetrymetadata): read span and log fields from the intrinsic map for related values
- the metadata field mapper resolves the span and log contexts to the
  intrinsic_attributes column written by signoz-otel-collector metadata
  migration 1002
- the related-values select reads a span or log field from the intrinsic
  map and falls back to attributes for rows written before the column
  existed; a key without a context is read from the intrinsic, resource or
  attributes map in that order
- bool span fields such as has_error are filtered by their "true"/"false"
  string form; numeric span fields still build no condition

Assisted-by: Claude Fable 5.1
Claude-Session: https://claude.ai/code/session_01NNQrizqt5XQ1N2NNAcuAi8
2026-09-07 19:28:03 +05:30
srikanthccv
ca1e063543 fix(telemetrytypes): floor the fields window start to the signal's metadata bucket
- traces and logs metadata rows are written per 6 h bucket, metrics rows
  per 24 h bucket; the start of a metrics or signal-less request is floored
  to the day so the bucket holding the start is not left out

Assisted-by: Claude Fable 5.1
Claude-Session: https://claude.ai/code/session_01NNQrizqt5XQ1N2NNAcuAi8
2026-09-07 19:24:20 +05:30
19 changed files with 994 additions and 110 deletions

View File

@@ -167,6 +167,26 @@ querier:
# to 0 to disable.
log_trace_id_window_padding: 5m
##################### TelemetryMetadata #####################
telemetrymetadata:
related_values:
# Server-side bound for the related-values query of the fields API; the values
# found by then are returned and the response is marked incomplete.
max_execution_time: 2s
# The related-values window is clamped to this length; the requested window
# keeps driving the all-values lookup.
max_window: 168h
# Related-values queries an org runs at once; further requests wait for a slot.
max_concurrency: 2
# max_threads for the related-values query; 0 keeps the server default.
max_threads: 0
# max_read_buffer_size_local_fs for the related-values query, in bytes; 0 keeps
# the server default.
read_buffer_size: 0
# Equality terms of the existing query whose value is checked against the
# all-values set before the scan; an absent value makes the related set empty.
max_existence_checks: 4
##################### TelemetryStore #####################
telemetrystore:
# Maximum number of idle connections in the connection pool.

View File

@@ -113,7 +113,7 @@ func NewTestManager(t *testing.T, testOpts *TestManagerOptions) *Manager {
}
// Create querier with test values
metadataStore := telemetrymetadata.NewTelemetryMetaStore(providerSettings, telemetryStore, flagger)
metadataStore := telemetrymetadata.NewTelemetryMetaStore(providerSettings, telemetryStore, flagger, telemetrymetadata.NewConfig())
cfg := statementbuilder.Config{}
ctx := context.Background()
traceStmtBuilder, err := tracesstatementbuilder.NewFactory(telemetryStore, metadataStore, flagger).New(ctx, providerSettings, cfg)

View File

@@ -4,6 +4,7 @@ import (
grammar "github.com/SigNoz/signoz/pkg/parser/filterquery/grammar"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/antlr4-go/antlr/v4"
"strconv"
)
// QueryStringToKeysSelectors converts a query string to a list of field key selectors
@@ -72,3 +73,61 @@ func QueryStringToKeysSelectors(query string) []*telemetrytypes.FieldKeySelector
return keys
}
// EqualityTerm is one `key = value` term of a filter expression.
type EqualityTerm struct {
Key *telemetrytypes.TelemetryFieldKey
Value string
}
// QueryStringEqualityTerms returns the `key = value` terms of a filter
// expression that is a conjunction: every term must hold for a row to match.
// ok is false when the expression contains OR or NOT, since a term may then
// be optional.
func QueryStringEqualityTerms(query string) (terms []EqualityTerm, ok bool) {
lexer := grammar.NewFilterQueryLexer(antlr.NewInputStream(query))
var lastKey *telemetrytypes.TelemetryFieldKey
equalsSeen := false
for {
tok := lexer.NextToken()
if tok.GetTokenType() == antlr.TokenEOF {
break
}
switch tok.GetTokenType() {
case grammar.FilterQueryLexerWS:
continue
case grammar.FilterQueryLexerOR, grammar.FilterQueryLexerNOT, grammar.FilterQueryLexerNOT_EQUALS, grammar.FilterQueryLexerNEQ:
return nil, false
case grammar.FilterQueryLexerKEY:
key := telemetrytypes.GetFieldKeyFromKeyText(tok.GetText())
lastKey = &key
equalsSeen = false
case grammar.FilterQueryLexerEQUALS:
equalsSeen = lastKey != nil
case grammar.FilterQueryLexerQUOTED_TEXT, grammar.FilterQueryLexerNUMBER, grammar.FilterQueryLexerBOOL:
if equalsSeen && lastKey != nil {
terms = append(terms, EqualityTerm{Key: lastKey, Value: literalText(tok.GetTokenType(), tok.GetText())})
}
lastKey = nil
equalsSeen = false
default:
lastKey = nil
equalsSeen = false
}
}
return terms, true
}
// literalText returns a literal as the where-clause visitor reads it: quoted
// text with its quotes and escapes removed, a number in its decimal form.
func literalText(tokenType int, text string) string {
switch tokenType {
case grammar.FilterQueryLexerQUOTED_TEXT:
return trimQuotes(text)
case grammar.FilterQueryLexerNUMBER:
if f, err := strconv.ParseFloat(text, 64); err == nil {
return strconv.FormatFloat(f, 'f', -1, 64)
}
}
return text
}

View File

@@ -133,3 +133,62 @@ func TestQueryToKeys(t *testing.T) {
}
}
}
func TestQueryStringEqualityTerms(t *testing.T) {
terms, ok := QueryStringEqualityTerms(`service.name = 'checkout' AND resource.k8s.namespace.name = "prod" AND http.status_code = 500 AND has_error = true`)
if !ok {
t.Fatalf("expected a conjunction to be accepted")
}
got := map[string]string{}
contexts := map[string]telemetrytypes.FieldContext{}
for _, term := range terms {
got[term.Key.Name] = term.Value
contexts[term.Key.Name] = term.Key.FieldContext
}
want := map[string]string{"service.name": "checkout", "k8s.namespace.name": "prod", "http.status_code": "500", "has_error": "true"}
if len(got) != len(want) {
t.Fatalf("expected %d terms, got %v", len(want), got)
}
for name, value := range want {
if got[name] != value {
t.Fatalf("expected %s = %q, got %q", name, value, got[name])
}
}
if contexts["k8s.namespace.name"] != telemetrytypes.FieldContextResource {
t.Fatalf("expected the resource prefix to set the context")
}
for _, query := range []string{
`service.name = 'checkout' OR service.name = 'cart'`,
`NOT service.name = 'checkout'`,
`service.name != 'checkout'`,
} {
if _, ok := QueryStringEqualityTerms(query); ok {
t.Fatalf("expected %q to be rejected", query)
}
}
terms, ok = QueryStringEqualityTerms(`service.name IN ('a', 'b') AND http.route LIKE '/api%' AND env = 'prod'`)
if !ok || len(terms) != 1 || terms[0].Key.Name != "env" {
t.Fatalf("expected only the equality term, got %v ok=%v", terms, ok)
}
terms, ok = QueryStringEqualityTerms(`resource.owner = 'O\'Reilly' AND http.status_code = 5e2 AND ratio = 1.50`)
if !ok || len(terms) != 3 {
t.Fatalf("expected three terms, got %v ok=%v", terms, ok)
}
for _, want := range []struct{ name, value string }{{"owner", "O'Reilly"}, {"http.status_code", "500"}, {"ratio", "1.5"}} {
found := false
for _, term := range terms {
if term.Key.Name == want.name {
found = true
if term.Value != want.value {
t.Fatalf("expected %s = %q, got %q", want.name, want.value, term.Value)
}
}
}
if !found {
t.Fatalf("expected a term for %s", want.name)
}
}
}

View File

@@ -40,6 +40,7 @@ import (
"github.com/SigNoz/signoz/pkg/sqlschema"
"github.com/SigNoz/signoz/pkg/sqlstore"
"github.com/SigNoz/signoz/pkg/statsreporter"
"github.com/SigNoz/signoz/pkg/telemetrymetadata"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/tokenizer"
"github.com/SigNoz/signoz/pkg/valuer"
@@ -97,6 +98,9 @@ type Config struct {
// Querier config
Querier querier.Config `mapstructure:"querier"`
// TelemetryMetadata config
TelemetryMetadata telemetrymetadata.Config `mapstructure:"telemetrymetadata"`
// Ruler config
Ruler ruler.Config `mapstructure:"ruler"`
@@ -166,6 +170,7 @@ func NewConfig(ctx context.Context, logger *slog.Logger, resolverConfig config.R
prometheus.NewConfigFactory(),
alertmanager.NewConfigFactory(),
querier.NewConfigFactory(),
telemetrymetadata.NewConfigFactory(),
ruler.NewConfigFactory(),
emailing.NewConfigFactory(),
sharder.NewConfigFactory(),

View File

@@ -125,7 +125,7 @@ func newQueryStack(
querier.BucketCache,
error,
) {
metadataStore := telemetrymetadata.NewTelemetryMetaStore(settings, telemetryStore, fl)
metadataStore := telemetrymetadata.NewTelemetryMetaStore(settings, telemetryStore, fl, config.TelemetryMetadata)
cfg := config.Querier.Config
traceStmtBuilder, err := tracesstatementbuilder.NewFactory(telemetryStore, metadataStore, fl).New(ctx, settings, cfg)

View File

@@ -3,6 +3,7 @@ package telemetrymetadata
import (
"context"
"fmt"
"strconv"
schema "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator"
"github.com/SigNoz/signoz/pkg/querybuilder"
@@ -91,19 +92,35 @@ func (c *conditionBuilder) conditionForKey(
return "", nil
}
if key.FieldDataType != telemetrytypes.FieldDataTypeString &&
key.FieldDataType != telemetrytypes.FieldDataTypeUnspecified {
// if the field data type is not string, we can't build a condition for related values
switch key.FieldDataType {
case telemetrytypes.FieldDataTypeString, telemetrytypes.FieldDataTypeUnspecified:
// the metadata maps hold strings only, so a numeric operand is
// compared in its decimal form instead of casting the map value
value = numericOperandToString(value)
case telemetrytypes.FieldDataTypeBool:
// bool fields are stored as the strings "true" and "false" in the
// intrinsic map, so the operand is compared in that form against the
// field as a string
value = boolOperandToString(value)
stringKey := *key
stringKey.FieldDataType = telemetrytypes.FieldDataTypeString
key = &stringKey
default:
// numeric fields are not stored in the metadata maps, so there is no
// condition to build for related values
return "", nil
}
fieldExpression, value = querybuilder.DataTypeCollisionHandledFieldName(key, value, fieldExpression, operator)
// key must exist to apply the main filter. for positive operators the
// absent-key rows are excluded (fallback false); for negative operators
// they are kept (fallback true) so rows legitimately lacking the key match.
keyMissingFallback := operator.IsNegativeOperator()
expr := `if(mapContains(%s, %s), %s, %t)`
// condition is a plain conjunction, which the skip indexes can analyse;
// for negative operators rows lacking the key are kept (fallback true)
// so rows legitimately lacking the key match.
expr := `(mapContains(%s, %s) AND %s)`
if operator.IsNegativeOperator() {
expr = `if(mapContains(%s, %s), %s, true)`
}
var cond string
@@ -141,23 +158,13 @@ func (c *conditionBuilder) conditionForKey(
if !ok {
return "", qbtypes.ErrInValues
}
// instead of using IN, we use `=` + `OR` to make use of index
conditions := []string{}
for _, value := range values {
conditions = append(conditions, sb.E(fieldExpression, value))
}
cond = sb.Or(conditions...)
cond = sb.In(fieldExpression, values...)
case qbtypes.FilterOperatorNotIn:
values, ok := value.([]any)
if !ok {
return "", qbtypes.ErrInValues
}
// instead of using NOT IN, we use `!=` + `AND` to make use of index
conditions := []string{}
for _, value := range values {
conditions = append(conditions, sb.NE(fieldExpression, value))
}
cond = sb.And(conditions...)
cond = sb.NotIn(fieldExpression, values...)
// exists and not exists
// in the query builder, `exists` and `not exists` are used for
@@ -177,5 +184,48 @@ func (c *conditionBuilder) conditionForKey(
}
}
return fmt.Sprintf(expr, columns[0].Name, sb.Var(key.Name), cond, keyMissingFallback), nil
return fmt.Sprintf(expr, columns[0].Name, sb.Var(key.Name), cond), nil
}
// numericOperandToString converts a numeric operand, or a list of them, to
// its decimal string form. Other values are returned unchanged.
func numericOperandToString(value any) any {
switch v := value.(type) {
case float64:
return strconv.FormatFloat(v, 'f', -1, 64)
case float32:
return strconv.FormatFloat(float64(v), 'f', -1, 32)
case int:
return strconv.Itoa(v)
case int64:
return strconv.FormatInt(v, 10)
case int32:
return strconv.FormatInt(int64(v), 10)
case uint64:
return strconv.FormatUint(v, 10)
case []any:
out := make([]any, len(v))
for i, item := range v {
out[i] = numericOperandToString(item)
}
return out
}
return value
}
// boolOperandToString converts a bool operand, or a list of them, to the
// "true"/"false" strings stored in the metadata maps. Other values are
// returned unchanged.
func boolOperandToString(value any) any {
switch v := value.(type) {
case bool:
return strconv.FormatBool(v)
case []any:
out := make([]any, len(v))
for i, item := range v {
out[i] = boolOperandToString(item)
}
return out
}
return value
}

View File

@@ -22,6 +22,7 @@ func TestConditionFor(t *testing.T) {
operator qbtypes.FilterOperator
value any
expectedSQL string
expectedArgs []any
expectedError error
}{
@@ -34,7 +35,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorILike,
value: "%admin%",
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) LIKE LOWER(?), false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND LOWER(attributes['user.id']) LIKE LOWER(?))",
expectedError: nil,
},
{
@@ -58,7 +59,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorEqual,
value: "admin",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] = ?, false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND attributes['user.id'] = ?)",
expectedError: nil,
},
{
@@ -82,7 +83,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorIn,
value: []any{"admin", "root"},
expectedSQL: "WHERE if(mapContains(attributes, ?), (attributes['user.id'] = ? OR attributes['user.id'] = ?), false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND attributes['user.id'] IN (?, ?))",
expectedError: nil,
},
{
@@ -94,7 +95,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorNotIn,
value: []any{"admin", "root"},
expectedSQL: "WHERE if(mapContains(attributes, ?), (attributes['user.id'] <> ? AND attributes['user.id'] <> ?), true)",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] NOT IN (?, ?), true)",
expectedError: nil,
},
{
@@ -106,7 +107,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorLike,
value: "%admin%",
expectedSQL: "WHERE if(mapContains(attributes, ?), attributes['user.id'] LIKE ?, false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND attributes['user.id'] LIKE ?)",
expectedError: nil,
},
{
@@ -130,7 +131,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorContains,
value: "admin",
expectedSQL: "WHERE if(mapContains(attributes, ?), LOWER(attributes['user.id']) LIKE LOWER(?), false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND LOWER(attributes['user.id']) LIKE LOWER(?))",
expectedError: nil,
},
{
@@ -154,7 +155,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorRegexp,
value: "adm.*",
expectedSQL: "WHERE if(mapContains(attributes, ?), match(attributes['user.id'], ?), false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND match(attributes['user.id'], ?))",
expectedError: nil,
},
{
@@ -169,6 +170,55 @@ func TestConditionFor(t *testing.T) {
expectedSQL: "WHERE if(mapContains(attributes, ?), NOT match(attributes['user.id'], ?), true)",
expectedError: nil,
},
{
name: "Equal operator - span field in the intrinsic map",
key: telemetrytypes.TelemetryFieldKey{
Name: "http_method",
FieldContext: telemetrytypes.FieldContextSpan,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorEqual,
value: "GET",
expectedSQL: "WHERE (mapContains(intrinsic_attributes, ?) AND intrinsic_attributes['http_method'] = ?)",
expectedError: nil,
},
{
name: "Equal operator - bool span field compares the string form",
key: telemetrytypes.TelemetryFieldKey{
Name: "has_error",
FieldContext: telemetrytypes.FieldContextSpan,
FieldDataType: telemetrytypes.FieldDataTypeBool,
},
operator: qbtypes.FilterOperatorEqual,
value: true,
expectedSQL: "WHERE (mapContains(intrinsic_attributes, ?) AND intrinsic_attributes['has_error'] = ?)",
expectedArgs: []any{"has_error", "true"},
expectedError: nil,
},
{
name: "Equal operator - numeric span field builds no condition",
key: telemetrytypes.TelemetryFieldKey{
Name: "duration_nano",
FieldContext: telemetrytypes.FieldContextSpan,
FieldDataType: telemetrytypes.FieldDataTypeNumber,
},
operator: qbtypes.FilterOperatorEqual,
value: 1000,
expectedSQL: "",
expectedError: nil,
},
{
name: "Equal operator - log severity in the intrinsic map",
key: telemetrytypes.TelemetryFieldKey{
Name: "severity_text",
FieldContext: telemetrytypes.FieldContextLog,
FieldDataType: telemetrytypes.FieldDataTypeString,
},
operator: qbtypes.FilterOperatorEqual,
value: "ERROR",
expectedSQL: "WHERE (mapContains(intrinsic_attributes, ?) AND intrinsic_attributes['severity_text'] = ?)",
expectedError: nil,
},
{
name: "Exists operator - positive fallback false",
key: telemetrytypes.TelemetryFieldKey{
@@ -178,7 +228,7 @@ func TestConditionFor(t *testing.T) {
},
operator: qbtypes.FilterOperatorExists,
value: nil,
expectedSQL: "WHERE if(mapContains(attributes, ?), mapContains(attributes, 'user.id') = ?, false)",
expectedSQL: "WHERE (mapContains(attributes, ?) AND mapContains(attributes, 'user.id') = ?)",
expectedError: nil,
},
{
@@ -205,9 +255,38 @@ func TestConditionFor(t *testing.T) {
assert.Equal(t, tc.expectedError, err)
} else {
require.NoError(t, err)
sql, _ := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, tc.expectedSQL)
if tc.expectedArgs != nil {
assert.Equal(t, tc.expectedArgs, args)
}
}
})
}
}
func TestConditionForComparesNumbersAsStrings(t *testing.T) {
ctx := context.Background()
conditionBuilder := NewConditionBuilder(NewFieldMapper())
key := telemetrytypes.TelemetryFieldKey{
Name: "http.status_code",
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
sb := sqlbuilder.NewSelectBuilder()
cond, _, err := conditionBuilder.ConditionFor(ctx, valuer.UUID{}, 0, 0, &key, map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {&key}}, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorEqual, float64(200), sb)
require.NoError(t, err)
sb.Where(cond...)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "WHERE (mapContains(attributes, ?) AND attributes['http.status_code'] = ?)")
assert.Equal(t, []any{"http.status_code", "200"}, args)
sb = sqlbuilder.NewSelectBuilder()
cond, _, err = conditionBuilder.ConditionFor(ctx, valuer.UUID{}, 0, 0, &key, map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {&key}}, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorIn, []any{float64(200), int64(404)}, sb)
require.NoError(t, err)
sb.Where(cond...)
sql, args = sb.BuildWithFlavor(sqlbuilder.ClickHouse)
assert.Contains(t, sql, "attributes['http.status_code'] IN (?, ?)")
assert.Equal(t, []any{"http.status_code", "200", "404"}, args)
}

View File

@@ -0,0 +1,82 @@
package telemetrymetadata
import (
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
)
// RelatedValuesConfig bounds the related-values lookup of the fields API,
// which scans the attributes metadata table for the requested window.
type RelatedValuesConfig struct {
// MaxExecutionTime bounds the related-values query on the server; when it
// is reached the query returns the values found so far and the response
// is marked incomplete.
MaxExecutionTime time.Duration `mapstructure:"max_execution_time"`
// MaxWindow clamps the window of the related-values query. The requested
// window keeps driving the all-values lookup.
MaxWindow time.Duration `mapstructure:"max_window"`
// MaxConcurrency caps the related-values queries an org runs at once;
// further requests wait for a slot.
MaxConcurrency int `mapstructure:"max_concurrency"`
// MaxThreads sets max_threads for the related-values query; 0 keeps the
// server default.
MaxThreads int `mapstructure:"max_threads"`
// ReadBufferSize sets max_read_buffer_size_local_fs for the related-values
// query, in bytes; 0 keeps the server default.
ReadBufferSize int `mapstructure:"read_buffer_size"`
// MaxExistenceChecks is the number of equality terms of the existing query
// whose value is looked up in the all-values set before the scan; a value
// absent from the window makes the related set empty without a scan.
MaxExistenceChecks int `mapstructure:"max_existence_checks"`
}
// Config is the configuration of the telemetry metadata store.
type Config struct {
RelatedValues RelatedValuesConfig `mapstructure:"related_values"`
}
func NewConfigFactory() factory.ConfigFactory {
return factory.NewConfigFactory(factory.MustNewName("telemetrymetadata"), newConfig)
}
func newConfig() factory.Config {
return NewConfig()
}
// NewConfig returns the default configuration.
func NewConfig() Config {
return Config{
RelatedValues: RelatedValuesConfig{
MaxExecutionTime: 2 * time.Second,
MaxWindow: 7 * 24 * time.Hour,
MaxConcurrency: 2,
MaxThreads: 0,
ReadBufferSize: 0,
MaxExistenceChecks: 4,
},
}
}
func (c Config) Validate() error {
if c.RelatedValues.MaxExecutionTime < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "related_values.max_execution_time must not be negative, got %v", c.RelatedValues.MaxExecutionTime)
}
if c.RelatedValues.MaxWindow < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "related_values.max_window must not be negative, got %v", c.RelatedValues.MaxWindow)
}
if c.RelatedValues.MaxConcurrency < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "related_values.max_concurrency must not be negative, got %v", c.RelatedValues.MaxConcurrency)
}
if c.RelatedValues.MaxThreads < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "related_values.max_threads must not be negative, got %v", c.RelatedValues.MaxThreads)
}
if c.RelatedValues.ReadBufferSize < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "related_values.read_buffer_size must not be negative, got %v", c.RelatedValues.ReadBufferSize)
}
if c.RelatedValues.MaxExistenceChecks < 0 {
return errors.NewInvalidInputf(errors.CodeInvalidInput, "related_values.max_existence_checks must not be negative, got %v", c.RelatedValues.MaxExistenceChecks)
}
return nil
}

View File

@@ -24,6 +24,13 @@ var (
KeyType: schema.LowCardinalityColumnType{ElementType: schema.ColumnTypeString},
ValueType: schema.ColumnTypeString,
}},
// intrinsic_attributes holds the span and log context fields: span
// name, kind, status, the calculated HTTP and database fields, and log
// severity. See metadata migration 1002 in signoz-otel-collector.
"intrinsic_attributes": {Name: "intrinsic_attributes", Type: schema.MapColumnType{
KeyType: schema.LowCardinalityColumnType{ElementType: schema.ColumnTypeString},
ValueType: schema.ColumnTypeString,
}},
}
)
@@ -46,6 +53,8 @@ func (m *fieldMapper) getColumn(_ context.Context, _, _ uint64, key *telemetryty
return []*schema.Column{attributeMetadataColumns["resource_attributes"]}, nil
case telemetrytypes.FieldContextAttribute:
return []*schema.Column{attributeMetadataColumns["attributes"]}, nil
case telemetrytypes.FieldContextSpan, telemetrytypes.FieldContextLog:
return []*schema.Column{attributeMetadataColumns["intrinsic_attributes"]}, nil
}
return nil, qbtypes.ErrColumnNotFound
}

View File

@@ -115,10 +115,28 @@ func TestGetColumn(t *testing.T) {
expectedError: nil,
},
{
name: "Log field - nonexistent",
name: "Span field - intrinsic map",
key: telemetrytypes.TelemetryFieldKey{
Name: "http_method",
FieldContext: telemetrytypes.FieldContextSpan,
},
expectedCol: attributeMetadataColumns["intrinsic_attributes"],
expectedError: nil,
},
{
name: "Log field - intrinsic map",
key: telemetrytypes.TelemetryFieldKey{
Name: "severity_text",
FieldContext: telemetrytypes.FieldContextLog,
},
expectedCol: attributeMetadataColumns["intrinsic_attributes"],
expectedError: nil,
},
{
name: "Metric field - no column",
key: telemetrytypes.TelemetryFieldKey{
Name: "nonexistent_field",
FieldContext: telemetrytypes.FieldContextLog,
FieldContext: telemetrytypes.FieldContextMetric,
},
expectedCol: nil,
expectedError: qbtypes.ErrColumnNotFound,
@@ -195,7 +213,7 @@ func TestGetFieldKeyName(t *testing.T) {
name: "Non-existent column",
key: telemetrytypes.TelemetryFieldKey{
Name: "nonexistent_field",
FieldContext: telemetrytypes.FieldContextLog,
FieldContext: telemetrytypes.FieldContextMetric,
},
expectedResult: "",
expectedError: qbtypes.ErrColumnNotFound,

View File

@@ -1,8 +1,11 @@
package telemetrymetadata
import (
"sync"
"context"
"fmt"
"golang.org/x/sync/semaphore"
"log/slog"
"strings"
"time"
@@ -69,6 +72,11 @@ type telemetryMetaStore struct {
conditionBuilder qbtypes.ConditionBuilder
fl flagger.Flagger
jsonColumnMetadata map[telemetrytypes.Signal]map[telemetrytypes.FieldContext]telemetrytypes.JSONColumnMetadata
config Config
relatedSlotsMu sync.Mutex
relatedSlots map[valuer.UUID]*semaphore.Weighted
}
func escapeForLike(s string) string {
@@ -79,6 +87,7 @@ func NewTelemetryMetaStore(
settings factory.ProviderSettings,
telemetrystore telemetrystore.TelemetryStore,
fl flagger.Flagger,
config Config,
) telemetrytypes.MetadataStore {
metadataSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetrymetadata")
@@ -120,6 +129,8 @@ func NewTelemetryMetaStore(
fl: fl,
fm: fm,
conditionBuilder: conditionBuilder,
config: config,
relatedSlots: map[valuer.UUID]*semaphore.Weighted{},
}
return t
@@ -1355,6 +1366,9 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
return nil, true, nil
}
cfg := t.config.RelatedValues
startMs, endMs, truncated := relatedValuesWindow(fieldValueSelector.StartUnixMilli, fieldValueSelector.EndUnixMilli, time.Now(), cfg.MaxWindow, telemetrytypes.MetadataBucketMilli(fieldValueSelector.Signal))
key := &telemetrytypes.TelemetryFieldKey{
Name: fieldValueSelector.Name,
Signal: fieldValueSelector.Signal,
@@ -1362,95 +1376,76 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
FieldDataType: fieldValueSelector.FieldDataType,
}
selectColumn, err := t.fm.FieldFor(ctx, orgID, 0, 0, key)
// the keys of the existing query and the requested key are looked up by
// their exact names, in the requested signal
keySelectors := querybuilder.QueryStringToKeysSelectors(fieldValueSelector.ExistingQuery)
keySelectors = append(keySelectors, &telemetrytypes.FieldKeySelector{
Name: fieldValueSelector.Name,
FieldContext: fieldValueSelector.FieldContext,
FieldDataType: fieldValueSelector.FieldDataType,
})
for _, keySelector := range keySelectors {
keySelector.Signal = fieldValueSelector.Signal
keySelector.SelectorMatchType = telemetrytypes.FieldSelectorMatchTypeExact
}
keys, _, err := t.GetKeysMulti(ctx, orgID, keySelectors)
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)
return nil, false, err
}
target := resolveRelatedTarget(fieldValueSelector, keys[fieldValueSelector.Name])
signal := relatedValuesSignal(fieldValueSelector.Signal, target)
// every query from here on, the existence checks included, runs in the
// org's related-values slot and under the related-values bounds
if slot := t.relatedValuesSlot(orgID); slot != nil {
if err := slot.Acquire(ctx, 1); err != nil {
return nil, false, ErrFailedToGetRelatedValues
}
defer slot.Release(1)
}
ctx = t.relatedValuesQueryContext(ctx, true)
if t.relatedValuesFilterAbsent(ctx, orgID, fieldValueSelector, signal, keys, startMs, endMs) {
return []string{}, !truncated, nil
}
selectColumn := t.relatedValuesSelectColumn(ctx, orgID, fieldValueSelector.Name, target)
sb := sqlbuilder.Select("DISTINCT " + selectColumn).From(t.relatedMetadataDBName + "." + t.relatedMetadataTblName)
if len(fieldValueSelector.ExistingQuery) != 0 {
keySelectors := querybuilder.QueryStringToKeysSelectors(fieldValueSelector.ExistingQuery)
for _, keySelector := range keySelectors {
keySelector.Signal = fieldValueSelector.Signal
}
keys, _, err := t.GetKeysMulti(ctx, orgID, keySelectors)
if err != nil {
return nil, false, err
}
whereClause, err := querybuilder.PrepareWhereClause(fieldValueSelector.ExistingQuery, querybuilder.FilterExprVisitorOpts{
Context: ctx,
Logger: t.logger,
FieldMapper: t.fm,
ConditionBuilder: t.conditionBuilder,
FieldKeys: keys,
})
if err != nil {
t.logger.WarnContext(ctx, "error parsing existing query for related values", errors.Attr(err))
}
if !whereClause.IsEmpty() {
sb.AddWhereClause(whereClause.WhereClause)
}
whereClause, err := querybuilder.PrepareWhereClause(fieldValueSelector.ExistingQuery, querybuilder.FilterExprVisitorOpts{
Context: ctx,
Logger: t.logger,
FieldMapper: t.fm,
ConditionBuilder: t.conditionBuilder,
FieldKeys: keys,
})
if err != nil {
t.logger.WarnContext(ctx, "error parsing existing query for related values", errors.Attr(err))
}
if !whereClause.IsEmpty() {
sb.AddWhereClause(whereClause.WhereClause)
}
if fieldValueSelector.StartUnixMilli != 0 {
sb.Where(sb.GE("unix_milli", fieldValueSelector.StartUnixMilli))
}
sb.Where(sb.GE("unix_milli", startMs))
sb.Where(sb.LE("unix_milli", endMs))
if fieldValueSelector.EndUnixMilli != 0 {
sb.Where(sb.LE("unix_milli", fieldValueSelector.EndUnixMilli))
}
// scope to the requested signal's rows;
if fieldValueSelector.Signal != telemetrytypes.SignalUnspecified {
sb.Where(sb.E("data_source", fieldValueSelector.Signal.StringValue()))
if signal != telemetrytypes.SignalUnspecified {
sb.Where(sb.E("data_source", signal.StringValue()))
}
if fieldValueSelector.Value != "" {
// the search text is matched in every map the key is read from
var conds []string
if fieldValueSelector.FieldContext != telemetrytypes.FieldContextAttribute &&
fieldValueSelector.FieldContext != telemetrytypes.FieldContextResource {
origContext := key.FieldContext
// search on attributes
key.FieldContext = telemetrytypes.FieldContextAttribute
attrConds, _, err := t.conditionBuilder.ConditionFor(ctx, orgID, 0, 0, key, map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}}, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorContains, fieldValueSelector.Value, sb)
if err == nil {
conds = append(conds, attrConds...)
}
// search on resource
key.FieldContext = telemetrytypes.FieldContextResource
resourceConds, _, err := t.conditionBuilder.ConditionFor(ctx, orgID, 0, 0, key, map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}}, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorContains, fieldValueSelector.Value, sb)
if err == nil {
conds = append(conds, resourceConds...)
}
key.FieldContext = origContext
} else {
keyConds, _, err := t.conditionBuilder.ConditionFor(ctx, orgID, 0, 0, key, map[string][]*telemetrytypes.TelemetryFieldKey{key.Name: {key}}, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorContains, fieldValueSelector.Value, sb)
for _, fieldContext := range relatedValuesContexts(fieldValueSelector.Name, target) {
searchKey := *key
searchKey.FieldContext = fieldContext
keyConds, _, err := t.conditionBuilder.ConditionFor(ctx, orgID, 0, 0, &searchKey, map[string][]*telemetrytypes.TelemetryFieldKey{searchKey.Name: {&searchKey}}, qbtypes.ConditionBuilderOptions{}, qbtypes.FilterOperatorContains, fieldValueSelector.Value, sb)
if err == nil {
conds = append(conds, keyConds...)
}
}
if len(conds) != 0 {
// the key may sit in the resource or attribute map (or both), so OR the
// two conditions — match if the key's value in either map contains searchText.
sb.Where(sb.Or(conds...))
}
}
@@ -1466,6 +1461,7 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
t.logger.DebugContext(ctx, "query for related values", slog.String("query", query), slog.Any("args", args))
started := time.Now()
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return nil, false, ErrFailedToGetRelatedValues
@@ -1490,8 +1486,11 @@ func (t *telemetryMetaStore) getRelatedValues(ctx context.Context, orgID valuer.
}
}
// hit the limit?
complete := rowCount <= limit
// hit the limit, the window was clamped, or the time bound cut the scan short?
complete := rowCount <= limit && !truncated
if cfg.MaxExecutionTime > 0 && time.Since(started) >= cfg.MaxExecutionTime {
complete = false
}
return attributeValues, complete, nil
}

View File

@@ -38,6 +38,7 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
instrumentationtest.New().ToProviderSettings(),
mockTelemetryStore,
flaggertest.New(t),
NewConfig(),
)
lookupKeys := []telemetrytypes.MetricMetadataLookupKey{

View File

@@ -0,0 +1,281 @@
package telemetrymetadata
import (
"context"
"fmt"
"log/slog"
"strconv"
"strings"
"time"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/huandu/go-sqlbuilder"
"golang.org/x/sync/semaphore"
)
// relatedValuesWindow returns the window of the related-values query. The
// requested window is clamped to maxWindow (an unset start becomes end minus
// maxWindow, an unset end becomes now) and the start is floored to the bucket
// that holds it. truncated reports that the clamp moved the start, so the
// result covers less than the request asked for.
func relatedValuesWindow(startMs, endMs int64, now time.Time, maxWindow time.Duration, bucketMs int64) (int64, int64, bool) {
if endMs == 0 {
endMs = now.UnixMilli()
}
truncated := false
if maxWindow > 0 && (startMs == 0 || endMs-startMs > maxWindow.Milliseconds()) {
startMs = endMs - maxWindow.Milliseconds()
truncated = true
}
if bucketMs > 0 {
startMs -= startMs % bucketMs
}
return startMs, endMs, truncated
}
// relatedTarget is where the requested key was found: the metadata contexts
// it resolves to and the signals it was seen in.
type relatedTarget struct {
contexts []telemetrytypes.FieldContext
signals map[telemetrytypes.Signal]struct{}
}
func metadataContextFor(fieldContext telemetrytypes.FieldContext) (telemetrytypes.FieldContext, bool) {
switch fieldContext {
case telemetrytypes.FieldContextResource, telemetrytypes.FieldContextAttribute,
telemetrytypes.FieldContextSpan, telemetrytypes.FieldContextLog:
return fieldContext, true
}
return telemetrytypes.FieldContextUnspecified, false
}
// resolveRelatedTarget derives the target from the requested context when one
// is given, otherwise from the keys the lookup returned for the name.
func resolveRelatedTarget(selector *telemetrytypes.FieldValueSelector, keys []*telemetrytypes.TelemetryFieldKey) relatedTarget {
target := relatedTarget{signals: map[telemetrytypes.Signal]struct{}{}}
if fieldContext, ok := metadataContextFor(selector.FieldContext); ok {
target.contexts = []telemetrytypes.FieldContext{fieldContext}
}
seen := map[telemetrytypes.FieldContext]struct{}{}
for _, key := range keys {
fieldContext, ok := metadataContextFor(key.FieldContext)
if !ok {
continue
}
if selector.FieldContext != telemetrytypes.FieldContextUnspecified && fieldContext != selector.FieldContext {
continue
}
if key.Signal != telemetrytypes.SignalUnspecified {
target.signals[key.Signal] = struct{}{}
}
if selector.FieldContext == telemetrytypes.FieldContextUnspecified {
if _, added := seen[fieldContext]; !added {
seen[fieldContext] = struct{}{}
target.contexts = append(target.contexts, fieldContext)
}
}
}
return target
}
// relatedValuesContexts returns the maps the key is read from and searched
// in: the maps it resolved to, or all three when nothing is known about it.
// The span name is also read from attributes for rows written before the
// intrinsic map existed.
func relatedValuesContexts(name string, target relatedTarget) []telemetrytypes.FieldContext {
contexts := target.contexts
if len(contexts) == 0 {
return []telemetrytypes.FieldContext{telemetrytypes.FieldContextSpan, telemetrytypes.FieldContextResource, telemetrytypes.FieldContextAttribute}
}
if name == "name" && len(contexts) == 1 && (contexts[0] == telemetrytypes.FieldContextSpan || contexts[0] == telemetrytypes.FieldContextLog) {
return append([]telemetrytypes.FieldContext{contexts[0]}, telemetrytypes.FieldContextAttribute)
}
return contexts
}
// relatedValuesSelectColumn returns the expression that reads the key from
// the metadata table: the single map it is read from, or a multiIf over
// several.
func (t *telemetryMetaStore) relatedValuesSelectColumn(ctx context.Context, orgID valuer.UUID, name string, target relatedTarget) string {
fieldFor := func(fieldContext telemetrytypes.FieldContext) string {
column, _ := t.fm.FieldFor(ctx, orgID, 0, 0, &telemetrytypes.TelemetryFieldKey{
Name: name,
FieldContext: fieldContext,
FieldDataType: telemetrytypes.FieldDataTypeString,
})
return column
}
contexts := relatedValuesContexts(name, target)
if len(contexts) == 1 {
return fieldFor(contexts[0])
}
args := []string{}
for _, fieldContext := range contexts[:len(contexts)-1] {
column := fieldFor(fieldContext)
args = append(args, fmt.Sprintf("notEmpty(%s), %s", column, column))
}
args = append(args, fieldFor(contexts[len(contexts)-1]))
return fmt.Sprintf("multiIf(%s)", strings.Join(args, ", "))
}
// relatedValuesSignal picks the data source to scan for a request without a
// signal: the one signal the key was seen in. A key seen in several signals
// keeps scanning all of them, since the values of one signal need not cover
// the others.
func relatedValuesSignal(requested telemetrytypes.Signal, target relatedTarget) telemetrytypes.Signal {
if requested != telemetrytypes.SignalUnspecified {
return requested
}
if len(target.signals) == 1 {
for signal := range target.signals {
return signal
}
}
return telemetrytypes.SignalUnspecified
}
// relatedValuesQueryContext applies the related-values bounds to the queries
// run with the context: the execution time, checked from the start, the
// thread count and the read buffer. With partialOnTimeout the values found by
// the bound are returned; without it the query fails on the bound.
func (t *telemetryMetaStore) relatedValuesQueryContext(ctx context.Context, partialOnTimeout bool) context.Context {
cfg := t.config.RelatedValues
settings := map[string]any{}
if cfg.MaxExecutionTime > 0 {
settings["max_execution_time"] = cfg.MaxExecutionTime.Seconds()
settings["timeout_before_checking_execution_speed"] = 0
if partialOnTimeout {
settings["timeout_overflow_mode"] = "break"
}
}
if cfg.MaxThreads > 0 {
settings["max_threads"] = cfg.MaxThreads
}
if cfg.ReadBufferSize > 0 {
settings["max_read_buffer_size_local_fs"] = cfg.ReadBufferSize
}
return ctxtypes.SetClickhouseSettings(ctx, settings)
}
// relatedValuesFilterAbsent reports whether an equality term of the existing
// query names a value the key does not have in the window, in which case no
// row can match and the scan is skipped. Only conjunctions of traces or logs
// filters on resource or attribute keys are checked, up to the configured
// number of terms.
func (t *telemetryMetaStore) relatedValuesFilterAbsent(ctx context.Context, orgID valuer.UUID, selector *telemetrytypes.FieldValueSelector, signal telemetrytypes.Signal, keys map[string][]*telemetrytypes.TelemetryFieldKey, startMs, endMs int64) bool {
maxChecks := t.config.RelatedValues.MaxExistenceChecks
if maxChecks <= 0 || (signal != telemetrytypes.SignalTraces && signal != telemetrytypes.SignalLogs) {
return false
}
terms, ok := querybuilder.QueryStringEqualityTerms(selector.ExistingQuery)
if !ok {
return false
}
checked := 0
for _, term := range terms {
if checked >= maxChecks {
break
}
fieldContext, ok := t.singleMapContext(term.Key, keys[term.Key.Name], signal)
if !ok {
continue
}
checked++
// the check fails on the time bound instead of returning a partial
// result, and a failed check proves nothing
exists, err := t.tagValueExists(t.relatedValuesQueryContext(ctx, false), signal, term.Key.Name, fieldContext, term.Value, startMs)
if err != nil {
t.logger.DebugContext(ctx, "failed to check filter value existence", slog.String("key", term.Key.Name), errors.Attr(err))
continue
}
if !exists {
t.logger.DebugContext(ctx, "related values skipped: filter value absent in window", slog.String("key", term.Key.Name), slog.String("value", term.Value))
return true
}
}
return false
}
// tagValueExists reports whether the tag table of the signal holds the value
// for the key, as a string or as the number it parses to, since the window
// start.
func (t *telemetryMetaStore) tagValueExists(ctx context.Context, signal telemetrytypes.Signal, name string, fieldContext telemetrytypes.FieldContext, value string, startMs int64) (bool, error) {
var table string
switch signal {
case telemetrytypes.SignalTraces:
table = t.tracesDBName + "." + t.tracesFieldsTblName
case telemetrytypes.SignalLogs:
table = t.logsDBName + "." + t.logsFieldsTblName
default:
return true, nil
}
sb := sqlbuilder.Select("1").From(table)
sb.Where(sb.E("tag_key", name))
sb.Where(sb.E("tag_type", fieldContext.TagType()))
valueConds := []string{sb.E("string_value", value)}
if number, err := strconv.ParseFloat(value, 64); err == nil {
valueConds = append(valueConds, sb.E("number_value", number))
}
sb.Where(sb.Or(valueConds...))
if startMs > 0 {
sb.Where(sb.GE("unix_milli", startMs))
}
sb.Limit(1)
query, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
rows, err := t.telemetrystore.ClickhouseDB().Query(ctx, query, args...)
if err != nil {
return false, err
}
defer rows.Close()
exists := rows.Next()
return exists, rows.Err()
}
// singleMapContext returns the resource or attribute context a filter key
// resolves to in the signal, when it resolves to exactly one.
func (t *telemetryMetaStore) singleMapContext(key *telemetrytypes.TelemetryFieldKey, resolved []*telemetrytypes.TelemetryFieldKey, signal telemetrytypes.Signal) (telemetrytypes.FieldContext, bool) {
if key.FieldContext == telemetrytypes.FieldContextResource || key.FieldContext == telemetrytypes.FieldContextAttribute {
return key.FieldContext, true
}
if key.FieldContext != telemetrytypes.FieldContextUnspecified {
return telemetrytypes.FieldContextUnspecified, false
}
found := telemetrytypes.FieldContextUnspecified
for _, candidate := range resolved {
if candidate.Signal != signal {
continue
}
if candidate.FieldContext != telemetrytypes.FieldContextResource && candidate.FieldContext != telemetrytypes.FieldContextAttribute {
return telemetrytypes.FieldContextUnspecified, false
}
if found != telemetrytypes.FieldContextUnspecified && found != candidate.FieldContext {
return telemetrytypes.FieldContextUnspecified, false
}
found = candidate.FieldContext
}
return found, found != telemetrytypes.FieldContextUnspecified
}
// relatedValuesSlot returns the org's concurrency limiter for related-values
// queries, or nil when no limit is configured.
func (t *telemetryMetaStore) relatedValuesSlot(orgID valuer.UUID) *semaphore.Weighted {
limit := t.config.RelatedValues.MaxConcurrency
if limit <= 0 {
return nil
}
t.relatedSlotsMu.Lock()
defer t.relatedSlotsMu.Unlock()
slot, ok := t.relatedSlots[orgID]
if !ok {
slot = semaphore.NewWeighted(int64(limit))
t.relatedSlots[orgID] = slot
}
return slot
}

View File

@@ -0,0 +1,164 @@
package telemetrymetadata
import (
"context"
"testing"
"time"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
"github.com/stretchr/testify/assert"
)
func TestRelatedValuesWindow(t *testing.T) {
// 2026-09-07T10:30:00Z
now := time.UnixMilli(1788777000000)
sixHours := int64(6 * 60 * 60 * 1000)
day := 24 * time.Hour
tests := []struct {
name string
start int64
end int64
maxWindow time.Duration
wantStart int64
wantEnd int64
wantTruncated bool
}{
{
name: "window inside the clamp keeps its end and is floored to the bucket",
start: now.Add(-3 * time.Hour).UnixMilli(),
end: now.UnixMilli(),
maxWindow: day,
wantStart: 1788760800000, // 06:00
wantEnd: now.UnixMilli(),
},
{
name: "window over the clamp keeps its end and loses its start",
start: now.Add(-30 * day).UnixMilli(),
end: now.UnixMilli(),
maxWindow: day,
wantStart: 1788674400000, // the day before, 10:30 floored to 06:00
wantEnd: now.UnixMilli(),
wantTruncated: true,
},
{
name: "unset bounds mean the last clamp window up to now",
maxWindow: day,
wantStart: 1788674400000,
wantEnd: now.UnixMilli(),
wantTruncated: true,
},
{
name: "no clamp keeps the requested start",
start: now.Add(-30 * day).UnixMilli(),
end: now.UnixMilli(),
wantStart: now.Add(-30*day).UnixMilli() - now.Add(-30*day).UnixMilli()%sixHours,
wantEnd: now.UnixMilli(),
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
start, end, truncated := relatedValuesWindow(tt.start, tt.end, now, tt.maxWindow, sixHours)
assert.Equal(t, tt.wantStart, start)
assert.Equal(t, tt.wantEnd, end)
assert.Equal(t, tt.wantTruncated, truncated)
})
}
}
func TestRelatedValuesSelectColumnFollowsTheResolvedContext(t *testing.T) {
store := &telemetryMetaStore{fm: NewFieldMapper()}
ctx := context.Background()
resourceKey := &telemetrytypes.TelemetryFieldKey{Name: "service.name", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextResource}
attributeKey := &telemetrytypes.TelemetryFieldKey{Name: "service.name", Signal: telemetrytypes.SignalLogs, FieldContext: telemetrytypes.FieldContextAttribute}
spanKey := &telemetrytypes.TelemetryFieldKey{Name: "http_method", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextSpan}
tests := []struct {
name string
selector *telemetrytypes.FieldValueSelector
keys []*telemetrytypes.TelemetryFieldKey
want string
signals int
}{
{
name: "resource key reads the resource map only",
selector: &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "service.name"}},
keys: []*telemetrytypes.TelemetryFieldKey{resourceKey},
want: "resource_attributes['service.name']",
signals: 1,
},
{
name: "requested context wins over the lookup",
selector: &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "service.name", FieldContext: telemetrytypes.FieldContextAttribute}},
keys: []*telemetrytypes.TelemetryFieldKey{resourceKey, attributeKey},
want: "attributes['service.name']",
signals: 1,
},
{
name: "key in two maps reads both",
selector: &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "service.name"}},
keys: []*telemetrytypes.TelemetryFieldKey{resourceKey, attributeKey},
want: "multiIf(notEmpty(resource_attributes['service.name']), resource_attributes['service.name'], attributes['service.name'])",
signals: 2,
},
{
name: "span field reads the intrinsic map",
selector: &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "http_method"}},
keys: []*telemetrytypes.TelemetryFieldKey{spanKey},
want: "intrinsic_attributes['http_method']",
signals: 1,
},
{
name: "span name also reads attributes for rows written before the intrinsic map",
selector: &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "name", FieldContext: telemetrytypes.FieldContextSpan}},
want: "multiIf(notEmpty(intrinsic_attributes['name']), intrinsic_attributes['name'], attributes['name'])",
},
{
name: "unknown key reads every map",
selector: &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "mystery"}},
want: "multiIf(notEmpty(intrinsic_attributes['mystery']), intrinsic_attributes['mystery'], notEmpty(resource_attributes['mystery']), resource_attributes['mystery'], attributes['mystery'])",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
target := resolveRelatedTarget(tt.selector, tt.keys)
assert.Equal(t, tt.want, store.relatedValuesSelectColumn(ctx, valuer.UUID{}, tt.selector.Name, target))
assert.Len(t, target.signals, tt.signals)
})
}
}
func TestRelatedValuesSignalScopesToTheOnlySignalSeen(t *testing.T) {
selector := &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "service.name"}}
target := resolveRelatedTarget(selector, []*telemetrytypes.TelemetryFieldKey{
{Name: "service.name", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextResource},
})
assert.Equal(t, telemetrytypes.SignalTraces, relatedValuesSignal(telemetrytypes.SignalUnspecified, target))
assert.Equal(t, telemetrytypes.SignalLogs, relatedValuesSignal(telemetrytypes.SignalLogs, target), "a requested signal is kept")
target = resolveRelatedTarget(selector, []*telemetrytypes.TelemetryFieldKey{
{Name: "service.name", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextResource},
{Name: "service.name", Signal: telemetrytypes.SignalLogs, FieldContext: telemetrytypes.FieldContextResource},
})
assert.Equal(t, telemetrytypes.SignalUnspecified, relatedValuesSignal(telemetrytypes.SignalUnspecified, target), "a key seen in two signals scans both")
}
func TestRelatedValuesContextsAreSharedBySelectAndSearch(t *testing.T) {
selector := &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "http_method"}}
target := resolveRelatedTarget(selector, []*telemetrytypes.TelemetryFieldKey{
{Name: "http_method", Signal: telemetrytypes.SignalTraces, FieldContext: telemetrytypes.FieldContextSpan},
})
assert.Equal(t, []telemetrytypes.FieldContext{telemetrytypes.FieldContextSpan}, relatedValuesContexts("http_method", target), "a span field is searched in the intrinsic map only")
selector = &telemetrytypes.FieldValueSelector{FieldKeySelector: &telemetrytypes.FieldKeySelector{Name: "name", FieldContext: telemetrytypes.FieldContextSpan}}
assert.Equal(t, []telemetrytypes.FieldContext{telemetrytypes.FieldContextSpan, telemetrytypes.FieldContextAttribute}, relatedValuesContexts("name", resolveRelatedTarget(selector, nil)))
}
func TestConfigValidate(t *testing.T) {
cfg := NewConfig()
assert.NoError(t, cfg.Validate())
cfg.RelatedValues.MaxConcurrency = -1
assert.Error(t, cfg.Validate())
}

View File

@@ -72,6 +72,10 @@ func (h *provider) BeforeQuery(ctx context.Context, _ *telemetrystore.QueryEvent
settings["result_overflow_mode"] = ctx.Value("result_overflow_mode")
}
for name, value := range ctxtypes.ClickhouseSettingsFromContext(ctx) {
settings[name] = value
}
// TODO(srikanthccv): enable it when the "Cannot read all data" issue is fixed
// https://github.com/ClickHouse/ClickHouse/issues/82283
settings["secondary_indices_enable_bulk_filtering"] = false

View File

@@ -6,9 +6,23 @@ type ctxKey string
const (
ClickhouseContextMaxThreadsKey ctxKey = "clickhouse_max_threads"
ClickhouseContextSettingsKey ctxKey = "clickhouse_settings"
)
// SetClickhouseMaxThreads stores the max threads value in context.
func SetClickhouseMaxThreads(ctx context.Context, maxThreads int) context.Context {
return context.WithValue(ctx, ClickhouseContextMaxThreadsKey, maxThreads)
}
// SetClickhouseSettings stores query settings in context that are applied on
// top of the configured defaults for the queries run with it.
func SetClickhouseSettings(ctx context.Context, settings map[string]any) context.Context {
return context.WithValue(ctx, ClickhouseContextSettingsKey, settings)
}
// ClickhouseSettingsFromContext returns the query settings stored with
// SetClickhouseSettings, or nil.
func ClickhouseSettingsFromContext(ctx context.Context) map[string]any {
settings, _ := ctx.Value(ClickhouseContextSettingsKey).(map[string]any)
return settings
}

View File

@@ -303,13 +303,24 @@ type PostableFieldValueParams struct {
ExistingQuery string `query:"existingQuery"`
}
// Metadata rows are written once per six-hour bucket at the bucket start, for
// every signal (the collector's default bucket). A window start is floored to
// the bucket holding it so that bucket's rows are not left out.
const metadataBucketMilli int64 = 6 * 60 * 60 * 1000
// MetadataBucketMilli returns the write bucket of the metadata table for the
// signal. Every signal uses the six-hour bucket until the collector writes a
// longer one for metrics.
func MetadataBucketMilli(_ Signal) int64 {
return metadataBucketMilli
}
func NewFieldKeySelectorFromPostableFieldKeysParams(params PostableFieldKeysParams) *FieldKeySelector {
var req FieldKeySelector
if params.StartUnixMilli != 0 {
req.StartUnixMilli = params.StartUnixMilli
// Round down to the nearest 6 hours (21600000 milliseconds)
req.StartUnixMilli -= req.StartUnixMilli % 21600000
req.StartUnixMilli -= req.StartUnixMilli % MetadataBucketMilli(params.Signal)
}
if params.EndUnixMilli != 0 {

View File

@@ -422,3 +422,32 @@ func TestNewFieldKeySelectorFromPostableFieldKeysParamsMetricNamespace(t *testin
t.Fatalf("expected metric name to remain empty, got %q", selector.MetricContext.MetricName)
}
}
func TestNewFieldKeySelectorFromPostableFieldKeysParamsFloorsStartToTheMetadataBucket(t *testing.T) {
// 2026-09-07T10:30:00Z
start := int64(1788777000000)
sixHourFloor := int64(1788760800000)
tests := []struct {
signal Signal
want int64
}{
{signal: SignalTraces, want: sixHourFloor},
{signal: SignalLogs, want: sixHourFloor},
{signal: SignalMetrics, want: sixHourFloor},
{signal: SignalUnspecified, want: sixHourFloor},
}
for _, tt := range tests {
t.Run(tt.signal.StringValue(), func(t *testing.T) {
selector := NewFieldKeySelectorFromPostableFieldKeysParams(PostableFieldKeysParams{Signal: tt.signal, StartUnixMilli: start})
if selector.StartUnixMilli != tt.want {
t.Fatalf("expected start %d, got %d", tt.want, selector.StartUnixMilli)
}
})
}
selector := NewFieldKeySelectorFromPostableFieldKeysParams(PostableFieldKeysParams{Signal: SignalTraces})
if selector.StartUnixMilli != 0 {
t.Fatalf("expected an unset start to stay unset, got %d", selector.StartUnixMilli)
}
}