From 94de5cf72b429b6425b83f5c7a2aee1fad9aabae Mon Sep 17 00:00:00 2001 From: Nikhil Soni Date: Thu, 1 Oct 2026 17:00:08 +0530 Subject: [PATCH] refactor(promote): move index creation into the metadata store --- pkg/modules/promote/implpromote/module.go | 42 +++--------------- .../promote/implpromote/module_test.go | 41 +++++++++-------- pkg/signoz/module.go | 2 +- pkg/telemetrymetadata/body_json_metadata.go | 16 +++++++ .../body_json_metadata_test.go | 44 +++++++++++++++++++ pkg/types/telemetrytypes/store.go | 4 ++ .../telemetrytypestest/metadata_store.go | 8 ++++ 7 files changed, 100 insertions(+), 57 deletions(-) diff --git a/pkg/modules/promote/implpromote/module.go b/pkg/modules/promote/implpromote/module.go index 21665f78cd..f78467af2b 100644 --- a/pkg/modules/promote/implpromote/module.go +++ b/pkg/modules/promote/implpromote/module.go @@ -7,25 +7,18 @@ import ( schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator" "github.com/SigNoz/signoz/pkg/errors" "github.com/SigNoz/signoz/pkg/modules/promote" - "github.com/SigNoz/signoz/pkg/telemetrystore" - "github.com/SigNoz/signoz/pkg/types/ctxtypes" - "github.com/SigNoz/signoz/pkg/types/instrumentationtypes" "github.com/SigNoz/signoz/pkg/types/promotetypes" "github.com/SigNoz/signoz/pkg/types/telemetrytypes" ) -var ( - CodeFailedToCreateIndex = errors.MustNewCode("failed_to_create_index_promoted_paths") - CodeFailedToQueryPromotedPaths = errors.MustNewCode("failed_to_query_promoted_paths") -) +var CodeFailedToQueryPromotedPaths = errors.MustNewCode("failed_to_query_promoted_paths") type module struct { - metadataStore telemetrytypes.MetadataStore - telemetryStore telemetrystore.TelemetryStore + metadataStore telemetrytypes.MetadataStore } -func NewModule(metadataStore telemetrytypes.MetadataStore, telemetrystore telemetrystore.TelemetryStore) promote.Module { - return &module{metadataStore: metadataStore, telemetryStore: telemetrystore} +func NewModule(metadataStore telemetrytypes.MetadataStore) promote.Module { + return &module{metadataStore: metadataStore} } func (m *module) ListPromotedPaths(ctx context.Context, filters promotetypes.ListPromotedPathsFilters) ([]promotetypes.PromotePath, error) { @@ -190,35 +183,10 @@ func (m *module) promotePaths(ctx context.Context, target promotetypes.Target, p } if len(indexes) > 0 { - if err := m.createIndexes(ctx, target, indexes); err != nil { + if err := m.metadataStore.CreateIndexes(ctx, target.IndexSource(), indexes); err != nil { return err } } return nil } - -func (m *module) createIndexes(ctx context.Context, target promotetypes.Target, indexes []schemamigrator.Index) error { - ctx = ctxtypes.NewContextWithCommentVals(ctx, map[string]string{ - instrumentationtypes.TelemetrySignal: target.Entry.Signal.StringValue(), - instrumentationtypes.CodeNamespace: "promote", - instrumentationtypes.CodeFunctionName: "createIndexes", - }) - if len(indexes) == 0 { - return nil - } - - for _, index := range indexes { - alterStmt := schemamigrator.AlterTableAddIndex{ - Database: target.DBName, - Table: target.LocalTableName, - Index: index, - } - op := alterStmt.OnCluster(m.telemetryStore.Cluster()) - if err := m.telemetryStore.ClickhouseDB().Exec(ctx, op.ToSQL()); err != nil { - return errors.WrapInternalf(err, CodeFailedToCreateIndex, "failed to create index") - } - } - - return nil -} diff --git a/pkg/modules/promote/implpromote/module_test.go b/pkg/modules/promote/implpromote/module_test.go index 903ef0a9a5..f59c5a4db2 100644 --- a/pkg/modules/promote/implpromote/module_test.go +++ b/pkg/modules/promote/implpromote/module_test.go @@ -2,12 +2,8 @@ package implpromote import ( "context" - "regexp" "testing" - sqlmock "github.com/DATA-DOG/go-sqlmock" - "github.com/SigNoz/signoz/pkg/telemetrystore" - "github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest" "github.com/SigNoz/signoz/pkg/types/promotetypes" "github.com/SigNoz/signoz/pkg/types/telemetrytypes" "github.com/SigNoz/signoz/pkg/types/telemetrytypes/telemetrytypestest" @@ -85,7 +81,7 @@ func TestPromotePaths(t *testing.T) { for _, testCase := range testCases { t.Run(testCase.name, func(t *testing.T) { store := telemetrytypestest.NewMockMetadataStore() - m := NewModule(store, nil) + m := NewModule(store) err := m.PromotePaths(ctx, testCase.paths...) if testCase.wantErr { @@ -112,10 +108,11 @@ func TestPromotePathsCreatesIndexes(t *testing.T) { ctx := context.Background() testCases := []struct { - name string - promoted map[string]bool - path *promotetypes.PromotePath - wantDDLColumn string + name string + promoted map[string]bool + path *promotetypes.PromotePath + wantName string + wantExprPart string }{ { name: "LogsNewPromotion_IndexesPromotedColumn", @@ -128,7 +125,8 @@ func TestPromotePathsCreatesIndexes(t *testing.T) { {FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1}, }, }, - wantDDLColumn: "dynamicElement(body_promoted.user.name", + wantName: "`body_promoted.user.name_String_ngrambf_v1`", + wantExprPart: "dynamicElement(body_promoted.user.name", }, { name: "LogsAlreadyPromoted_IndexesPromotedColumn", @@ -141,7 +139,8 @@ func TestPromotePathsCreatesIndexes(t *testing.T) { {FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1}, }, }, - wantDDLColumn: "dynamicElement(body_promoted.user.name", + wantName: "`body_promoted.user.name_String_ngrambf_v1`", + wantExprPart: "dynamicElement(body_promoted.user.name", }, { name: "LogsUnpromotedPath_IndexesBaseColumn", @@ -153,7 +152,8 @@ func TestPromotePathsCreatesIndexes(t *testing.T) { {FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1}, }, }, - wantDDLColumn: "dynamicElement(body_v2.user.name", + wantName: "`body_v2.user.name_String_ngrambf_v1`", + wantExprPart: "dynamicElement(body_v2.user.name", }, { name: "TracesNewPromotion_IndexesPromotedColumn", @@ -166,7 +166,8 @@ func TestPromotePathsCreatesIndexes(t *testing.T) { {FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1}, }, }, - wantDDLColumn: "`attributes_promoted.http.method_String_ngrambf_v1` attributes_promoted.http.method::String", + wantName: "`attributes_promoted.http.method_String_ngrambf_v1`", + wantExprPart: "attributes_promoted.http.method::String", }, { name: "TracesUnpromotedPath_IndexesBaseColumn", @@ -178,22 +179,24 @@ func TestPromotePathsCreatesIndexes(t *testing.T) { {FieldDataType: telemetrytypes.FieldDataTypeString, Type: "ngrambf_v1(4, 1024, 2, 0)", Granularity: 1}, }, }, - wantDDLColumn: "`attributes.http.method_String_ngrambf_v1` attributes.http.method::String", + wantName: "`attributes.http.method_String_ngrambf_v1`", + wantExprPart: "attributes.http.method::String", }, } for _, testCase := range testCases { t.Run(testCase.name, func(t *testing.T) { - ts := telemetrystoretest.New(telemetrystore.Config{}, sqlmock.QueryMatcherRegexp) store := telemetrytypestest.NewMockMetadataStore() if testCase.promoted != nil { store.PromotedPathsMap = testCase.promoted } - m := NewModule(store, ts) + m := NewModule(store) - ts.Mock().ExpectExec("ADD INDEX (.+)" + regexp.QuoteMeta(testCase.wantDDLColumn)).WillReturnError(nil) require.NoError(t, m.PromotePaths(ctx, testCase.path)) - assert.NoError(t, ts.Mock().ExpectationsWereMet()) + require.Len(t, store.CreatedIndexes, 1) + assert.Equal(t, testCase.wantName, store.CreatedIndexes[0].Name) + assert.Contains(t, store.CreatedIndexes[0].Expression, testCase.wantExprPart) + assert.Equal(t, "ngrambf_v1(4, 1024, 2, 0)", store.CreatedIndexes[0].Type) }) } } @@ -382,7 +385,7 @@ func TestListPromotedPaths(t *testing.T) { store := telemetrytypestest.NewMockMetadataStore() store.PromotedPathsMap = testCase.promoted store.LogsJSONIndexes = testCase.indexes - m := NewModule(store, nil) + m := NewModule(store) paths, err := m.ListPromotedPaths(ctx, testCase.filters) require.NoError(t, err) diff --git a/pkg/signoz/module.go b/pkg/signoz/module.go index 7a7e78b693..96e5c9d39c 100644 --- a/pkg/signoz/module.go +++ b/pkg/signoz/module.go @@ -155,7 +155,7 @@ func NewModules( MetricsExplorer: implmetricsexplorer.NewModule(telemetryStore, telemetryMetadataStore, cache, ruleStore, dashboard, fl, providerSettings, config.MetricsExplorer), MetricReductionRule: metricReductionRule, InfraMonitoring: implinframonitoring.NewModule(telemetryStore, telemetryMetadataStore, querier, fl, providerSettings, config.InfraMonitoring), - Promote: implpromote.NewModule(telemetryMetadataStore, telemetryStore), + Promote: implpromote.NewModule(telemetryMetadataStore), ServiceAccount: serviceAccount, ServiceAccountGetter: serviceAccountGetter, LogsPipeline: impllogspipeline.NewModule(sqlstore), diff --git a/pkg/telemetrymetadata/body_json_metadata.go b/pkg/telemetrymetadata/body_json_metadata.go index 8d91fd74f5..61837ec788 100644 --- a/pkg/telemetrymetadata/body_json_metadata.go +++ b/pkg/telemetrymetadata/body_json_metadata.go @@ -25,6 +25,7 @@ var ( CodeFailLoadPromotedPaths = errors.MustNewCode("fail_load_promoted_paths") CodeFailCheckPathPromoted = errors.MustNewCode("fail_check_path_promoted") CodeFailLoadLogsJSONIndexes = errors.MustNewCode("fail_load_logs_json_indexes") + CodeFailCreateJSONIndexes = errors.MustNewCode("fail_create_json_indexes") CodeFailListJSONValues = errors.MustNewCode("fail_list_json_values") CodeFailScanJSONValue = errors.MustNewCode("fail_scan_json_value") CodeFailScanVariant = errors.MustNewCode("fail_scan_variant") @@ -234,6 +235,21 @@ func (t *telemetryMetaStore) ListJSONIndexes(ctx context.Context, source telemet return indexes, nil } +func (t *telemetryMetaStore) CreateIndexes(ctx context.Context, source telemetrytypes.JSONIndexSource, indexes []schemamigrator.Index) error { + ctx = withTelemetryContext(ctx, source.Signal, "CreateIndexes") + for _, index := range indexes { + op := schemamigrator.AlterTableAddIndex{ + Database: source.DBName, + Table: source.LocalTableName, + Index: index, + }.OnCluster(t.telemetrystore.Cluster()) + if err := t.telemetrystore.ClickhouseDB().Exec(ctx, op.ToSQL()); err != nil { + return errors.WrapInternalf(err, CodeFailCreateJSONIndexes, "failed to create index") + } + } + return nil +} + // TODO(Piyush): Remove this if not used in future. func (t *telemetryMetaStore) ListJSONValues(ctx context.Context, path string, limit int) (*telemetrytypes.TelemetryFieldValues, bool, error) { ctx = withTelemetryContext(ctx, telemetrytypes.SignalLogs, "ListJSONValues") diff --git a/pkg/telemetrymetadata/body_json_metadata_test.go b/pkg/telemetrymetadata/body_json_metadata_test.go index 8ab1f7d92d..804783c69f 100644 --- a/pkg/telemetrymetadata/body_json_metadata_test.go +++ b/pkg/telemetrymetadata/body_json_metadata_test.go @@ -1,12 +1,20 @@ package telemetrymetadata import ( + "context" "fmt" + "regexp" "testing" + sqlmock "github.com/DATA-DOG/go-sqlmock" + schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator" "github.com/SigNoz/signoz-otel-collector/constants" + "github.com/SigNoz/signoz/pkg/errors" "github.com/SigNoz/signoz/pkg/querybuilder" "github.com/SigNoz/signoz/pkg/telemetryschema/logstelemetryschema" + "github.com/SigNoz/signoz/pkg/telemetrystore" + "github.com/SigNoz/signoz/pkg/telemetrystore/telemetrystoretest" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -57,3 +65,39 @@ func TestBuildListLogsJSONIndexesQuery(t *testing.T) { }) } } + +func TestCreateIndexes(t *testing.T) { + index := schemamigrator.Index{ + Name: "`body_promoted.user.name_String_ngrambf_v1`", + Expression: "lower(assumeNotNull(dynamicElement(body_promoted.user.name, 'String')))", + Type: "ngrambf_v1(4, 1024, 2, 0)", + Granularity: 1, + } + + testCases := []struct { + name string + execErr error + wantErr bool + }{ + {name: "AddIndex_ExecutedOnSourceTable"}, + {name: "ExecError_Returned", execErr: errors.NewInternalf(errors.CodeInternal, "boom"), wantErr: true}, + } + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + ts := telemetrystoretest.New(telemetrystore.Config{}, sqlmock.QueryMatcherRegexp) + store := &telemetryMetaStore{telemetrystore: ts} + + ts.Mock().ExpectExec("ALTER TABLE (.+)" + logstelemetryschema.LogsV2LocalTableName + "(.+) ADD INDEX IF NOT EXISTS " + regexp.QuoteMeta(index.Name)). + WillReturnError(testCase.execErr) + + err := store.CreateIndexes(context.Background(), logsBodyIndexSource, []schemamigrator.Index{index}) + if testCase.wantErr { + assert.Error(t, err) + return + } + require.NoError(t, err) + assert.NoError(t, ts.Mock().ExpectationsWereMet()) + }) + } +} diff --git a/pkg/types/telemetrytypes/store.go b/pkg/types/telemetrytypes/store.go index 61dfbcac5c..6bc8d53337 100644 --- a/pkg/types/telemetrytypes/store.go +++ b/pkg/types/telemetrytypes/store.go @@ -3,6 +3,7 @@ package telemetrytypes import ( "context" + schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator" "github.com/SigNoz/signoz/pkg/types/metrictypes" "github.com/SigNoz/signoz/pkg/valuer" ) @@ -37,6 +38,9 @@ type MetadataStore interface { // ListJSONIndexes lists the per-path JSON skip indexes of the given source. ListJSONIndexes(ctx context.Context, source JSONIndexSource, filters ...string) ([]TelemetryFieldKeySkipIndex, error) + // CreateIndexes adds per-path JSON skip indexes to the source's table. + CreateIndexes(ctx context.Context, source JSONIndexSource, indexes []schemamigrator.Index) error + // GetPromotedPaths lists the promoted paths recorded in the column // evolution table for the entry's signal, column and field context. GetPromotedPaths(ctx context.Context, entry EvolutionEntry, paths ...string) (map[string]bool, error) diff --git a/pkg/types/telemetrytypes/telemetrytypestest/metadata_store.go b/pkg/types/telemetrytypes/telemetrytypestest/metadata_store.go index 3fda163d4b..e83aa07e88 100644 --- a/pkg/types/telemetrytypes/telemetrytypestest/metadata_store.go +++ b/pkg/types/telemetrytypes/telemetrytypestest/metadata_store.go @@ -4,6 +4,7 @@ import ( "context" "strings" + schemamigrator "github.com/SigNoz/signoz-otel-collector/cmd/signozschemamigrator/schema_migrator" "github.com/SigNoz/signoz/pkg/types/metrictypes" "github.com/SigNoz/signoz/pkg/types/telemetrytypes" "github.com/SigNoz/signoz/pkg/valuer" @@ -20,6 +21,7 @@ type MockMetadataStore struct { ReducedMap map[string]bool PromotedPathsMap map[string]bool LogsJSONIndexes []telemetrytypes.TelemetryFieldKeySkipIndex + CreatedIndexes []schemamigrator.Index ColumnEvolutionMetadataMap map[string][]*telemetrytypes.EvolutionEntry LookupKeysMap map[telemetrytypes.MetricMetadataLookupKey]int64 // StaticFields holds signal-specific intrinsic field definitions (e.g. logstelemetryschema.IntrinsicFields). @@ -36,6 +38,7 @@ func NewMockMetadataStore() *MockMetadataStore { TypeMap: make(map[string]metrictypes.Type), PromotedPathsMap: make(map[string]bool), LogsJSONIndexes: []telemetrytypes.TelemetryFieldKeySkipIndex{}, + CreatedIndexes: []schemamigrator.Index{}, ColumnEvolutionMetadataMap: make(map[string][]*telemetrytypes.EvolutionEntry), LookupKeysMap: make(map[telemetrytypes.MetricMetadataLookupKey]int64), StaticFields: make(map[string]telemetrytypes.TelemetryFieldKey), @@ -384,6 +387,11 @@ func (m *MockMetadataStore) ListJSONIndexes(ctx context.Context, source telemetr return indexes, nil } +func (m *MockMetadataStore) CreateIndexes(_ context.Context, _ telemetrytypes.JSONIndexSource, indexes []schemamigrator.Index) error { + m.CreatedIndexes = append(m.CreatedIndexes, indexes...) + return nil +} + func (m *MockMetadataStore) updateColumnEvolutionMetadataForKeys(_ context.Context, keysToUpdate []*telemetrytypes.TelemetryFieldKey) map[string][]*telemetrytypes.EvolutionEntry { var metadataKeySelectors []*telemetrytypes.EvolutionSelector