diff --git a/sei-cosmos/storev2/rootmulti/snapshot_metrics_test.go b/sei-cosmos/storev2/rootmulti/snapshot_metrics_test.go new file mode 100644 index 0000000000..3752d00aea --- /dev/null +++ b/sei-cosmos/storev2/rootmulti/snapshot_metrics_test.go @@ -0,0 +1,69 @@ +package rootmulti + +import ( + "bytes" + "context" + "testing" + + protoio "github.com/gogo/protobuf/io" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" +) + +// collectTotalNumKeys collects once from reader and returns the +// iavl_total_num_keys value of every store_name data point. +func collectTotalNumKeys(t *testing.T, reader *sdkmetric.ManualReader) map[string]int64 { + t.Helper() + var collected metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &collected)) + + values := map[string]int64{} + for _, scope := range collected.ScopeMetrics { + for _, m := range scope.Metrics { + if m.Name != "iavl_total_num_keys" { + continue + } + gauge, ok := m.Data.(metricdata.Gauge[int64]) + require.True(t, ok, "unexpected data type %T", m.Data) + for _, point := range gauge.DataPoints { + name, ok := point.Attributes.Value("store_name") + require.True(t, ok, "data point has no store_name attribute") + values[name.AsString()] = point.Value + } + } + } + return values +} + +// TestSnapshotReportsZeroKeysForEmptyStore pins that a store exported with no +// nodes is reported with zero keys rather than left unreported. After the EVM +// migration the memiavl evm store is exported empty; a gauge that is never +// re-recorded keeps exporting its pre-migration value. +func TestSnapshotReportsZeroKeysForEmptyStore(t *testing.T) { + reader := sdkmetric.NewManualReader() + provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + prev := otel.GetMeterProvider() + otel.SetMeterProvider(provider) + t.Cleanup(func() { + otel.SetMeterProvider(prev) + require.NoError(t, provider.Shutdown(context.Background())) + }) + + cfg := evmMigratedConfig() + evmData := newEVMTestData(0x56) + store, storeKeys := newTestRootMulti(t, t.TempDir(), cfg) + for block := 1; block <= 3; block++ { + simulateBlock(t, store, storeKeys, block, evmData) + } + + var buf bytes.Buffer + require.NoError(t, store.Snapshot(3, protoio.NewDelimitedWriter(&buf))) + require.NoError(t, store.Close()) + + totals := collectTotalNumKeys(t, reader) + require.Contains(t, totals, "evm", "empty memiavl evm store must still be reported") + require.Zero(t, totals["evm"]) + require.Positive(t, totals["flatkv"], "migrated evm keys are reported under flatkv") +} diff --git a/sei-cosmos/storev2/rootmulti/store.go b/sei-cosmos/storev2/rootmulti/store.go index 99579ffc4f..1444093f74 100644 --- a/sei-cosmos/storev2/rootmulti/store.go +++ b/sei-cosmos/storev2/rootmulti/store.go @@ -1364,6 +1364,9 @@ func (rs *Store) Snapshot(height uint64, protoWriter protoio.Writer) error { return err } currentStoreName = item + keySizePerStore[item] = 0 + valueSizePerStore[item] = 0 + numKeysPerStore[item] = 0 default: return fmt.Errorf("unknown item type %T", item) } diff --git a/sei-db/state_db/sc/composite/store.go b/sei-db/state_db/sc/composite/store.go index 461dcb497e..e08a8c618a 100644 --- a/sei-db/state_db/sc/composite/store.go +++ b/sei-db/state_db/sc/composite/store.go @@ -481,8 +481,14 @@ func (cs *CompositeCommitStore) resolveCurrentWriteMode(closeIdleFlatKV bool) er // mode has been resolved. func (cs *CompositeCommitStore) buildRouter() error { routerCtx, cancel := context.WithCancel(cs.ctx) + var options []migration.RouterOption + if cs.derived { + // A derived store sees the migration state at its own height; publishing it would overwrite the + // live store's gauges. + options = append(options, migration.WithoutTelemetry()) + } router, err := migration.BuildRouter( - routerCtx, cs.currentWriteMode, cs.memIAVL, cs.flatKV, int(cs.migrationBatchSize.Load())) + routerCtx, cs.currentWriteMode, cs.memIAVL, cs.flatKV, int(cs.migrationBatchSize.Load()), options...) if err != nil { cancel() return fmt.Errorf("failed to build router: %w", err) diff --git a/sei-db/state_db/sc/composite/store_metrics_test.go b/sei-db/state_db/sc/composite/store_metrics_test.go new file mode 100644 index 0000000000..31907a08eb --- /dev/null +++ b/sei-db/state_db/sc/composite/store_metrics_test.go @@ -0,0 +1,122 @@ +package composite + +import ( + "context" + "fmt" + "testing" + + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" + + "github.com/sei-protocol/sei-chain/sei-db/common/keys" + "github.com/sei-protocol/sei-chain/sei-db/config" + "github.com/sei-protocol/sei-chain/sei-db/proto" + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/migration" + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" +) + +const migrationVersionInstrument = "seidb_migration_version" + +// installManualMeterProvider routes the global OTel meter to a ManualReader +// for the duration of the test. +func installManualMeterProvider(t *testing.T) *sdkmetric.ManualReader { + t.Helper() + reader := sdkmetric.NewManualReader() + provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + prev := otel.GetMeterProvider() + otel.SetMeterProvider(provider) + t.Cleanup(func() { + otel.SetMeterProvider(prev) + require.NoError(t, provider.Shutdown(context.Background())) + }) + return reader +} + +// collectMigrationVersions collects once from reader and returns every +// seidb_migration_version data point value. +func collectMigrationVersions(t *testing.T, reader *sdkmetric.ManualReader) []int64 { + t.Helper() + var collected metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &collected)) + + var values []int64 + for _, scope := range collected.ScopeMetrics { + for _, m := range scope.Metrics { + if m.Name != migrationVersionInstrument { + continue + } + gauge, ok := m.Data.(metricdata.Gauge[int64]) + require.True(t, ok, "unexpected data type %T", m.Data) + for _, point := range gauge.DataPoints { + values = append(values, point.Value) + } + } + } + return values +} + +// TestLoadVersionReadOnlyDoesNotReportMigrationVersion pins that a read-only +// handle opened at a mid-migration height leaves the process-wide migration +// version gauge at the live manager's value. State-sync snapshot exports open +// such handles after the migration completes; they must not report the +// pre-completion version they observe. +func TestLoadVersionReadOnlyDoesNotReportMigrationVersion(t *testing.T) { + dir := t.TempDir() + + // Phase 1: seed evm/ keys in MemiavlOnly so the migration spans several blocks. + v0Cfg := config.DefaultStateCommitConfig() + v0Cfg.WriteMode = types.MemiavlOnly + cs1, err := NewCompositeCommitStore(t.Context(), dir, v0Cfg) + require.NoError(t, err) + require.NoError(t, cs1.Initialize([]string{keys.BankStoreKey, keys.EVMStoreKey})) + require.NoError(t, cs1.LoadLatest()) + pairs := make([]*proto.KVPair, 0, 8) + for i := 0; i < 8; i++ { + pairs = append(pairs, &proto.KVPair{Key: []byte(fmt.Sprintf("evm_key_%d", i)), Value: []byte("v")}) + } + require.NoError(t, cs1.ApplyChangeSets([]*proto.NamedChangeSet{ + {Name: keys.EVMStoreKey, Changeset: proto.ChangeSet{Pairs: pairs}}, + })) + _, err = cs1.Commit() + require.NoError(t, err) + require.NoError(t, cs1.Close()) + + reader := installManualMeterProvider(t) + + // Phase 2: migrate one key per block until the completion block lands. + migrateCfg := config.DefaultStateCommitConfig() + migrateCfg.WriteMode = types.MigrateEVM + cs2, err := NewCompositeCommitStore(t.Context(), dir, migrateCfg) + require.NoError(t, err) + require.NoError(t, cs2.SetMigrationBatchSize(1)) + require.NoError(t, cs2.Initialize([]string{keys.BankStoreKey, keys.EVMStoreKey})) + require.NoError(t, cs2.LoadLatest()) + defer cs2.Close() + + var midMigrationVersion int64 + for block := 0; block < 32; block++ { + require.NoError(t, cs2.ApplyChangeSets([]*proto.NamedChangeSet{ + {Name: keys.BankStoreKey, Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{ + {Key: []byte(fmt.Sprintf("bank_%d", block)), Value: []byte("v")}, + }}}, + })) + _, err = cs2.Commit() + require.NoError(t, err) + if _, done := cs2.flatKV.Get(migration.MigrationStore, []byte(migration.MigrationVersionKey)); done { + break + } + midMigrationVersion = cs2.Version() + } + require.NotZero(t, midMigrationVersion, "migration must span more than one block") + require.Equal(t, []int64{int64(migration.Version1_MigrateEVM)}, collectMigrationVersions(t, reader), + "precondition: live manager reports the completed migration version") + + ro, err := cs2.LoadVersionReadOnly(midMigrationVersion) + require.NoError(t, err) + defer func() { _ = ro.Close() }() + + require.Equal(t, []int64{int64(migration.Version1_MigrateEVM)}, collectMigrationVersions(t, reader), + "read-only handle at a mid-migration height must not overwrite the live migration version") +} diff --git a/sei-db/state_db/sc/migration/router_builder.go b/sei-db/state_db/sc/migration/router_builder.go index fabba1c0e1..b4940b88aa 100644 --- a/sei-db/state_db/sc/migration/router_builder.go +++ b/sei-db/state_db/sc/migration/router_builder.go @@ -12,6 +12,26 @@ import ( "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" ) +// RouterOption adjusts how BuildRouter assembles a router. +type RouterOption func(*routerOptions) + +type routerOptions struct { + telemetry bool +} + +// WithoutTelemetry makes the router's MigrationManager keep its migration +// metrics in process, without publishing them on the process-wide OTel instruments. +func WithoutTelemetry() RouterOption { + return func(o *routerOptions) { o.telemetry = false } +} + +func (o routerOptions) migrationMetrics(ctx context.Context, targetVersion uint64) *MigrationMetrics { + if !o.telemetry { + return newLocalMigrationMetrics() + } + return NewMigrationMetrics(ctx, targetVersion) +} + // Builds a router for the given migration write mode. A router is responsible for splitting // reads/writes between the memiavl and flatkv backends. func BuildRouter( @@ -21,7 +41,12 @@ func BuildRouter( flatKV flatkv.Store, // If this router will be doing data migration, this is the number of keys to migrate in each batch. migrationBatchSize int, + options ...RouterOption, ) (Router, error) { + opts := routerOptions{telemetry: true} + for _, apply := range options { + apply(&opts) + } switch writeMode { case types.MemiavlOnly: @@ -31,7 +56,7 @@ func BuildRouter( } return router, nil case types.MigrateEVM: - router, err := buildMigrateEVMRouter(ctx, memIAVL, flatKV, migrationBatchSize) + router, err := buildMigrateEVMRouter(ctx, memIAVL, flatKV, migrationBatchSize, opts) if err != nil { return nil, fmt.Errorf("buildMigrateEVMRouter: %w", err) } @@ -51,7 +76,7 @@ func BuildRouter( } return threadSafe, nil case types.MigrateAllButBank: - router, err := buildMigrateAllButBankRouter(ctx, memIAVL, flatKV, migrationBatchSize) + router, err := buildMigrateAllButBankRouter(ctx, memIAVL, flatKV, migrationBatchSize, opts) if err != nil { return nil, fmt.Errorf("buildMigrateAllButBankRouter: %w", err) } @@ -71,7 +96,7 @@ func BuildRouter( } return threadSafe, nil case types.MigrateBank: - router, err := buildMigrateBankRouter(ctx, memIAVL, flatKV, migrationBatchSize) + router, err := buildMigrateBankRouter(ctx, memIAVL, flatKV, migrationBatchSize, opts) if err != nil { return nil, fmt.Errorf("buildMigrateBankRouter: %w", err) } @@ -149,6 +174,7 @@ func buildMigrateEVMRouter( memIAVL *memiavl.CommitStore, flatKV flatkv.Store, migrationBatchSize int, + opts routerOptions, ) (Router, error) { if memIAVL == nil { @@ -167,7 +193,7 @@ func buildMigrateEVMRouter( buildFlatKVReader(flatKV), buildFlatKVWriter(flatKV), NewMemiavlMigrationIterator(memIAVL.GetDB(), []string{keys.EVMStoreKey}), - NewMigrationMetrics(ctx, Version1_MigrateEVM), + opts.migrationMetrics(ctx, Version1_MigrateEVM), ) if err != nil { return nil, fmt.Errorf("NewMigrationManager: %w", err) @@ -276,6 +302,7 @@ func buildMigrateAllButBankRouter( memIAVL *memiavl.CommitStore, flatKV flatkv.Store, migrationBatchSize int, + opts routerOptions, ) (Router, error) { if memIAVL == nil { @@ -299,7 +326,7 @@ func buildMigrateAllButBankRouter( buildFlatKVReader(flatKV), buildFlatKVWriter(flatKV), NewMemiavlMigrationIterator(memIAVL.GetDB(), allModulesButEvmAndBank), - NewMigrationMetrics(ctx, Version2_MigrateAllButBank), + opts.migrationMetrics(ctx, Version2_MigrateAllButBank), ) if err != nil { return nil, fmt.Errorf("NewMigrationManager: %w", err) @@ -401,6 +428,7 @@ func buildMigrateBankRouter( memIAVL *memiavl.CommitStore, flatKV flatkv.Store, migrationBatchSize int, + opts routerOptions, ) (Router, error) { if memIAVL == nil { @@ -426,7 +454,7 @@ func buildMigrateBankRouter( buildFlatKVReader(flatKV), buildFlatKVWriter(flatKV), NewMemiavlMigrationIterator(memIAVL.GetDB(), []string{keys.BankStoreKey}), - NewMigrationMetrics(ctx, Version3_FlatKVOnly), + opts.migrationMetrics(ctx, Version3_FlatKVOnly), ) if err != nil { return nil, fmt.Errorf("NewMigrationManager: %w", err)