Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
69 changes: 69 additions & 0 deletions sei-cosmos/storev2/rootmulti/snapshot_metrics_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
3 changes: 3 additions & 0 deletions sei-cosmos/storev2/rootmulti/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -1463,6 +1463,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)
}
Expand Down
8 changes: 7 additions & 1 deletion sei-db/state_db/sc/composite/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -536,8 +536,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)
Expand Down
122 changes: 122 additions & 0 deletions sei-db/state_db/sc/composite/store_metrics_test.go
Original file line number Diff line number Diff line change
@@ -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(cs1.Version() + 1)
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(cs2.Version() + 1)
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")
}
40 changes: 34 additions & 6 deletions sei-db/state_db/sc/migration/router_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -21,7 +41,12 @@ func BuildRouter(
flatKV gigatypes.LiveStateStore,
// 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:
Expand All @@ -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)
}
Expand All @@ -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)
}
Expand All @@ -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)
}
Expand Down Expand Up @@ -149,6 +174,7 @@ func buildMigrateEVMRouter(
memIAVL *memiavl.CommitStore,
flatKV gigatypes.LiveStateStore,
migrationBatchSize int,
opts routerOptions,
) (Router, error) {

if memIAVL == nil {
Expand All @@ -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)
Expand Down Expand Up @@ -276,6 +302,7 @@ func buildMigrateAllButBankRouter(
memIAVL *memiavl.CommitStore,
flatKV gigatypes.LiveStateStore,
migrationBatchSize int,
opts routerOptions,
) (Router, error) {

if memIAVL == nil {
Expand All @@ -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)
Expand Down Expand Up @@ -401,6 +428,7 @@ func buildMigrateBankRouter(
memIAVL *memiavl.CommitStore,
flatKV gigatypes.LiveStateStore,
migrationBatchSize int,
opts routerOptions,
) (Router, error) {

if memIAVL == nil {
Expand All @@ -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)
Expand Down
Loading