mirror of
https://github.com/SigNoz/signoz.git
synced 2026-09-15 16:00:41 +01:00
Compare commits
6 Commits
feat/sqlco
...
feat/relat
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e0de0fcf6f | ||
|
|
c456041089 | ||
|
|
96128e9f87 | ||
|
|
16a0c2375d | ||
|
|
14df7ff810 | ||
|
|
ca1e063543 |
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
82
pkg/telemetrymetadata/config.go
Normal file
82
pkg/telemetrymetadata/config.go
Normal 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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -38,6 +38,7 @@ func TestGetFirstSeenFromMetricMetadata(t *testing.T) {
|
||||
instrumentationtest.New().ToProviderSettings(),
|
||||
mockTelemetryStore,
|
||||
flaggertest.New(t),
|
||||
NewConfig(),
|
||||
)
|
||||
|
||||
lookupKeys := []telemetrytypes.MetricMetadataLookupKey{
|
||||
|
||||
281
pkg/telemetrymetadata/related_values.go
Normal file
281
pkg/telemetrymetadata/related_values.go
Normal 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
|
||||
}
|
||||
164
pkg/telemetrymetadata/related_values_test.go
Normal file
164
pkg/telemetrymetadata/related_values_test.go
Normal 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())
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user