fix: address comments

This commit is contained in:
nityanandagohain
2026-07-14 18:53:30 +05:30
parent d502d12ac3
commit 31efe177a4
14 changed files with 1463 additions and 1101 deletions

View File

@@ -493,6 +493,10 @@ func (q *querier) QueryRawStream(ctx context.Context, orgID valuer.UUID, req *qb
if query.Type == qbtypes.QueryTypeBuilder {
switch spec := query.Spec.(type) {
case qbtypes.QueryBuilderQuery[qbtypes.LogAggregation]:
if spec.Source == telemetrytypes.SourceAI {
client.Error <- errors.NewInvalidInputf(errors.CodeInvalidInput, "source \"ai\" is only supported for the traces signal, not logs")
return
}
event.FilterApplied = spec.Filter != nil && spec.Filter.Expression != ""
default:
// return if it's not log aggregation

View File

@@ -101,8 +101,6 @@ func newProvider(
aiTraceStmtBuilder := telemetryai.NewAITraceStatementBuilder(
settings,
telemetryMetadataStore,
traceFieldMapper,
traceConditionBuilder,
aiBaseCondition,
traceStmtBuilder,
telemetryStore,

View File

@@ -114,20 +114,27 @@ func (s *filterSplitter) route(atom antlr.ParserRuleContext) {
}
// classifyKeys reports whether a subtree references trace-level and/or span-level keys.
// A key is trace-level only when it has no explicit field context and its name — after
// the optional user-facing `trace.` prefix is stripped — is a known aggregate. Any
// explicit context (`span.`, `resource.`, …) is span-level.
// A key is trace-level when it has no explicit field context and its name — after the
// optional user-facing `trace.` prefix is stripped — is a known aggregate, or when it
// carries the trace field context explicitly (`tracefield.`, which Normalize parses
// into FieldContextTrace; no span field mapper resolves that context, so it can only
// mean a trace-level aggregate — an unknown name is then rejected by the aggregate
// validation with a targeted error instead of failing as an unknown span field). Any
// other explicit context (`span.`, `resource.`, …) is span-level.
func classifyKeys(node antlr.Tree, aggregateNames map[string]struct{}) (isTrace, isSpan bool) {
kc, ok := node.(*grammar.KeyContext)
if ok {
key := telemetrytypes.GetFieldKeyFromKeyText(kc.GetText())
if key.FieldContext == telemetrytypes.FieldContextUnspecified {
switch key.FieldContext {
case telemetrytypes.FieldContextUnspecified:
// `trace.` is the user-facing prefix for trace-level aggregates. It is not a
// registered field context, so it stays on the name; strip it before matching.
name := strings.TrimPrefix(key.Name, telemetrytypes.FieldContextTrace.StringValue()+".")
_, isTrace = aggregateNames[name]
isSpan = !isTrace
} else {
case telemetrytypes.FieldContextTrace:
isTrace = true
default:
isSpan = true
}
return

View File

@@ -45,10 +45,19 @@ func TestSplitFilterForAggregates(t *testing.T) {
having: "trace.completion_tokens > 1000",
},
{
// `tracefield.` is not supported: it has an explicit context, so it is span-level.
name: "tracefield prefix is span-level",
query: "tracefield.completion_tokens > 1000",
span: "tracefield.completion_tokens > 1000",
// `tracefield.` is the explicit trace field context (Normalize parses it into
// FieldContextTrace), so it marks a trace-level aggregate like `trace.`.
name: "agg only tracefield prefix",
query: "tracefield.completion_tokens > 1000",
having: "tracefield.completion_tokens > 1000",
},
{
// an unknown name under the explicit trace context still routes trace-level,
// so the aggregate validation rejects it with a targeted error instead of the
// span path failing on an unknown field.
name: "unknown agg under tracefield prefix stays trace-level",
query: "tracefield.not_an_aggregate > 1000",
having: "tracefield.not_an_aggregate > 1000",
},
{

View File

@@ -3,31 +3,18 @@ package telemetryai
import (
"strings"
scopedtraces "github.com/SigNoz/signoz/pkg/telemetryscopedtraces"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// BaseConditionProvider defines which spans (and traces) are in scope, decoupled
// from the query. Swap it to redefine "an AI trace" without touching the builder.
//
// It only *declares* the gate (a grammar expression + its field keys); the builder
// resolves those keys through the field mapper, so all attribute access is
// materialization/evolution aware — no hardcoded map lookups.
type BaseConditionProvider interface {
// FilterExpression is the grammar-level (EXISTS) gate, resolved via the visitor.
FilterExpression() string
// FieldKeys are the gate's keys: registered in metadata and used to build the
// per-span mask (OR of resolved EXISTS conditions).
FieldKeys() []*telemetrytypes.TelemetryFieldKey
}
// genAIBaseConditionProvider: an AI trace has >=1 gen_ai LLM, tool, or agent span.
type genAIBaseConditionProvider struct {
keys []string
}
var _ BaseConditionProvider = (*genAIBaseConditionProvider)(nil)
var _ scopedtraces.BaseConditionProvider = (*genAIBaseConditionProvider)(nil)
func NewGenAIBaseConditionProvider() BaseConditionProvider {
func NewGenAIBaseConditionProvider() scopedtraces.BaseConditionProvider {
return &genAIBaseConditionProvider{
keys: []string{telemetrytypes.GenAIRequestModel, telemetrytypes.GenAIToolName, telemetrytypes.GenAIAgentName},
}
@@ -42,8 +29,8 @@ func (p *genAIBaseConditionProvider) FilterExpression() string {
}
func (p *genAIBaseConditionProvider) FieldKeys() []*telemetrytypes.TelemetryFieldKey {
// Definitions (context/type) come from GenAIFieldDefinitions so they can't drift
// from the canonical semconv keys; copy to take the address.
// Definitions come from GenAIFieldDefinitions so they can't drift from the
// canonical semconv keys; copy to take the address.
keys := make([]*telemetrytypes.TelemetryFieldKey, 0, len(p.keys))
for _, k := range p.keys {
def := telemetrytypes.GenAIFieldDefinitions[k]

View File

@@ -1,196 +1,20 @@
package telemetryai
import (
"fmt"
"strings"
scopedtraces "github.com/SigNoz/signoz/pkg/telemetryscopedtraces"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/huandu/go-sqlbuilder"
)
// TraceColumn is one per-trace output column: an alias, whether it can be sorted
// on, and the Aggregate that renders its SQL.
type TraceColumn struct {
Alias string
// Orderable marks a column usable in ORDER BY and the aggregate filter. All-span
// aggregates (span_count, duration_nano, …) are display-only and set false.
Orderable bool
// SpanLevel marks a column that merely surfaces a real span/resource attribute
// (service.name, input/output messages). It is excluded from AggregateAliases so a
// filter on that attribute is applied span-level, not treated as a trace aggregate.
SpanLevel bool
Expr Aggregate
}
// Aggregate renders one column's SQL against the query resolver, and lists the
// attribute keys it references so the builder can pre-fetch them from metadata.
// Build one with the constructors below (Intrinsic, CountExists, Reduce, …); the
// zero value is not usable.
type Aggregate struct {
keys []*telemetrytypes.TelemetryFieldKey
render func(r aggResolver) (expr string, args []any, err error)
}
// aggResolver hands each aggregate the two field-mapper primitives it may need — an
// EXISTS predicate and a resolved value expression — plus the gen_ai gate mask. The
// builder populates it per query (see resolveColumns).
type aggResolver struct {
exists func(key *telemetrytypes.TelemetryFieldKey) (string, []any, error)
value func(key *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType) (string, []any, error)
maskExpr string
maskArgs []any
}
// AggFunc is a ClickHouse aggregate function name.
type AggFunc string
const (
AggSum AggFunc = "sum"
AggMax AggFunc = "max"
AggMin AggFunc = "min"
)
// PickDirection selects the earliest (argMin) or latest (argMax) span by ordering.
type PickDirection int
const (
PickLatest PickDirection = iota
PickEarliest
)
// Intrinsic emits fixed intrinsic-column SQL verbatim (escaped once).
func Intrinsic(text string) Aggregate {
return Aggregate{render: func(aggResolver) (string, []any, error) {
return sqlbuilder.Escape(text), nil, nil
}}
}
// CountExists renders countIf(<key> EXISTS) — counts spans carrying key.
func CountExists(key *telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: keysOf(key), render: func(r aggResolver) (string, []any, error) {
cond, args, err := r.exists(key)
return fmt.Sprintf("countIf(%s)", cond), args, err
}}
}
// Reduce renders <fn>(<value>) over a resolved numeric attribute value.
func Reduce(fn AggFunc, valueKey *telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: keysOf(valueKey), render: func(r aggResolver) (string, []any, error) {
v, args, err := r.value(valueKey, telemetrytypes.FieldDataTypeFloat64)
return fmt.Sprintf("%s(%s)", fn, v), args, err
}}
}
// ScopedReduce renders <fn>If(<valueExpr>, <gate mask>) over a fixed value expression.
func ScopedReduce(fn AggFunc, valueExpr string) Aggregate {
return Aggregate{render: func(r aggResolver) (string, []any, error) {
return fmt.Sprintf("%sIf(%s, %s)", fn, valueExpr, r.maskExpr), append([]any{}, r.maskArgs...), nil
}}
}
// ScopedToKey renders <fn>If(<valueExpr>, <scopeKey> EXISTS) — a fixed value
// aggregated only over spans carrying scopeKey (e.g. max LLM latency).
func ScopedToKey(fn AggFunc, valueExpr string, scopeKey *telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: keysOf(scopeKey), render: func(r aggResolver) (string, []any, error) {
cond, args, err := r.exists(scopeKey)
return fmt.Sprintf("%sIf(%s, %s)", fn, valueExpr, cond), args, err
}}
}
// PickBy renders argMinIf/argMaxIf(<value>, <orderExpr>, <value> EXISTS) — the value
// from the earliest/latest span that carries it.
func PickBy(valueKey *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType, orderExpr string, dir PickDirection) Aggregate {
fn := "argMaxIf"
if dir == PickEarliest {
fn = "argMinIf"
}
return Aggregate{keys: keysOf(valueKey), render: func(r aggResolver) (string, []any, error) {
v, vargs, err := r.value(valueKey, dt)
if err != nil {
return "", nil, err
}
cond, cargs, err := r.exists(valueKey)
return fmt.Sprintf("%s(%s, %s, %s)", fn, v, orderExpr, cond), append(vargs, cargs...), err
}}
}
// UniqCount renders uniqIf(<value>, <value> EXISTS) — distinct count of an attribute.
func UniqCount(valueKey *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType) Aggregate {
return Aggregate{keys: keysOf(valueKey), render: func(r aggResolver) (string, []any, error) {
v, vargs, err := r.value(valueKey, dt)
if err != nil {
return "", nil, err
}
cond, cargs, err := r.exists(valueKey)
return fmt.Sprintf("uniqIf(%s, %s)", v, cond), append(vargs, cargs...), err
}}
}
// PredicateCount renders countIf(<predicate>) over a fixed boolean predicate.
func PredicateCount(predicate string) Aggregate {
return Aggregate{render: func(aggResolver) (string, []any, error) {
return fmt.Sprintf("countIf(%s)", sqlbuilder.Escape(predicate)), nil, nil
}}
}
// SumOfKeys renders sum(<v1>) + sum(<v2>) + … over several resolved numeric attributes.
func SumOfKeys(dt telemetrytypes.FieldDataType, valueKeys ...*telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: valueKeys, render: func(r aggResolver) (string, []any, error) {
parts := make([]string, 0, len(valueKeys))
var args []any
for _, k := range valueKeys {
v, vargs, err := r.value(k, dt)
if err != nil {
return "", nil, err
}
parts = append(parts, fmt.Sprintf("sum(%s)", v))
args = append(args, vargs...)
}
return strings.Join(parts, " + "), args, nil
}}
}
func keysOf(k *telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
return []*telemetrytypes.TelemetryFieldKey{k}
}
// --- Provider --------------------------------------------------------------
// The columns a trace list computes, decoupled from selection and topology.
type ColumnProvider interface {
Columns() []TraceColumn
// DefaultOrderAlias is sorted by (desc) when the query gives no order.
DefaultOrderAlias() string
// AggregateAliases are the computed per-trace column names, used to classify a filter
// key as trace-level vs span-level. Excludes SpanLevel columns (see TraceColumn.SpanLevel).
AggregateAliases() []string
}
// CommonTraceColumns are domain-neutral intrinsic columns any trace list can reuse.
// All aggregate over every span, so none is Orderable — output-only (see TraceColumn.Orderable).
func CommonTraceColumns() []TraceColumn {
return []TraceColumn{
{Alias: "start_time", Expr: Intrinsic("min(timestamp)")},
{Alias: "end_time", Expr: Intrinsic("max(timestamp)")},
{Alias: "duration_nano", Expr: Intrinsic("(max(toUnixTimestamp64Nano(timestamp) + duration_nano) - min(toUnixTimestamp64Nano(timestamp)))")},
{Alias: "span_count", Expr: Intrinsic("count()")},
{Alias: "root_span_name", Expr: Intrinsic("anyIf(name, parent_span_id = '')")},
{Alias: "service.name", SpanLevel: true, Expr: Intrinsic("any(resource_string_service$$name)")},
}
}
// --- gen_ai provider -------------------------------------------------------
// Adds AI/LLM per-trace metrics on top of the common columns.
// genAIColumnProvider adds AI/LLM per-trace metrics on top of the common columns.
type genAIColumnProvider struct{}
var _ ColumnProvider = (*genAIColumnProvider)(nil)
var _ scopedtraces.ColumnProvider = (*genAIColumnProvider)(nil)
func NewGenAIColumnProvider() ColumnProvider {
func NewGenAIColumnProvider() scopedtraces.ColumnProvider {
return &genAIColumnProvider{}
}
func (genAIColumnProvider) Columns() []TraceColumn {
func (genAIColumnProvider) Columns() []scopedtraces.TraceColumn {
defs := telemetrytypes.GenAIFieldDefinitions
reqModel := defs[telemetrytypes.GenAIRequestModel]
toolName := defs[telemetrytypes.GenAIToolName]
@@ -201,39 +25,34 @@ func (genAIColumnProvider) Columns() []TraceColumn {
outMsg := defs[telemetrytypes.GenAIOutputMessages]
str := telemetrytypes.FieldDataTypeString
return append(CommonTraceColumns(),
// LLM calls only (request model present), not the full gate (which includes tool/agent).
TraceColumn{Alias: "llm_call_count", Orderable: true, Expr: CountExists(&reqModel)},
// tool / distinct-tool activity across the trace's tool spans.
TraceColumn{Alias: "tool_call_count", Orderable: true, Expr: CountExists(&toolName)},
TraceColumn{Alias: "distinct_tool_count", Orderable: true, Expr: UniqCount(&toolName, str)},
// tokens live only on LLM spans, so a plain sum is correct without scoping to the gate.
TraceColumn{Alias: "input_tokens", Orderable: true, Expr: Reduce(AggSum, &inTok)},
TraceColumn{Alias: "output_tokens", Orderable: true, Expr: Reduce(AggSum, &outTok)},
TraceColumn{Alias: "total_tokens", Orderable: true, Expr: SumOfKeys(telemetrytypes.FieldDataTypeFloat64, &inTok, &outTok)},
// per-span total cost attached by the SigNoz LLM pricing processor.
TraceColumn{Alias: "estimated_cost_usd", Orderable: true, Expr: Reduce(AggSum, &cost)},
// slowest single LLM call in the trace (duration over request.model spans).
// duration_nano is table-qualified so it binds to the physical column, not the
// output alias `duration_nano` (which is itself an aggregate) — else ClickHouse
// errors with "aggregate function found inside another aggregate function".
TraceColumn{Alias: "max_llm_latency_ns", Orderable: true, Expr: ScopedToKey(AggMax, spanTable()+".duration_nano", &reqModel)},
// error spans across the whole trace (any span), so display-only.
TraceColumn{Alias: "error_count", Expr: PredicateCount("has_error = true")},
// last gen_ai span (LLM/tool/agent), so scope the max to the gate mask.
TraceColumn{Alias: "last_activity_time", Orderable: true, Expr: ScopedReduce(AggMax, "timestamp")},
// display previews: first call's input (the prompt), last call's output (the answer).
TraceColumn{Alias: "input", SpanLevel: true, Expr: PickBy(&inMsg, str, "timestamp", PickEarliest)},
TraceColumn{Alias: "output", SpanLevel: true, Expr: PickBy(&outMsg, str, "timestamp", PickLatest)},
return append(scopedtraces.CommonTraceColumns(),
// LLM calls only (request model present), not the full gate.
scopedtraces.TraceColumn{Alias: "llm_call_count", Orderable: true, Expr: scopedtraces.CountExists(&reqModel)},
scopedtraces.TraceColumn{Alias: "tool_call_count", Orderable: true, Expr: scopedtraces.CountExists(&toolName)},
scopedtraces.TraceColumn{Alias: "distinct_tool_count", Orderable: true, Expr: scopedtraces.UniqCount(&toolName, str)},
// tokens live only on LLM spans, so a plain sum needs no gate scoping.
scopedtraces.TraceColumn{Alias: "input_tokens", Orderable: true, Expr: scopedtraces.Reduce(scopedtraces.AggSum, &inTok)},
scopedtraces.TraceColumn{Alias: "output_tokens", Orderable: true, Expr: scopedtraces.Reduce(scopedtraces.AggSum, &outTok)},
scopedtraces.TraceColumn{Alias: "total_tokens", Orderable: true, Expr: scopedtraces.SumOfKeys(telemetrytypes.FieldDataTypeFloat64, &inTok, &outTok)},
// per-span cost attached by the SigNoz LLM pricing processor.
scopedtraces.TraceColumn{Alias: "estimated_cost_usd", Orderable: true, Expr: scopedtraces.Reduce(scopedtraces.AggSum, &cost)},
// slowest single LLM call in the trace.
scopedtraces.TraceColumn{Alias: "max_llm_latency_ns", Orderable: true, Expr: scopedtraces.ScopedToKeyColumn(scopedtraces.AggMax, "duration_nano", &reqModel)},
// errors across the whole trace (any span), so display-only.
scopedtraces.TraceColumn{Alias: "error_count", Expr: scopedtraces.PredicateCount("has_error = true")},
// timestamp of the last gen_ai span (LLM/tool/agent), hence gate-scoped.
scopedtraces.TraceColumn{Alias: "last_activity_time", Orderable: true, Expr: scopedtraces.ScopedReduce(scopedtraces.AggMax, "timestamp")},
// previews: first call's input (the prompt), last call's output (the answer).
scopedtraces.TraceColumn{Alias: "input", SpanLevel: true, Expr: scopedtraces.PickBy(&inMsg, str, "timestamp", scopedtraces.PickEarliest)},
scopedtraces.TraceColumn{Alias: "output", SpanLevel: true, Expr: scopedtraces.PickBy(&outMsg, str, "timestamp", scopedtraces.PickLatest)},
)
}
func (genAIColumnProvider) DefaultOrderAlias() string { return "last_activity_time" }
func (p genAIColumnProvider) AggregateAliases() []string {
// Derived from Columns() so a new column can't be forgotten. SpanLevel columns
// surface a real attribute, so a filter on them is applied span-level rather than
// treated as a trace aggregate — skip those.
// Derived from Columns() so a new column can't be forgotten; SpanLevel columns
// are filtered span-level, so skip them.
cols := p.Columns()
aliases := make([]string, 0, len(cols))
for _, c := range cols {

View File

@@ -1,799 +1,26 @@
package telemetryai
import (
"context"
"fmt"
"log/slog"
"sort"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/telemetryresourcefilter"
scopedtraces "github.com/SigNoz/signoz/pkg/telemetryscopedtraces"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrytraces"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/huandu/go-sqlbuilder"
)
var (
ErrUnsupportedRequestType = errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported request type for source=ai")
)
// scopedTraceStatementBuilder builds a trace list scoped to a span-selection
// category. Topology is fixed; selection (BaseConditionProvider) and columns
// (ColumnProvider) are pluggable, so a new category is a new pair of
// providers, not new topology.
type scopedTraceStatementBuilder struct {
logger *slog.Logger
metadataStore telemetrytypes.MetadataStore
fm qbtypes.FieldMapper
cb qbtypes.ConditionBuilder
baseCond BaseConditionProvider
columnProvider ColumnProvider
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation]
resourceFilterResolver *telemetryresourcefilter.ResourceFingerprintResolver[qbtypes.TraceAggregation]
skipResourceFingerprintEnabled bool
}
var _ qbtypes.StatementBuilder[qbtypes.TraceAggregation] = (*scopedTraceStatementBuilder)(nil)
// NewScopedTraceStatementBuilder wires the generic trace-list builder. The trace
// builder is reused for the span-list (raw) path.
func NewScopedTraceStatementBuilder(
settings factory.ProviderSettings,
metadataStore telemetrytypes.MetadataStore,
fieldMapper qbtypes.FieldMapper,
conditionBuilder qbtypes.ConditionBuilder,
baseCond BaseConditionProvider,
columnProvider ColumnProvider,
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation],
telemetryStore telemetrystore.TelemetryStore,
fl flagger.Flagger,
skipResourceFingerprintEnable bool,
skipResourceFingerprintThreshold uint64,
) *scopedTraceStatementBuilder {
aiSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryai")
// Resource-fingerprint prune over the traces resource table, mirroring the standard
// trace builder — the list scans the same span index, so the same resource table applies.
resourceFilterResolver := telemetryresourcefilter.NewResolver[qbtypes.TraceAggregation](
settings,
telemetrytraces.DBName,
telemetrytraces.TracesResourceV3TableName,
telemetrytypes.SignalTraces,
telemetrytypes.SourceUnspecified,
metadataStore,
nil,
fl,
telemetryStore,
skipResourceFingerprintThreshold,
)
return &scopedTraceStatementBuilder{
logger: aiSettings.Logger(),
metadataStore: metadataStore,
fm: fieldMapper,
cb: conditionBuilder,
baseCond: baseCond,
columnProvider: columnProvider,
traceStmtBuilder: traceStmtBuilder,
resourceFilterResolver: resourceFilterResolver,
skipResourceFingerprintEnabled: skipResourceFingerprintEnable,
}
}
// NewAITraceStatementBuilder is the scoped builder with the gen_ai gate + AI columns.
// NewAITraceStatementBuilder wires the generic scoped-trace builder with the gen_ai
// gate and AI columns. This package holds only gen_ai domain knowledge; the query
// topology lives in telemetryscopedtraces.
func NewAITraceStatementBuilder(
settings factory.ProviderSettings,
metadataStore telemetrytypes.MetadataStore,
fieldMapper qbtypes.FieldMapper,
conditionBuilder qbtypes.ConditionBuilder,
baseCond BaseConditionProvider,
baseCond scopedtraces.BaseConditionProvider,
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation],
telemetryStore telemetrystore.TelemetryStore,
fl flagger.Flagger,
skipResourceFingerprintEnable bool,
skipResourceFingerprintThreshold uint64,
) *scopedTraceStatementBuilder {
return NewScopedTraceStatementBuilder(settings, metadataStore, fieldMapper, conditionBuilder, baseCond, NewGenAIColumnProvider(), traceStmtBuilder, telemetryStore, fl, skipResourceFingerprintEnable, skipResourceFingerprintThreshold)
}
func (b *scopedTraceStatementBuilder) Build(
ctx context.Context,
start uint64,
end uint64,
requestType qbtypes.RequestType,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
switch requestType {
case qbtypes.RequestTypeTrace:
return b.buildTraceListQuery(ctx, querybuilder.ToNanoSecs(start), querybuilder.ToNanoSecs(end), query, variables)
case qbtypes.RequestTypeRaw:
return b.buildDelegated(ctx, start, end, requestType, query, variables)
default:
return nil, ErrUnsupportedRequestType
}
}
// buildDelegated ANDs the base gate into the user filter and delegates to the
// standard trace builder (the span-list / raw path).
func (b *scopedTraceStatementBuilder) buildDelegated(
ctx context.Context,
start, end uint64,
requestType qbtypes.RequestType,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
gate := b.baseCond.FilterExpression()
expr := gate
if query.Filter != nil && strings.TrimSpace(query.Filter.Expression) != "" {
expr = fmt.Sprintf("(%s) AND (%s)", gate, query.Filter.Expression)
}
// shallow copy; only Filter is replaced, caller's query untouched
gated := query
gated.Filter = &qbtypes.Filter{Expression: expr}
return b.traceStmtBuilder.Build(ctx, start, end, requestType, gated, variables)
}
// buildTraceListQuery is the map for the whole file. It resolves the columns, then
// wires the CTE pipeline that was benchmarked (see ai-qb-handoff.md): a single
// windowed pass picks the top-N traces, then a bucket-pruned pass enriches only those.
// The helpers appear in this file in the order they run here.
//
// RESOLVE (turn keys/columns into SQL, field-mapper aware)
// fetchKeys → metadata for the keys we reference
// resolveMask → the "is a gen_ai span" predicate (OR of EXISTS) [existsExpr]
// resolveColumns → per-trace column SQL: intrinsics + resolved aggregates
// resolveListOrders→ which resolved columns to ORDER BY
// splitFilter → span-level predicate + trace-level HAVING expression
//
// BUILD (compose the CTE pipeline)
// matched [buildMatchedCTE] ONE windowed, mask-pruned GROUP BY trace_id pass over
// │ the span index that applies the gate (+ span filter as
// │ countIf existence), the trace-level HAVING, ORDER BY and
// │ LIMIT/OFFSET in a single scan → top-N trace_ids + their
// │ gen_ai-scoped ranking metrics. No giant gate id-set.
// ▼
// ranked [buildRankedCTE] per-trace [start,end] bounds for those N traces, read
// │ from the small distributed_trace_summary table.
// ▼
// buckets [buildBucketsCTE] the exact ts_bucket_start values those N traces touch,
// │ so the enrichment scan is primary-key pruned.
// ▼
// enrichment[buildEnrichmentSelect] all per-trace columns for the N traces, scanning
// only their buckets (full trace, not window-clipped).
//
// Only gen_ai-scoped aggregates (tokens, llm activity, llm_call_count) are computable in
// the mask-pruned `matched` pass, so only those are orderable / usable in the aggregate
// filter. All-span columns (span_count, duration_nano, …) are output-only.
//
// start/end are nanoseconds.
func (b *scopedTraceStatementBuilder) buildTraceListQuery(
ctx context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
startBucket := start/querybuilder.NsToSeconds - querybuilder.BucketAdjustment
endBucket := end / querybuilder.NsToSeconds
limit := query.Limit
if limit <= 0 {
limit = 100
}
// Resolve the gate keys + columns once; every attribute access below goes through
// the field mapper (materialization/evolution aware), never a hardcoded map lookup.
keys, err := b.fetchKeys(ctx)
if err != nil {
return nil, err
}
maskExpr, maskArgs, err := b.resolveMask(ctx, start, end, keys)
if err != nil {
return nil, err
}
resolved, err := b.resolveColumns(ctx, start, end, keys, maskExpr, maskArgs)
if err != nil {
return nil, err
}
orders, err := b.resolveListOrders(query.Order, resolved)
if err != nil {
return nil, err
}
orderableSet := orderableAliasSet(resolved)
// Resource-fingerprint prune: when the filter references resource attributes, build a
// __resource_filter CTE and narrow the `matched` scan by resource_fingerprint, exactly
// like the standard trace builder. `skipResourceFilter` then drops those resource keys
// from the span predicate so they aren't re-applied on the span index.
resourceFrag, resourceArgs, resourcePred, skipResourceFilter, err := b.maybeAttachResourceFilter(ctx, query, start, end, variables)
if err != nil {
return nil, err
}
// Split the user filter: span-level predicate + trace-level HAVING expression.
fp, err := b.splitFilter(ctx, query, b.aggregateAliasSet(), orderableSet, start, end, skipResourceFilter, variables)
if err != nil {
return nil, err
}
// matched → ranked → buckets → enrichment
matchedFrag, matchedArgs, err := b.buildMatchedCTE(start, end, startBucket, endBucket, resolved, orders, orderableSet, maskExpr, maskArgs, fp, resourcePred, limit, query.Offset)
if err != nil {
return nil, err
}
rankedFrag, rankedArgs := b.buildRankedCTE(start, end)
bucketsFrag := buildBucketsCTE()
mainSQL, mainArgs := b.buildEnrichmentSelect(resolved, orders)
cteFragments := []string{matchedFrag, rankedFrag, bucketsFrag}
cteArgs := [][]any{matchedArgs, rankedArgs, nil}
// __resource_filter must precede `matched`, which references it.
if resourceFrag != "" {
cteFragments = append([]string{resourceFrag}, cteFragments...)
cteArgs = append([][]any{resourceArgs}, cteArgs...)
}
finalSQL := querybuilder.CombineCTEs(cteFragments) + mainSQL + " SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000"
finalArgs := querybuilder.PrependArgs(cteArgs, mainArgs)
return &qbtypes.Statement{
Query: finalSQL,
Args: finalArgs,
Warnings: fp.warnings,
WarningsDocURL: fp.warningsURL,
}, nil
}
// maybeAttachResourceFilter mirrors the standard trace builder: if the user filter
// references resource attributes, it builds a __resource_filter CTE (fingerprints
// matching those resource conditions) and returns the predicate that narrows the span
// scan by resource_fingerprint. `skipResourceFilter` tells the caller to drop those
// resource keys from the span predicate so they aren't re-applied on the span index.
// When the resolver is unset or the filter contains no resource conditions, it returns
// empty fragments and leaves the resource keys inline (skipResourceFilter=false).
func (b *scopedTraceStatementBuilder) maybeAttachResourceFilter(
ctx context.Context,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
start, end uint64,
variables map[string]qbtypes.VariableItem,
) (cteFrag string, cteArgs []any, fingerprintPred string, skipResourceFilter bool, err error) {
if b.resourceFilterResolver == nil {
return "", nil, "", false, nil
}
if b.skipResourceFingerprintEnabled {
decision, err := b.resourceFilterResolver.Resolve(ctx, query, start, end, variables)
if err != nil {
return "", nil, "", true, err
}
switch decision {
case qbtypes.ResourceFilterResolveKindNoOp:
return "", nil, "", true, nil
case qbtypes.ResourceFilterResolveKindFallback:
return "", nil, "", false, nil
}
}
stmt, err := b.resourceFilterResolver.StatementBuilder().Build(
ctx, start, end, qbtypes.RequestTypeRaw, query, variables,
)
if err != nil {
return "", nil, "", true, err
}
if stmt == nil {
return "", nil, "", true, nil
}
return fmt.Sprintf("__resource_filter AS (%s)", stmt.Query), stmt.Args,
"resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter)", true, nil
}
// ---------------------------------------------------------------------------
// RESOLVE — turn keys/columns into field-mapper-aware SQL
// ---------------------------------------------------------------------------
func (b *scopedTraceStatementBuilder) fetchKeys(ctx context.Context) (map[string][]*telemetrytypes.TelemetryFieldKey, error) {
fields := b.resolverFieldKeys()
selectors := make([]*telemetrytypes.FieldKeySelector, 0, len(fields))
for _, k := range fields {
selectors = append(selectors, &telemetrytypes.FieldKeySelector{
Name: k.Name,
Signal: k.Signal,
FieldContext: k.FieldContext,
})
}
keys, _, err := b.metadataStore.GetKeysMulti(ctx, selectors)
return keys, err
}
func (b *scopedTraceStatementBuilder) resolverFieldKeys() []*telemetrytypes.TelemetryFieldKey {
seen := make(map[string]struct{})
var out []*telemetrytypes.TelemetryFieldKey
add := func(k *telemetrytypes.TelemetryFieldKey) {
if k == nil {
return
}
if _, dup := seen[k.Name]; dup {
return
}
seen[k.Name] = struct{}{}
out = append(out, k)
}
for _, k := range b.baseCond.FieldKeys() {
add(k)
}
for _, c := range b.columnProvider.Columns() {
for _, k := range c.Expr.keys {
add(k)
}
}
return out
}
// resolveMask builds the per-span in-scope mask: OR of resolved EXISTS predicates
// over the base condition's field keys.
func (b *scopedTraceStatementBuilder) resolveMask(ctx context.Context, start, end uint64, keys map[string][]*telemetrytypes.TelemetryFieldKey) (string, []any, error) {
fieldKeys := b.baseCond.FieldKeys()
parts := make([]string, 0, len(fieldKeys))
var args []any
for _, key := range fieldKeys {
e, a, err := b.existsExpr(ctx, start, end, keys, key)
if err != nil {
return "", nil, err
}
parts = append(parts, e)
args = append(args, a...)
}
return "(" + strings.Join(parts, " OR ") + ")", args, nil
}
// existsExpr resolves a field-mapper-aware EXISTS predicate for key (materialized
// column when present, else the map). Escaped once so it round-trips when embedded
// in an outer builder.
func (b *scopedTraceStatementBuilder) existsExpr(ctx context.Context, start, end uint64, keys map[string][]*telemetrytypes.TelemetryFieldKey, key *telemetrytypes.TelemetryFieldKey) (string, []any, error) {
resolvedKey := key
cands := keys[key.Name]
if len(cands) == 0 {
cands = []*telemetrytypes.TelemetryFieldKey{key}
} else {
resolvedKey = cands[0]
}
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := b.cb.ConditionFor(ctx, start, end, resolvedKey, cands, qbtypes.FilterOperatorExists, nil, sb)
if err != nil {
return "", nil, err
}
sb.Where(conds[0])
expr, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
expr = strings.TrimPrefix(expr, "WHERE ")
return sqlbuilder.Escape(expr), args, nil
}
// resolvedColumn is a column whose attribute access has been resolved to
// SQL via the field mapper. expr is escaped once, ready to embed in an outer SELECT.
type resolvedColumn struct {
alias string
expr string
args []any
orderable bool
}
// resolveColumns turns the declarative columns into SQL, resolving all
// attribute access through the field mapper.
func (b *scopedTraceStatementBuilder) resolveColumns(ctx context.Context, start, end uint64, keys map[string][]*telemetrytypes.TelemetryFieldKey, maskExpr string, maskArgs []any) ([]resolvedColumn, error) {
r := aggResolver{
exists: func(key *telemetrytypes.TelemetryFieldKey) (string, []any, error) {
return b.existsExpr(ctx, start, end, keys, key)
},
value: func(key *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType) (string, []any, error) {
// Resolve to the metadata variant, which carries Materialized, so a promoted
// attribute uses its materialized column instead of map access. The static
// GenAIFieldDefinitions key always has Materialized=false and a concrete
// context, so CollisionHandledFinalExpr would otherwise never consult metadata.
// Mirrors existsExpr.
if cands := keys[key.Name]; len(cands) > 0 {
key = cands[0]
}
return querybuilder.CollisionHandledFinalExpr(ctx, start, end, key, b.fm, b.cb, keys, dt, nil, false)
},
maskExpr: maskExpr,
maskArgs: maskArgs,
}
cols := b.columnProvider.Columns()
out := make([]resolvedColumn, 0, len(cols))
for _, c := range cols {
expr, args, err := c.Expr.render(r)
if err != nil {
return nil, err
}
out = append(out, resolvedColumn{alias: c.Alias, expr: expr, args: args, orderable: c.Orderable})
}
return out, nil
}
// listOrder resolves a sort key to an aggregate-column alias + direction. Both the
// matched CTE and the enrichment select that alias, so both ORDER BY it.
type listOrder struct {
alias string
direction string
}
// resolveListOrders maps order keys to the resolved orderable columns; non-orderable
// columns are rejected. Defaults to the column provider's default order.
func (b *scopedTraceStatementBuilder) resolveListOrders(order []qbtypes.OrderBy, resolved []resolvedColumn) ([]listOrder, error) {
byAlias := make(map[string]resolvedColumn, len(resolved))
orderable := make([]string, 0, len(resolved))
for _, rc := range resolved {
byAlias[rc.alias] = rc
if rc.orderable {
orderable = append(orderable, rc.alias)
}
}
if len(order) == 0 {
return []listOrder{{alias: b.columnProvider.DefaultOrderAlias(), direction: "DESC"}}, nil
}
orders := make([]listOrder, 0, len(order))
for _, o := range order {
direction := "DESC"
if o.Direction == qbtypes.OrderDirectionAsc {
direction = "ASC"
}
rc, ok := byAlias[o.Key.Name]
if !ok || !rc.orderable {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput,
"unsupported order key %q for the trace list; orderable keys: %s", o.Key.Name, strings.Join(orderable, ", "))
}
orders = append(orders, listOrder{alias: rc.alias, direction: direction})
}
return orders, nil
}
// filterParts is the split of the user filter into a resolved span-level predicate
// (used both to widen the matched WHERE prune and as a countIf existence in HAVING)
// and a trace-level HAVING expression.
type filterParts struct {
spanPred string
spanArgs []any
hasSpanFilter bool
havingExpr string
warnings []string
warningsURL string
}
// splitFilter partitions query.Filter into a span-level predicate and a trace-level
// HAVING expression; an explicit query.Having is ANDed onto the latter. The trace-level
// expression is validated against the aggregates computable in the matched pass.
func (b *scopedTraceStatementBuilder) splitFilter(ctx context.Context, query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation], classifySet, orderableSet map[string]struct{}, start, end uint64, skipResourceFilter bool, variables map[string]qbtypes.VariableItem) (filterParts, error) {
var fp filterParts
if query.Filter != nil && strings.TrimSpace(query.Filter.Expression) != "" {
spanExpr, traceExpr, err := querybuilder.SplitFilterForAggregates(query.Filter.Expression, classifySet)
if err != nil {
return fp, err
}
fp.havingExpr = traceExpr
if strings.TrimSpace(spanExpr) != "" {
pred, args, warnings, url, err := b.resolveSpanPredicate(ctx, start, end, spanExpr, skipResourceFilter, variables)
if err != nil {
return fp, err
}
// pred can resolve to empty when the span-level keys were only resource
// attributes now handled by the __resource_filter CTE (skipResourceFilter).
if strings.TrimSpace(pred) != "" {
fp.spanPred, fp.spanArgs, fp.hasSpanFilter = pred, args, true
}
fp.warnings, fp.warningsURL = warnings, url
}
}
if query.Having != nil && strings.TrimSpace(query.Having.Expression) != "" {
if fp.havingExpr != "" {
fp.havingExpr = fmt.Sprintf("(%s) AND (%s)", fp.havingExpr, query.Having.Expression)
} else {
fp.havingExpr = query.Having.Expression
}
}
if err := validateAggregateFilter(fp.havingExpr, orderableSet); err != nil {
return fp, err
}
return fp, nil
}
// resolveSpanPredicate resolves a span-level filter expression to a bare boolean SQL
// predicate (escaped) + args via the field mapper.
func (b *scopedTraceStatementBuilder) resolveSpanPredicate(ctx context.Context, start, end uint64, expr string, skipResourceFilter bool, variables map[string]qbtypes.VariableItem) (string, []any, []string, string, error) {
selectors := querybuilder.QueryStringToKeysSelectors(expr)
for i := range selectors {
selectors[i].Signal = telemetrytypes.SignalTraces
}
keys, _, err := b.metadataStore.GetKeysMulti(ctx, selectors)
if err != nil {
return "", nil, nil, "", err
}
prepared, err := querybuilder.PrepareWhereClause(expr, querybuilder.FilterExprVisitorOpts{
Context: ctx,
Logger: b.logger,
FieldMapper: b.fm,
ConditionBuilder: b.cb,
FieldKeys: keys,
SkipResourceFilter: skipResourceFilter,
Variables: variables,
StartNs: start,
EndNs: end,
})
if err != nil {
return "", nil, nil, "", err
}
if prepared.IsEmpty() {
return "", nil, nil, "", nil
}
sb := sqlbuilder.NewSelectBuilder()
sb.Select("1")
sb.AddWhereClause(prepared.WhereClause)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
pred := sql[strings.Index(sql, "WHERE ")+len("WHERE "):]
return sqlbuilder.Escape(pred), args, prepared.Warnings, prepared.WarningsDocURL, nil
}
// buildMatchedCTE builds `matched`: the single windowed, mask-pruned GROUP BY trace_id
// scan that fuses gate + span filter + aggregate HAVING + ORDER BY + LIMIT/OFFSET,
// selecting only the aggregate aliases the ORDER BY / HAVING reference.
func (b *scopedTraceStatementBuilder) buildMatchedCTE(start, end, startBucket, endBucket uint64, resolved []resolvedColumn, orders []listOrder, orderableSet map[string]struct{}, maskExpr string, maskArgs []any, fp filterParts, resourcePred string, limit, offset int) (string, []any, error) {
sb := sqlbuilder.NewSelectBuilder()
// SELECT trace_id + only the aggregates ORDER BY / HAVING reference (as aliases).
needed := neededMatchedAliases(orders, fp.havingExpr, orderableSet)
selects := []string{"trace_id"}
for _, rc := range resolved {
if _, ok := needed[rc.alias]; !ok {
continue
}
selects = append(selects, embedExpr(sb, rc.expr, rc.args)+" AS "+quoteAlias(rc.alias))
}
sb.Select(selects...)
sb.From(spanTable())
// WHERE: window + coarse prune to gen_ai spans (widened so span-filter spans are
// visible for the countIf existence check below).
win := windowWhere(sb, start, end, startBucket, endBucket)
prune := "(" + embedExpr(sb, maskExpr, maskArgs)
if fp.hasSpanFilter {
prune += " OR " + embedExpr(sb, fp.spanPred, fp.spanArgs)
}
prune += ")"
where := append(win, prune)
// Resource-fingerprint prune: restrict the scan to spans on matching resources.
if resourcePred != "" {
where = append(where, resourcePred)
}
sb.Where(where...)
sb.GroupBy("trace_id")
// HAVING: the gate + span-existence checks are only needed once the WHERE has been
// widened by a span filter; otherwise WHERE = mask already enforces the gate.
var having []string
if fp.hasSpanFilter {
having = append(having, "countIf("+embedExpr(sb, maskExpr, maskArgs)+") > 0")
having = append(having, "countIf("+embedExpr(sb, fp.spanPred, fp.spanArgs)+") > 0")
}
if strings.TrimSpace(fp.havingExpr) != "" {
hv, err := b.buildHaving(fp.havingExpr, orderableSet)
if err != nil {
return "", nil, err
}
if hv != "" {
having = append(having, hv)
}
}
if len(having) > 0 {
sb.Having(strings.Join(having, " AND "))
}
sb.OrderBy(orderClause(orders)...)
sb.Limit(limit)
if offset > 0 {
sb.Offset(offset)
}
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return fmt.Sprintf("matched AS (%s)", sql), args, nil
}
// buildRankedCTE builds `ranked`: per-trace [start,end] bounds for the matched traces,
// read from the small trace-summary table (used to derive the bucket prune).
func (b *scopedTraceStatementBuilder) buildRankedCTE(start, end uint64) (string, []any) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("trace_id", "min(start) AS t_start", "max(end) AS t_end")
sb.From(summaryTable())
sb.Where(
"trace_id GLOBAL IN (SELECT trace_id FROM matched)",
"end >= fromUnixTimestamp64Nano("+sb.Var(start)+")",
"start < fromUnixTimestamp64Nano("+sb.Var(end)+")",
)
sb.GroupBy("trace_id")
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return fmt.Sprintf("ranked AS (%s)", sql), args
}
// buildBucketsCTE builds `buckets`: the exact ts_bucket_start values the matched traces
// span, so the enrichment scan is pruned to those primary-key buckets. No args.
func buildBucketsCTE() string {
adj := querybuilder.BucketAdjustment // 30-min bucket width in seconds
return fmt.Sprintf("buckets AS (SELECT DISTINCT b AS ts_bucket FROM ranked "+
"ARRAY JOIN range("+
"toUInt64(intDiv(toUnixTimestamp(t_start), %d) * %d - %d), "+
"toUInt64(intDiv(toUnixTimestamp(t_end), %d) * %d + %d), "+
"%d) AS b)", adj, adj, adj, adj, adj, adj, adj)
}
// buildEnrichmentSelect builds the final SELECT: all per-trace columns for the matched
// traces, scanning only their buckets (full trace, not window-clipped). SELECT-expr
// args lead; the WHERE / ORDER BY carry none.
func (b *scopedTraceStatementBuilder) buildEnrichmentSelect(resolved []resolvedColumn, orders []listOrder) (string, []any) {
sb := sqlbuilder.NewSelectBuilder()
selects, selectArgs := selectAllColumns(resolved)
sb.Select(selects...)
sb.From(spanTable())
sb.Where(
"ts_bucket_start GLOBAL IN (SELECT ts_bucket FROM buckets)",
"trace_id GLOBAL IN (SELECT trace_id FROM ranked)",
)
sb.GroupBy("trace_id")
sb.OrderBy(orderClause(orders)...)
sql, builtArgs := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return sql, append(append([]any{}, selectArgs...), builtArgs...)
}
// buildHaving rewrites a trace-level HAVING expression against the aggregate column
// aliases computable in the matched pass; bare, trace., and tracefield. forms all map
// to the selected alias.
func (b *scopedTraceStatementBuilder) buildHaving(havingExpr string, orderableSet map[string]struct{}) (string, error) {
columnMap := make(map[string]string, len(orderableSet)*3)
for a := range orderableSet {
columnMap[a] = quoteAlias(a)
columnMap["trace."+a] = quoteAlias(a)
columnMap["tracefield."+a] = quoteAlias(a)
}
return querybuilder.NewHavingExpressionRewriter().Rewrite(havingExpr, columnMap)
}
// ---------------------------------------------------------------------------
// Small shared SQL-builder utilities
// ---------------------------------------------------------------------------
// spanTable is the fully-qualified span index table.
func spanTable() string {
return fmt.Sprintf("%s.%s", telemetrytraces.DBName, telemetrytraces.SpanIndexV3TableName)
}
// summaryTable is the fully-qualified trace-summary table.
func summaryTable() string {
return fmt.Sprintf("%s.%s", telemetrytraces.DBName, telemetrytraces.TraceSummaryTableName)
}
// aggregateAliasSet is the set of all trace-level (computed) column aliases, used to
// classify which filter keys are trace-level vs span-level.
func (b *scopedTraceStatementBuilder) aggregateAliasSet() map[string]struct{} {
set := make(map[string]struct{}, len(b.columnProvider.AggregateAliases()))
for _, a := range b.columnProvider.AggregateAliases() {
set[a] = struct{}{}
}
return set
}
// orderableAliasSet is the subset of aggregate aliases computable in the matched pass
// (gen_ai-scoped): the only ones usable for ORDER BY and the aggregate filter.
func orderableAliasSet(resolved []resolvedColumn) map[string]struct{} {
set := make(map[string]struct{})
for _, rc := range resolved {
if rc.orderable {
set[rc.alias] = struct{}{}
}
}
return set
}
// neededMatchedAliases is the minimal set of aggregate aliases the matched pass must
// select: those referenced by ORDER BY plus those referenced in the aggregate HAVING
// (bare / trace. / tracefield. forms). Anything else is left to the enrichment scan.
func neededMatchedAliases(orders []listOrder, havingExpr string, orderableSet map[string]struct{}) map[string]struct{} {
needed := make(map[string]struct{})
for _, o := range orders {
needed[o.alias] = struct{}{}
}
for _, sel := range querybuilder.QueryStringToKeysSelectors(havingExpr) {
name := strings.TrimPrefix(strings.TrimPrefix(sel.Name, "trace."), "tracefield.")
if _, ok := orderableSet[name]; ok {
needed[name] = struct{}{}
}
}
return needed
}
// validateAggregateFilter rejects a trace-level filter that references an aggregate not
// computable in the matched pass (e.g. span_count, duration_nano), with a clear message.
func validateAggregateFilter(havingExpr string, orderableSet map[string]struct{}) error {
if strings.TrimSpace(havingExpr) == "" {
return nil
}
allowed := make([]string, 0, len(orderableSet))
for a := range orderableSet {
allowed = append(allowed, a)
}
sort.Strings(allowed)
for _, sel := range querybuilder.QueryStringToKeysSelectors(havingExpr) {
name := strings.TrimPrefix(strings.TrimPrefix(sel.Name, "trace."), "tracefield.")
if _, ok := orderableSet[name]; !ok {
return errors.NewInvalidInputf(errors.CodeInvalidInput,
"aggregate %q cannot be used in an AI trace-list filter; filterable aggregates: %s", name, strings.Join(allowed, ", "))
}
}
return nil
}
// embedExpr inlines a resolved (escaped) expr carrying `?` placeholders into sb by
// replacing each `?` with a builder Var, so go-sqlbuilder tracks the args in appearance
// order and un-escapes the expr at Build time.
func embedExpr(sb *sqlbuilder.SelectBuilder, expr string, args []any) string {
var out strings.Builder
ai := 0
for i := 0; i < len(expr); i++ {
if expr[i] == '?' && ai < len(args) {
out.WriteString(sb.Var(args[ai]))
ai++
continue
}
out.WriteByte(expr[i])
}
return out.String()
}
// windowWhere binds the shared time-window predicates to sb and returns them, so a
// caller can add its own predicate in the same Where call.
func windowWhere(sb *sqlbuilder.SelectBuilder, start, end, startBucket, endBucket uint64) []string {
return []string{
sb.GE("timestamp", fmt.Sprintf("%d", start)),
sb.L("timestamp", fmt.Sprintf("%d", end)),
sb.GE("ts_bucket_start", startBucket),
sb.LE("ts_bucket_start", endBucket),
}
}
// orderClause renders the ORDER BY terms (by column alias) + the trace_id tiebreak.
func orderClause(orders []listOrder) []string {
out := make([]string, 0, len(orders)+1)
for _, o := range orders {
out = append(out, fmt.Sprintf("%s %s", quoteAlias(o.alias), o.direction))
}
return append(out, "trace_id DESC")
}
// selectAllColumns renders `expr AS alias` for every resolved column and returns their
// field-mapper args in select order.
func selectAllColumns(resolved []resolvedColumn) ([]string, []any) {
selects := []string{"trace_id"}
var args []any
for _, rc := range resolved {
selects = append(selects, rc.expr+" AS "+quoteAlias(rc.alias))
args = append(args, rc.args...)
}
return selects, args
}
// quoteAlias backticks an alias that carries characters special to the SQL builder.
func quoteAlias(alias string) string {
if strings.ContainsAny(alias, ".$`") {
return "`" + alias + "`"
}
return alias
) qbtypes.StatementBuilder[qbtypes.TraceAggregation] {
return scopedtraces.NewScopedTraceStatementBuilder(settings, metadataStore, baseCond, NewGenAIColumnProvider(), traceStmtBuilder, telemetryStore, fl, skipResourceFingerprintEnable, skipResourceFingerprintThreshold)
}

View File

@@ -9,6 +9,7 @@ import (
"github.com/SigNoz/signoz/pkg/flagger/flaggertest"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/querybuilder"
scopedtraces "github.com/SigNoz/signoz/pkg/telemetryscopedtraces"
"github.com/SigNoz/signoz/pkg/telemetrytraces"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
@@ -64,14 +65,14 @@ const (
testEndMs = uint64(1747983448000)
)
func newTestBuilder(t *testing.T) *scopedTraceStatementBuilder {
func newTestBuilder(t *testing.T) qbtypes.StatementBuilder[qbtypes.TraceAggregation] {
return newTestBuilderWithKeys(t, otelKeysMap())
}
// newTestBuilderWithKeys mirrors the production wiring in signozquerier's provider.
// The gen_ai keys are seeded via keysMap here; in production the metadata store
// surfaces them itself (enrichWithGenAIKeys).
func newTestBuilderWithKeys(t *testing.T, keysMap map[string][]*telemetrytypes.TelemetryFieldKey) *scopedTraceStatementBuilder {
func newTestBuilderWithKeys(t *testing.T, keysMap map[string][]*telemetrytypes.TelemetryFieldKey) qbtypes.StatementBuilder[qbtypes.TraceAggregation] {
t.Helper()
settings := instrumentationtest.New().ToProviderSettings()
fm := telemetrytraces.NewFieldMapper()
@@ -98,8 +99,6 @@ func newTestBuilderWithKeys(t *testing.T, keysMap map[string][]*telemetrytypes.T
return NewAITraceStatementBuilder(
settings,
metadataStore,
fm,
cb,
baseCond,
traceStmtBuilder,
nil, // telemetryStore: only used by the skip-fingerprint count query, which is disabled here
@@ -220,7 +219,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, mapContains(attributes_string, 'gen_ai.request.model') = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -295,7 +294,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists = true, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists = true, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(attribute_number_gen_ai$$usage$$input_tokens_exists = true, toFloat64(attribute_number_gen_ai$$usage$$input_tokens), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, attribute_string_gen_ai$$request$$model_exists = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -369,7 +368,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, mapContains(attributes_string, 'gen_ai.request.model') = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -439,7 +438,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, mapContains(attributes_string, 'gen_ai.request.model') = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -510,7 +509,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, mapContains(attributes_string, 'gen_ai.request.model') = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -596,7 +595,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, mapContains(attributes_string, 'gen_ai.request.model') = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -674,7 +673,7 @@ SELECT trace_id,
uniqIf(multiIf(mapContains(attributes_string, 'gen_ai.tool.name') = true, attributes_string['gen_ai.tool.name'], NULL), mapContains(attributes_string, 'gen_ai.tool.name') = true) AS distinct_tool_count,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) AS input_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS output_tokens,
sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)) + sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)) AS total_tokens,
coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.input_tokens') = true, toFloat64(attributes_number['gen_ai.usage.input_tokens']), NULL)), 0) + coalesce(sum(multiIf(mapContains(attributes_number, 'gen_ai.usage.output_tokens') = true, toFloat64(attributes_number['gen_ai.usage.output_tokens']), NULL)), 0) AS total_tokens,
sum(multiIf(mapContains(attributes_number, '_signoz.gen_ai.total_cost') = true, toFloat64(attributes_number['_signoz.gen_ai.total_cost']), NULL)) AS estimated_cost_usd,
maxIf(signoz_traces.distributed_signoz_index_v3.duration_nano, mapContains(attributes_string, 'gen_ai.request.model') = true) AS max_llm_latency_ns,
countIf(has_error = true) AS error_count,
@@ -769,23 +768,8 @@ func TestBuild_TraceList_ResourcePlusSpanPlusAggregateFilter(t *testing.T) {
require.Contains(t, got, "output_tokens")
}
// With the resolver unset (nil), the resource filter falls back to being applied inline
// on the span index — no fingerprint CTE — so existing behavior is preserved.
func TestBuild_TraceList_ResourceFilter_NoResolver(t *testing.T) {
b := newTestBuilderWithKeys(t, resourceKeysMap())
b.resourceFilterResolver = nil
stmt, err := b.Build(context.Background(), testStartMs, testEndMs, qbtypes.RequestTypeTrace,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces, Source: telemetrytypes.SourceAI,
Filter: &qbtypes.Filter{Expression: "resource.service.name = 'checkout'"},
Limit: 10,
}, nil)
require.NoError(t, err)
got := renderSQL(t, stmt)
require.NotContains(t, got, "__resource_filter")
require.Contains(t, got, "resources_string['service.name']")
}
// The resolver-unset (nil) fallback is covered in pkg/telemetryscopedtraces, which
// can construct that builder state directly.
// Trace-level and span-level predicates may not be OR-combined.
func TestBuild_TraceList_TraceOrSpanMixRejected(t *testing.T) {
@@ -877,5 +861,92 @@ func TestBuild_UnsupportedRequestType(t *testing.T) {
},
}
_, err := b.Build(context.Background(), testStartMs, testEndMs, qbtypes.RequestTypeDistribution, query, nil)
require.ErrorIs(t, err, ErrUnsupportedRequestType)
require.ErrorIs(t, err, scopedtraces.ErrUnsupportedRequestType)
}
// A gate key ingested under several data types (e.g. string + number from a
// misbehaving SDK) contributes ALL variants to the mask, OR-combined — not just
// the first — matching the standard visitor's EXISTS handling.
func TestBuild_TraceList_MultiVariantGateKey(t *testing.T) {
keys := otelKeysMap()
keys[telemetrytypes.GenAIToolName] = append(keys[telemetrytypes.GenAIToolName], &telemetrytypes.TelemetryFieldKey{
Name: telemetrytypes.GenAIToolName,
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeFloat64,
})
b := newTestBuilderWithKeys(t, keys)
stmt, err := b.Build(context.Background(), testStartMs, testEndMs, qbtypes.RequestTypeTrace,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces, Source: telemetrytypes.SourceAI, Limit: 10,
}, nil)
require.NoError(t, err)
got := renderSQL(t, stmt)
require.Contains(t, got, "mapContains(attributes_string, 'gen_ai.tool.name') = true OR mapContains(attributes_number, 'gen_ai.tool.name') = true")
}
// `tracefield.` is the explicit trace field context, so in a filter it marks a
// trace-level aggregate exactly like the user-facing `trace.` prefix — same statement,
// and the same targeted rejection for a non-filterable aggregate. (Filter and Having
// accept the same forms; the splitter used to misroute tracefield. as span-level.)
func TestBuild_TraceList_TracefieldPrefixMatchesTracePrefix(t *testing.T) {
b := newTestBuilder(t)
build := func(expr string) (*qbtypes.Statement, error) {
return b.Build(context.Background(), testStartMs, testEndMs, qbtypes.RequestTypeTrace,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces, Source: telemetrytypes.SourceAI,
Filter: &qbtypes.Filter{Expression: expr},
Limit: 20,
}, nil)
}
viaTrace, err := build("trace.output_tokens > 1000")
require.NoError(t, err)
viaTracefield, err := build("tracefield.output_tokens > 1000")
require.NoError(t, err)
require.Equal(t, viaTrace.Query, viaTracefield.Query)
require.Equal(t, viaTrace.Args, viaTracefield.Args)
// output-only aggregate under tracefield. gets the aggregate rejection, not an
// unknown-span-field failure.
_, err = build("tracefield.span_count > 3")
require.Error(t, err)
require.Contains(t, err.Error(), "cannot be used")
}
// Query variables in a trace-level condition are substituted into the HAVING (the
// span path binds them via PrepareWhereClause; the HAVING is a text rewrite).
func TestBuild_TraceList_VariableInAggregateFilter(t *testing.T) {
b := newTestBuilder(t)
build := func(expr string, vars map[string]qbtypes.VariableItem) (*qbtypes.Statement, error) {
return b.Build(context.Background(), testStartMs, testEndMs, qbtypes.RequestTypeTrace,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces, Source: telemetrytypes.SourceAI,
Filter: &qbtypes.Filter{Expression: expr},
Limit: 20,
}, vars)
}
// scalar variable -> literal in HAVING
stmt, err := build("trace.output_tokens > $threshold",
map[string]qbtypes.VariableItem{"threshold": {Value: 700}})
require.NoError(t, err)
require.Contains(t, stmt.Query, "HAVING output_tokens > 700")
// list variable with IN
stmt, err = build("trace.llm_call_count IN $counts",
map[string]qbtypes.VariableItem{"counts": {Value: []any{1, 2}}})
require.NoError(t, err)
require.Contains(t, stmt.Query, "HAVING llm_call_count IN")
// dynamic __all__ -> condition dropped, no HAVING at all
stmt, err = build("trace.output_tokens > $threshold",
map[string]qbtypes.VariableItem{"threshold": {Type: qbtypes.DynamicVariableType, Value: "__all__"}})
require.NoError(t, err)
require.NotContains(t, stmt.Query, "HAVING")
// unresolved variable -> rejected, not compared as a literal
_, err = build("trace.output_tokens > $missing", map[string]qbtypes.VariableItem{"other": {Value: 1}})
require.Error(t, err)
}

View File

@@ -19,6 +19,7 @@ import (
"github.com/SigNoz/signoz/pkg/telemetrymetrics"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrytraces"
"github.com/SigNoz/signoz/pkg/types/authtypes"
"github.com/SigNoz/signoz/pkg/types/ctxtypes"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
"github.com/SigNoz/signoz/pkg/types/instrumentationtypes"
@@ -1190,6 +1191,24 @@ func enrichWithIntrinsicMetricKeys(keys map[string][]*telemetrytypes.TelemetryFi
return keys
}
// genAIEnrichmentEnabled reports whether the org in ctx has AI observability enabled.
// The static gen_ai key definitions are surfaced (autocomplete + query-time resolution
// before any gen_ai data is ingested) only for those orgs, so other tenants don't see
// gen_ai keys in trace autocomplete. Contexts without claims (internal paths) resolve
// to false — an org relying on enrichment has no gen_ai data ingested, so those paths
// had nothing to resolve anyway.
func (t *telemetryMetaStore) genAIEnrichmentEnabled(ctx context.Context) bool {
claims, err := authtypes.ClaimsFromContext(ctx)
if err != nil {
return false
}
orgID, err := valuer.NewUUID(claims.OrgID)
if err != nil {
return false
}
return t.fl.BooleanOrEmpty(ctx, flagger.FeatureEnableAIObservability, featuretypes.NewFlaggerEvaluationContext(orgID))
}
// enrichWithGenAIKeys adds keys that can be queried for GenAI signals, even though they have not been ingested yet.
func enrichWithGenAIKeys(keys map[string][]*telemetrytypes.TelemetryFieldKey, selectors []*telemetrytypes.FieldKeySelector) map[string][]*telemetrytypes.TelemetryFieldKey {
for _, selector := range selectors {
@@ -1296,7 +1315,9 @@ func (t *telemetryMetaStore) GetKeys(ctx context.Context, fieldKeySelector *tele
applyBackwardCompatibleKeys(mapOfKeys)
mapOfKeys = enrichWithIntrinsicMetricKeys(mapOfKeys, selectors)
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, selectors)
if t.genAIEnrichmentEnabled(ctx) {
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, selectors)
}
return mapOfKeys, complete, nil
}
@@ -1375,7 +1396,9 @@ func (t *telemetryMetaStore) GetKeysMulti(ctx context.Context, fieldKeySelectors
applyBackwardCompatibleKeys(mapOfKeys)
mapOfKeys = enrichWithIntrinsicMetricKeys(mapOfKeys, fieldKeySelectors)
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, fieldKeySelectors)
if t.genAIEnrichmentEnabled(ctx) {
mapOfKeys = enrichWithGenAIKeys(mapOfKeys, fieldKeySelectors)
}
return mapOfKeys, complete, nil
}

View File

@@ -0,0 +1,17 @@
package telemetryscopedtraces
import (
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
)
// BaseConditionProvider defines which spans are in scope. It only declares the gate
// (a filter expression + its field keys); the builder resolves the keys through the
// field mapper, so attribute access stays materialization-aware.
type BaseConditionProvider interface {
// FilterExpression is the grammar-level (EXISTS) gate, used on the delegated
// span-list path.
FilterExpression() string
// FieldKeys are the gate's keys, used to build the per-span mask
// (OR of resolved EXISTS conditions).
FieldKeys() []*telemetrytypes.TelemetryFieldKey
}

View File

@@ -0,0 +1,181 @@
package telemetryscopedtraces
import (
"fmt"
"strings"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/huandu/go-sqlbuilder"
)
// TraceColumn is one per-trace output column.
type TraceColumn struct {
Alias string
// Orderable columns can be used in ORDER BY and the aggregate filter. All-span
// aggregates (span_count, duration_nano, …) are display-only and set false.
Orderable bool
// SpanLevel columns surface a real span/resource attribute (service.name,
// input/output messages); a filter on them is applied span-level, so they are
// excluded from AggregateAliases.
SpanLevel bool
Expr Aggregate
}
// Aggregate renders one column's SQL and lists the attribute keys it references so
// the builder can pre-fetch their metadata. Build one with the constructors below;
// the zero value is not usable.
type Aggregate struct {
keys []*telemetrytypes.TelemetryFieldKey
render func(r aggResolver) (expr string, args []any, err error)
}
// aggResolver hands each aggregate the field-mapper primitives it may need — an
// EXISTS predicate, a resolved value expression, and the gate mask. Populated per
// query by resolveColumns.
type aggResolver struct {
exists func(key *telemetrytypes.TelemetryFieldKey) (string, []any, error)
value func(key *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType) (string, []any, error)
maskExpr string
maskArgs []any
}
// AggFunc is a ClickHouse aggregate function name.
type AggFunc string
const (
AggSum AggFunc = "sum"
AggMax AggFunc = "max"
AggMin AggFunc = "min"
)
// PickDirection selects the earliest (argMin) or latest (argMax) span by ordering.
type PickDirection int
const (
PickLatest PickDirection = iota
PickEarliest
)
// Intrinsic emits fixed intrinsic-column SQL verbatim (escaped once).
func Intrinsic(text string) Aggregate {
return Aggregate{render: func(aggResolver) (string, []any, error) {
return sqlbuilder.Escape(text), nil, nil
}}
}
// CountExists renders countIf(<key> EXISTS) — counts spans carrying key.
func CountExists(key *telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: keysOf(key), render: func(r aggResolver) (string, []any, error) {
cond, args, err := r.exists(key)
return fmt.Sprintf("countIf(%s)", cond), args, err
}}
}
// Reduce renders <fn>(<value>) over a resolved numeric attribute value.
func Reduce(fn AggFunc, valueKey *telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: keysOf(valueKey), render: func(r aggResolver) (string, []any, error) {
v, args, err := r.value(valueKey, telemetrytypes.FieldDataTypeFloat64)
return fmt.Sprintf("%s(%s)", fn, v), args, err
}}
}
// ScopedReduce renders <fn>If(<valueExpr>, <gate mask>) over a fixed value expression.
func ScopedReduce(fn AggFunc, valueExpr string) Aggregate {
return Aggregate{render: func(r aggResolver) (string, []any, error) {
return fmt.Sprintf("%sIf(%s, %s)", fn, valueExpr, r.maskExpr), append([]any{}, r.maskArgs...), nil
}}
}
// ScopedToKeyColumn renders <fn>If(<column>, <scopeKey> EXISTS) — a physical
// span-index column aggregated over spans carrying scopeKey (e.g. max LLM latency).
// Providers pass the bare column name; it is table-qualified here so it binds to the
// physical column and not a same-named output alias, which ClickHouse would reject
// as an aggregate inside an aggregate.
func ScopedToKeyColumn(fn AggFunc, column string, scopeKey *telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: keysOf(scopeKey), render: func(r aggResolver) (string, []any, error) {
cond, args, err := r.exists(scopeKey)
return fmt.Sprintf("%sIf(%s.%s, %s)", fn, spanTable(), column, cond), args, err
}}
}
// PickBy renders argMinIf/argMaxIf(<value>, <orderExpr>, <value> EXISTS) — the value
// from the earliest/latest span that carries it.
func PickBy(valueKey *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType, orderExpr string, dir PickDirection) Aggregate {
fn := "argMaxIf"
if dir == PickEarliest {
fn = "argMinIf"
}
return Aggregate{keys: keysOf(valueKey), render: func(r aggResolver) (string, []any, error) {
v, vargs, err := r.value(valueKey, dt)
if err != nil {
return "", nil, err
}
cond, cargs, err := r.exists(valueKey)
return fmt.Sprintf("%s(%s, %s, %s)", fn, v, orderExpr, cond), append(vargs, cargs...), err
}}
}
// UniqCount renders uniqIf(<value>, <value> EXISTS) — distinct count of an attribute.
func UniqCount(valueKey *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType) Aggregate {
return Aggregate{keys: keysOf(valueKey), render: func(r aggResolver) (string, []any, error) {
v, vargs, err := r.value(valueKey, dt)
if err != nil {
return "", nil, err
}
cond, cargs, err := r.exists(valueKey)
return fmt.Sprintf("uniqIf(%s, %s)", v, cond), append(vargs, cargs...), err
}}
}
// PredicateCount renders countIf(<predicate>) over a fixed boolean predicate.
func PredicateCount(predicate string) Aggregate {
return Aggregate{render: func(aggResolver) (string, []any, error) {
return fmt.Sprintf("countIf(%s)", sqlbuilder.Escape(predicate)), nil, nil
}}
}
// SumOfKeys renders coalesce(sum(<v1>), 0) + coalesce(sum(<v2>), 0) + … over several
// numeric attributes. Coalesced because a key absent from every span sums to NULL and
// NULL + n = NULL — a trace with only output tokens would otherwise total NULL.
func SumOfKeys(dt telemetrytypes.FieldDataType, valueKeys ...*telemetrytypes.TelemetryFieldKey) Aggregate {
return Aggregate{keys: valueKeys, render: func(r aggResolver) (string, []any, error) {
parts := make([]string, 0, len(valueKeys))
var args []any
for _, k := range valueKeys {
v, vargs, err := r.value(k, dt)
if err != nil {
return "", nil, err
}
parts = append(parts, fmt.Sprintf("coalesce(sum(%s), 0)", v))
args = append(args, vargs...)
}
return strings.Join(parts, " + "), args, nil
}}
}
func keysOf(k *telemetrytypes.TelemetryFieldKey) []*telemetrytypes.TelemetryFieldKey {
return []*telemetrytypes.TelemetryFieldKey{k}
}
// ColumnProvider supplies the columns a trace list computes.
type ColumnProvider interface {
Columns() []TraceColumn
// DefaultOrderAlias is sorted by (desc) when the query gives no order.
DefaultOrderAlias() string
// AggregateAliases are the computed per-trace column names, used to classify a
// filter key as trace-level vs span-level. Excludes SpanLevel columns.
AggregateAliases() []string
}
// CommonTraceColumns are domain-neutral columns any trace list can reuse. All
// aggregate over every span, so none is Orderable.
func CommonTraceColumns() []TraceColumn {
return []TraceColumn{
{Alias: "start_time", Expr: Intrinsic("min(timestamp)")},
{Alias: "end_time", Expr: Intrinsic("max(timestamp)")},
{Alias: "duration_nano", Expr: Intrinsic("(max(toUnixTimestamp64Nano(timestamp) + duration_nano) - min(toUnixTimestamp64Nano(timestamp)))")},
{Alias: "span_count", Expr: Intrinsic("count()")},
{Alias: "root_span_name", Expr: Intrinsic("anyIf(name, parent_span_id = '')")},
{Alias: "service.name", SpanLevel: true, Expr: Intrinsic("any(resource_string_service$$name)")},
}
}

View File

@@ -0,0 +1,817 @@
package telemetryscopedtraces
import (
"context"
"fmt"
"log/slog"
"sort"
"strings"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/flagger"
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/telemetryresourcefilter"
"github.com/SigNoz/signoz/pkg/telemetrystore"
"github.com/SigNoz/signoz/pkg/telemetrytraces"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
qbvariables "github.com/SigNoz/signoz/pkg/variables"
"github.com/huandu/go-sqlbuilder"
)
var (
ErrUnsupportedRequestType = errors.NewInvalidInputf(errors.CodeInvalidInput, "unsupported request type for the scoped trace builder")
)
// scopedTraceStatementBuilder builds a trace list scoped to one span category
// (e.g. gen_ai spans). The query shape is fixed; BaseConditionProvider decides which
// spans are in scope and ColumnProvider decides the per-trace columns, so a new
// category only needs a new pair of providers.
type scopedTraceStatementBuilder struct {
logger *slog.Logger
metadataStore telemetrytypes.MetadataStore
fm qbtypes.FieldMapper
cb qbtypes.ConditionBuilder
baseCond BaseConditionProvider
columnProvider ColumnProvider
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation]
resourceFilterResolver *telemetryresourcefilter.ResourceFingerprintResolver[qbtypes.TraceAggregation]
skipResourceFingerprintEnabled bool
}
var _ qbtypes.StatementBuilder[qbtypes.TraceAggregation] = (*scopedTraceStatementBuilder)(nil)
// NewScopedTraceStatementBuilder wires the generic trace-list builder. The field
// mapper / condition builder are built here, not injected — the list always scans the
// telemetrytraces span index. traceStmtBuilder (the delegate for the span-list path)
// is injected because the provider already has the canonical instance.
func NewScopedTraceStatementBuilder(
settings factory.ProviderSettings,
metadataStore telemetrytypes.MetadataStore,
baseCond BaseConditionProvider,
columnProvider ColumnProvider,
traceStmtBuilder qbtypes.StatementBuilder[qbtypes.TraceAggregation],
telemetryStore telemetrystore.TelemetryStore,
fl flagger.Flagger,
skipResourceFingerprintEnable bool,
skipResourceFingerprintThreshold uint64,
) qbtypes.StatementBuilder[qbtypes.TraceAggregation] {
scopedSettings := factory.NewScopedProviderSettings(settings, "github.com/SigNoz/signoz/pkg/telemetryscopedtraces")
fieldMapper := telemetrytraces.NewFieldMapper()
conditionBuilder := telemetrytraces.NewConditionBuilder(fieldMapper)
// Same resource-fingerprint prune as the standard trace builder — the list scans
// the same span index.
resourceFilterResolver := telemetryresourcefilter.NewResolver[qbtypes.TraceAggregation](
settings,
telemetrytraces.DBName,
telemetrytraces.TracesResourceV3TableName,
telemetrytypes.SignalTraces,
telemetrytypes.SourceUnspecified,
metadataStore,
nil,
fl,
telemetryStore,
skipResourceFingerprintThreshold,
)
return &scopedTraceStatementBuilder{
logger: scopedSettings.Logger(),
metadataStore: metadataStore,
fm: fieldMapper,
cb: conditionBuilder,
baseCond: baseCond,
columnProvider: columnProvider,
traceStmtBuilder: traceStmtBuilder,
resourceFilterResolver: resourceFilterResolver,
skipResourceFingerprintEnabled: skipResourceFingerprintEnable,
}
}
func (b *scopedTraceStatementBuilder) Build(
ctx context.Context,
start uint64,
end uint64,
requestType qbtypes.RequestType,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
switch requestType {
case qbtypes.RequestTypeTrace:
return b.buildTraceListQuery(ctx, querybuilder.ToNanoSecs(start), querybuilder.ToNanoSecs(end), query, variables)
case qbtypes.RequestTypeRaw:
return b.buildDelegated(ctx, start, end, requestType, query, variables)
default:
return nil, ErrUnsupportedRequestType
}
}
// buildDelegated ANDs the base gate into the user filter and delegates to the
// standard trace builder (the span-list / raw path).
func (b *scopedTraceStatementBuilder) buildDelegated(
ctx context.Context,
start, end uint64,
requestType qbtypes.RequestType,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
gate := b.baseCond.FilterExpression()
expr := gate
if query.Filter != nil && strings.TrimSpace(query.Filter.Expression) != "" {
expr = fmt.Sprintf("(%s) AND (%s)", gate, query.Filter.Expression)
}
// shallow copy; only Filter is replaced, caller's query untouched
gated := query
gated.Filter = &qbtypes.Filter{Expression: expr}
return b.traceStmtBuilder.Build(ctx, start, end, requestType, gated, variables)
}
// buildTraceListQuery wires the CTE pipeline (benchmarked, see ai-qb-handoff.md): one
// windowed pass picks the top-N traces, then a bucket-pruned pass enriches only those.
// Helpers appear in this file in the order they run. start/end are nanoseconds.
//
// RESOLVE (keys/columns → SQL via the field mapper)
// fetchKeys metadata for every key we reference
// resolveMask the "span is in scope" predicate (OR of EXISTS)
// resolveColumns per-trace column SQL
// resolveListOrders which columns to ORDER BY
// splitFilter span-level predicate + trace-level HAVING
//
// BUILD
// matched one windowed, mask-pruned GROUP BY trace_id scan fusing gate + span
// │ filter + HAVING + ORDER BY + LIMIT/OFFSET → the top-N trace_ids
// ▼
// ranked [start,end] bounds of those traces, from the small summary table
// ▼
// buckets the ts_bucket_start values they touch, to prune the next scan
// ▼
// enrichment every per-trace column for those traces over their full extent
// (not window-clipped), scanning only their buckets
//
// Only Orderable columns are computable in the mask-pruned matched pass, so only they
// can be ordered or filtered on; all-span columns (span_count, …) are output-only.
func (b *scopedTraceStatementBuilder) buildTraceListQuery(
ctx context.Context,
start, end uint64,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
variables map[string]qbtypes.VariableItem,
) (*qbtypes.Statement, error) {
startBucket := start/querybuilder.NsToSeconds - querybuilder.BucketAdjustment
endBucket := end / querybuilder.NsToSeconds
limit := query.Limit
if limit <= 0 {
limit = 100
}
// Resolve keys and columns once; all attribute access goes through the field mapper.
keys, err := b.fetchKeys(ctx)
if err != nil {
return nil, err
}
maskExpr, maskArgs, err := b.resolveMask(ctx, start, end, keys)
if err != nil {
return nil, err
}
resolved, err := b.resolveColumns(ctx, start, end, keys, maskExpr, maskArgs)
if err != nil {
return nil, err
}
orders, err := b.resolveListOrders(query.Order, resolved)
if err != nil {
return nil, err
}
orderableSet := orderableAliasSet(resolved)
// If the filter references resource attributes, add a __resource_filter CTE and
// narrow the matched scan by resource_fingerprint; skipResourceFilter then drops
// those keys from the span predicate so they aren't applied twice.
resourceFrag, resourceArgs, resourcePred, skipResourceFilter, err := b.maybeAttachResourceFilter(ctx, query, start, end, variables)
if err != nil {
return nil, err
}
// Split the user filter: span-level predicate + trace-level HAVING expression.
fp, err := b.splitFilter(ctx, query, b.aggregateAliasSet(), orderableSet, start, end, skipResourceFilter, variables)
if err != nil {
return nil, err
}
// matched → ranked → buckets → enrichment
matchedFrag, matchedArgs, err := b.buildMatchedCTE(start, end, startBucket, endBucket, resolved, orders, orderableSet, maskExpr, maskArgs, fp, resourcePred, limit, query.Offset)
if err != nil {
return nil, err
}
rankedFrag, rankedArgs := b.buildRankedCTE(start, end)
bucketsFrag := buildBucketsCTE()
mainSQL, mainArgs := b.buildEnrichmentSelect(resolved, orders)
cteFragments := []string{matchedFrag, rankedFrag, bucketsFrag}
cteArgs := [][]any{matchedArgs, rankedArgs, nil}
// __resource_filter must precede `matched`, which references it.
if resourceFrag != "" {
cteFragments = append([]string{resourceFrag}, cteFragments...)
cteArgs = append([][]any{resourceArgs}, cteArgs...)
}
finalSQL := querybuilder.CombineCTEs(cteFragments) + mainSQL + " SETTINGS distributed_product_mode='allow', max_memory_usage=10000000000"
finalArgs := querybuilder.PrependArgs(cteArgs, mainArgs)
return &qbtypes.Statement{
Query: finalSQL,
Args: finalArgs,
Warnings: fp.warnings,
WarningsDocURL: fp.warningsURL,
}, nil
}
// maybeAttachResourceFilter builds the __resource_filter CTE (fingerprints matching
// the filter's resource conditions) and the predicate narrowing the span scan by
// resource_fingerprint, mirroring the standard trace builder. With no resolver or no
// resource conditions it returns empty fragments and the resource keys stay in the
// span predicate (skipResourceFilter=false).
func (b *scopedTraceStatementBuilder) maybeAttachResourceFilter(
ctx context.Context,
query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation],
start, end uint64,
variables map[string]qbtypes.VariableItem,
) (cteFrag string, cteArgs []any, fingerprintPred string, skipResourceFilter bool, err error) {
if b.resourceFilterResolver == nil {
return "", nil, "", false, nil
}
if b.skipResourceFingerprintEnabled {
decision, err := b.resourceFilterResolver.Resolve(ctx, query, start, end, variables)
if err != nil {
return "", nil, "", true, err
}
switch decision {
case qbtypes.ResourceFilterResolveKindNoOp:
return "", nil, "", true, nil
case qbtypes.ResourceFilterResolveKindFallback:
return "", nil, "", false, nil
}
}
stmt, err := b.resourceFilterResolver.StatementBuilder().Build(
ctx, start, end, qbtypes.RequestTypeRaw, query, variables,
)
if err != nil {
return "", nil, "", true, err
}
if stmt == nil {
return "", nil, "", true, nil
}
return fmt.Sprintf("__resource_filter AS (%s)", stmt.Query), stmt.Args,
"resource_fingerprint GLOBAL IN (SELECT fingerprint FROM __resource_filter)", true, nil
}
// ---------------------------------------------------------------------------
// RESOLVE — turn keys/columns into field-mapper-aware SQL
// ---------------------------------------------------------------------------
func (b *scopedTraceStatementBuilder) fetchKeys(ctx context.Context) (map[string][]*telemetrytypes.TelemetryFieldKey, error) {
fields := b.resolverFieldKeys()
selectors := make([]*telemetrytypes.FieldKeySelector, 0, len(fields))
for _, k := range fields {
selectors = append(selectors, &telemetrytypes.FieldKeySelector{
Name: k.Name,
Signal: k.Signal,
FieldContext: k.FieldContext,
})
}
keys, _, err := b.metadataStore.GetKeysMulti(ctx, selectors)
return keys, err
}
func (b *scopedTraceStatementBuilder) resolverFieldKeys() []*telemetrytypes.TelemetryFieldKey {
seen := make(map[string]struct{})
var out []*telemetrytypes.TelemetryFieldKey
add := func(k *telemetrytypes.TelemetryFieldKey) {
if k == nil {
return
}
if _, dup := seen[k.Name]; dup {
return
}
seen[k.Name] = struct{}{}
out = append(out, k)
}
for _, k := range b.baseCond.FieldKeys() {
add(k)
}
for _, c := range b.columnProvider.Columns() {
for _, k := range c.Expr.keys {
add(k)
}
}
return out
}
// resolveMask builds the per-span in-scope mask: OR of resolved EXISTS predicates
// over the base condition's field keys.
func (b *scopedTraceStatementBuilder) resolveMask(ctx context.Context, start, end uint64, keys map[string][]*telemetrytypes.TelemetryFieldKey) (string, []any, error) {
fieldKeys := b.baseCond.FieldKeys()
parts := make([]string, 0, len(fieldKeys))
var args []any
for _, key := range fieldKeys {
e, a, err := b.existsExpr(ctx, start, end, keys, key)
if err != nil {
return "", nil, err
}
parts = append(parts, e)
args = append(args, a...)
}
return "(" + strings.Join(parts, " OR ") + ")", args, nil
}
// existsExpr resolves an EXISTS predicate for key via the field mapper (materialized
// column when present, else map access). Escaped once so it can be embedded in an
// outer builder.
func (b *scopedTraceStatementBuilder) existsExpr(ctx context.Context, start, end uint64, keys map[string][]*telemetrytypes.TelemetryFieldKey, key *telemetrytypes.TelemetryFieldKey) (string, []any, error) {
resolvedKey := key
cands := keys[key.Name]
if len(cands) == 0 {
cands = []*telemetrytypes.TelemetryFieldKey{key}
} else {
resolvedKey = cands[0]
}
sb := sqlbuilder.NewSelectBuilder()
conds, _, err := b.cb.ConditionFor(ctx, start, end, resolvedKey, cands, qbtypes.FilterOperatorExists, nil, sb)
if err != nil {
return "", nil, err
}
// One condition per candidate variant (a key can be ingested under several data
// types); OR them all, like the visitor does for EXISTS.
if len(conds) == 1 {
sb.Where(conds[0])
} else {
sb.Where(sb.Or(conds...))
}
expr, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
expr = strings.TrimPrefix(expr, "WHERE ")
return sqlbuilder.Escape(expr), args, nil
}
// resolvedColumn is a column resolved to SQL via the field mapper; expr is escaped
// once, ready to embed in an outer SELECT.
type resolvedColumn struct {
alias string
expr string
args []any
orderable bool
}
// resolveColumns turns the declarative columns into SQL, resolving all
// attribute access through the field mapper.
func (b *scopedTraceStatementBuilder) resolveColumns(ctx context.Context, start, end uint64, keys map[string][]*telemetrytypes.TelemetryFieldKey, maskExpr string, maskArgs []any) ([]resolvedColumn, error) {
r := aggResolver{
exists: func(key *telemetrytypes.TelemetryFieldKey) (string, []any, error) {
return b.existsExpr(ctx, start, end, keys, key)
},
value: func(key *telemetrytypes.TelemetryFieldKey, dt telemetrytypes.FieldDataType) (string, []any, error) {
// Use the metadata variant, which carries Materialized — a provider's static
// definition never does, so a promoted attribute would otherwise fall back
// to map access. Mirrors existsExpr.
if cands := keys[key.Name]; len(cands) > 0 {
key = cands[0]
}
return querybuilder.CollisionHandledFinalExpr(ctx, start, end, key, b.fm, b.cb, keys, dt, nil, false)
},
maskExpr: maskExpr,
maskArgs: maskArgs,
}
cols := b.columnProvider.Columns()
out := make([]resolvedColumn, 0, len(cols))
for _, c := range cols {
expr, args, err := c.Expr.render(r)
if err != nil {
return nil, err
}
out = append(out, resolvedColumn{alias: c.Alias, expr: expr, args: args, orderable: c.Orderable})
}
return out, nil
}
// listOrder is a sort key resolved to a column alias + direction; both the matched
// CTE and the enrichment ORDER BY it.
type listOrder struct {
alias string
direction string
}
// resolveListOrders maps order keys to the resolved orderable columns; non-orderable
// columns are rejected. Defaults to the column provider's default order.
func (b *scopedTraceStatementBuilder) resolveListOrders(order []qbtypes.OrderBy, resolved []resolvedColumn) ([]listOrder, error) {
byAlias := make(map[string]resolvedColumn, len(resolved))
orderable := make([]string, 0, len(resolved))
for _, rc := range resolved {
byAlias[rc.alias] = rc
if rc.orderable {
orderable = append(orderable, rc.alias)
}
}
if len(order) == 0 {
return []listOrder{{alias: b.columnProvider.DefaultOrderAlias(), direction: "DESC"}}, nil
}
orders := make([]listOrder, 0, len(order))
for _, o := range order {
direction := "DESC"
if o.Direction == qbtypes.OrderDirectionAsc {
direction = "ASC"
}
rc, ok := byAlias[o.Key.Name]
if !ok || !rc.orderable {
return nil, errors.NewInvalidInputf(errors.CodeInvalidInput,
"unsupported order key %q for the trace list; orderable keys: %s", o.Key.Name, strings.Join(orderable, ", "))
}
orders = append(orders, listOrder{alias: rc.alias, direction: direction})
}
return orders, nil
}
// filterParts is the user filter split into a span-level predicate (widens the
// matched WHERE prune and becomes a countIf existence check in HAVING) and a
// trace-level HAVING expression.
type filterParts struct {
spanPred string
spanArgs []any
hasSpanFilter bool
havingExpr string
warnings []string
warningsURL string
}
// splitFilter splits query.Filter into a span-level predicate and a trace-level
// HAVING expression (an explicit query.Having is ANDed onto the latter), then
// validates the trace-level part against the matched-pass aggregates.
func (b *scopedTraceStatementBuilder) splitFilter(ctx context.Context, query qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation], classifySet, orderableSet map[string]struct{}, start, end uint64, skipResourceFilter bool, variables map[string]qbtypes.VariableItem) (filterParts, error) {
var fp filterParts
if query.Filter != nil && strings.TrimSpace(query.Filter.Expression) != "" {
spanExpr, traceExpr, err := querybuilder.SplitFilterForAggregates(query.Filter.Expression, classifySet)
if err != nil {
return fp, err
}
fp.havingExpr = traceExpr
if strings.TrimSpace(spanExpr) != "" {
pred, args, warnings, url, err := b.resolveSpanPredicate(ctx, start, end, spanExpr, skipResourceFilter, variables)
if err != nil {
return fp, err
}
// pred is empty when the span-level keys were all resource attributes
// already handled by the __resource_filter CTE.
if strings.TrimSpace(pred) != "" {
fp.spanPred, fp.spanArgs, fp.hasSpanFilter = pred, args, true
}
fp.warnings, fp.warningsURL = warnings, url
}
}
if query.Having != nil && strings.TrimSpace(query.Having.Expression) != "" {
if fp.havingExpr != "" {
fp.havingExpr = fmt.Sprintf("(%s) AND (%s)", fp.havingExpr, query.Having.Expression)
} else {
fp.havingExpr = query.Having.Expression
}
}
// The span predicate binds variables via PrepareWhereClause; the HAVING is a plain
// text rewrite, so substitute variables here (list/IN quoting, __all__ drops the
// condition) before validating.
if strings.TrimSpace(fp.havingExpr) != "" && len(variables) > 0 {
replaced, err := qbvariables.ReplaceVariablesInExpression(fp.havingExpr, variables)
if err != nil {
return fp, err
}
fp.havingExpr = replaced
}
if err := validateAggregateFilter(fp.havingExpr, orderableSet); err != nil {
return fp, err
}
return fp, nil
}
// resolveSpanPredicate resolves a span-level filter expression to a bare boolean
// SQL predicate + args via the field mapper.
func (b *scopedTraceStatementBuilder) resolveSpanPredicate(ctx context.Context, start, end uint64, expr string, skipResourceFilter bool, variables map[string]qbtypes.VariableItem) (string, []any, []string, string, error) {
selectors := querybuilder.QueryStringToKeysSelectors(expr)
for i := range selectors {
selectors[i].Signal = telemetrytypes.SignalTraces
}
keys, _, err := b.metadataStore.GetKeysMulti(ctx, selectors)
if err != nil {
return "", nil, nil, "", err
}
prepared, err := querybuilder.PrepareWhereClause(expr, querybuilder.FilterExprVisitorOpts{
Context: ctx,
Logger: b.logger,
FieldMapper: b.fm,
ConditionBuilder: b.cb,
FieldKeys: keys,
SkipResourceFilter: skipResourceFilter,
Variables: variables,
StartNs: start,
EndNs: end,
})
if err != nil {
return "", nil, nil, "", err
}
if prepared.IsEmpty() {
return "", nil, nil, "", nil
}
sb := sqlbuilder.NewSelectBuilder()
sb.Select("1")
sb.AddWhereClause(prepared.WhereClause)
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
pred := sql[strings.Index(sql, "WHERE ")+len("WHERE "):]
return sqlbuilder.Escape(pred), args, prepared.Warnings, prepared.WarningsDocURL, nil
}
// buildMatchedCTE builds `matched`: the single windowed GROUP BY trace_id scan that
// fuses gate + span filter + HAVING + ORDER BY + LIMIT/OFFSET, selecting only the
// aliases the ORDER BY / HAVING reference.
func (b *scopedTraceStatementBuilder) buildMatchedCTE(start, end, startBucket, endBucket uint64, resolved []resolvedColumn, orders []listOrder, orderableSet map[string]struct{}, maskExpr string, maskArgs []any, fp filterParts, resourcePred string, limit, offset int) (string, []any, error) {
sb := sqlbuilder.NewSelectBuilder()
// SELECT trace_id + only the aggregates ORDER BY / HAVING reference (as aliases).
needed := neededMatchedAliases(orders, fp.havingExpr, orderableSet)
selects := []string{"trace_id"}
for _, rc := range resolved {
if _, ok := needed[rc.alias]; !ok {
continue
}
colExpr, err := embedExpr(sb, rc.expr, rc.args)
if err != nil {
return "", nil, err
}
selects = append(selects, colExpr+" AS "+quoteAlias(rc.alias))
}
sb.Select(selects...)
sb.From(spanTable())
// WHERE: window + prune to in-scope spans, widened by the span filter so its
// spans survive for the countIf existence check below.
win := windowWhere(sb, start, end, startBucket, endBucket)
mask, err := embedExpr(sb, maskExpr, maskArgs)
if err != nil {
return "", nil, err
}
prune := "(" + mask
if fp.hasSpanFilter {
spanPred, err := embedExpr(sb, fp.spanPred, fp.spanArgs)
if err != nil {
return "", nil, err
}
prune += " OR " + spanPred
}
prune += ")"
where := append(win, prune)
if resourcePred != "" {
where = append(where, resourcePred)
}
sb.Where(where...)
sb.GroupBy("trace_id")
// HAVING: the gate/span existence checks are only needed when the WHERE was
// widened by a span filter; otherwise the mask alone already enforces the gate.
var having []string
if fp.hasSpanFilter {
havingMask, err := embedExpr(sb, maskExpr, maskArgs)
if err != nil {
return "", nil, err
}
havingPred, err := embedExpr(sb, fp.spanPred, fp.spanArgs)
if err != nil {
return "", nil, err
}
having = append(having, "countIf("+havingMask+") > 0")
having = append(having, "countIf("+havingPred+") > 0")
}
if strings.TrimSpace(fp.havingExpr) != "" {
hv, err := b.buildHaving(fp.havingExpr, orderableSet)
if err != nil {
return "", nil, err
}
if hv != "" {
having = append(having, hv)
}
}
if len(having) > 0 {
sb.Having(strings.Join(having, " AND "))
}
sb.OrderBy(orderClause(orders)...)
sb.Limit(limit)
if offset > 0 {
sb.Offset(offset)
}
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return fmt.Sprintf("matched AS (%s)", sql), args, nil
}
// buildRankedCTE builds `ranked`: [start,end] bounds per matched trace, read from the
// small trace-summary table.
func (b *scopedTraceStatementBuilder) buildRankedCTE(start, end uint64) (string, []any) {
sb := sqlbuilder.NewSelectBuilder()
sb.Select("trace_id", "min(start) AS t_start", "max(end) AS t_end")
sb.From(summaryTable())
sb.Where(
"trace_id GLOBAL IN (SELECT trace_id FROM matched)",
"end >= fromUnixTimestamp64Nano("+sb.Var(start)+")",
"start < fromUnixTimestamp64Nano("+sb.Var(end)+")",
)
sb.GroupBy("trace_id")
sql, args := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return fmt.Sprintf("ranked AS (%s)", sql), args
}
// buildBucketsCTE builds `buckets`: the ts_bucket_start values the matched traces
// span, so the enrichment scan is primary-key pruned. No args.
func buildBucketsCTE() string {
adj := querybuilder.BucketAdjustment // 30-min bucket width in seconds
return fmt.Sprintf("buckets AS (SELECT DISTINCT b AS ts_bucket FROM ranked "+
"ARRAY JOIN range("+
"toUInt64(intDiv(toUnixTimestamp(t_start), %d) * %d - %d), "+
"toUInt64(intDiv(toUnixTimestamp(t_end), %d) * %d + %d), "+
"%d) AS b)", adj, adj, adj, adj, adj, adj, adj)
}
// buildEnrichmentSelect builds the final SELECT: every per-trace column for the
// matched traces over their full extent, scanning only their buckets.
//
// Accepted discrepancy: matched ranks/paginates on window-clipped values, this pass
// recomputes and ORDER BYs full-trace values, so a trace with activity outside the
// window can sort differently than it ranked. Page membership is unaffected
// (LIMIT/OFFSET runs only in matched); rows still sort by the values the user sees.
// Ordering by matched's values instead would re-run the matched scan (ClickHouse
// re-executes a CTE per reference) without fixing the visible cross-page artifact.
func (b *scopedTraceStatementBuilder) buildEnrichmentSelect(resolved []resolvedColumn, orders []listOrder) (string, []any) {
sb := sqlbuilder.NewSelectBuilder()
selects, selectArgs := selectAllColumns(resolved)
sb.Select(selects...)
sb.From(spanTable())
sb.Where(
"ts_bucket_start GLOBAL IN (SELECT ts_bucket FROM buckets)",
"trace_id GLOBAL IN (SELECT trace_id FROM ranked)",
)
sb.GroupBy("trace_id")
sb.OrderBy(orderClause(orders)...)
sql, builtArgs := sb.BuildWithFlavor(sqlbuilder.ClickHouse)
return sql, append(append([]any{}, selectArgs...), builtArgs...)
}
// buildHaving rewrites a trace-level HAVING expression to the matched-pass column
// aliases; bare, trace., and tracefield. forms all map to the same alias.
func (b *scopedTraceStatementBuilder) buildHaving(havingExpr string, orderableSet map[string]struct{}) (string, error) {
columnMap := make(map[string]string, len(orderableSet)*3)
for a := range orderableSet {
columnMap[a] = quoteAlias(a)
columnMap["trace."+a] = quoteAlias(a)
columnMap["tracefield."+a] = quoteAlias(a)
}
return querybuilder.NewHavingExpressionRewriter().Rewrite(havingExpr, columnMap)
}
// ---------------------------------------------------------------------------
// Small shared SQL-builder utilities
// ---------------------------------------------------------------------------
// spanTable is the fully-qualified span index table.
func spanTable() string {
return fmt.Sprintf("%s.%s", telemetrytraces.DBName, telemetrytraces.SpanIndexV3TableName)
}
// summaryTable is the fully-qualified trace-summary table.
func summaryTable() string {
return fmt.Sprintf("%s.%s", telemetrytraces.DBName, telemetrytraces.TraceSummaryTableName)
}
// aggregateAliasSet is every trace-level column alias, used to classify filter keys
// as trace-level vs span-level.
func (b *scopedTraceStatementBuilder) aggregateAliasSet() map[string]struct{} {
set := make(map[string]struct{}, len(b.columnProvider.AggregateAliases()))
for _, a := range b.columnProvider.AggregateAliases() {
set[a] = struct{}{}
}
return set
}
// orderableAliasSet is the subset of aliases computable in the matched pass — the
// only ones usable in ORDER BY and the aggregate filter.
func orderableAliasSet(resolved []resolvedColumn) map[string]struct{} {
set := make(map[string]struct{})
for _, rc := range resolved {
if rc.orderable {
set[rc.alias] = struct{}{}
}
}
return set
}
// neededMatchedAliases is the minimal alias set the matched pass must select: those
// in ORDER BY plus those in the aggregate HAVING. Everything else is left to the
// enrichment scan.
func neededMatchedAliases(orders []listOrder, havingExpr string, orderableSet map[string]struct{}) map[string]struct{} {
needed := make(map[string]struct{})
for _, o := range orders {
needed[o.alias] = struct{}{}
}
for _, sel := range querybuilder.QueryStringToKeysSelectors(havingExpr) {
name := strings.TrimPrefix(strings.TrimPrefix(sel.Name, "trace."), "tracefield.")
if _, ok := orderableSet[name]; ok {
needed[name] = struct{}{}
}
}
return needed
}
// validateAggregateFilter rejects a trace-level filter referencing an aggregate not
// computable in the matched pass (e.g. span_count, duration_nano).
func validateAggregateFilter(havingExpr string, orderableSet map[string]struct{}) error {
if strings.TrimSpace(havingExpr) == "" {
return nil
}
allowed := make([]string, 0, len(orderableSet))
for a := range orderableSet {
allowed = append(allowed, a)
}
sort.Strings(allowed)
for _, sel := range querybuilder.QueryStringToKeysSelectors(havingExpr) {
name := strings.TrimPrefix(strings.TrimPrefix(sel.Name, "trace."), "tracefield.")
if _, ok := orderableSet[name]; !ok {
return errors.NewInvalidInputf(errors.CodeInvalidInput,
"aggregate %q cannot be used in the trace-list filter; filterable aggregates: %s", name, strings.Join(allowed, ", "))
}
}
return nil
}
// embedExpr inlines a resolved expr into sb, replacing each `?` placeholder with a
// builder Var so the args are tracked in appearance order. Resolved exprs carry
// values only as bound args, so every `?` is a placeholder; a count mismatch would
// silently shift args into the wrong slots — error out instead.
func embedExpr(sb *sqlbuilder.SelectBuilder, expr string, args []any) (string, error) {
if n := strings.Count(expr, "?"); n != len(args) {
return "", errors.NewInternalf(errors.CodeInternal,
"scoped trace builder: %d placeholders != %d args embedding %q", n, len(args), expr)
}
var out strings.Builder
ai := 0
for i := 0; i < len(expr); i++ {
if expr[i] == '?' {
out.WriteString(sb.Var(args[ai]))
ai++
continue
}
out.WriteByte(expr[i])
}
return out.String(), nil
}
// windowWhere binds the time-window predicates to sb and returns them so the caller
// can add its own predicates in the same Where call.
func windowWhere(sb *sqlbuilder.SelectBuilder, start, end, startBucket, endBucket uint64) []string {
return []string{
sb.GE("timestamp", fmt.Sprintf("%d", start)),
sb.L("timestamp", fmt.Sprintf("%d", end)),
sb.GE("ts_bucket_start", startBucket),
sb.LE("ts_bucket_start", endBucket),
}
}
// orderClause renders the ORDER BY terms plus the trace_id tiebreak.
func orderClause(orders []listOrder) []string {
out := make([]string, 0, len(orders)+1)
for _, o := range orders {
out = append(out, fmt.Sprintf("%s %s", quoteAlias(o.alias), o.direction))
}
return append(out, "trace_id DESC")
}
// selectAllColumns renders `expr AS alias` for every resolved column, args in select
// order.
func selectAllColumns(resolved []resolvedColumn) ([]string, []any) {
selects := []string{"trace_id"}
var args []any
for _, rc := range resolved {
selects = append(selects, rc.expr+" AS "+quoteAlias(rc.alias))
args = append(args, rc.args...)
}
return selects, args
}
// quoteAlias backticks an alias containing characters special to the SQL builder.
func quoteAlias(alias string) string {
if strings.ContainsAny(alias, ".$`") {
return "`" + alias + "`"
}
return alias
}

View File

@@ -0,0 +1,139 @@
package telemetryscopedtraces
import (
"context"
"fmt"
"strings"
"testing"
"github.com/SigNoz/signoz/pkg/factory"
"github.com/SigNoz/signoz/pkg/instrumentation/instrumentationtest"
"github.com/SigNoz/signoz/pkg/telemetrytraces"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest"
"github.com/huandu/go-sqlbuilder"
"github.com/stretchr/testify/require"
)
// The full-pipeline golden tests live in pkg/telemetryai, which exercises this
// builder through its production provider pair. The tests here cover only what
// needs the package internals: builder states not reachable through the
// constructor.
// stubGate scopes to spans carrying a single attribute key.
type stubGate struct {
key *telemetrytypes.TelemetryFieldKey
}
func (s stubGate) FilterExpression() string { return s.key.Name + " EXISTS" }
func (s stubGate) FieldKeys() []*telemetrytypes.TelemetryFieldKey {
return []*telemetrytypes.TelemetryFieldKey{s.key}
}
// stubColumns is the common columns plus one orderable scoped aggregate, the
// minimum a column provider must supply (a default order key).
type stubColumns struct {
key *telemetrytypes.TelemetryFieldKey
}
func (s stubColumns) Columns() []TraceColumn {
return append(CommonTraceColumns(),
TraceColumn{Alias: "scoped_span_count", Orderable: true, Expr: CountExists(s.key)})
}
func (s stubColumns) DefaultOrderAlias() string { return "scoped_span_count" }
func (s stubColumns) AggregateAliases() []string {
aliases := make([]string, 0)
for _, c := range s.Columns() {
if !c.SpanLevel {
aliases = append(aliases, c.Alias)
}
}
return aliases
}
// renderSQL substitutes bound args into the `?` placeholders so assertions can
// match the statement as one literal SQL string.
func renderSQL(t *testing.T, stmt *qbtypes.Statement) string {
t.Helper()
var b strings.Builder
argi := 0
for i := 0; i < len(stmt.Query); i++ {
if stmt.Query[i] == '?' {
require.Less(t, argi, len(stmt.Args), "more ? than args in query")
if s, ok := stmt.Args[argi].(string); ok {
b.WriteString("'" + s + "'")
} else {
fmt.Fprintf(&b, "%v", stmt.Args[argi])
}
argi++
continue
}
b.WriteByte(stmt.Query[i])
}
require.Equal(t, len(stmt.Args), argi, "arg count does not match number of placeholders")
return b.String()
}
// With the resolver unset (nil), the resource filter falls back to being applied inline
// on the span index — no fingerprint CTE — so existing behavior is preserved. The nil
// state is not reachable through the constructor, hence the direct struct literal.
func TestBuild_TraceList_ResourceFilter_NoResolver(t *testing.T) {
gateKey := &telemetrytypes.TelemetryFieldKey{
Name: "scope.marker",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextAttribute,
FieldDataType: telemetrytypes.FieldDataTypeString,
}
mockMetadataStore := telemetrytypestest.NewMockMetadataStore()
mockMetadataStore.KeysMap = map[string][]*telemetrytypes.TelemetryFieldKey{
gateKey.Name: {gateKey},
"service.name": {{
Name: "service.name",
Signal: telemetrytypes.SignalTraces,
FieldContext: telemetrytypes.FieldContextResource,
FieldDataType: telemetrytypes.FieldDataTypeString,
}},
}
settings := factory.NewScopedProviderSettings(instrumentationtest.New().ToProviderSettings(), "github.com/SigNoz/signoz/pkg/telemetryscopedtraces")
fm := telemetrytraces.NewFieldMapper()
b := &scopedTraceStatementBuilder{
logger: settings.Logger(),
metadataStore: mockMetadataStore,
fm: fm,
cb: telemetrytraces.NewConditionBuilder(fm),
baseCond: stubGate{key: gateKey},
columnProvider: stubColumns{key: gateKey},
resourceFilterResolver: nil,
}
stmt, err := b.Build(context.Background(), 1747947419000, 1747983448000, qbtypes.RequestTypeTrace,
qbtypes.QueryBuilderQuery[qbtypes.TraceAggregation]{
Signal: telemetrytypes.SignalTraces,
Filter: &qbtypes.Filter{Expression: "resource.service.name = 'checkout'"},
Limit: 10,
}, nil)
require.NoError(t, err)
got := renderSQL(t, stmt)
require.NotContains(t, got, "__resource_filter")
require.Contains(t, got, "resources_string['service.name']")
}
// embedExpr treats every `?` byte as a placeholder; a count/args mismatch (an expr
// carrying a literal `?`, or a dropped arg) must fail loudly instead of silently
// shifting every subsequent arg into the wrong placeholder.
func TestEmbedExpr_PlaceholderArgMismatch(t *testing.T) {
sb := sqlbuilder.NewSelectBuilder()
out, err := embedExpr(sb, "x = ? AND y = ?", []any{1, 2})
require.NoError(t, err)
require.Equal(t, 2, strings.Count(out, "$"), "both placeholders bound as builder vars")
_, err = embedExpr(sb, "x = ? AND y LIKE 'a?b'", []any{1})
require.Error(t, err, "literal ? in the expr must not pass as a placeholder")
_, err = embedExpr(sb, "x = ?", []any{1, 2})
require.Error(t, err, "extra args must not be silently dropped")
}

View File

@@ -33,7 +33,7 @@ def _ai_trace(
now: datetime,
service: str,
user: str,
in_tokens: int,
in_tokens: int | None,
out_tokens: int,
cost: float,
model: str = "gpt-4o-mini",
@@ -41,7 +41,8 @@ def _ai_trace(
error: bool = False,
environment: str = "production",
) -> list[Traces]:
"""A minimal AI trace: root span + one LLM span with gen_ai attributes."""
"""A minimal AI trace: root span + one LLM span with gen_ai attributes.
in_tokens=None omits the input-tokens attribute entirely (not zero)."""
trace_id = TraceIdGenerator.trace_id()
root_id = TraceIdGenerator.span_id()
llm_id = TraceIdGenerator.span_id()
@@ -59,6 +60,16 @@ def _ai_trace(
resources=resources,
attributes={"http.request.method": "POST"},
)
attributes = {
"gen_ai.request.model": model,
"gen_ai.system": "openai",
"gen_ai.user.id": user,
# numeric values land in attributes_number
"gen_ai.usage.output_tokens": out_tokens,
"_signoz.gen_ai.total_cost": cost,
}
if in_tokens is not None:
attributes["gen_ai.usage.input_tokens"] = in_tokens
llm = Traces(
timestamp=now - timedelta(seconds=4),
duration=timedelta(seconds=llm_duration_s),
@@ -69,15 +80,7 @@ def _ai_trace(
kind=TracesKind.SPAN_KIND_CLIENT,
status_code=(TracesStatusCode.STATUS_CODE_ERROR if error else TracesStatusCode.STATUS_CODE_OK),
resources=resources,
attributes={
"gen_ai.request.model": model,
"gen_ai.system": "openai",
"gen_ai.user.id": user,
# numeric values land in attributes_number
"gen_ai.usage.input_tokens": in_tokens,
"gen_ai.usage.output_tokens": out_tokens,
"_signoz.gen_ai.total_cost": cost,
},
attributes=attributes,
)
return [root, llm]
@@ -223,7 +226,10 @@ def test_ai_list_having_aggregate_filter(
"""
Aggregate filter written in the SAME filter box: the span-level predicate narrows
to the service, the trace-level `output_tokens > 100` keeps the large-token
trace and drops the small one (split internally into WHERE + HAVING).
trace and drops the small one (split internally into WHERE + HAVING). All three
spellings of a trace-level aggregate — bare, `trace.`, `tracefield.` — behave
identically (unit tests pin them to byte-identical SQL; this covers the wiring
once end-to-end). An output-only aggregate is rejected under any spelling.
"""
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
service = "ai-it-having"
@@ -237,19 +243,32 @@ def test_ai_list_having_aggregate_filter(
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms, end_ms = _window_ms(now)
query = BuilderQuery(
for spelling in ("output_tokens", "trace.output_tokens", "tracefield.output_tokens"):
query = BuilderQuery(
signal="traces",
source="ai",
name="A",
filter_expression=f"service.name = '{service}' AND {spelling} > 100",
limit=10,
)
response = make_query_request(signoz, token, start_ms, end_ms, [query.to_dict()], request_type="trace")
assert response.status_code == HTTPStatus.OK, f"{spelling}: {response.text}"
body = json.dumps(response.json())
assert large_id in body, f"{spelling}: trace with 500 out-tokens should pass > 100"
assert small_id not in body, f"{spelling}: trace with 20 out-tokens should be filtered out by HAVING"
# output-only aggregate gets the targeted rejection, also under the explicit context.
bad = BuilderQuery(
signal="traces",
source="ai",
name="A",
filter_expression=f"service.name = '{service}' AND output_tokens > 100",
filter_expression="tracefield.span_count > 3",
limit=10,
)
response = make_query_request(signoz, token, start_ms, end_ms, [query.to_dict()], request_type="trace")
assert response.status_code == HTTPStatus.OK, response.text
body = json.dumps(response.json())
assert large_id in body, "trace with 500 out-tokens should pass output_tokens > 100"
assert small_id not in body, "trace with 20 out-tokens should be filtered out by HAVING"
response = make_query_request(signoz, token, start_ms, end_ms, [bad.to_dict()], request_type="trace")
assert response.status_code == HTTPStatus.BAD_REQUEST, response.text
assert "cannot be used" in response.text
def test_ai_list_order_limit_offset(
@@ -387,39 +406,6 @@ def test_ai_list_having_or_aggregates(
assert small_id not in body
def test_ai_list_having_trace_context_prefix(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
"""The `trace.` context prefix on an aggregate column works like the bare name."""
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
service = "ai-it-trace-ctx"
small = _ai_trace(now=now, service=service, user="a", in_tokens=10, out_tokens=20, cost=0.1)
large = _ai_trace(now=now, service=service, user="b", in_tokens=10, out_tokens=500, cost=0.2)
small_id, large_id = small[0].trace_id, large[0].trace_id
insert_traces(small + large)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms, end_ms = _window_ms(now)
query = BuilderQuery(
signal="traces",
source="ai",
name="A",
filter_expression=f"service.name = '{service}' AND trace.output_tokens > 100",
limit=10,
)
response = make_query_request(signoz, token, start_ms, end_ms, [query.to_dict()], request_type="trace")
assert response.status_code == HTTPStatus.OK, response.text
body = json.dumps(response.json())
assert large_id in body
assert small_id not in body
def test_ai_list_resource_filter_isolates_by_fingerprint(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
@@ -587,6 +573,83 @@ def test_ai_list_rejects_order_by_span_attribute(
assert "order key" in response.text
def test_ai_list_total_tokens_output_only(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
"""
A trace whose LLM span carries only output tokens (no input-tokens attribute at
all) must still total: total_tokens is coalesce(sum(in),0)+coalesce(sum(out),0),
since sum over an absent attribute is NULL and NULL + n = NULL in ClickHouse.
"""
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
service = "ai-it-total-coalesce"
insert_traces(_ai_trace(now=now, service=service, user="a", in_tokens=None, out_tokens=300, cost=0.1))
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms, end_ms = _window_ms(now)
query = BuilderQuery(
signal="traces",
source="ai",
name="A",
filter_expression=f"service.name = '{service}'",
limit=10,
)
response = make_query_request(signoz, token, start_ms, end_ms, [query.to_dict()], request_type="trace")
assert response.status_code == HTTPStatus.OK, response.text
rows = response.json()["data"]["data"]["results"][0]["rows"]
assert len(rows) == 1, f"expected one trace, got: {rows}"
data = rows[0]["data"]
assert data["input_tokens"] is None, data # attribute absent -> NULL, not 0
assert data["output_tokens"] == 300, data
assert data["total_tokens"] == 300, f"total must coalesce the missing input side: {data}"
def test_ai_list_variable_in_aggregate_filter(
signoz: types.SigNoz,
create_user_admin: None, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
insert_traces: Callable[[list[Traces]], None],
) -> None:
"""A query variable in a trace-level condition is substituted into the HAVING."""
now = datetime.now(tz=UTC).replace(second=0, microsecond=0)
service = "ai-it-having-var"
small = _ai_trace(now=now, service=service, user="a", in_tokens=10, out_tokens=20, cost=0.1)
large = _ai_trace(now=now, service=service, user="b", in_tokens=10, out_tokens=500, cost=0.2)
small_id, large_id = small[0].trace_id, large[0].trace_id
insert_traces(small + large)
token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
start_ms, end_ms = _window_ms(now)
query = BuilderQuery(
signal="traces",
source="ai",
name="A",
filter_expression=f"service.name = '{service}' AND trace.output_tokens > $threshold",
limit=10,
)
response = make_query_request(
signoz,
token,
start_ms,
end_ms,
[query.to_dict()],
request_type="trace",
variables={"threshold": {"type": "custom", "value": 100}},
)
assert response.status_code == HTTPStatus.OK, response.text
body = json.dumps(response.json())
assert large_id in body
assert small_id not in body
def _ai_trace_two_llm(*, now: datetime, service: str) -> list[Traces]:
"""Root + two LLM spans at different times, each with distinct input/output messages."""
trace_id = TraceIdGenerator.trace_id()