diff --git a/.github/workflows/sei-db-tests.yml b/.github/workflows/sei-db-tests.yml index f559d39964..906ed9b1af 100644 --- a/.github/workflows/sei-db-tests.yml +++ b/.github/workflows/sei-db-tests.yml @@ -66,7 +66,7 @@ jobs: -run '^$' \ -bench . \ -benchtime=1x \ - ./sei-db/state_db/bench/... + ./sei-db/bench/... coverage: name: Coverage diff --git a/sei-db/bench/cryptosim/config/basic-config.json b/sei-db/bench/cryptosim/config/basic-config.json index c97250a157..b53a5f97df 100644 --- a/sei-db/bench/cryptosim/config/basic-config.json +++ b/sei-db/bench/cryptosim/config/basic-config.json @@ -1,8 +1,7 @@ { "Comment": "Basic configuration for the cryptosim benchmark. Intended for basic correctness/sanity testing.", - "Backend": "FlatKV", "StateStoreConfig": { - "Enable": true, + "Enable": false, "DBDirectory": "", "Backend": "pebbledb", "AsyncWriteBuffer": 100, @@ -10,8 +9,16 @@ "PruneIntervalSeconds": 600, "ImportNumWorkers": 1, "KeepLastVersion": true, - "UseDefaultComparer": false, - "EVMDBDirectory": "" + "UseDefaultComparer": false + }, + "CheckpointConfig": { + "TimeInterval": 600000000000, + "BlockInterval": 0 + }, + "PruningConfig": { + "RollbackWindow": 1000, + "LookbackWindow": 0, + "PruneInterval": 300000000000 }, "CannedRandomSize": 1073741824, "ConstantThreadCount": 0, diff --git a/sei-db/bench/cryptosim/config/historical-offload-kafka.json b/sei-db/bench/cryptosim/config/historical-offload-kafka.json deleted file mode 100644 index 16e9d44f32..0000000000 --- a/sei-db/bench/cryptosim/config/historical-offload-kafka.json +++ /dev/null @@ -1,25 +0,0 @@ -{ - "Comment": "Sample cryptosim config for offchain pipeline testing. Replace HistoricalOffload.Kafka.Brokers, Topic, and Region with your environment values.", - "Backend": "SSHistoricalOffload", - "HistoricalOffload": { - "Provider": "kafka", - "Kafka": { - "Brokers": [ - "b-1.example.kafka.amazonaws.com:9098", - "b-2.example.kafka.amazonaws.com:9098" - ], - "Topic": "historical-offload", - "Region": "us-east-1", - "TLSEnabled": true, - "SASLMechanism": "aws-msk-iam" - } - }, - "ConsoleUpdateIntervalSeconds": 5, - "DataDir": "data/historical-offload-kafka", - "EnableSuspension": false, - "MetricsAddr": "", - "MaxRuntimeSeconds": 1200, - "DisableTransactionReads": true, - "LogDir": "logs/historical-offload-kafka", - "MaxTPS": 150000 -} diff --git a/sei-db/bench/cryptosim/config/ss-composite-config.json b/sei-db/bench/cryptosim/config/ss-composite-config.json deleted file mode 100644 index d0332599ea..0000000000 --- a/sei-db/bench/cryptosim/config/ss-composite-config.json +++ /dev/null @@ -1,6 +0,0 @@ -{ - "Comment": "State store benchmark using composite SS with EVM sub-stores (split_write + EVM-first read).", - "Backend": "SSComposite", - "DataDir": "data", - "LogDir": "logs" -} diff --git a/sei-db/bench/cryptosim/config/ss-composite-pebbledb-write-only.json b/sei-db/bench/cryptosim/config/ss-composite-pebbledb-write-only.json deleted file mode 100644 index c353124fab..0000000000 --- a/sei-db/bench/cryptosim/config/ss-composite-pebbledb-write-only.json +++ /dev/null @@ -1,49 +0,0 @@ -{ - "Comment": "Write-only SSComposite benchmark using PebbleDB.", - "Backend": "SSComposite", - "StateStoreConfig": { - "Enable": true, - "DBDirectory": "", - "Backend": "pebbledb", - "AsyncWriteBuffer": 100, - "KeepRecent": 100000, - "PruneIntervalSeconds": 600, - "ImportNumWorkers": 1, - "KeepLastVersion": true, - "UseDefaultComparer": false, - "EVMDBDirectory": "" - }, - "CannedRandomSize": 1073741824, - "ConstantThreadCount": 0, - "ConsoleUpdateIntervalSeconds": 5, - "ConsoleUpdateIntervalTransactions": 1000000, - "DataDir": "data/pebble20", - "EnableSuspension": false, - "Erc20ContractSize": 2048, - "Erc20InteractionsPerAccount": 10, - "Erc20StorageSlotSize": 32, - "ExecutorQueueSize": 1024, - "HotAccountProbability": 0.1, - "HotErc20ContractProbability": 0.5, - "HotErc20ContractSetSize": 100, - "MetricsAddr": "", - "MinimumNumberOfColdAccounts": 1000000, - "MinimumNumberOfDormantAccounts": 1000000, - "MinimumNumberOfErc20Contracts": 10000, - "NewAccountDormancyProbability": 1.0, - "NewAccountProbability": 0.001, - "NumberOfHotAccounts": 100, - "PaddedAccountSize": 32, - "Seed": 1337, - "SetupUpdateIntervalCount": 100000, - "ThreadsPerCore": 2.0, - "TransactionsPerBlock": 1024, - "MaxRuntimeSeconds": 1200, - "TransactionMetricsSampleRate": 0.001, - "BackgroundMetricsScrapeInterval": 60, - "BlockChannelCapacity": 8, - "DisableTransactionReads": true, - "LogDir": "logs/pebble20", - "LogLevel": "info", - "MaxTPS": 150000 -} diff --git a/sei-db/bench/cryptosim/config/ss-composite-rocksdb-write-only.json b/sei-db/bench/cryptosim/config/ss-composite-rocksdb-write-only.json deleted file mode 100644 index 9339f56958..0000000000 --- a/sei-db/bench/cryptosim/config/ss-composite-rocksdb-write-only.json +++ /dev/null @@ -1,49 +0,0 @@ -{ - "Comment": "Write-only SSComposite benchmark using RocksDB.", - "Backend": "SSComposite", - "StateStoreConfig": { - "Enable": true, - "DBDirectory": "", - "Backend": "rocksdb", - "AsyncWriteBuffer": 100, - "KeepRecent": 100000, - "PruneIntervalSeconds": 600, - "ImportNumWorkers": 1, - "KeepLastVersion": true, - "UseDefaultComparer": false, - "EVMDBDirectory": "" - }, - "CannedRandomSize": 1073741824, - "ConstantThreadCount": 0, - "ConsoleUpdateIntervalSeconds": 5, - "ConsoleUpdateIntervalTransactions": 1000000, - "DataDir": "data/rocks20", - "EnableSuspension": false, - "Erc20ContractSize": 2048, - "Erc20InteractionsPerAccount": 10, - "Erc20StorageSlotSize": 32, - "ExecutorQueueSize": 1024, - "HotAccountProbability": 0.1, - "HotErc20ContractProbability": 0.5, - "HotErc20ContractSetSize": 100, - "MetricsAddr": "", - "MinimumNumberOfColdAccounts": 1000000, - "MinimumNumberOfDormantAccounts": 1000000, - "MinimumNumberOfErc20Contracts": 10000, - "NewAccountDormancyProbability": 1.0, - "NewAccountProbability": 0.001, - "NumberOfHotAccounts": 100, - "PaddedAccountSize": 32, - "Seed": 1337, - "SetupUpdateIntervalCount": 100000, - "ThreadsPerCore": 2.0, - "TransactionsPerBlock": 1024, - "MaxRuntimeSeconds": 1200, - "TransactionMetricsSampleRate": 0.001, - "BackgroundMetricsScrapeInterval": 60, - "BlockChannelCapacity": 8, - "DisableTransactionReads": true, - "LogDir": "logs/rocks20", - "LogLevel": "info", - "MaxTPS": 150000 -} diff --git a/sei-db/bench/cryptosim/cryptosim.go b/sei-db/bench/cryptosim/cryptosim.go index b48ec0e945..08a78d3b40 100644 --- a/sei-db/bench/cryptosim/cryptosim.go +++ b/sei-db/bench/cryptosim/cryptosim.go @@ -3,14 +3,17 @@ package cryptosim import ( "context" "fmt" + "path/filepath" "runtime" "time" - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" + "golang.org/x/time/rate" + "github.com/sei-protocol/sei-chain/sei-db/common/keys" crand "github.com/sei-protocol/sei-chain/sei-db/common/rand" "github.com/sei-protocol/sei-chain/sei-db/common/utils" - "golang.org/x/time/rate" + "github.com/sei-protocol/sei-chain/sei-db/controller" + "github.com/sei-protocol/sei-chain/sei-db/state_db/giga" ) const ( @@ -133,23 +136,50 @@ func NewCryptoSim( fmt.Printf("Running cryptosim benchmark from data directory: %s\n", config.DataDir) fmt.Printf("Logs are being routed to: %s\n", config.LogDir) - var dbConfig any - switch config.Backend { - case wrappers.FlatKV: - dbConfig = config.FlatKVConfig - case wrappers.SSComposite, wrappers.CompositeDual_SSComposite: - dbConfig = config.StateStoreConfig - case wrappers.SSHistoricalOffload: - dbConfig = config.HistoricalOffload + config.FlatKVConfig.DataDir = config.DataDir + config.StateStoreConfig.EVMDBDirectory = filepath.Join( + config.DataDir, "state_store", "evm", config.StateStoreConfig.Backend) + + // Every store the state DB opens is pruned by the collector started below, so each one stands its + // own pruner down. This is the same handover bootstrap.GigaStorageManager performs for a node. + config.FlatKVConfig.ExternalPruning = true + config.StateStoreConfig.ExternalPruning = true + + // giga.NewStateDB is the node's own entry point, and the only one that leaves the state WAL + // outside the live state DB: it opens the WAL itself and writes each block to it ahead of the + // commit. A live state DB opened directly would own its WAL and write it inline instead. + db, err := giga.NewStateDB(ctx, config.FlatKVConfig, config.StateStoreConfig, config.CheckpointConfig) + if err != nil { + cancel() + return nil, fmt.Errorf("failed to open the state DB: %w", err) } - db, err := wrappers.NewDBImpl(ctx, config.Backend, config.DataDir, dbConfig) + // Nothing the state DB opens prunes itself on this path: the state WAL, and the historical state + // DB when it is enabled, shrink only when a collector tells them to. + garbageCollector, err := controller.NewStorageGarbageCollector( + ctx, config.PruningConfig, db.PrunableStores()) if err != nil { cancel() - return nil, fmt.Errorf("failed to create database: %w", err) + if closeErr := db.Close(); closeErr != nil { + fmt.Printf("failed to close the state DB during error recovery: %v\n", closeErr) + } + return nil, fmt.Errorf("failed to start the storage garbage collector: %w", err) } - metrics := NewCryptosimMetrics(ctx, db.GetPhaseTimer(), config) + // Every construction failure past this point releases through here. The state DB holds the state + // WAL directory's exclusive lock, so a handle left open makes an in-process retry fail to open the + // WAL rather than only leaking descriptors. + releaseStorage := func() { + cancel() + if closeErr := garbageCollector.Close(); closeErr != nil { + fmt.Printf("failed to close the garbage collector during error recovery: %v\n", closeErr) + } + if closeErr := db.Close(); closeErr != nil { + fmt.Printf("failed to close the state DB during error recovery: %v\n", closeErr) + } + } + + metrics := NewCryptosimMetrics(ctx, db.SC().GetPhaseTimer(), config) // Server start deferred until after DataGenerator loads DB state and sets gauges, // avoiding rate() spikes when restarting with a preserved DB. @@ -160,24 +190,15 @@ func NewCryptoSim( start := time.Now() - database, err := NewDatabase(config, db, metrics, 0) + database, err := NewDatabase(config, db, garbageCollector, metrics) if err != nil { - cancel() - if closeErr := db.Close(); closeErr != nil { - fmt.Printf("failed to close database during error recovery: %v\n", closeErr) - } + releaseStorage() return nil, fmt.Errorf("failed to create database: %w", err) } - dataGenerator, err := NewDataGenerator(config, database, rand, metrics) - if err != nil { - cancel() - if closeErr := db.Close(); closeErr != nil { - fmt.Printf("failed to close database during error recovery: %v\n", closeErr) - } - return nil, fmt.Errorf("failed to create data generator: %w", err) - } - database.nextBlockNumber = dataGenerator.InitialNextBlockNumber() + fmt.Printf("Next block number: %s.\n", int64Commas(database.nextBlockNumber)) + + dataGenerator := NewDataGenerator(config, database, rand, metrics) threadCount := int(config.ThreadsPerCore)*runtime.NumCPU() + config.ConstantThreadCount if threadCount < 1 { threadCount = 1 @@ -195,7 +216,7 @@ func NewCryptoSim( recieptsChan = make(chan *block, config.RecieptChannelCapacity) _, err := NewRecieptStoreSimulator(ctx, config, recieptsChan, metrics, rand.Clone(false)) if err != nil { - cancel() + releaseStorage() return nil, fmt.Errorf("failed to create receipt store simulator: %w", err) } metrics.startReceiptChannelDepthSampling(recieptsChan, config.BackgroundMetricsScrapeInterval) @@ -230,6 +251,7 @@ func NewCryptoSim( err = c.setup() if err != nil { + releaseStorage() return nil, fmt.Errorf("failed to setup benchmark: %w", err) } diff --git a/sei-db/bench/cryptosim/cryptosim_config.go b/sei-db/bench/cryptosim/cryptosim_config.go index 922614479d..96797cdad5 100644 --- a/sei-db/bench/cryptosim/cryptosim_config.go +++ b/sei-db/bench/cryptosim/cryptosim_config.go @@ -7,7 +7,6 @@ import ( "path/filepath" "strings" - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" "github.com/sei-protocol/sei-chain/sei-db/config" flatkvConfig "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/config" ) @@ -86,7 +85,7 @@ type CryptoSimConfig struct { // How many blocks the benchmark may run ahead of block hashing. Databases hash committed blocks // asynchronously, and the benchmark takes one block's hash per block committed once it is this far // ahead — so a block's hash must arrive no later than this many blocks after it was committed, and - // the benchmark waits when it does not. A database that publishes no block hashes waits on nothing. + // the benchmark waits when it does not. HashLagBlocks int // The directory to store the benchmark data. @@ -100,16 +99,14 @@ type CryptoSimConfig struct { // in undefined behavior, don't change the size unless you are starting a new run from scratch. CannedRandomSize int - // The backend to use for the benchmark database. - Backend wrappers.DBType + // Configures the historical state DB. + StateStoreConfig config.StateStoreConfig - // StateStoreConfig controls SS-backed benchmark backends such as SSComposite. - // The default preserves the benchmark SS defaults: pebbledb, async buffer 100. - StateStoreConfig *config.StateStoreConfig + // Configures the cadence the state DB checkpoints both halves of state on. + CheckpointConfig config.CheckpointConfig - // HistoricalOffload configures the transport used by the - // SSHistoricalOffload backend. - HistoricalOffload *wrappers.HistoricalOffloadConfig + // Configures the prune cycle that enforces retention across the state DB's stores. + PruningConfig *config.StorageGarbageCollectorConfig // This field is ignored, but allows for a comment to be added to the config file. // Something, something, why in the name of all things holy doesn't json support comments? @@ -163,7 +160,7 @@ type CryptoSimConfig struct { // If true, the log directory will be deleted on a clean shutdown. DeleteLogDirOnShutdown bool - // Configures the FlatKV database. Ignored if Backend is not "FlatKV". + // Configures the live state DB. FlatKVConfig *flatkvConfig.Config // The capacity of the channel that holds blocks awaiting execution. @@ -243,6 +240,11 @@ func DefaultCryptoSimConfig() *CryptoSimConfig { // Note: if you add new fields or modify default values, be sure to keep config/basic-config.json in sync. // That file should contain every available config set to its default value, as a reference. + ssConfig := config.DefaultStateStoreConfig() + // Nothing in the benchmark reads the historical state DB, so a run pays to write it only when + // the config asks for it. + ssConfig.Enable = false + cfg := &CryptoSimConfig{ NumberOfHotAccounts: 100, MinimumNumberOfColdAccounts: 1_000_000, @@ -262,8 +264,9 @@ func DefaultCryptoSimConfig() *CryptoSimConfig { HashLagBlocks: 32, Seed: 1337, CannedRandomSize: 1024 * 1024 * 1024, // 1GB - Backend: wrappers.FlatKV, - StateStoreConfig: wrappers.DefaultBenchStateStoreConfig(), + StateStoreConfig: ssConfig, + CheckpointConfig: config.DefaultCheckpointConfig(), + PruningConfig: config.DefaultStorageGarbageCollectorConfig(), ConsoleUpdateIntervalSeconds: 1, ConsoleUpdateIntervalTransactions: 1_000_000, SetupUpdateIntervalCount: 100_000, @@ -412,19 +415,14 @@ func (c *CryptoSimConfig) Validate() error { return fmt.Errorf("ReceiptLogFilterMaxBlockRange must be >= ReceiptLogFilterMinBlockRange (got %d < %d)", c.ReceiptLogFilterMaxBlockRange, c.ReceiptLogFilterMinBlockRange) } - if c.StateStoreConfig == nil { - return fmt.Errorf("StateStoreConfig is required") - } switch c.StateStoreConfig.Backend { case config.PebbleDBBackend, config.RocksDBBackend: default: return fmt.Errorf("StateStoreConfig.Backend must be one of %q or %q (got %q)", config.PebbleDBBackend, config.RocksDBBackend, c.StateStoreConfig.Backend) } - if c.Backend == wrappers.SSHistoricalOffload { - if err := c.HistoricalOffload.Validate(); err != nil { - return err - } + if err := c.PruningConfig.Validate(); err != nil { + return fmt.Errorf("PruningConfig is invalid: %w", err) } switch strings.ToLower(c.LogLevel) { case "debug", "info", "warn", "error": diff --git a/sei-db/bench/cryptosim/cryptosim_config_test.go b/sei-db/bench/cryptosim/cryptosim_config_test.go index 808532bb07..ad09b24bef 100644 --- a/sei-db/bench/cryptosim/cryptosim_config_test.go +++ b/sei-db/bench/cryptosim/cryptosim_config_test.go @@ -7,7 +7,6 @@ import ( "github.com/stretchr/testify/require" - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" "github.com/sei-protocol/sei-chain/sei-db/config" ) @@ -16,7 +15,6 @@ func TestLoadConfigFromFile_StateStoreConfigOverridePreservesBenchmarkDefaults(t configPath := filepath.Join(t.TempDir(), "cryptosim.json") err := os.WriteFile(configPath, []byte(`{ - "Backend": "SSComposite", "StateStoreConfig": { "Backend": "rocksdb" }, @@ -27,10 +25,9 @@ func TestLoadConfigFromFile_StateStoreConfigOverridePreservesBenchmarkDefaults(t cfg, err := LoadConfigFromFile(configPath) require.NoError(t, err) - require.Equal(t, wrappers.SSComposite, cfg.Backend) require.Equal(t, config.RocksDBBackend, cfg.StateStoreConfig.Backend) require.Equal(t, config.DefaultSSAsyncBuffer, cfg.StateStoreConfig.AsyncWriteBuffer) - require.True(t, cfg.StateStoreConfig.EVMSplit) + require.Equal(t, config.DefaultSSKeepRecent, cfg.StateStoreConfig.KeepRecent) } func TestLoadConfigFromFile_InvalidStateStoreBackend(t *testing.T) { @@ -74,7 +71,6 @@ func TestLoadConfigFromFile_DisableTransactionReadsOverride(t *testing.T) { configPath := filepath.Join(t.TempDir(), "cryptosim.json") err := os.WriteFile(configPath, []byte(`{ - "Backend": "NoOp", "DisableTransactionReads": true, "DataDir": "data", "LogDir": "logs" @@ -83,6 +79,5 @@ func TestLoadConfigFromFile_DisableTransactionReadsOverride(t *testing.T) { cfg, err := LoadConfigFromFile(configPath) require.NoError(t, err) - require.Equal(t, wrappers.NoOp, cfg.Backend) require.True(t, cfg.DisableTransactionReads) } diff --git a/sei-db/bench/cryptosim/data_generator.go b/sei-db/bench/cryptosim/data_generator.go index c896bbe9c3..3523c9b214 100644 --- a/sei-db/bench/cryptosim/data_generator.go +++ b/sei-db/bench/cryptosim/data_generator.go @@ -13,8 +13,6 @@ const ( accountIdCounterKey = "accountIdCounterKey" // Used to store the next ERC20 contract ID in the database. erc20IdCounterKey = "erc20IdCounterKey" - // Used to store the next block number in the database. - blockNumberCounterKey = "blockNumberCounterKey" // Use the code hash as a proxy. There is currently no mechanism to force FlatKV to update the account balance // field, and code hash keys will cause the account DB to get updated, which is the important part for this @@ -32,10 +30,6 @@ type DataGenerator struct { // The next ERC20 contract ID to be used when creating a new ERC20 contract. nextErc20ContractID int64 - // The next block number at startup time. Not updated after initialization; - // the block builder tracks the ongoing value. - initialNextBlockNumber uint64 - // The random number generator. rand *crand.CannedRandom @@ -67,12 +61,9 @@ func NewDataGenerator( database *Database, rand *crand.CannedRandom, metrics *CryptosimMetrics, -) (*DataGenerator, error) { +) *DataGenerator { - nextAccountIDBinary, found, err := database.Get(AccountIDCounterKey()) - if err != nil { - return nil, fmt.Errorf("failed to read account counter: %w", err) - } + nextAccountIDBinary, found := database.Get(AccountIDCounterKey()) var nextAccountID int64 if found { //nolint:gosec // G115 - persisted counter value, overflow acceptable @@ -84,10 +75,7 @@ func NewDataGenerator( cold := min(int64(config.MinimumNumberOfColdAccounts), max(0, nextAccountID-1-hot)) metrics.SetTotalNumberOfAccounts(nextAccountID, hot, cold) - nextErc20ContractIDBinary, found, err := database.Get(Erc20IDCounterKey()) - if err != nil { - return nil, fmt.Errorf("failed to read ERC20 contract counter: %w", err) - } + nextErc20ContractIDBinary, found := database.Get(Erc20IDCounterKey()) var nextErc20ContractID int64 if found { //nolint:gosec // G115 - persisted counter value, overflow acceptable @@ -97,17 +85,6 @@ func NewDataGenerator( fmt.Printf("There are currently %s ERC20 contracts in the database.\n", int64Commas(nextErc20ContractID)) metrics.SetTotalNumberOfERC20Contracts(nextErc20ContractID) - nextBlockNumberBinary, found, err := database.Get(BlockNumberCounterKey()) - if err != nil { - return nil, fmt.Errorf("failed to read block number counter: %w", err) - } - var nextBlockNumber uint64 - if found { - nextBlockNumber = binary.BigEndian.Uint64(nextBlockNumberBinary) - } - - fmt.Printf("Next block number: %s.\n", int64Commas(int64(nextBlockNumber))) //nolint:gosec - feeCollectionAddress := keys.BuildEVMKey( accountKeyPrefix, rand.Address(accountPrefix, 0, keys.AddressLen), @@ -117,14 +94,13 @@ func NewDataGenerator( config: config, nextAccountID: nextAccountID, nextErc20ContractID: nextErc20ContractID, - initialNextBlockNumber: nextBlockNumber, rand: rand, feeCollectionAddress: feeCollectionAddress, database: database, highestSafeAccountIDInBlock: nextAccountID - 1, numberOfColdAccounts: int64(config.MinimumNumberOfColdAccounts), metrics: metrics, - }, nil + } } // Get the next account ID to be used when creating a new account. This is also the total number of accounts @@ -149,11 +125,6 @@ func (d *DataGenerator) NextErc20ContractID() int64 { return d.nextErc20ContractID } -// Get the next block number as it was at startup time. -func (d *DataGenerator) InitialNextBlockNumber() uint64 { - return d.initialNextBlockNumber -} - // Creates a new account and optionally writes it to the database. Returns the address of the new // account and whether it is a cold account (vs dormant). func (d *DataGenerator) CreateNewAccount( diff --git a/sei-db/bench/cryptosim/database.go b/sei-db/bench/cryptosim/database.go index 7aca0f7cff..5400288cfa 100644 --- a/sei-db/bench/cryptosim/database.go +++ b/sei-db/bench/cryptosim/database.go @@ -2,10 +2,13 @@ package cryptosim import ( "encoding/binary" + "errors" "fmt" - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" + "github.com/sei-protocol/sei-chain/sei-db/common/keys" + "github.com/sei-protocol/sei-chain/sei-db/controller" "github.com/sei-protocol/sei-chain/sei-db/proto" + gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" ) // Encapsulates the database for the cryptosim benchmark. @@ -14,7 +17,14 @@ type Database struct { config *CryptoSimConfig // The database implementation to use for the benchmark. - db wrappers.DBWrapper + db gigatypes.StateDB + + // Enforces retention across the stores the database opened. + garbageCollector *controller.StorageGarbageCollector + + // A read-only view of the most recently committed block, which every read that misses the + // current batch is served from. Replaced after each commit. + view gigatypes.StateView // The total number of transactions executed by the benchmark since it last started. transactionCount int64 @@ -22,8 +32,8 @@ type Database struct { // A count of the number of transactions in the current batch. transactionsInCurrentBlock int64 - // The next block number to be persisted. Tracked internally and incremented after each finalized block. - nextBlockNumber uint64 + // The block number the next commit lands on. Incremented after each finalized block. + nextBlockNumber int64 // The current batch of key-value pairs waiting to be committed. Represents changes we are accumulating // as part of a simulated "block". Stored as value []byte; converted to NamedChangeSet when applied to the DB. @@ -35,36 +45,38 @@ type Database struct { // The metrics for the benchmark. metrics *CryptosimMetrics - // Takes one block hash per block committed, so that the benchmark cannot outrun hashing. Nil when - // the database publishes no block hashes. + // Takes one block hash per block committed, so that the benchmark cannot outrun hashing. hashes *blockHashWaiter } // Creates a new database for the cryptosim benchmark. func NewDatabase( config *CryptoSimConfig, - db wrappers.DBWrapper, + db gigatypes.StateDB, + garbageCollector *controller.StorageGarbageCollector, metrics *CryptosimMetrics, - initialNextBlockNumber uint64, ) (*Database, error) { + // The view is both what reads are served from and where the starting height comes from: the + // store accepts only the block after the one it opened at. + view := db.OpenView() database := &Database{ - config: config, - db: db, - batch: NewSyncMap[string, []byte](), - metrics: metrics, - nextBlockNumber: initialNextBlockNumber, + config: config, + db: db, + garbageCollector: garbageCollector, + view: view, + batch: NewSyncMap[string, []byte](), + metrics: metrics, + nextBlockNumber: view.GetBlockHeight() + 1, } // Registered here because this is before the first block is committed, and that is the only place // a listener can be sure of being handed every block's hash. waiter := newBlockHashWaiter(config.HashLagBlocks, metrics) - registered, err := db.RegisterHashListener(waiter.listen) - if err != nil { + if _, err := db.RegisterHashListener(waiter.listen); err != nil { + view.Close() return nil, fmt.Errorf("failed to register a block hash listener: %w", err) } - if registered { - database.hashes = waiter - } + database.hashes = waiter return database, nil } @@ -81,20 +93,11 @@ func (d *Database) Put(key []byte, value []byte) error { // // This method is safe to call concurrently with other calls to Put() and Get(). Is not thread // safe with FinalizeBlock(). -func (d *Database) Get(key []byte) ([]byte, bool, error) { +func (d *Database) Get(key []byte) ([]byte, bool) { if value, found := d.batch.Get(string(key)); found { - return value, true, nil - } - - value, found, err := d.db.Read(key) - if err != nil { - return nil, false, fmt.Errorf("failed to read from database: %w", err) - } - if found { - return value, true, nil + return value, true } - - return nil, false, nil + return d.view.Get(keys.EVMStoreKey, key) } // Signal that a transaction has been added to the current block. @@ -152,7 +155,7 @@ func (d *Database) FinalizeBlock( changeSets := make([]*proto.NamedChangeSet, 0, d.transactionsInCurrentBlock+3) for key, value := range d.batch.Iterator() { changeSets = append(changeSets, &proto.NamedChangeSet{ - Name: wrappers.EVMStoreName, + Name: keys.EVMStoreKey, Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{{Key: []byte(key), Value: value}}}, }) } @@ -163,7 +166,7 @@ func (d *Database) FinalizeBlock( //nolint:gosec // G115 - nextAccountID is benchmark counter, overflow acceptable binary.BigEndian.PutUint64(nonceValue, uint64(nextAccountID)) changeSets = append(changeSets, &proto.NamedChangeSet{ - Name: wrappers.EVMStoreName, + Name: keys.EVMStoreKey, Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{ {Key: AccountIDCounterKey(), Value: nonceValue}, }}, @@ -174,49 +177,30 @@ func (d *Database) FinalizeBlock( //nolint:gosec // G115 - nextErc20ContractID is benchmark counter, overflow acceptable binary.BigEndian.PutUint64(erc20ContractIDValue, uint64(nextErc20ContractID)) changeSets = append(changeSets, &proto.NamedChangeSet{ - Name: wrappers.EVMStoreName, + Name: keys.EVMStoreKey, Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{ {Key: Erc20IDCounterKey(), Value: erc20ContractIDValue}, }}, }) - // Persist the block number counter in every batch. - blockNumberValue := make([]byte, 8) - binary.BigEndian.PutUint64(blockNumberValue, d.nextBlockNumber) - changeSets = append(changeSets, &proto.NamedChangeSet{ - Name: wrappers.EVMStoreName, - Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{ - {Key: BlockNumberCounterKey(), Value: blockNumberValue}, - }}, - }) - d.nextBlockNumber++ - - entry := &proto.ChangelogEntry{ - Version: d.db.Version() + 1, - Changesets: changeSets, - } - err := d.db.ApplyChangeSets(entry) - if err != nil { - return fmt.Errorf("failed to apply change sets: %w", err) - } + blockNum := d.nextBlockNumber d.metrics.ReportBlockFinalized(d.transactionsInCurrentBlock) d.transactionsInCurrentBlock = 0 // One commit per block: that is the store contract, so the benchmark must not batch. d.metrics.SetMainThreadPhase("committing") - if _, err := d.db.Commit(); err != nil { - return fmt.Errorf("failed to commit: %w", err) + if err := d.db.CommitStateChanges(blockNum, changeSets); err != nil { + return fmt.Errorf("failed to commit block %d: %w", blockNum, err) } + d.nextBlockNumber++ d.metrics.ReportDBCommit() + d.reopenView() // Committing a block is not finishing it: the hash of a block committed a bounded number of // blocks ago is taken here, and waited for when hashing has fallen behind execution. - if d.hashes != nil { - if err := d.hashes.awaitBlock(); err != nil { - return fmt.Errorf("failed to obtain a block hash after committing block %d: %w", - d.db.Version(), err) - } + if err := d.hashes.awaitBlock(); err != nil { + return fmt.Errorf("failed to obtain a block hash after committing block %d: %w", blockNum, err) } d.metrics.SetMainThreadPhase("executing") @@ -224,32 +208,48 @@ func (d *Database) FinalizeBlock( return nil } +// reopenView replaces the read view with one over the block just committed. A view never observes +// writes made after it was opened, so without this every read would keep answering from the height +// the benchmark started at. +func (d *Database) reopenView() { + d.view.Close() + d.view = d.db.OpenView() +} + // Close the database and release any resources. func (d *Database) Close(nextAccountID int64, nextErc20ContractID int64) error { fmt.Printf("Committing final batch.\n") + // A failed final commit still has to release the stores below: they hold the state WAL directory's + // exclusive lock, which an in-process retry needs back. + var errs error if err := d.FinalizeBlock(nextAccountID, nextErc20ContractID); err != nil { - return fmt.Errorf("failed to commit batch: %w", err) - } - - fmt.Printf("Closing database.\n") - err := d.db.Close() - if err != nil { - return fmt.Errorf("failed to close database: %w", err) + errs = errors.Join(errs, fmt.Errorf("failed to commit batch: %w", err)) } - return nil + return errors.Join(errs, d.CloseWithoutFinalizing()) } // Close the database and release any resources without finalizing the last batch. func (d *Database) CloseWithoutFinalizing() error { fmt.Printf("Closing database.\n") - err := d.db.Close() - if err != nil { - return fmt.Errorf("failed to close database: %w", err) + + var errs error + + // The collector prunes the stores closed below, so it stops before them. A failure to stop it + // does not skip those closes: every failure here is collected and reported together. + if err := d.garbageCollector.Close(); err != nil { + errs = errors.Join(errs, fmt.Errorf("failed to close the storage garbage collector: %w", err)) } - return nil + // The view holds a reference into the store, which cannot release it while the view is open. + d.view.Close() + + if err := d.db.Close(); err != nil { + errs = errors.Join(errs, fmt.Errorf("failed to close database: %w", err)) + } + + return errs } // Set the function that flushes the executors. This setter is required to break a circular dependency. diff --git a/sei-db/bench/cryptosim/historical_offload_test.go b/sei-db/bench/cryptosim/historical_offload_test.go deleted file mode 100644 index 7e79cd4570..0000000000 --- a/sei-db/bench/cryptosim/historical_offload_test.go +++ /dev/null @@ -1,76 +0,0 @@ -package cryptosim - -import ( - "testing" - - "github.com/stretchr/testify/require" - - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" -) - -func TestValidateHistoricalOffloadRequiresConfigForHistoricalOffloadBackend(t *testing.T) { - cfg := DefaultCryptoSimConfig() - cfg.DataDir = t.TempDir() - cfg.LogDir = t.TempDir() - cfg.Backend = wrappers.SSHistoricalOffload - - err := cfg.Validate() - require.ErrorContains(t, err, "historical offload config is required") -} - -func TestValidateHistoricalOffloadRequiresKafkaConfig(t *testing.T) { - cfg := DefaultCryptoSimConfig() - cfg.DataDir = t.TempDir() - cfg.LogDir = t.TempDir() - cfg.Backend = wrappers.SSHistoricalOffload - cfg.HistoricalOffload = &wrappers.HistoricalOffloadConfig{ - Provider: "kafka", - } - - err := cfg.Validate() - require.Error(t, err) - require.Contains(t, err.Error(), "historical offload kafka config is required") -} - -func TestValidateHistoricalOffloadKafkaAcceptsMinimalValidConfig(t *testing.T) { - cfg := DefaultCryptoSimConfig() - cfg.DataDir = t.TempDir() - cfg.LogDir = t.TempDir() - cfg.Backend = wrappers.SSHistoricalOffload - cfg.HistoricalOffload = &wrappers.HistoricalOffloadConfig{ - Provider: "kafka", - Kafka: &wrappers.KafkaHistoricalOffloadConfig{ - Brokers: []string{"localhost:9092"}, - Topic: "historical-offload", - }, - } - - require.NoError(t, cfg.Validate()) - require.Equal(t, "cryptosim-historical-offload", cfg.HistoricalOffload.Kafka.ClientID) - require.Equal(t, "none", cfg.HistoricalOffload.Kafka.RequiredAcks) - require.Nil(t, cfg.HistoricalOffload.Kafka.Async) - require.Equal(t, 1000, cfg.HistoricalOffload.Kafka.BatchSize) - require.Equal(t, 4<<20, cfg.HistoricalOffload.Kafka.BatchBytes) -} - -func TestValidateHistoricalOffloadKafkaIAMRequiresRegion(t *testing.T) { - cfg := DefaultCryptoSimConfig() - cfg.DataDir = t.TempDir() - cfg.LogDir = t.TempDir() - cfg.Backend = wrappers.SSHistoricalOffload - cfg.HistoricalOffload = &wrappers.HistoricalOffloadConfig{ - Provider: "kafka", - Kafka: &wrappers.KafkaHistoricalOffloadConfig{ - Brokers: []string{"localhost:9098"}, - Topic: "historical-offload", - TLSEnabled: true, - SASLMechanism: "aws-msk-iam", - }, - } - - err := cfg.Validate() - require.ErrorContains(t, err, "region is required") - - cfg.HistoricalOffload.Kafka.Region = "eu-central-1" - require.NoError(t, cfg.Validate()) -} diff --git a/sei-db/bench/cryptosim/transaction.go b/sei-db/bench/cryptosim/transaction.go index 762d2586ab..04647e66db 100644 --- a/sei-db/bench/cryptosim/transaction.go +++ b/sei-db/bench/cryptosim/transaction.go @@ -103,9 +103,7 @@ func (txn *transaction) Execute( phaseTimer.SetPhase("read_erc20") // Read the simulated ERC20 contract. - if _, _, err := database.Get(txn.erc20Contract); err != nil { - return fmt.Errorf("failed to get ERC20 contract: %w", err) - } + database.Get(txn.erc20Contract) // Read the following: // - the sender's native balance / nonce / codehash @@ -120,39 +118,29 @@ func (txn *transaction) Execute( // Technically, we are just requesting to read the codehash, but internally the codehash is bundled with // the nonce and balance, so all of this data will be read from low level storage, even if it isn't being // returned to the caller. - if _, _, err := database.Get(txn.srcAccount); err != nil { - return fmt.Errorf("failed to get source account: %w", err) - } + database.Get(txn.srcAccount) phaseTimer.SetPhase("read_dst_account") // Read the receiver's native balance / nonce / codehash. - if _, _, err := database.Get(txn.dstAccount); err != nil { - return fmt.Errorf("failed to get destination account: %w", err) - } + database.Get(txn.dstAccount) phaseTimer.SetPhase("read_src_account_slot") // Read the sender's storage slot for the ERC20 contract. // We don't care if the value isn't in the DB yet, since we don't pre-populate the database with storage slots. - if _, _, err := database.Get(txn.srcAccountSlot); err != nil { - return fmt.Errorf("failed to get source account slot: %w", err) - } + database.Get(txn.srcAccountSlot) phaseTimer.SetPhase("read_dst_account_slot") // Read the receiver's storage slot for the ERC20 contract. // We don't care if the value isn't in the DB yet, since we don't pre-populate the database with storage slots. - if _, _, err := database.Get(txn.dstAccountSlot); err != nil { - return fmt.Errorf("failed to get destination account slot: %w", err) - } + database.Get(txn.dstAccountSlot) phaseTimer.SetPhase("read_fee_collection_account") // Read the fee collection account's native balance. - if _, _, err := database.Get(feeCollectionAddress); err != nil { - return fmt.Errorf("failed to get fee collection account: %w", err) - } + database.Get(feeCollectionAddress) } phaseTimer.SetPhase("update_balances") diff --git a/sei-db/bench/cryptosim/transaction_test.go b/sei-db/bench/cryptosim/transaction_test.go index 1469932491..047988120f 100644 --- a/sei-db/bench/cryptosim/transaction_test.go +++ b/sei-db/bench/cryptosim/transaction_test.go @@ -5,52 +5,41 @@ import ( "github.com/stretchr/testify/require" - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" - commonmetrics "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - scTypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" + "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/lthash" ) -type readTrackingWrapper struct { - readCalls int -} +// readTrackingView counts the reads a transaction issues. It embeds StateView without implementing +// it, so any method the transaction is not expected to call panics on the nil interface rather than +// answering with a zero value. +type readTrackingView struct { + gigatypes.StateView -func (r *readTrackingWrapper) ApplyChangeSets(_ *proto.ChangelogEntry) error { - return nil -} - -func (r *readTrackingWrapper) Read(_ []byte) ([]byte, bool, error) { - r.readCalls++ - return nil, false, nil + // Reads served, in call order. + readCalls int } -func (r *readTrackingWrapper) Commit() (int64, error) { - return 0, nil -} +func (v *readTrackingView) GetBlockHeight() int64 { return 0 } -func (r *readTrackingWrapper) Close() error { - return nil +func (v *readTrackingView) Get(_ string, _ []byte) ([]byte, bool) { + v.readCalls++ + return nil, false } -func (r *readTrackingWrapper) Version() int64 { - return 0 -} +func (v *readTrackingView) Close() {} -func (r *readTrackingWrapper) LoadLatest() error { - return nil -} +// readTrackingStateDB serves every view from one readTrackingView, so a test can count the reads +// made through it. +type readTrackingStateDB struct { + gigatypes.StateDB -func (r *readTrackingWrapper) Importer(_ int64) (scTypes.Importer, error) { - return nil, nil + view *readTrackingView } -func (r *readTrackingWrapper) GetPhaseTimer() *commonmetrics.PhaseTimer { - return nil -} +func (s *readTrackingStateDB) OpenView() gigatypes.StateView { return s.view } -func (r *readTrackingWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil +func (s *readTrackingStateDB) RegisterHashListener(_ gigatypes.HashListener) (lthash.BlockHash, error) { + return lthash.BlockHash{}, nil } func TestTransactionExecuteSkipsReadsWhenDisabled(t *testing.T) { @@ -59,8 +48,8 @@ func TestTransactionExecuteSkipsReadsWhenDisabled(t *testing.T) { cfg := DefaultCryptoSimConfig() cfg.DisableTransactionReads = true - wrapper := &readTrackingWrapper{} - db, err := NewDatabase(cfg, wrapper, nil, 0) + stateDB := &readTrackingStateDB{view: &readTrackingView{}} + db, err := NewDatabase(cfg, stateDB, nil, nil) require.NoError(t, err) txn := &transaction{ @@ -77,10 +66,10 @@ func TestTransactionExecuteSkipsReadsWhenDisabled(t *testing.T) { } require.NoError(t, txn.Execute(db, []byte("fee"), nil)) - require.Zero(t, wrapper.readCalls) + require.Zero(t, stateDB.view.readCalls) - _, found, err := db.Get([]byte("src")) - require.NoError(t, err) + // The write the transaction made is in the batch, so it is served without reaching the view. + _, found := db.Get([]byte("src")) require.True(t, found) } @@ -89,5 +78,4 @@ func TestDefaultCryptoSimConfigDisablesTransactionReadsByDefaultFalse(t *testing cfg := DefaultCryptoSimConfig() require.False(t, cfg.DisableTransactionReads) - require.Equal(t, wrappers.FlatKV, cfg.Backend) } diff --git a/sei-db/bench/cryptosim/util.go b/sei-db/bench/cryptosim/util.go index 9eab5713ba..17d3594dbb 100644 --- a/sei-db/bench/cryptosim/util.go +++ b/sei-db/bench/cryptosim/util.go @@ -28,11 +28,6 @@ func Erc20IDCounterKey() []byte { return keys.BuildEVMKey(keys.EVMKeyCode, paddedCounterKey(erc20IdCounterKey)) } -// Get the key for the block number counter in the database. -func BlockNumberCounterKey() []byte { - return keys.BuildEVMKey(keys.EVMKeyCode, paddedCounterKey(blockNumberCounterKey)) -} - // paddedCounterKey pads the string to AddressLen bytes for use with EVM key builders. func paddedCounterKey(s string) []byte { b := make([]byte, keys.AddressLen) diff --git a/sei-db/bench/wrappers/combined_wrapper.go b/sei-db/bench/wrappers/combined_wrapper.go deleted file mode 100644 index 368f96516a..0000000000 --- a/sei-db/bench/wrappers/combined_wrapper.go +++ /dev/null @@ -1,78 +0,0 @@ -package wrappers - -import ( - "sync/atomic" - - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - dbTypes "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - scTypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -var _ DBWrapper = (*combinedWrapper)(nil) - -// combinedWrapper drives both a State Commit (SC) and State Store (SS) backend -// from the same changeset stream, mirroring production where SC and SS receive -// identical writes. -type combinedWrapper struct { - sc DBWrapper - ss dbTypes.StateStore - ssVersion atomic.Int64 -} - -func NewCombinedWrapper(sc DBWrapper, ss dbTypes.StateStore) DBWrapper { - w := &combinedWrapper{sc: sc, ss: ss} - w.ssVersion.Store(ss.GetLatestVersion()) - return w -} - -func (c *combinedWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - if err := c.sc.ApplyChangeSets(entry); err != nil { - return err - } - c.ssVersion.Store(entry.Version) - return c.ss.ApplyChangesetAsync(entry.Version, entry.Changesets) -} - -func (c *combinedWrapper) Read(key []byte) (data []byte, found bool, err error) { - return c.sc.Read(key) -} - -func (c *combinedWrapper) Commit() (int64, error) { - if _, err := c.sc.Commit(); err != nil { - return 0, err - } - return c.ssVersion.Load(), nil -} - -func (c *combinedWrapper) Close() error { - scErr := c.sc.Close() - ssErr := c.ss.Close() - if scErr != nil { - return scErr - } - return ssErr -} - -func (c *combinedWrapper) Version() int64 { - return c.ssVersion.Load() -} - -func (c *combinedWrapper) LoadLatest() error { - return c.sc.LoadLatest() -} - -func (c *combinedWrapper) Importer(version int64) (scTypes.Importer, error) { - return c.sc.Importer(version) -} - -// RegisterHashListener forwards to the SC backend. The SS backend commits the same changesets but -// computes no block hash of its own. -func (c *combinedWrapper) RegisterHashListener(listener gigatypes.HashListener) (bool, error) { - return c.sc.RegisterHashListener(listener) -} - -func (c *combinedWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} diff --git a/sei-db/bench/wrappers/composite_wrapper.go b/sei-db/bench/wrappers/composite_wrapper.go deleted file mode 100644 index 8f462b06a6..0000000000 --- a/sei-db/bench/wrappers/composite_wrapper.go +++ /dev/null @@ -1,65 +0,0 @@ -package wrappers - -import ( - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/composite" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -var _ DBWrapper = (*compositeWrapper)(nil) - -// compositeWrapper wraps a composite commit store to implement the DBWrapper interface. -type compositeWrapper struct { - base *composite.CompositeCommitStore -} - -// NewCompositeWrapper creates a new compositeWrapper with a given composite commit store. -func NewCompositeWrapper(store *composite.CompositeCommitStore) DBWrapper { - return &compositeWrapper{ - base: store, - } -} - -func (c *compositeWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - return c.base.ApplyChangeSets(entry.Changesets) -} - -func (c *compositeWrapper) Commit() (int64, error) { - // The benchmark wrapper interface carries no height, so the next one is derived here. That is - // sound only because nothing in the benchmark path takes a block's hash before committing it. - return c.base.Commit(c.base.Version() + 1) -} - -func (c *compositeWrapper) LoadLatest() error { - return c.base.LoadLatest() -} - -func (c *compositeWrapper) Version() int64 { - return c.base.Version() -} - -func (c *compositeWrapper) Importer(version int64) (types.Importer, error) { - return c.base.Importer(version) -} - -func (c *compositeWrapper) Close() error { - return c.base.Close() -} - -func (c *compositeWrapper) Read(key []byte) (data []byte, found bool, err error) { - store := c.base.GetChildStoreByName(EVMStoreName) - data = store.Get(key) - return data, data != nil, nil -} - -// RegisterHashListener reports that this DB publishes no block hashes. The composite store consumes -// flatKV's hashes itself, in order to answer Cosmos synchronously, so it admits no second consumer. -func (c *compositeWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil -} - -func (c *compositeWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} diff --git a/sei-db/bench/wrappers/db_implementations.go b/sei-db/bench/wrappers/db_implementations.go deleted file mode 100644 index c17e2b3af2..0000000000 --- a/sei-db/bench/wrappers/db_implementations.go +++ /dev/null @@ -1,228 +0,0 @@ -package wrappers - -import ( - "context" - "fmt" - "path/filepath" - - commonevm "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/state_db/sc/composite" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv" - flatkvConfig "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/config" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/memiavl" - sctypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" - ssComposite "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/composite" -) - -const EVMStoreName = commonevm.EVMStoreKey - -type DBType string - -const ( - NoOp DBType = "NoOp" - MemIAVL DBType = "MemIAVL" - FlatKV DBType = "FlatKV" - CompositeDual DBType = "CompositeDual" - CompositeSplit DBType = "CompositeSplit" - CompositeCosmos DBType = "CompositeCosmos" - - SSComposite DBType = "SSComposite" - SSHistoricalOffload DBType = "SSHistoricalOffload" - CompositeDual_SSComposite DBType = "CompositeDual+SSComposite" -) - -func DefaultBenchStateStoreConfig() *config.StateStoreConfig { - cfg := config.DefaultStateStoreConfig() - cfg.AsyncWriteBuffer = config.DefaultSSAsyncBuffer - cfg.EVMSplit = true - return &cfg -} - -// DefaultBenchMemIAVLConfig returns the memiavl config the benchmarks open -// with by default. Note AsyncCommitBuffer=10: Commit() returns once the WAL -// write is enqueued, not once it is durable. -func DefaultBenchMemIAVLConfig() memiavl.Config { - cfg := memiavl.DefaultConfig() - cfg.AsyncCommitBuffer = 10 - cfg.SnapshotInterval = 1000 - cfg.SnapshotMinTimeInterval = 60 - return cfg -} - -func newMemIAVLCommitStore(dbDir string, cfg *memiavl.Config) (DBWrapper, error) { - if cfg == nil { - defaultCfg := DefaultBenchMemIAVLConfig() - cfg = &defaultCfg - } - fmt.Printf("Opening memIAVL from directory %s\n", dbDir) - cs := memiavl.NewCommitStore(dbDir, *cfg) - if err := cs.Initialize([]string{EVMStoreName}); err != nil { - return nil, fmt.Errorf("memiavl Initialize: %w", err) - } - _, err := cs.LoadVersion(0, false) - if err != nil { - if closeErr := cs.Close(); closeErr != nil { - fmt.Printf("failed to close commit store during error recovery: %v\n", closeErr) - } - return nil, fmt.Errorf("failed to load version: %w", err) - } - return NewMemIAVLWrapper(cs), nil -} - -func newFlatKVCommitStore(ctx context.Context, dbDir string, config *flatkvConfig.Config) (DBWrapper, error) { - if config == nil { - config = flatkvConfig.DefaultConfig() - } - config.DataDir = dbDir - - fmt.Printf("Opening flatKV from directory %s\n", dbDir) - stateWAL, err := flatkv.OpenStateWAL(config) - if err != nil { - return nil, fmt.Errorf("failed to open FlatKV state WAL: %w", err) - } - cs, err := flatkv.NewCommitStore(ctx, config, stateWAL) - if err != nil { - _ = stateWAL.Close() - return nil, fmt.Errorf("failed to create FlatKV commit store: %w", err) - } - if err := cs.LoadLatest(); err != nil { - if closeErr := cs.Close(); closeErr != nil { - fmt.Printf("failed to close commit store during error recovery: %v\n", closeErr) - } - return nil, fmt.Errorf("failed to load version: %w", err) - } - return NewFlatKVWrapper(cs), nil -} - -func newCompositeCommitStore(ctx context.Context, dbDir string, writeMode sctypes.WriteMode) (DBWrapper, error) { - cfg := config.DefaultStateCommitConfig() - cfg.WriteMode = writeMode - cfg.MemIAVLConfig.AsyncCommitBuffer = 10 - cfg.MemIAVLConfig.SnapshotInterval = 100 - - cs, err := composite.NewCompositeCommitStore(ctx, dbDir, cfg) - if err != nil { - return nil, fmt.Errorf("failed to create composite commit store: %w", err) - } - if err := cs.CleanupCrashArtifacts(); err != nil { - return nil, fmt.Errorf("failed to cleanup crash artifacts: %w", err) - } - if err := cs.Initialize([]string{EVMStoreName}); err != nil { - return nil, fmt.Errorf("composite Initialize: %w", err) - } - - if err := cs.LoadLatest(); err != nil { - if closeErr := cs.Close(); closeErr != nil { - fmt.Printf("failed to close commit store during error recovery: %v\n", closeErr) - } - return nil, fmt.Errorf("failed to load version: %w", err) - } - - return NewCompositeWrapper(cs), nil -} - -func openSSComposite(dir string, cfg config.StateStoreConfig) (*ssComposite.CompositeStateStore, error) { - return ssComposite.NewCompositeStateStore(cfg, dir) -} - -func newSSCompositeStateStore(dbDir string, ssConfig *config.StateStoreConfig) (DBWrapper, error) { - if ssConfig == nil { - ssConfig = DefaultBenchStateStoreConfig() - } - fmt.Printf("Opening composite state store from directory %s\n", dbDir) - store, err := openSSComposite(dbDir, *ssConfig) - if err != nil { - return nil, fmt.Errorf("failed to open composite state store: %w", err) - } - return NewStateStoreWrapper(store), nil -} - -func newCombinedCompositeDualSSComposite( - ctx context.Context, - dbDir string, - ssConfig *config.StateStoreConfig, -) (DBWrapper, error) { - if ssConfig == nil { - ssConfig = DefaultBenchStateStoreConfig() - } - - fmt.Printf("Opening CompositeDual (SC) + Composite (SS) from directory %s\n", dbDir) - sc, err := newCompositeCommitStore(ctx, filepath.Join(dbDir, "sc"), sctypes.TestOnlyDualWrite) - if err != nil { - return nil, fmt.Errorf("failed to create SC store: %w", err) - } - ss, err := openSSComposite(filepath.Join(dbDir, "ss"), *ssConfig) - if err != nil { - _ = sc.Close() - return nil, fmt.Errorf("failed to create SS store: %w", err) - } - return NewCombinedWrapper(sc, ss), nil -} - -// backendConfig converts the untyped config NewDBImpl is handed into the one its backend expects. -// -// A nil config is ordinary rather than exceptional: runBenchmark passes nil for every backend, and -// each constructor supplies its own default. So nil is passed straight through, and only a config of -// the wrong type is an error. Asserting without this — dbConfig.(*T) on a nil interface — panics -// before the constructor can apply that default, which is how three bench backends came to crash at -// startup. -func backendConfig[T any](dbType DBType, dbConfig any) (*T, error) { - if dbConfig == nil { - return nil, nil - } - typed, ok := dbConfig.(*T) - if !ok { - var want T - return nil, fmt.Errorf("invalid %s config type %T, want *%T", dbType, dbConfig, want) - } - return typed, nil -} - -// NewDBImpl instantiates a new empty DBWrapper based on the given DBType. -func NewDBImpl(ctx context.Context, dbType DBType, dataDir string, dbConfig any) (DBWrapper, error) { - switch dbType { - case NoOp: - return NewNoOpWrapper(), nil - case MemIAVL: - cfg, err := backendConfig[memiavl.Config](dbType, dbConfig) - if err != nil { - return nil, err - } - return newMemIAVLCommitStore(dataDir, cfg) - case FlatKV: - cfg, err := backendConfig[flatkvConfig.Config](dbType, dbConfig) - if err != nil { - return nil, err - } - return newFlatKVCommitStore(ctx, dataDir, cfg) - case CompositeDual: - return newCompositeCommitStore(ctx, dataDir, sctypes.TestOnlyDualWrite) - case CompositeSplit: - return newCompositeCommitStore(ctx, dataDir, sctypes.EVMMigrated) - case CompositeCosmos: - return newCompositeCommitStore(ctx, dataDir, sctypes.MemiavlOnly) - case SSComposite: - cfg, err := backendConfig[config.StateStoreConfig](dbType, dbConfig) - if err != nil { - return nil, err - } - return newSSCompositeStateStore(dataDir, cfg) - case SSHistoricalOffload: - // No default: the stream needs brokers only the caller knows, so a missing config is - // reported by HistoricalOffloadConfig.Validate rather than invented here. - cfg, err := backendConfig[HistoricalOffloadConfig](dbType, dbConfig) - if err != nil { - return nil, err - } - return newSSHistoricalOffloadStateStore(ctx, dataDir, cfg) - case CompositeDual_SSComposite: - cfg, err := backendConfig[config.StateStoreConfig](dbType, dbConfig) - if err != nil { - return nil, err - } - return newCombinedCompositeDualSSComposite(ctx, dataDir, cfg) - default: - return nil, fmt.Errorf("unsupported DB type: %s", dbType) - } -} diff --git a/sei-db/bench/wrappers/db_wrapper.go b/sei-db/bench/wrappers/db_wrapper.go deleted file mode 100644 index 5efe9cf0d7..0000000000 --- a/sei-db/bench/wrappers/db_wrapper.go +++ /dev/null @@ -1,46 +0,0 @@ -package wrappers - -import ( - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -// This benchmarking utility is capable of benchmarking a DB that implements this interface. -type DBWrapper interface { - // ApplyChangeSets applies a versioned changelog entry. SC-backed wrappers buffer - // entry changesets until Commit, while SS-backed wrappers can use entry.Version - // to persist at the benchmark-assigned version immediately. - ApplyChangeSets(entry *proto.ChangelogEntry) error - - // Read reads the value for the given key. - Read(key []byte) (data []byte, found bool, err error) - - // Commit persists buffered writes and advances the version. - Commit() (int64, error) - - // Close releases any resources held by the DB. - Close() error - - // Version returns the latest committed version. - Version() int64 - - // LoadLatest opens the DB at its latest committed version. Benchmarks only ever run against the tip, so - // there is no way to ask for a historical version. - LoadLatest() error - - // Importer return an importer which load snapshot data into the database - Importer(version int64) (types.Importer, error) - - // Get the phase timer used to measure time spent in various phases of execution. Useful for metrics - // integration with external phases of execution. - // - // If the underlying DB does not support phase timers, return nil. - GetPhaseTimer() *metrics.PhaseTimer - - // RegisterHashListener subscribes listener to the hash of each block this DB commits, one per - // block in block order. It reports false when the DB publishes no block hashes, in which case - // there is nothing for a benchmark to wait on. - RegisterHashListener(listener gigatypes.HashListener) (registered bool, err error) -} diff --git a/sei-db/bench/wrappers/flatkv_wrapper.go b/sei-db/bench/wrappers/flatkv_wrapper.go deleted file mode 100644 index 97591cd146..0000000000 --- a/sei-db/bench/wrappers/flatkv_wrapper.go +++ /dev/null @@ -1,86 +0,0 @@ -package wrappers - -import ( - "fmt" - - "github.com/sei-protocol/sei-chain/sei-db/common/keys" - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -var _ DBWrapper = (*flatKVWrapper)(nil) - -// flatKVWrapper wraps a flatkv commit store to implement the DBWrapper interface. -// FlatKV persists exactly one block per Commit, so benchmarks must commit every -// block. Several -// ApplyChangeSets calls may still precede one Commit as long as they all target the -// same height; Commit() consults PendingVersion() to find that height. -type flatKVWrapper struct { - base gigatypes.LiveStateStore -} - -// NewFlatKVWrapper creates a new flatKVWrapper with a given flatkv store. -func NewFlatKVWrapper(store gigatypes.LiveStateStore) DBWrapper { - return &flatKVWrapper{ - base: store, - } -} - -func (f *flatKVWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - version := entry.Version - if version <= 0 { - version = f.nextVersion() - } - return f.base.ApplyChangeSets(version, entry.Changesets) -} - -func (f *flatKVWrapper) Commit() (int64, error) { - version := f.base.PendingVersion() - if version == 0 { - version = f.base.Version() + 1 - } - return f.base.Commit(version) -} - -func (f *flatKVWrapper) LoadLatest() error { - return f.base.LoadLatest() -} - -func (f *flatKVWrapper) Version() int64 { - return f.base.Version() -} - -// nextVersion computes the height for the next ApplyChangeSets call: one past the -// committed version. It deliberately ignores PendingVersion() — a pending block's -// writes may be extended at its own height, never continued at the next one. -func (f *flatKVWrapper) nextVersion() int64 { - return f.base.Version() + 1 -} - -func (f *flatKVWrapper) Importer(version int64) (types.Importer, error) { - return f.base.Importer(version) -} - -func (f *flatKVWrapper) Close() error { - return f.base.Close() -} - -func (f *flatKVWrapper) Read(key []byte) (data []byte, found bool, err error) { - val, ok := f.base.Get(keys.EVMStoreKey, key) - return val, ok, nil -} - -// RegisterHashListener subscribes listener to flatKV's block hashes. The hash the store returns is -// dropped: a benchmark waits on the blocks it is about to commit, not the one already behind it. -func (f *flatKVWrapper) RegisterHashListener(listener gigatypes.HashListener) (bool, error) { - if _, err := f.base.RegisterHashListener(listener); err != nil { - return false, fmt.Errorf("register a hash listener on flatkv: %w", err) - } - return true, nil -} - -func (f *flatKVWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return f.base.GetPhaseTimer() -} diff --git a/sei-db/bench/wrappers/flatkv_wrapper_test.go b/sei-db/bench/wrappers/flatkv_wrapper_test.go deleted file mode 100644 index 2e8722eb46..0000000000 --- a/sei-db/bench/wrappers/flatkv_wrapper_test.go +++ /dev/null @@ -1,58 +0,0 @@ -package wrappers - -import ( - "testing" - - "github.com/stretchr/testify/require" - - "github.com/sei-protocol/sei-chain/sei-db/proto" -) - -func flatKVEntry(version int64, value byte) *proto.ChangelogEntry { - return &proto.ChangelogEntry{ - Version: version, - Changesets: []*proto.NamedChangeSet{{ - Name: EVMStoreName, - Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{{Key: []byte("key"), Value: []byte{value}}}}, - }}, - } -} - -// TestFlatKVWrapperCommitsOneBlockPerCommit drives the cryptosim -// Database.FinalizeBlock pattern against a real state -// WAL: each block is applied at Version()+1 and committed immediately. It runs -// several cycles because the WAL only rejects a non-contiguous block number on -// the commit after the first one. -func TestFlatKVWrapperCommitsOneBlockPerCommit(t *testing.T) { - wrapper, err := NewDBImpl(t.Context(), FlatKV, t.TempDir(), nil) - require.NoError(t, err) - defer func() { require.NoError(t, wrapper.Close()) }() - - for block := 1; block <= 5; block++ { - require.NoError(t, - wrapper.ApplyChangeSets(flatKVEntry(wrapper.Version()+1, byte(block))), "block %d", block) - committed, err := wrapper.Commit() - require.NoError(t, err, "block %d", block) - require.Equal(t, int64(block), committed) - require.Equal(t, int64(block), wrapper.Version()) - } -} - -// TestFlatKVWrapperRejectsSecondBlockBeforeCommit pins the barrier against -// batched blocks: FlatKV persists exactly one block per Commit, and the refusal -// lands on the offending ApplyChangeSets rather than on a later Commit. -func TestFlatKVWrapperRejectsSecondBlockBeforeCommit(t *testing.T) { - wrapper, err := NewDBImpl(t.Context(), FlatKV, t.TempDir(), nil) - require.NoError(t, err) - defer func() { require.NoError(t, wrapper.Close()) }() - - require.NoError(t, wrapper.ApplyChangeSets(flatKVEntry(1, 0x01))) - - err = wrapper.ApplyChangeSets(flatKVEntry(2, 0x02)) - require.ErrorContains(t, err, "flatkv: apply version 2 must be committed version 0 plus one") - - // The rejected call left the pending block intact and committable. - committed, err := wrapper.Commit() - require.NoError(t, err) - require.Equal(t, int64(1), committed) -} diff --git a/sei-db/bench/wrappers/historical_offload_wrapper.go b/sei-db/bench/wrappers/historical_offload_wrapper.go deleted file mode 100644 index ff50078c46..0000000000 --- a/sei-db/bench/wrappers/historical_offload_wrapper.go +++ /dev/null @@ -1,187 +0,0 @@ -package wrappers - -import ( - "context" - "fmt" - "io" - "strings" - "sync/atomic" - "time" - - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - scTypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" - "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/offload" -) - -var _ DBWrapper = (*historicalOffloadWrapper)(nil) - -type HistoricalOffloadConfig struct { - Provider string - Kafka *KafkaHistoricalOffloadConfig -} - -type KafkaHistoricalOffloadConfig struct { - Brokers []string - Topic string - ClientID string - Region string - Async *bool - RequiredAcks string - Compression string - BatchSize int - BatchTimeoutMS int - BatchBytes int - TLSEnabled bool - SASLMechanism string -} - -type historicalOffloadWrapper struct { - stream offload.Stream - version atomic.Int64 -} - -func (c *HistoricalOffloadConfig) Validate() error { - if c == nil { - return fmt.Errorf("historical offload config is required") - } - switch strings.ToLower(c.Provider) { - case "kafka": - if c.Kafka == nil { - return fmt.Errorf("historical offload kafka config is required when provider is kafka") - } - c.Kafka.applyDefaults() - return c.Kafka.validate() - default: - return fmt.Errorf("unsupported historical offload provider %q", c.Provider) - } -} - -func (c *KafkaHistoricalOffloadConfig) applyDefaults() { - if c.ClientID == "" { - c.ClientID = "cryptosim-historical-offload" - } - if c.RequiredAcks == "" { - c.RequiredAcks = "none" - } - if c.Compression == "" { - c.Compression = "snappy" - } - if c.BatchSize == 0 { - c.BatchSize = 1000 - } - if c.BatchTimeoutMS == 0 { - c.BatchTimeoutMS = 50 - } - if c.BatchBytes == 0 { - c.BatchBytes = 4 << 20 - } -} - -func (c *KafkaHistoricalOffloadConfig) validate() error { - cfg := offload.KafkaConfig{ - Brokers: c.Brokers, - Topic: c.Topic, - ClientID: c.ClientID, - Region: c.Region, - Async: c.asyncValue(), - RequiredAcks: c.RequiredAcks, - Compression: c.Compression, - BatchSize: c.BatchSize, - BatchTimeout: time.Duration(c.BatchTimeoutMS) * time.Millisecond, - BatchBytes: c.BatchBytes, - TLSEnabled: c.TLSEnabled, - SASLMechanism: c.SASLMechanism, - } - return cfg.Validate() -} - -func (c *KafkaHistoricalOffloadConfig) asyncValue() bool { - return c.Async == nil || *c.Async -} - -func newHistoricalOffloadStream(cfg *HistoricalOffloadConfig) (offload.Stream, error) { - if err := cfg.Validate(); err != nil { - return nil, err - } - kafkaCfg := *cfg.Kafka - kafkaCfg.applyDefaults() - return offload.NewKafkaStream(offload.KafkaConfig{ - Brokers: append([]string(nil), kafkaCfg.Brokers...), - Topic: kafkaCfg.Topic, - ClientID: kafkaCfg.ClientID, - Region: kafkaCfg.Region, - Async: kafkaCfg.asyncValue(), - RequiredAcks: kafkaCfg.RequiredAcks, - Compression: kafkaCfg.Compression, - BatchSize: kafkaCfg.BatchSize, - BatchTimeout: time.Duration(kafkaCfg.BatchTimeoutMS) * time.Millisecond, - BatchBytes: kafkaCfg.BatchBytes, - TLSEnabled: kafkaCfg.TLSEnabled, - SASLMechanism: kafkaCfg.SASLMechanism, - }) -} - -func newSSHistoricalOffloadStateStore(_ context.Context, dbDir string, cfg *HistoricalOffloadConfig) (DBWrapper, error) { - fmt.Printf("Opening historical offload stream from directory %s\n", dbDir) - stream, err := newHistoricalOffloadStream(cfg) - if err != nil { - return nil, fmt.Errorf("failed to create historical offload stream: %w", err) - } - return NewHistoricalOffloadWrapper(stream), nil -} - -func NewHistoricalOffloadWrapper(stream offload.Stream) DBWrapper { - return &historicalOffloadWrapper{stream: stream} -} - -func (h *historicalOffloadWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - ack, err := h.stream.Publish(context.Background(), entry) - if err != nil { - return err - } - if !ack.Accepted { - return fmt.Errorf("historical offload publish was not acknowledged at version %d", entry.Version) - } - h.version.Store(entry.Version) - return nil -} - -func (h *historicalOffloadWrapper) Read(_ []byte) (data []byte, found bool, err error) { - return nil, false, nil -} - -func (h *historicalOffloadWrapper) Commit() (int64, error) { - return h.version.Load(), nil -} - -func (h *historicalOffloadWrapper) Close() error { - var streamErr error - if closer, ok := h.stream.(io.Closer); ok { - streamErr = closer.Close() - } - return streamErr -} - -func (h *historicalOffloadWrapper) Version() int64 { - return h.version.Load() -} - -func (h *historicalOffloadWrapper) LoadLatest() error { - return nil -} - -func (h *historicalOffloadWrapper) Importer(_ int64) (scTypes.Importer, error) { - return nil, fmt.Errorf("import not supported for historical offload wrapper") -} - -// RegisterHashListener reports that this DB publishes no block hashes. An offload stream computes -// none. -func (h *historicalOffloadWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil -} - -func (h *historicalOffloadWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} diff --git a/sei-db/bench/wrappers/historical_offload_wrapper_test.go b/sei-db/bench/wrappers/historical_offload_wrapper_test.go deleted file mode 100644 index d4a185dd57..0000000000 --- a/sei-db/bench/wrappers/historical_offload_wrapper_test.go +++ /dev/null @@ -1,80 +0,0 @@ -package wrappers - -import ( - "context" - "testing" - - "github.com/stretchr/testify/require" - - "github.com/sei-protocol/sei-chain/sei-db/proto" - "github.com/sei-protocol/sei-chain/sei-db/state_db/ss/offload" -) - -type mockOffloadStream struct { - closeCalls int - calls int - entries []*proto.ChangelogEntry -} - -func (m *mockOffloadStream) Publish(_ context.Context, entry *proto.ChangelogEntry) (offload.Ack, error) { - m.calls++ - m.entries = append(m.entries, entry) - return offload.Ack{Accepted: true}, nil -} - -func (m *mockOffloadStream) Close() error { - m.closeCalls++ - return nil -} - -func TestHistoricalOffloadWrapperPublishesWithoutLocalWrites(t *testing.T) { - stream := &mockOffloadStream{} - wrapper := NewHistoricalOffloadWrapper(stream) - - entry := &proto.ChangelogEntry{ - Version: 5, - Changesets: []*proto.NamedChangeSet{{Name: EVMStoreName}}, - } - - err := wrapper.ApplyChangeSets(entry) - require.NoError(t, err) - require.Equal(t, 1, stream.calls) - require.Len(t, stream.entries, 1) - require.Equal(t, entry.Version, stream.entries[0].Version) - require.Equal(t, entry.Changesets, stream.entries[0].Changesets) - require.Equal(t, int64(5), wrapper.Version()) - - data, found, err := wrapper.Read([]byte("ignored")) - require.NoError(t, err) - require.Nil(t, data) - require.False(t, found) - - version, err := wrapper.Commit() - require.NoError(t, err) - require.Equal(t, int64(5), version) - - require.NoError(t, wrapper.Close()) - require.Equal(t, 1, stream.closeCalls) -} - -func TestHistoricalOffloadWrapperImporterUnsupported(t *testing.T) { - wrapper := NewHistoricalOffloadWrapper(&mockOffloadStream{}) - - importer, err := wrapper.Importer(1) - require.Nil(t, importer) - require.Error(t, err) -} - -func TestHistoricalOffloadWrapperGetPhaseTimerIsNil(t *testing.T) { - wrapper := NewHistoricalOffloadWrapper(&mockOffloadStream{}) - require.Nil(t, wrapper.GetPhaseTimer()) -} - -func TestHistoricalOffloadConfigValidateRequiresKafkaConfig(t *testing.T) { - cfg := &HistoricalOffloadConfig{Provider: "kafka"} - err := cfg.Validate() - require.ErrorContains(t, err, "historical offload kafka config is required") -} - -var _ DBWrapper = NewHistoricalOffloadWrapper(&mockOffloadStream{}) -var _ offload.Stream = (*mockOffloadStream)(nil) diff --git a/sei-db/bench/wrappers/memiavl_wrapper.go b/sei-db/bench/wrappers/memiavl_wrapper.go deleted file mode 100644 index 8b43849b71..0000000000 --- a/sei-db/bench/wrappers/memiavl_wrapper.go +++ /dev/null @@ -1,71 +0,0 @@ -package wrappers - -import ( - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/memiavl" - "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -var _ DBWrapper = (*memIAVLWrapper)(nil) - -// A light wrapper around a memiavl commit store to implement the DBWrapper interface. -type memIAVLWrapper struct { - base *memiavl.CommitStore -} - -// NewMemIAVLWrapper creates a new memIAVLWrapper with a given memiavl commit store. -func NewMemIAVLWrapper(commitStore *memiavl.CommitStore) DBWrapper { - return &memIAVLWrapper{ - base: commitStore, - } -} - -func (m *memIAVLWrapper) Commit() (int64, error) { - // The benchmark wrapper interface carries no height, so the next one is derived here. That is - // sound only because nothing in the benchmark path takes a block's hash before committing it. - return m.base.Commit(m.base.Version() + 1) -} - -func (m *memIAVLWrapper) LoadLatest() error { - // memiavl's Committer signature is pinned; (0, false) is its load-latest-writable path. - _, err := m.base.LoadVersion(0, false) - return err -} - -func (m *memIAVLWrapper) Version() int64 { - return m.base.Version() -} - -func (m *memIAVLWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - return m.base.ApplyChangeSets(entry.Changesets) -} - -func (m *memIAVLWrapper) Importer(version int64) (types.Importer, error) { - // Close DB first to release lock - if err := m.Close(); err != nil { - return nil, err - } - return m.base.Importer(version) -} - -func (m *memIAVLWrapper) Close() error { - return m.base.Close() -} - -func (m *memIAVLWrapper) Read(key []byte) (data []byte, found bool, err error) { - store := m.base.GetChildStoreByName(EVMStoreName) - data = store.Get(key) - return data, data != nil, nil -} - -// RegisterHashListener reports that this DB publishes no block hashes. memIAVL's root is a -// Cosmos-layer aggregation over its per-module hashes rather than a hash the store hands out. -func (m *memIAVLWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil -} - -func (m *memIAVLWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} diff --git a/sei-db/bench/wrappers/noop_wrapper.go b/sei-db/bench/wrappers/noop_wrapper.go deleted file mode 100644 index b01c54a591..0000000000 --- a/sei-db/bench/wrappers/noop_wrapper.go +++ /dev/null @@ -1,62 +0,0 @@ -package wrappers - -import ( - "fmt" - "sync/atomic" - - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - scTypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -var _ DBWrapper = (*noOpWrapper)(nil) - -// noOpWrapper lets the benchmark measure its own overhead without DB read/write cost. -type noOpWrapper struct { - version atomic.Int64 -} - -func NewNoOpWrapper() DBWrapper { - return &noOpWrapper{} -} - -func (n *noOpWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - n.version.Store(entry.Version) - return nil -} - -func (n *noOpWrapper) Read(_ []byte) ([]byte, bool, error) { - return nil, false, nil -} - -func (n *noOpWrapper) Commit() (int64, error) { - return n.version.Load(), nil -} - -func (n *noOpWrapper) Close() error { - return nil -} - -func (n *noOpWrapper) Version() int64 { - return n.version.Load() -} - -// LoadLatest is a no-op: the tracked version already is this store's latest, since only ApplyChangeSets moves it. -func (n *noOpWrapper) LoadLatest() error { - return nil -} - -func (n *noOpWrapper) Importer(_ int64) (scTypes.Importer, error) { - return nil, fmt.Errorf("import not supported for no-op wrapper") -} - -// RegisterHashListener reports that this DB publishes no block hashes. A store that persists nothing -// hashes nothing. -func (n *noOpWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil -} - -func (n *noOpWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} diff --git a/sei-db/bench/wrappers/state_store_wrapper.go b/sei-db/bench/wrappers/state_store_wrapper.go deleted file mode 100644 index c6f8411bbb..0000000000 --- a/sei-db/bench/wrappers/state_store_wrapper.go +++ /dev/null @@ -1,78 +0,0 @@ -package wrappers - -import ( - "fmt" - "sync/atomic" - - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - dbTypes "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - scTypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -var _ DBWrapper = (*stateStoreWrapper)(nil) - -// stateStoreWrapper adapts a versioned StateStore (SS layer) to the DBWrapper -// interface used by the cryptosim benchmark. Each ApplyChangeSets call maps to -// a single ApplyChangesetAsync at the benchmark-provided version. The SS layer -// persists on every apply, so Commit is a no-op. -type stateStoreWrapper struct { - base dbTypes.StateStore - version atomic.Int64 -} - -func NewStateStoreWrapper(store dbTypes.StateStore) DBWrapper { - w := &stateStoreWrapper{ - base: store, - } - w.version.Store(store.GetLatestVersion()) - return w -} - -func (s *stateStoreWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - s.version.Store(entry.Version) - return s.base.ApplyChangesetAsync(entry.Version, entry.Changesets) -} - -func (s *stateStoreWrapper) Read(key []byte) (data []byte, found bool, err error) { - version := s.version.Load() - if version == 0 { - return nil, false, nil - } - val, err := s.base.Get(EVMStoreName, version, key) - if err != nil { - return nil, false, err - } - return val, val != nil, nil -} - -func (s *stateStoreWrapper) Commit() (int64, error) { - return s.version.Load(), nil -} - -func (s *stateStoreWrapper) Close() error { - return s.base.Close() -} - -func (s *stateStoreWrapper) Version() int64 { - return s.version.Load() -} - -func (s *stateStoreWrapper) LoadLatest() error { - return nil -} - -func (s *stateStoreWrapper) Importer(_ int64) (scTypes.Importer, error) { - return nil, fmt.Errorf("import not supported for state store wrapper") -} - -// RegisterHashListener reports that this DB publishes no block hashes. The historical state store -// computes none. -func (s *stateStoreWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil -} - -func (s *stateStoreWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} diff --git a/sei-db/bench/wrappers/state_store_wrapper_test.go b/sei-db/bench/wrappers/state_store_wrapper_test.go deleted file mode 100644 index a2596131d6..0000000000 --- a/sei-db/bench/wrappers/state_store_wrapper_test.go +++ /dev/null @@ -1,80 +0,0 @@ -package wrappers - -import ( - "bytes" - "testing" - - commonevm "github.com/sei-protocol/sei-chain/sei-db/common/keys" - "github.com/sei-protocol/sei-chain/sei-db/proto" - "github.com/stretchr/testify/require" -) - -func TestStateStoreWrapperApplyChangesetsAsyncPreservesHistoricalState(t *testing.T) { - dataDir := t.TempDir() - - store, err := openSSComposite(dataDir, *DefaultBenchStateStoreConfig()) - require.NoError(t, err) - - wrapper := NewStateStoreWrapper(store) - - keyV1AndV2 := commonevm.BuildEVMKey(commonevm.EVMKeyNonce, bytes.Repeat([]byte{0x11}, 20)) - keyV2Only := commonevm.BuildEVMKey(commonevm.EVMKeyCodeHash, bytes.Repeat([]byte{0x22}, 20)) - - require.NoError(t, wrapper.ApplyChangeSets(changelogEntry(1, []*proto.NamedChangeSet{ - { - Name: EVMStoreName, - Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{ - {Key: keyV1AndV2, Value: []byte("value-v1")}, - }}, - }, - }))) - - version, err := wrapper.Commit() - require.NoError(t, err) - require.Equal(t, int64(1), version) - - require.NoError(t, wrapper.ApplyChangeSets(changelogEntry(2, []*proto.NamedChangeSet{ - { - Name: EVMStoreName, - Changeset: proto.ChangeSet{Pairs: []*proto.KVPair{ - {Key: keyV1AndV2, Value: []byte("value-v2")}, - {Key: keyV2Only, Value: []byte("value-v2-only")}, - }}, - }, - }))) - - version, err = wrapper.Commit() - require.NoError(t, err) - require.Equal(t, int64(2), version) - - require.NoError(t, wrapper.Close()) - - reopened, err := openSSComposite(dataDir, *DefaultBenchStateStoreConfig()) - require.NoError(t, err) - t.Cleanup(func() { - require.NoError(t, reopened.Close()) - }) - - historical, err := reopened.Get(EVMStoreName, 1, keyV1AndV2) - require.NoError(t, err) - require.Equal(t, []byte("value-v1"), historical) - - missingAtV1, err := reopened.Get(EVMStoreName, 1, keyV2Only) - require.NoError(t, err) - require.Nil(t, missingAtV1) - - latest, err := reopened.Get(EVMStoreName, 2, keyV1AndV2) - require.NoError(t, err) - require.Equal(t, []byte("value-v2"), latest) - - latestOnly, err := reopened.Get(EVMStoreName, 2, keyV2Only) - require.NoError(t, err) - require.Equal(t, []byte("value-v2-only"), latestOnly) -} - -func changelogEntry(version int64, changesets []*proto.NamedChangeSet) *proto.ChangelogEntry { - return &proto.ChangelogEntry{ - Version: version, - Changesets: changesets, - } -} diff --git a/sei-db/bench/wrappers/wrappers_test.go b/sei-db/bench/wrappers/wrappers_test.go deleted file mode 100644 index 19a9ecda59..0000000000 --- a/sei-db/bench/wrappers/wrappers_test.go +++ /dev/null @@ -1,230 +0,0 @@ -package wrappers - -import ( - "testing" - - "github.com/stretchr/testify/require" - dbm "github.com/tendermint/tm-db" - - "github.com/sei-protocol/sei-chain/sei-db/common/metrics" - dbTypes "github.com/sei-protocol/sei-chain/sei-db/db_engine/types" - "github.com/sei-protocol/sei-chain/sei-db/proto" - gigatypes "github.com/sei-protocol/sei-chain/sei-db/state_db/giga/types" - scTypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -type mockDBWrapper struct { - appliedEntries []*proto.ChangelogEntry - commitVersion int64 - applyErr error - commitErr error -} - -func (m *mockDBWrapper) ApplyChangeSets(entry *proto.ChangelogEntry) error { - m.appliedEntries = append(m.appliedEntries, entry) - return m.applyErr -} - -func (m *mockDBWrapper) Read(_ []byte) ([]byte, bool, error) { - return nil, false, nil -} - -func (m *mockDBWrapper) Commit() (int64, error) { - return m.commitVersion, m.commitErr -} - -func (m *mockDBWrapper) Close() error { - return nil -} - -func (m *mockDBWrapper) Version() int64 { - return m.commitVersion -} - -func (m *mockDBWrapper) LoadLatest() error { - return nil -} - -func (m *mockDBWrapper) Importer(_ int64) (scTypes.Importer, error) { - return nil, nil -} - -func (m *mockDBWrapper) GetPhaseTimer() *metrics.PhaseTimer { - return nil -} - -type mockStateStore struct { - latestVersion int64 - asyncVersion int64 - asyncChanges []*proto.NamedChangeSet - asyncCalls int - syncVersion int64 - syncChanges []*proto.NamedChangeSet - syncCalls int -} - -func (m *mockStateStore) Get(_ string, _ int64, _ []byte) ([]byte, error) { - return nil, nil -} - -func (m *mockStateStore) Has(_ string, _ int64, _ []byte) (bool, error) { - return false, nil -} - -func (m *mockStateStore) Iterator(_ string, _ int64, _, _ []byte) (dbm.Iterator, error) { - return nil, nil -} - -func (m *mockStateStore) ReverseIterator(_ string, _ int64, _, _ []byte) (dbm.Iterator, error) { - return nil, nil -} - -func (m *mockStateStore) RawIterate(_ string, _ func([]byte, []byte, int64) bool) (bool, error) { - return false, nil -} - -func (m *mockStateStore) GetLatestVersion() int64 { - return m.latestVersion -} - -func (m *mockStateStore) SetLatestVersion(version int64) error { - m.latestVersion = version - return nil -} - -func (m *mockStateStore) GetEarliestVersion() int64 { - return 0 -} - -func (m *mockStateStore) SetEarliestVersion(_ int64, _ bool) error { - return nil -} - -func (m *mockStateStore) ApplyChangesetSync(version int64, changesets []*proto.NamedChangeSet) error { - m.syncCalls++ - m.syncVersion = version - m.syncChanges = changesets - m.latestVersion = version - return nil -} - -func (m *mockStateStore) ApplyChangesetAsync(version int64, changesets []*proto.NamedChangeSet) error { - m.asyncCalls++ - m.asyncVersion = version - m.asyncChanges = changesets - m.latestVersion = version - return nil -} - -func (m *mockStateStore) Prune(_ int64) error { - return nil -} - -func (m *mockStateStore) Import(_ int64, _ <-chan dbTypes.SnapshotNode) error { - return nil -} - -func (m *mockStateStore) Close() error { - return nil -} - -func (m *mockDBWrapper) RegisterHashListener(_ gigatypes.HashListener) (bool, error) { - return false, nil -} - -func TestCombinedWrapperApplyChangeSetsUsesAsyncSS(t *testing.T) { - sc := &mockDBWrapper{commitVersion: 7} - ss := &mockStateStore{latestVersion: 7} - wrapper := NewCombinedWrapper(sc, ss) - - changesets := []*proto.NamedChangeSet{{Name: EVMStoreName}} - entry := &proto.ChangelogEntry{ - Version: 8, - Changesets: changesets, - } - - err := wrapper.ApplyChangeSets(entry) - require.NoError(t, err) - require.Len(t, sc.appliedEntries, 1) - require.Same(t, entry, sc.appliedEntries[0]) - require.Equal(t, 1, ss.asyncCalls) - require.Equal(t, entry.Version, ss.asyncVersion) - require.Equal(t, changesets, ss.asyncChanges) - require.Zero(t, ss.syncCalls) - require.Equal(t, entry.Version, wrapper.Version()) -} - -func TestStateStoreWrapperApplyChangeSetsUsesEntryVersion(t *testing.T) { - store := &mockStateStore{latestVersion: 11} - wrapper := NewStateStoreWrapper(store) - - changesets := []*proto.NamedChangeSet{{Name: EVMStoreName}} - entry := &proto.ChangelogEntry{ - Version: 15, - Changesets: changesets, - } - - err := wrapper.ApplyChangeSets(entry) - require.NoError(t, err) - require.Equal(t, 1, store.asyncCalls) - require.Equal(t, entry.Version, store.asyncVersion) - require.Equal(t, changesets, store.asyncChanges) - require.Zero(t, store.syncCalls) - require.Equal(t, entry.Version, wrapper.Version()) - version, err := wrapper.Commit() - require.NoError(t, err) - require.Equal(t, entry.Version, version) -} - -func TestNoOpWrapperTracksVersionWithoutReadsOrWrites(t *testing.T) { - wrapper := NewNoOpWrapper() - - entry := &proto.ChangelogEntry{ - Version: 9, - Changesets: []*proto.NamedChangeSet{{Name: EVMStoreName}}, - } - - require.NoError(t, wrapper.ApplyChangeSets(entry)) - require.Equal(t, int64(9), wrapper.Version()) - - data, found, err := wrapper.Read([]byte("key")) - require.NoError(t, err) - require.Nil(t, data) - require.False(t, found) - - version, err := wrapper.Commit() - require.NoError(t, err) - require.Equal(t, int64(9), version) -} - -// runBenchmark passes nil for every backend, so each of these opened with a panic before their -// config was made optional. -func TestNewDBImplUsesDefaultConfigWhenNil(t *testing.T) { - for _, dbType := range []DBType{MemIAVL, FlatKV, SSComposite, CompositeDual_SSComposite} { - t.Run(string(dbType), func(t *testing.T) { - wrapper, err := NewDBImpl(t.Context(), dbType, t.TempDir(), nil) - require.NoError(t, err) - require.NoError(t, wrapper.Close()) - }) - } -} - -// The offload stream needs brokers that only the caller knows, so this backend has no default to -// fall back on and must say so rather than panic. -func TestNewDBImplSSHistoricalOffloadReportsMissingConfig(t *testing.T) { - wrapper, err := NewDBImpl(t.Context(), SSHistoricalOffload, t.TempDir(), nil) - require.Error(t, err) - require.Nil(t, wrapper) - require.ErrorContains(t, err, "historical offload config is required") -} - -func TestNewDBImplRejectsInvalidConfigType(t *testing.T) { - for _, dbType := range []DBType{MemIAVL, FlatKV, SSComposite, SSHistoricalOffload, CompositeDual_SSComposite} { - t.Run(string(dbType), func(t *testing.T) { - wrapper, err := NewDBImpl(t.Context(), dbType, t.TempDir(), "invalid") - require.Error(t, err) - require.Nil(t, wrapper) - require.ErrorContains(t, err, "invalid "+string(dbType)+" config type string") - }) - } -} diff --git a/sei-db/state_db/bench/README.md b/sei-db/state_db/bench/README.md deleted file mode 100644 index ceab839090..0000000000 --- a/sei-db/state_db/bench/README.md +++ /dev/null @@ -1,73 +0,0 @@ -# Benchmarks - -This package contains benchmarks for the state DB commit store. - -## Run benchmarks - -From the repo root: - -- Standard benchmarks: - - `go test ./sei-db/state_db/bench -run ^$ -bench . -benchmem` -- Run a single benchmark: - - `go test ./sei-db/state_db/bench -run ^$ -bench BenchmarkMemIAVLWriteWithDifferentBlockSize -benchmem` - -## CI - -CI does not measure anything here. `sei-db-tests.yml` runs every benchmark once -with `-benchtime=1x`, which catches a benchmark that no longer runs but produces -no usable timings. Treat a green CI run as "the benchmarks still work", and take -real numbers from a local run on a quiet machine. - -A benchmark that cannot run without external infrastructure does not belong in -this package, because that smoke step will fail on it. - -## Snapshot pre-population - -`TestScenario.SnapshotPath` loads a Cosmos SDK state sync snapshot into the -database before the timed region, so throughput is measured against a -realistically sized tree instead of an empty one. Point it at the directory -holding the numbered chunk files (`0`, `1`, `2`, …), typically -`/data/snapshots///`. - -## Define new scenarios - -Benchmarks are configured via `TestScenario`: - -- `Name`: scenario name used for the sub-benchmark -- `TotalKeys`: total number of keys to write across all blocks -- `NumBlocks`: number of blocks to commit -- `DuplicateRatio`: fraction of keys that are updates instead of inserts -- `Backend`: database backend (`wrappers.MemIAVL`, `wrappers.FlatKV`, - `wrappers.CompositeCosmos`, `wrappers.CompositeSplit`, `wrappers.CompositeDual`) -- `Distribution`: per-block key distribution function -- `SnapshotPath`: (optional) path to a state sync snapshot chunks directory; - when set, the snapshot is imported via the native `Committer.Importer` path - before the benchmark begins - -Example: - -```go -scenario := TestScenario{ - Name: "bursty_updates", - TotalKeys: 100_000, - NumBlocks: 10_000, - DuplicateRatio: 0.25, - Distribution: BurstyDistribution(1, 10, 5, 3), - Backend: wrappers.MemIAVL, -} -``` - -## Add a new distribution - -Define a new `KeyDistribution` in `helper.go`: - -```go -func MyDistribution(numBlocks, totalKeys, block int64) int64 { - // return the number of keys for this block - return totalKeys / numBlocks -} -``` - -Then set it on a `TestScenario` in `bench_sc_test.go`. A scenario that leaves -`Distribution` unset gets `EvenDistribution`, so a named distribution scenario -that forgets the field silently measures the even case. diff --git a/sei-db/state_db/bench/bench_sc_test.go b/sei-db/state_db/bench/bench_sc_test.go deleted file mode 100644 index 05a89e1fc5..0000000000 --- a/sei-db/state_db/bench/bench_sc_test.go +++ /dev/null @@ -1,275 +0,0 @@ -package bench - -import ( - "testing" - - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" -) - -// Using MemIAVL, Tests throughput with fixed total keys, -// varying keysPerBlock and numBlocks to find optimal block size. -func BenchmarkMemIAVLWriteWithDifferentBlockSize(b *testing.B) { - const totalKeys int64 = 100_000 - - scenarios := []TestScenario{ - { - Name: "1_key_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys, // 1 key per block - Backend: wrappers.MemIAVL, - }, - { - Name: "2_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 2, // 2 keys per block - Backend: wrappers.MemIAVL, - }, - { - Name: "10_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 10, // 10 keys per block - Backend: wrappers.MemIAVL, - }, - { - Name: "20_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 20, // 20 keys per block - Backend: wrappers.MemIAVL, - }, - { - Name: "100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, // 100 keys per block - Backend: wrappers.MemIAVL, - }, - { - Name: "200_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 200, // 200 keys per block - Backend: wrappers.MemIAVL, - }, - } - - for _, scenario := range scenarios { - b.Run(scenario.Name, func(b *testing.B) { - runBenchmark(b, scenario) - }) - } -} - -// Using FlatKV, Tests throughput with fixed total keys, -// varying keysPerBlock and numBlocks to find optimal block size. -func BenchmarkFlatKVWriteWithDifferentBlockSize(b *testing.B) { - // Note: FlatKV is currently behaving more slowly than expected, and so - // the total number of keys/blocks is reduced by a factor of 1000 compared to the equivalent MemIAVL benchmarks. - const totalKeys int64 = 100_000 - - scenarios := []TestScenario{ - { - Name: "100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, - Backend: wrappers.FlatKV, - }, - { - Name: "200_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 200, - Backend: wrappers.FlatKV, - }, - { - Name: "1000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 1000, - Backend: wrappers.FlatKV, - }, - { - Name: "2000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 2000, - Backend: wrappers.FlatKV, - }, - } - - for _, scenario := range scenarios { - b.Run(scenario.Name, func(b *testing.B) { - runBenchmark(b, scenario) - }) - } -} - -// Using composite backends, tests throughput with fixed total keys, -// varying keysPerBlock and numBlocks. Uses the same reduced scale as FlatKV. -func BenchmarkCompositeWriteWithDifferentBlockSize(b *testing.B) { - // Note: FlatKV is currently behaving more slowly than expected, and so - // the total number of keys/blocks is reduced by a factor of 1000 compared to the equivalent MemIAVL benchmarks. - const totalKeys int64 = 100_000 - - scenarios := []TestScenario{ - { - Name: "cosmos/100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, - Backend: wrappers.CompositeCosmos, - }, - { - Name: "split/100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, - Backend: wrappers.CompositeSplit, - }, - { - Name: "dual/100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, - Backend: wrappers.CompositeDual, - }, - { - Name: "cosmos/200_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 200, - Backend: wrappers.CompositeCosmos, - }, - { - Name: "split/200_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 200, - Backend: wrappers.CompositeSplit, - }, - { - Name: "dual/200_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 200, - Backend: wrappers.CompositeDual, - }, - { - Name: "cosmos/1000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 1000, - Backend: wrappers.CompositeCosmos, - }, - { - Name: "split/1000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 1000, - Backend: wrappers.CompositeSplit, - }, - { - Name: "dual/1000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 1000, - Backend: wrappers.CompositeDual, - }, - { - Name: "cosmos/2000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 2000, - Backend: wrappers.CompositeCosmos, - }, - { - Name: "split/2000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 2000, - Backend: wrappers.CompositeSplit, - }, - { - Name: "dual/2000_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 2000, - Backend: wrappers.CompositeDual, - }, - } - - for _, scenario := range scenarios { - b.Run(scenario.Name, func(b *testing.B) { - runBenchmark(b, scenario) - }) - } -} - -// Compares throughput across key distributions with the Composite backend. -func BenchmarkCompositeWriteWithDifferentKeyDistributions(b *testing.B) { - // Note: FlatKV is currently behaving more slowly than expected, and so - // the total number of keys/blocks is reduced by a factor of 10 compared to the equivalent MemIAVL benchmarks. - const ( - totalKeys int64 = 100_000 - numBlocks int64 = 1_000 - ) - - // The even distribution is covered by BenchmarkCompositeWriteWithDifferentBlockSize - // at 100 keys per block, which is the same scenario. - scenarios := []TestScenario{ - // Bursty distribution - { - Name: "cosmos/bursty_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeCosmos, - Distribution: BurstyDistribution(1, 10, 5, 3), - }, - { - Name: "split/bursty_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeSplit, - Distribution: BurstyDistribution(1, 10, 5, 3), - }, - { - Name: "dual/bursty_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeDual, - Distribution: BurstyDistribution(1, 10, 5, 3), - }, - // Normal distribution - { - Name: "cosmos/normal_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeCosmos, - Distribution: NormalDistribution(1, 0.2), - }, - { - Name: "split/normal_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeSplit, - Distribution: NormalDistribution(1, 0.2), - }, - { - Name: "dual/normal_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeDual, - Distribution: NormalDistribution(1, 0.2), - }, - // Ramp distribution - { - Name: "cosmos/ramp_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeCosmos, - Distribution: RampDistribution(0.5, 1.5), - }, - { - Name: "split/ramp_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeSplit, - Distribution: RampDistribution(0.5, 1.5), - }, - { - Name: "dual/ramp_distribution", - TotalKeys: totalKeys, - NumBlocks: numBlocks, - Backend: wrappers.CompositeDual, - Distribution: RampDistribution(0.5, 1.5), - }, - } - - for _, scenario := range scenarios { - b.Run(scenario.Name, func(b *testing.B) { - runBenchmark(b, scenario) - }) - } -} diff --git a/sei-db/state_db/bench/bench_ss_test.go b/sei-db/state_db/bench/bench_ss_test.go deleted file mode 100644 index 299c6dfc4d..0000000000 --- a/sei-db/state_db/bench/bench_ss_test.go +++ /dev/null @@ -1,45 +0,0 @@ -package bench - -import ( - "testing" - - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" -) - -func BenchmarkSSCompositeWrite(b *testing.B) { - const totalKeys int64 = 10_000 - - scenarios := []TestScenario{ - { - Name: "100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, - Backend: wrappers.SSComposite, - }, - } - - for _, scenario := range scenarios { - b.Run(scenario.Name, func(b *testing.B) { - runBenchmark(b, scenario) - }) - } -} - -func BenchmarkCombinedCompositeDualSSCompositeWrite(b *testing.B) { - const totalKeys int64 = 10_000 - - scenarios := []TestScenario{ - { - Name: "100_keys_per_block", - TotalKeys: totalKeys, - NumBlocks: totalKeys / 100, - Backend: wrappers.CompositeDual_SSComposite, - }, - } - - for _, scenario := range scenarios { - b.Run(scenario.Name, func(b *testing.B) { - runBenchmark(b, scenario) - }) - } -} diff --git a/sei-db/state_db/bench/helper.go b/sei-db/state_db/bench/helper.go deleted file mode 100644 index d880f05ce8..0000000000 --- a/sei-db/state_db/bench/helper.go +++ /dev/null @@ -1,375 +0,0 @@ -package bench - -import ( - "crypto/rand" - "crypto/sha256" - "encoding/binary" - "fmt" - "io" - "math" - mrand "math/rand" - "os" - "path/filepath" - "strconv" - "testing" - "time" - - "github.com/stretchr/testify/require" - - "github.com/sei-protocol/sei-chain/sei-cosmos/snapshots" - snapshottypes "github.com/sei-protocol/sei-chain/sei-cosmos/snapshots/types" - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" - commonevm "github.com/sei-protocol/sei-chain/sei-db/common/keys" - "github.com/sei-protocol/sei-chain/sei-db/proto" - sctypes "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/types" -) - -const ( - // EVMStoreName simulates the EVM module store - EVMStoreName = commonevm.EVMStoreKey - - // KeySize EVM storage key: 0x03 prefix + 20-byte address + 32-byte slot = 53 bytes - KeySize = 53 - ValueSize = 32 -) - -// TestScenario bundles benchmark parameters and distribution. -type TestScenario struct { - Name string - TotalKeys int64 - NumBlocks int64 - DuplicateRatio float64 // 0.0 = all inserts, 1.0 = all updates - // The database backend to use for the benchmark. - Backend wrappers.DBType - Distribution KeyDistribution - - // SnapshotPath, when set, points to a state sync snapshot chunks directory - // (e.g. "/data/snapshots///") containing numbered - // chunk files (0, 1, 2, ...). Before the benchmark begins, the snapshot is - // imported into the database via the native Committer.Importer path as a - // preparation stage. - SnapshotPath string -} - -// KeyDistribution defines how many keys to generate per block. -type KeyDistribution func(numBlocks, totalKeys, block int64) int64 - -// EvenDistribution generates same number of keys on each block. -func EvenDistribution(numBlocks, totalKeys, _ int64) int64 { - if numBlocks <= 0 || totalKeys < numBlocks { - return 0 - } - return totalKeys / numBlocks -} - -// BurstyDistribution emits periodic bursts with optional jitter. -// Example: base=100 keys/block, burstEvery=5, burstMultiplier=3 => -// blocks 0,5,10... emit 300 keys; other blocks emit 100 keys (then +/- jitter). -func BurstyDistribution(seed int64, burstEvery, burstMultiplier, maxJitter int64) KeyDistribution { - rng := mrand.New(mrand.NewSource(seed)) - return func(numBlocks, totalKeys, block int64) int64 { - if numBlocks <= 0 { - return 0 - } - keysPerBlock := totalKeys / numBlocks - count := keysPerBlock - if burstEvery > 0 && block%burstEvery == 0 { - count *= burstMultiplier - } - if maxJitter > 0 { - count += rng.Int63n(2*maxJitter+1) - maxJitter - } - if count < 0 { - return 0 - } - return count - } -} - -// NormalDistribution samples keys per block from a normal distribution. -// Example: totalKeys=1000, numBlocks=10, stddevFactor=0.2 => -// mean=100 keys, stddev=20 keys -func NormalDistribution(seed int64, stddevFactor float64) KeyDistribution { - rng := mrand.New(mrand.NewSource(seed)) - return func(numBlocks, totalKeys, _ int64) int64 { - if numBlocks <= 0 { - return 0 - } - mean := float64(totalKeys) / float64(numBlocks) - stddev := mean * stddevFactor - if stddev <= 0 { - return int64(mean) - } - count := int64(mean + rng.NormFloat64()*stddev) - if count < 0 { - return 0 - } - return count - } -} - -// RampDistribution linearly ramps keysPerBlock by a factor over the run. -// Example: totalKeys=1000, numBlocks=10, startFactor=0.5, endFactor=1.5 => -// per-block base=100; block 0 ~50 keys, block 9 ~150 keys (linearly interpolated). -func RampDistribution(startFactor, endFactor float64) KeyDistribution { - return func(numBlocks, totalKeys, block int64) int64 { - if numBlocks <= 1 { - return int64(float64(totalKeys) * endFactor) - } - keysPerBlock := totalKeys / numBlocks - t := float64(block) / float64(numBlocks-1) - factor := startFactor + t*(endFactor-startFactor) - count := int64(float64(keysPerBlock) * factor) - if count < 0 { - return 0 - } - return count - } -} - -// startChangesetGenerator streams per-block changesets based on the scenario distribution. -func startChangesetGenerator(scenario TestScenario) <-chan *proto.NamedChangeSet { - if scenario.Distribution == nil { - scenario.Distribution = EvenDistribution - } - duplicateRatio := scenario.DuplicateRatio - if duplicateRatio < 0 { - duplicateRatio = 0 - } - if duplicateRatio > 1 { - duplicateRatio = 1 - } - rng := mrand.New(mrand.NewSource(1)) - out := make(chan *proto.NamedChangeSet) - go func() { - defer close(out) - var uniqueCounter int64 - for i := range scenario.NumBlocks { - numKeysInBlock := scenario.Distribution(scenario.NumBlocks, scenario.TotalKeys, i) - if numKeysInBlock < 0 { - numKeysInBlock = 0 - } - kvPairs := make([]*proto.KVPair, int(numKeysInBlock)) - duplicateCount := int64(float64(numKeysInBlock) * duplicateRatio) - for j := range kvPairs { - var keyIndex int64 - if int64(j) < duplicateCount && uniqueCounter > 0 { - keyIndex = rng.Int63n(uniqueCounter) - } else { - keyIndex = uniqueCounter - uniqueCounter++ - } - key := keyFromIndex(keyIndex) - val := make([]byte, ValueSize) - if _, err := rand.Read(val); err != nil { - panic(fmt.Sprintf("failed to generate random value: %v", err)) - } - kvPairs[j] = &proto.KVPair{Key: key, Value: val} - } - cs := &proto.NamedChangeSet{ - Name: EVMStoreName, - Changeset: proto.ChangeSet{Pairs: kvPairs}, - } - out <- cs - } - }() - return out -} - -func keyFromIndex(index int64) []byte { - key := make([]byte, KeySize) - key[0] = 0x03 - var input [9]byte - if index < 0 { - panic(fmt.Sprintf("negative key index: %d", index)) - } - binary.LittleEndian.PutUint64(input[1:], uint64(index)) //nolint:gosec // index validated non-negative above - sum1 := sha256.Sum256(input[:]) - input[0] = 1 //nolint:gosec - sum2 := sha256.Sum256(input[:]) - copy(key[1:], sum1[:]) - copy(key[1+len(sum1):], sum2[:len(key)-1-len(sum1)]) - return key -} - -// parseSnapshotHeight extracts the block height from a state sync snapshot -// chunks directory path. The expected layout is ///, -// so the height is the parent of the format directory. -func parseSnapshotHeight(chunksDir string) (int64, error) { - heightStr := filepath.Base(filepath.Dir(filepath.Clean(chunksDir))) - h, err := strconv.ParseInt(heightStr, 10, 64) - if err != nil { - return 0, fmt.Errorf("parse snapshot height from path %q: %w", chunksDir, err) - } - if h <= 0 || h > math.MaxUint32 { - return 0, fmt.Errorf("snapshot height %d out of range", h) - } - return h, nil -} - -// openSnapshotStream opens the numbered chunk files in chunksDir and returns a -// StreamReader that decompresses and demuxes the protobuf item stream. -func openSnapshotStream(chunksDir string) (*snapshots.StreamReader, error) { - if _, err := os.Stat(filepath.Join(chunksDir, "0")); err != nil { - return nil, fmt.Errorf("no chunk files found in %s: %w", chunksDir, err) - } - - chunks := make(chan io.ReadCloser) - go func() { - defer close(chunks) - for i := 0; ; i++ { - path := filepath.Join(chunksDir, strconv.Itoa(i)) - f, err := os.Open(filepath.Clean(path)) - if err != nil { - if os.IsNotExist(err) { - return - } - pr, pw := io.Pipe() - _ = pw.CloseWithError(fmt.Errorf("open chunk %d: %w", i, err)) - chunks <- pr - return - } - chunks <- f - } - }() - - return snapshots.NewStreamReader(chunks) -} - -// importSnapshot reads a state sync snapshot from chunksDir and feeds every -// item through the given Importer (AddModule / AddNode). This is the same -// import path used by the real state sync restore logic. -// Returns the total number of leaf keys imported. -func importSnapshot(chunksDir string, importer sctypes.Importer) error { - streamReader, err := openSnapshotStream(chunksDir) - if err != nil { - return fmt.Errorf("create stream reader: %w", err) - } - defer func() { - _ = streamReader.Close() - }() - - var ( - totalKeys int64 - startTime = time.Now() - currModule = "" - ) - - for { - var item snapshottypes.SnapshotItem - err := streamReader.ReadMsg(&item) - if err == io.EOF { - break - } - if err != nil { - return fmt.Errorf("read snapshot item: %w", err) - } - - switch i := item.Item.(type) { - case *snapshottypes.SnapshotItem_Store: - currModule = i.Store.Name - if currModule == commonevm.EVMStoreKey { - if err := importer.AddModule(i.Store.Name); err != nil { - return fmt.Errorf("add module %s: %w", i.Store.Name, err) - } - fmt.Printf("[Snapshot] Importing store: %s\n", i.Store.Name) - } else { - fmt.Printf("[Snapshot] Skipping store: %s\n", i.Store.Name) - } - case *snapshottypes.SnapshotItem_IAVL: - if currModule != commonevm.EVMStoreKey { - continue - } - if i.IAVL.Height > math.MaxInt8 { - return fmt.Errorf("node height %d exceeds int8", i.IAVL.Height) - } - node := &sctypes.SnapshotNode{ - Key: i.IAVL.Key, - Value: i.IAVL.Value, - Height: int8(i.IAVL.Height), //nolint:gosec - Version: i.IAVL.Version, - } - if node.Height == 0 && node.Value == nil { - node.Value = []byte{} - } - importer.AddNode(node) - if node.Height == 0 { - totalKeys++ - if totalKeys%1_000_000 == 0 { - elapsed := time.Since(startTime).Seconds() - fmt.Printf("[Snapshot] keys=%d, keys/sec=%.0f, elapsed=%.2fs\n", - totalKeys, float64(totalKeys)/elapsed, elapsed) - } - } - default: - break - } - } - - elapsed := time.Since(startTime).Seconds() - fmt.Printf("[Snapshot] Import Done: keys=%d, keys/sec=%.0f, elapsed=%.2fs\n", - totalKeys, float64(totalKeys)/elapsed, elapsed) - - return importer.Close() -} - -// runBenchmark runs the scenario against a freshly opened backend, reporting -// keys/sec and elapsed seconds as custom benchmark metrics. -func runBenchmark(b *testing.B, scenario TestScenario) { - if scenario.Distribution == nil { - scenario.Distribution = EvenDistribution - } - - b.ResetTimer() - b.ReportAllocs() - - for range b.N { - func() { - dbDir := b.TempDir() - b.StopTimer() - cs, err := wrappers.NewDBImpl(b.Context(), scenario.Backend, dbDir, nil) - require.NoError(b, err) - - // Load snapshot if available - if scenario.SnapshotPath != "" { - snapshotHeight, err := parseSnapshotHeight(scenario.SnapshotPath) - require.NoError(b, err) - importer, err := cs.Importer(snapshotHeight) - require.NoError(b, err) - err = importSnapshot(scenario.SnapshotPath, importer) - require.NoError(b, err) - err = cs.LoadLatest() - require.NoError(b, err) - } - changesetChannel := startChangesetGenerator(scenario) - - baseVersion := cs.Version() - b.StartTimer() - fmt.Printf("Opening DB with base version %d\n", baseVersion) - - for block := int64(1); block < scenario.NumBlocks; block++ { - changeset, ok := <-changesetChannel - if !ok { - break - } - entry := &proto.ChangelogEntry{ - Version: baseVersion + block, - Changesets: []*proto.NamedChangeSet{changeset}, - } - err := cs.ApplyChangeSets(entry) - require.NoError(b, err) - version, err := cs.Commit() - require.NoError(b, err) - require.Equal(b, baseVersion+block, version) - } - closeErr := cs.Close() // close to make sure all data got flushed - require.NoError(b, closeErr) - - b.StopTimer() - - elapsed := b.Elapsed().Seconds() - b.ReportMetric(float64(scenario.TotalKeys)/elapsed, "keys/sec") - b.ReportMetric(elapsed, "seconds") - }() - } -} diff --git a/sei-db/state_db/bench/writeset.go b/sei-db/state_db/bench/writeset.go deleted file mode 100644 index 0bf2a9a063..0000000000 --- a/sei-db/state_db/bench/writeset.go +++ /dev/null @@ -1,280 +0,0 @@ -package bench - -import ( - "context" - "encoding/hex" - "encoding/json" - "fmt" - "os" - "strings" - "time" - - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" - "github.com/sei-protocol/sei-chain/sei-db/common/keys" - "github.com/sei-protocol/sei-chain/sei-db/proto" - flatkvConfig "github.com/sei-protocol/sei-chain/sei-db/state_db/sc/flatkv/config" -) - -// This file implements the write-set replay adapter: it parses a captured -// write-set file (typically derived from a debug_traceCall prestateTracer -// diff) into per-block changesets and replays them through a storage-engine -// wrapper, timing ApplyChangeSets and Commit separately. -// -// v1 scope: EVM-module keys only (storage/code/nonce/codehash/raw). Bank -// (balance) changesets require the bank store layout and are deliberately -// out of scope; see the gas-repricing storage doc. - -// WriteSetEntryKind enumerates the supported key kinds in a write-set file. -const ( - WriteKindStorage = "storage" // requires address + slot - WriteKindCode = "code" // requires address - WriteKindNonce = "nonce" // requires address - WriteKindCodeHash = "codehash" // requires address - WriteKindRaw = "raw" // requires key (full store key, hex) -) - -// WriteSetEntry is one captured write. Hex fields accept an optional 0x prefix. -type WriteSetEntry struct { - // Kind is one of the WriteKind* constants. - Kind string `json:"kind"` - // Address is the 20-byte EVM address (storage/code/nonce/codehash kinds). - Address string `json:"address,omitempty"` - // Slot is the 32-byte storage slot (storage kind only). - Slot string `json:"slot,omitempty"` - // Key is the full raw store key (raw kind only). - Key string `json:"key,omitempty"` - // Value is the new value. Ignored when Delete is true. - Value string `json:"value,omitempty"` - // Delete marks a deletion instead of a write. - Delete bool `json:"delete,omitempty"` -} - -// WriteSetBlock groups the writes that commit together as one block. -type WriteSetBlock struct { - Writes []WriteSetEntry `json:"writes"` -} - -// WriteSet is the top-level write-set file format. -type WriteSet struct { - // Module is the store the writes belong to. Only "evm" is supported in v1; - // empty defaults to "evm". - Module string `json:"module,omitempty"` - Blocks []WriteSetBlock `json:"blocks"` -} - -// LoadWriteSet reads and validates a write-set file. -func LoadWriteSet(path string) (*WriteSet, error) { - data, err := os.ReadFile(path) //nolint:gosec // benchmark input path supplied by the operator - if err != nil { - return nil, fmt.Errorf("read write-set file: %w", err) - } - var ws WriteSet - if err := json.Unmarshal(data, &ws); err != nil { - return nil, fmt.Errorf("parse write-set file: %w", err) - } - if err := ws.Validate(); err != nil { - return nil, err - } - return &ws, nil -} - -// Validate checks module support and per-entry field consistency. -func (ws *WriteSet) Validate() error { - if ws.Module != "" && ws.Module != keys.EVMStoreKey { - return fmt.Errorf("unsupported module %q: v1 replay supports only %q", ws.Module, keys.EVMStoreKey) - } - if len(ws.Blocks) == 0 { - return fmt.Errorf("write set has no blocks") - } - for bi, block := range ws.Blocks { - for wi, w := range block.Writes { - if _, err := buildEntryKey(w); err != nil { - return fmt.Errorf("block %d write %d: %w", bi, wi, err) - } - if !w.Delete { - if _, err := decodeEntryValue(w); err != nil { - return fmt.Errorf("block %d write %d: %w", bi, wi, err) - } - } - } - } - return nil -} - -// TotalKeys returns the total number of writes across all blocks. -func (ws *WriteSet) TotalKeys() int { - total := 0 - for _, b := range ws.Blocks { - total += len(b.Writes) - } - return total -} - -// BlockChangesets converts one block into the NamedChangeSet slice consumed by -// DBWrapper.ApplyChangeSets. -func (ws *WriteSet) BlockChangesets(blockIdx int) ([]*proto.NamedChangeSet, error) { - block := ws.Blocks[blockIdx] - pairs := make([]*proto.KVPair, 0, len(block.Writes)) - for wi, w := range block.Writes { - key, err := buildEntryKey(w) - if err != nil { - return nil, fmt.Errorf("block %d write %d: %w", blockIdx, wi, err) - } - pair := &proto.KVPair{Key: key, Delete: w.Delete} - if !w.Delete { - value, err := decodeEntryValue(w) - if err != nil { - return nil, fmt.Errorf("block %d write %d: %w", blockIdx, wi, err) - } - pair.Value = value - } - pairs = append(pairs, pair) - } - return []*proto.NamedChangeSet{{ - Name: keys.EVMStoreKey, - Changeset: proto.ChangeSet{Pairs: pairs}, - }}, nil -} - -// buildEntryKey builds the raw store key for a write-set entry. -func buildEntryKey(w WriteSetEntry) ([]byte, error) { - switch w.Kind { - case WriteKindStorage: - addr, err := decodeHexField("address", w.Address, keys.AddressLen) - if err != nil { - return nil, err - } - slot, err := decodeHexField("slot", w.Slot, 32) - if err != nil { - return nil, err - } - return keys.BuildEVMKey(keys.EVMKeyStorage, append(addr, slot...)), nil - case WriteKindCode, WriteKindNonce, WriteKindCodeHash: - addr, err := decodeHexField("address", w.Address, keys.AddressLen) - if err != nil { - return nil, err - } - kind := map[string]keys.EVMKeyKind{ - WriteKindCode: keys.EVMKeyCode, - WriteKindNonce: keys.EVMKeyNonce, - WriteKindCodeHash: keys.EVMKeyCodeHash, - }[w.Kind] - return keys.BuildEVMKey(kind, addr), nil - case WriteKindRaw: - key, err := decodeHexField("key", w.Key, 0) - if err != nil { - return nil, err - } - if len(key) == 0 { - return nil, fmt.Errorf("raw write has empty key") - } - return key, nil - default: - return nil, fmt.Errorf("unknown write kind %q", w.Kind) - } -} - -// valueLenForKind returns the exact byte length a kind's value must have, or 0 -// when the length is unconstrained (code and raw). The fixed widths mirror the -// FlatKV apply path (vtype.ParseNonce/ParseCodeHash/ParseStorageValue), so a -// wrong-length value in a hand-authored write-set file is rejected up front by -// Validate rather than only failing later inside ApplyChangeSets on FlatKV -// (memiavl stores raw bytes and would silently accept it, breaking the -// same-write-set-across-backends premise of the benchmark). -// -// Raw entries are the escape hatch this check cannot cover: their keys are -// opaque here, so hand-authored raw entries must target key families FlatKV -// does not width-check (legacy prefixes such as 0x09 codesize). A raw key -// aliasing an optimized family (e.g. 0x0a nonce) with a wrong-width value -// passes Validate, replays on memiavl, and hard-fails on FlatKV. -func valueLenForKind(kind string) int { - switch kind { - case WriteKindNonce: - return 8 - case WriteKindStorage, WriteKindCodeHash: - return 32 - default: // WriteKindCode, WriteKindRaw: unconstrained - return 0 - } -} - -// decodeEntryValue decodes a write entry's value, enforcing the fixed width its -// kind requires (see valueLenForKind). -func decodeEntryValue(w WriteSetEntry) ([]byte, error) { - return decodeHexField("value", w.Value, valueLenForKind(w.Kind)) -} - -// decodeHexField decodes a hex field, tolerating a 0x prefix. wantLen of 0 -// disables the length check. An empty string decodes to nil. -func decodeHexField(name, value string, wantLen int) ([]byte, error) { - trimmed := strings.TrimPrefix(value, "0x") - decoded, err := hex.DecodeString(trimmed) - if err != nil { - return nil, fmt.Errorf("field %s: invalid hex %q: %w", name, value, err) - } - if wantLen > 0 && len(decoded) != wantLen { - return nil, fmt.Errorf("field %s: expected %d bytes, got %d", name, wantLen, len(decoded)) - } - return decoded, nil -} - -// OpenReplayWrapper opens a fresh DBWrapper for a replay run, supplying the -// explicit default config that the FlatKV wrapper factory requires. -// -// memiavl is opened with AsyncCommitBuffer=0 (synchronous WAL write) rather -// than the shared bench default of 10. With the async buffer, memiavl's -// Commit() returns once the WAL entry is enqueued, while FlatKV's Commit() -// waits for its WAL write — the reported commit_ns/key would compare enqueue -// latency against write latency. Neither backend fsyncs, so with a -// synchronous WAL write on both sides the durability semantics match. -func OpenReplayWrapper(ctx context.Context, backend wrappers.DBType, dbDir string) (wrappers.DBWrapper, error) { - var dbConfig any - switch backend { - case wrappers.FlatKV: - dbConfig = flatkvConfig.DefaultConfig() - case wrappers.MemIAVL: - cfg := wrappers.DefaultBenchMemIAVLConfig() - cfg.AsyncCommitBuffer = 0 - dbConfig = &cfg - } - return wrappers.NewDBImpl(ctx, backend, dbDir, dbConfig) -} - -// ReplayResult reports a replay run with apply and commit timed separately. -type ReplayResult struct { - Blocks int - Keys int - ApplyDuration time.Duration - CommitDuration time.Duration -} - -// ReplayWriteSet replays the write set through the wrapper, one block per -// version, timing ApplyChangeSets and Commit separately. The wrapper must be -// freshly opened (or snapshot-loaded); replay starts at wrapper.Version()+1. -func ReplayWriteSet(wrapper wrappers.DBWrapper, ws *WriteSet) (ReplayResult, error) { - result := ReplayResult{Blocks: len(ws.Blocks), Keys: ws.TotalKeys()} - baseVersion := wrapper.Version() - for i := range ws.Blocks { - changesets, err := ws.BlockChangesets(i) - if err != nil { - return result, err - } - entry := &proto.ChangelogEntry{ - Version: baseVersion + int64(i) + 1, - Changesets: changesets, - } - - applyStart := time.Now() - if err := wrapper.ApplyChangeSets(entry); err != nil { - return result, fmt.Errorf("apply block %d: %w", i, err) - } - result.ApplyDuration += time.Since(applyStart) - - commitStart := time.Now() - if _, err := wrapper.Commit(); err != nil { - return result, fmt.Errorf("commit block %d: %w", i, err) - } - result.CommitDuration += time.Since(commitStart) - } - return result, nil -} diff --git a/sei-db/state_db/bench/writeset_bench_test.go b/sei-db/state_db/bench/writeset_bench_test.go deleted file mode 100644 index 634a95e038..0000000000 --- a/sei-db/state_db/bench/writeset_bench_test.go +++ /dev/null @@ -1,102 +0,0 @@ -package bench - -import ( - "fmt" - "os" - "testing" - "time" - - "github.com/stretchr/testify/require" - - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" -) - -// BenchmarkWriteSetReplay replays a captured write-set file against both -// storage backends, timing ApplyChangeSets and Commit separately. -// -// Inputs (environment variables): -// -// TRACE_PATH path to a prestateTracer diffMode JSON file (the raw -// {"pre","post"} result or a whole JSON-RPC response), -// converted on the fly; takes precedence over WRITESET_PATH. -// WRITESET_PATH path to a write-set JSON file (see writeset.go). Raw -// tracer output is not accepted here; use TRACE_PATH. -// SNAPSHOT_PATH optional state sync snapshot chunks directory imported -// before the timed region (same as the other benchmarks). -// -// Example: -// -// TRACE_PATH=/tmp/sstore_trace.json go test ./sei-db/state_db/bench \ -// -run '^$' -bench '^BenchmarkWriteSetReplay$' -benchtime=5x -func BenchmarkWriteSetReplay(b *testing.B) { - ws := loadBenchWriteSet(b) - - for _, backend := range []wrappers.DBType{wrappers.MemIAVL, wrappers.FlatKV} { - b.Run(string(backend), func(b *testing.B) { - // Accumulate across iterations and report once: b.ReportMetric keeps - // only the last value for a given unit, so reporting inside the loop - // would surface a single iteration instead of an average over b.N. - var totalApply, totalCommit time.Duration - var totalKeys int - for range b.N { - result := runWriteSetReplay(b, backend, ws) - totalApply += result.ApplyDuration - totalCommit += result.CommitDuration - totalKeys += result.Keys - } - if totalKeys == 0 { - return - } - keys := float64(totalKeys) - b.ReportMetric(totalApply.Seconds()/keys*1e9, "apply_ns/key") - b.ReportMetric(totalCommit.Seconds()/keys*1e9, "commit_ns/key") - }) - } -} - -func loadBenchWriteSet(b *testing.B) *WriteSet { - if tracePath := os.Getenv("TRACE_PATH"); tracePath != "" { - converted, err := ConvertPrestateDiffFile(tracePath) - require.NoError(b, err) - if converted.SkippedBalanceChanges > 0 { - b.Logf("skipped %d balance change(s): bank-module replay is out of scope", - converted.SkippedBalanceChanges) - } - return converted.WriteSet - } - if wsPath := os.Getenv("WRITESET_PATH"); wsPath != "" { - ws, err := LoadWriteSet(wsPath) - require.NoError(b, err) - return ws - } - b.Skip("set TRACE_PATH or WRITESET_PATH to run the write-set replay benchmark") - return nil -} - -func runWriteSetReplay(b *testing.B, backend wrappers.DBType, ws *WriteSet) ReplayResult { - b.StopTimer() - dbDir := b.TempDir() - wrapper, err := OpenReplayWrapper(b.Context(), backend, dbDir) - require.NoError(b, err) - defer func() { - require.NoError(b, wrapper.Close()) - }() - - if snapshotPath := os.Getenv("SNAPSHOT_PATH"); snapshotPath != "" { - snapshotHeight, err := parseSnapshotHeight(snapshotPath) - require.NoError(b, err) - importer, err := wrapper.Importer(snapshotHeight) - require.NoError(b, err) - require.NoError(b, importSnapshot(snapshotPath, importer)) - require.NoError(b, wrapper.LoadLatest()) - } - - b.StartTimer() - result, err := ReplayWriteSet(wrapper, ws) - b.StopTimer() - require.NoError(b, err) - - fmt.Printf("[Replay %s] blocks=%d keys=%d apply=%s commit=%s\n", - backend, result.Blocks, result.Keys, result.ApplyDuration, result.CommitDuration) - return result -} diff --git a/sei-db/state_db/bench/writeset_convert.go b/sei-db/state_db/bench/writeset_convert.go deleted file mode 100644 index 1323b5e730..0000000000 --- a/sei-db/state_db/bench/writeset_convert.go +++ /dev/null @@ -1,262 +0,0 @@ -package bench - -import ( - "encoding/binary" - "encoding/hex" - "encoding/json" - "fmt" - "os" - "sort" - "strings" - - ethcrypto "github.com/ethereum/go-ethereum/crypto" -) - -// This file converts a debug_traceCall prestateTracer diffMode result -// ({"pre": {...}, "post": {...}}) into a WriteSet for replay. -// -// Mapping rules (v1): -// - storage slots present in post -> storage write with the post value -// - storage slots present in pre, not post -> storage delete (slot zeroed) -// - nonce changed -> nonce write (8-byte big-endian) -// - code changed -> code write + codehash write -// (keccak256 of the code) + codesize raw write (0x09||addr, 8-byte length), -// mirroring x/evm's deploy path -// - account present in pre, absent in post (SELFDESTRUCT) -> deletes for the -// account's nonce, code, codehash, and codesize keys, plus the storage -// deletes from the rule above. Only slots present in pre can be deleted; -// a real account wipe also removes slots the trace never touched, so -// self-destruct replays are a lower bound on the true delete volume. -// - balance changes are bank-module writes and are NOT converted; they are -// counted in SkippedBalanceChanges so callers can see what was dropped -// (a removed account's balance zeroing is counted the same way) -// -// Known v1 fidelity gap: deploying to a previously-unassociated address also -// writes the Sei<->EVM address mapping (two raw keys, 0x01||evm and 0x02||sei) -// and creates a Sei account, via x/evm's SetCode -> SetAddressMapping path. Those -// writes are NOT emitted here because a prestate trace does not reveal the prior -// association state, so we cannot tell whether the mapping write actually fired -// (the same reason balance changes are skipped). New-contract deploy replays -// therefore slightly undercount apply/commit cost; emitting them conditionally -// is left to a future revision. -// -// Addresses and slots are emitted in sorted order so conversion output is -// deterministic for a given trace. - -// codeSizeKeyPrefix mirrors x/evm/types.CodeSizeKeyPrefix. It routes to the -// legacy key family, which the sei-db keys package intentionally does not -// enumerate, so the byte is duplicated here the same way keys/evm.go -// duplicates the other prefixes. -var codeSizeKeyPrefix = []byte{0x09} - -// prestateAccount is one account entry in a prestateTracer result. -type prestateAccount struct { - Balance string `json:"balance,omitempty"` - Nonce *uint64 `json:"nonce,omitempty"` - Code string `json:"code,omitempty"` - Storage map[string]string `json:"storage,omitempty"` -} - -// prestateDiff is the diffMode payload of a prestateTracer trace. -type prestateDiff struct { - Pre map[string]prestateAccount `json:"pre"` - Post map[string]prestateAccount `json:"post"` -} - -// ConvertResult carries the converted write set plus conversion statistics. -type ConvertResult struct { - WriteSet *WriteSet - // SkippedBalanceChanges counts balance changes that were not converted - // because balances live in the bank module (out of v1 scope): accounts - // with a post-state balance, plus removed accounts whose pre-state - // balance was zeroed. - SkippedBalanceChanges int -} - -// ConvertPrestateDiffFile reads a prestateTracer diffMode JSON file (either -// the raw {"pre","post"} object or a JSON-RPC response with that object under -// "result") and converts it into a single-block WriteSet. -func ConvertPrestateDiffFile(path string) (*ConvertResult, error) { - data, err := os.ReadFile(path) //nolint:gosec // benchmark input path supplied by the operator - if err != nil { - return nil, fmt.Errorf("read trace file: %w", err) - } - return ConvertPrestateDiff(data) -} - -// ConvertPrestateDiff converts prestateTracer diffMode JSON bytes into a -// single-block WriteSet. -func ConvertPrestateDiff(data []byte) (*ConvertResult, error) { - var rpcEnvelope struct { - Result json.RawMessage `json:"result"` - } - if err := json.Unmarshal(data, &rpcEnvelope); err == nil && len(rpcEnvelope.Result) > 0 { - data = rpcEnvelope.Result - } - - var diff prestateDiff - if err := json.Unmarshal(data, &diff); err != nil { - return nil, fmt.Errorf("parse prestate diff: %w", err) - } - if diff.Post == nil { - return nil, fmt.Errorf("trace has no post state; was the tracer run with diffMode=true?") - } - - result := &ConvertResult{} - var writes []WriteSetEntry - - for _, addr := range sortedKeys(diff.Post) { - post := diff.Post[addr] - pre := diff.Pre[addr] - - writes = append(writes, convertStorage(addr, pre, post)...) - - if post.Nonce != nil { - nonce := make([]byte, 8) - binary.BigEndian.PutUint64(nonce, *post.Nonce) - writes = append(writes, WriteSetEntry{ - Kind: WriteKindNonce, - Address: addr, - Value: hex.EncodeToString(nonce), - }) - } - - if post.Code != "" && post.Code != pre.Code { - codeWrites, err := convertCode(addr, post.Code) - if err != nil { - return nil, err - } - writes = append(writes, codeWrites...) - } - - if post.Balance != "" { - result.SkippedBalanceChanges++ - } - } - - // Slots that were zeroed appear in pre but not post: emit deletes. Membership - // is tested on the normalized (padded) slot key because the write pass - // normalizes the post slot the same way; comparing raw hex could miss a match - // when pre and post encode the same slot differently (padded vs unpadded, 0x - // prefix, case), emitting a spurious delete that clobbers the write — deletes - // are appended after writes, and both engines apply last-write-wins per key. - for _, addr := range sortedKeys(diff.Pre) { - pre := diff.Pre[addr] - post, inPost := diff.Post[addr] - postSlots := make(map[string]struct{}, len(post.Storage)) - for slot := range post.Storage { - postSlots[padTo32(slot)] = struct{}{} - } - for _, slot := range sortedKeys(pre.Storage) { - if _, stillSet := postSlots[padTo32(slot)]; !stillSet { - writes = append(writes, WriteSetEntry{ - Kind: WriteKindStorage, - Address: addr, - Slot: padTo32(slot), - Delete: true, - }) - } - } - if !inPost { - removalWrites, err := convertAccountRemoval(addr, pre) - if err != nil { - return nil, err - } - writes = append(writes, removalWrites...) - if pre.Balance != "" { - result.SkippedBalanceChanges++ - } - } - } - - if len(writes) == 0 { - return nil, fmt.Errorf("trace produced no convertible writes") - } - result.WriteSet = &WriteSet{Blocks: []WriteSetBlock{{Writes: writes}}} - if err := result.WriteSet.Validate(); err != nil { - return nil, fmt.Errorf("converted write set is invalid: %w", err) - } - return result, nil -} - -// convertStorage emits writes for every slot present in the post state. -// Slots and values are padded to 32 bytes: x/evm writes fixed 32-byte values, -// and some tracers emit unpadded hex. -func convertStorage(addr string, _, post prestateAccount) []WriteSetEntry { - writes := make([]WriteSetEntry, 0, len(post.Storage)) - for _, slot := range sortedKeys(post.Storage) { - writes = append(writes, WriteSetEntry{ - Kind: WriteKindStorage, - Address: addr, - Slot: padTo32(slot), - Value: padTo32(post.Storage[slot]), - }) - } - return writes -} - -// convertCode emits the code, codehash, and codesize writes that a contract -// deployment produces. It deliberately omits the Sei<->EVM address-mapping and -// account-creation writes that SetCode also performs for a previously -// unassociated address; see the fidelity-gap note in the file header. -func convertCode(addr, codeHex string) ([]WriteSetEntry, error) { - code, err := hex.DecodeString(strings.TrimPrefix(codeHex, "0x")) - if err != nil { - return nil, fmt.Errorf("address %s: invalid code hex: %w", addr, err) - } - addrBytes, err := decodeHexField("address", addr, 20) - if err != nil { - return nil, err - } - size := make([]byte, 8) - binary.BigEndian.PutUint64(size, uint64(len(code))) - return []WriteSetEntry{ - {Kind: WriteKindCode, Address: addr, Value: hex.EncodeToString(code)}, - {Kind: WriteKindCodeHash, Address: addr, Value: hex.EncodeToString(ethcrypto.Keccak256(code))}, - {Kind: WriteKindRaw, Key: hex.EncodeToString(append(codeSizeKeyPrefix, addrBytes...)), Value: hex.EncodeToString(size)}, - }, nil -} - -// convertAccountRemoval emits the account-level deletes for an address that is -// present in pre but absent from post (the diffMode shape a SELFDESTRUCT -// produces): nonce, and — when the account had code — code, codehash, and -// codesize. Storage-slot deletes are handled by the caller's delete pass; see -// the file header for why slots absent from pre cannot be deleted. -func convertAccountRemoval(addr string, pre prestateAccount) ([]WriteSetEntry, error) { - var writes []WriteSetEntry - if pre.Nonce != nil { - writes = append(writes, WriteSetEntry{Kind: WriteKindNonce, Address: addr, Delete: true}) - } - if pre.Code != "" { - addrBytes, err := decodeHexField("address", addr, 20) - if err != nil { - return nil, err - } - writes = append(writes, - WriteSetEntry{Kind: WriteKindCode, Address: addr, Delete: true}, - WriteSetEntry{Kind: WriteKindCodeHash, Address: addr, Delete: true}, - WriteSetEntry{Kind: WriteKindRaw, Key: hex.EncodeToString(append(codeSizeKeyPrefix, addrBytes...)), Delete: true}, - ) - } - return writes, nil -} - -// padTo32 left-pads a hex slot to 32 bytes, tolerating a 0x prefix. Some -// tracers emit unpadded slot keys; store keys are always 32 bytes. -func padTo32(slot string) string { - trimmed := strings.TrimPrefix(slot, "0x") - if len(trimmed) >= 64 { - return trimmed - } - return strings.Repeat("0", 64-len(trimmed)) + trimmed -} - -// sortedKeys returns the map's keys in sorted order for deterministic output. -func sortedKeys[V any](m map[string]V) []string { - out := make([]string, 0, len(m)) - for k := range m { - out = append(out, k) - } - sort.Strings(out) - return out -} diff --git a/sei-db/state_db/bench/writeset_test.go b/sei-db/state_db/bench/writeset_test.go deleted file mode 100644 index 1d2630d761..0000000000 --- a/sei-db/state_db/bench/writeset_test.go +++ /dev/null @@ -1,287 +0,0 @@ -package bench - -import ( - "encoding/hex" - "os" - "testing" - - "github.com/stretchr/testify/require" - - "github.com/sei-protocol/sei-chain/sei-db/bench/wrappers" -) - -// prestateFixture is a real debug_traceCall prestateTracer diffMode response -// captured from pacific-1 for bytecode 0x602a60005500 -// (PUSH1 0x2a; PUSH1 0; SSTORE; STOP) with a code state override. -const prestateFixture = `{ - "jsonrpc": "2.0", - "id": 1, - "result": { - "post": { - "0x0000000000000000000000000000000000000001": {"nonce": 1}, - "0x1000000000000000000000000000000000000001": { - "storage": { - "0x0000000000000000000000000000000000000000000000000000000000000000": - "0x000000000000000000000000000000000000000000000000000000000000002a" - } - } - }, - "pre": { - "0x0000000000000000000000000000000000000001": {"balance": "0x27147114878000"}, - "0x1000000000000000000000000000000000000001": {"balance": "0x0", "code": "0x602a60005500"} - } - } -}` - -func TestConvertPrestateDiff(t *testing.T) { - converted, err := ConvertPrestateDiff([]byte(prestateFixture)) - require.NoError(t, err) - ws := converted.WriteSet - require.Len(t, ws.Blocks, 1) - - byKind := map[string]int{} - for _, w := range ws.Blocks[0].Writes { - byKind[w.Kind]++ - } - require.Equal(t, 1, byKind[WriteKindStorage], "one SSTORE slot") - require.Equal(t, 1, byKind[WriteKindNonce], "sender nonce bump") - - changesets, err := ws.BlockChangesets(0) - require.NoError(t, err) - require.Len(t, changesets, 1) - require.Equal(t, "evm", changesets[0].Name) - - var sawStorageKey bool - for _, pair := range changesets[0].Changeset.Pairs { - if pair.Key[0] == 0x03 { - sawStorageKey = true - require.Len(t, pair.Key, 53, "storage key is 0x03||addr||slot") - require.Len(t, pair.Value, 32, "storage value padded to 32 bytes") - require.Equal(t, byte(0x2a), pair.Value[31]) - } - } - require.True(t, sawStorageKey) -} - -func TestConvertPrestateDiffEmitsDeletes(t *testing.T) { - trace := `{ - "pre": {"0x1000000000000000000000000000000000000001": {"storage": {"0x01": "0x2a"}}}, - "post": {"0x1000000000000000000000000000000000000001": {"nonce": 1}} - }` - converted, err := ConvertPrestateDiff([]byte(trace)) - require.NoError(t, err) - - var deletes int - for _, w := range converted.WriteSet.Blocks[0].Writes { - if w.Delete { - deletes++ - require.Equal(t, WriteKindStorage, w.Kind) - require.Len(t, w.Slot, 64, "slot padded to 32 bytes") - } - } - require.Equal(t, 1, deletes, "slot zeroed in post emits a delete") -} - -func TestConvertPrestateDiffSelfDestruct(t *testing.T) { - // An account present in pre but absent from post is the diffMode shape a - // SELFDESTRUCT produces: the account-level keys must be deleted too, not - // just the storage slots seen in pre. - trace := `{ - "pre": {"0x3000000000000000000000000000000000000003": { - "balance": "0x2a", "nonce": 1, "code": "0x602a60005500", - "storage": {"0x01": "0x2a"}}}, - "post": {} - }` - converted, err := ConvertPrestateDiff([]byte(trace)) - require.NoError(t, err) - require.Equal(t, 1, converted.SkippedBalanceChanges, - "removed account's balance zeroing is counted as skipped") - - deletesByKind := map[string]int{} - for _, w := range converted.WriteSet.Blocks[0].Writes { - require.True(t, w.Delete, "a removed account produces only deletes") - deletesByKind[w.Kind]++ - } - require.Equal(t, map[string]int{ - WriteKindStorage: 1, - WriteKindNonce: 1, - WriteKindCode: 1, - WriteKindCodeHash: 1, - WriteKindRaw: 1, // codesize (0x09||addr) - }, deletesByKind) - - // The delete-only write set must build valid changesets. - changesets, err := converted.WriteSet.BlockChangesets(0) - require.NoError(t, err) - for _, pair := range changesets[0].Changeset.Pairs { - require.True(t, pair.Delete) - } - - // Deleting keys that were never written must replay cleanly on both - // backends (a fresh DB has none of the removed account's keys). - for _, backend := range []wrappers.DBType{wrappers.MemIAVL, wrappers.FlatKV} { - t.Run(string(backend), func(t *testing.T) { - wrapper, err := OpenReplayWrapper(t.Context(), backend, t.TempDir()) - require.NoError(t, err) - defer func() { - require.NoError(t, wrapper.Close()) - }() - _, err = ReplayWriteSet(wrapper, converted.WriteSet) - require.NoError(t, err) - }) - } -} - -func TestConvertPrestateDiffNoSpuriousDeleteOnEncodingMismatch(t *testing.T) { - // The same slot is unpadded in pre but padded in post. The delete pass must - // normalize both before comparing, or it emits a delete that clobbers the - // write (last-write-wins), losing the updated value. - trace := `{ - "pre": {"0x1000000000000000000000000000000000000001": {"storage": {"0x01": "0x2a"}}}, - "post": {"0x1000000000000000000000000000000000000001": {"storage": { - "0x0000000000000000000000000000000000000000000000000000000000000001": "0x2b"}}} - }` - converted, err := ConvertPrestateDiff([]byte(trace)) - require.NoError(t, err) - - for _, w := range converted.WriteSet.Blocks[0].Writes { - require.False(t, w.Delete, "same slot written in post must not also be deleted") - } - - changesets, err := converted.WriteSet.BlockChangesets(0) - require.NoError(t, err) - var storageWrites int - for _, pair := range changesets[0].Changeset.Pairs { - if pair.Key[0] == 0x03 { - storageWrites++ - require.False(t, pair.Delete) - require.Equal(t, byte(0x2b), pair.Value[31], "updated value survives") - } - } - require.Equal(t, 1, storageWrites) -} - -func TestConvertPrestateDiffCodeDeployment(t *testing.T) { - trace := `{ - "pre": {"0x2000000000000000000000000000000000000002": {}}, - "post": {"0x2000000000000000000000000000000000000002": {"code": "0x602a60005500", "nonce": 1}} - }` - converted, err := ConvertPrestateDiff([]byte(trace)) - require.NoError(t, err) - - byKind := map[string]WriteSetEntry{} - for _, w := range converted.WriteSet.Blocks[0].Writes { - byKind[w.Kind] = w - } - require.Contains(t, byKind, WriteKindCode) - require.Contains(t, byKind, WriteKindCodeHash) - require.Contains(t, byKind, WriteKindRaw, "codesize write") - - codeHash, err := hex.DecodeString(byKind[WriteKindCodeHash].Value) - require.NoError(t, err) - require.Len(t, codeHash, 32) - - rawKey, err := hex.DecodeString(byKind[WriteKindRaw].Key) - require.NoError(t, err) - require.Equal(t, byte(0x09), rawKey[0], "codesize key prefix") - require.Len(t, rawKey, 21) -} - -func TestConvertPrestateDiffRequiresDiffMode(t *testing.T) { - _, err := ConvertPrestateDiff([]byte(`{"pre": {}}`)) - require.ErrorContains(t, err, "diffMode") -} - -func TestReplayWriteSetOnBothBackends(t *testing.T) { - converted, err := ConvertPrestateDiff([]byte(prestateFixture)) - require.NoError(t, err) - ws := converted.WriteSet - - for _, backend := range []wrappers.DBType{wrappers.MemIAVL, wrappers.FlatKV} { - t.Run(string(backend), func(t *testing.T) { - wrapper, err := OpenReplayWrapper(t.Context(), backend, t.TempDir()) - require.NoError(t, err) - defer func() { - require.NoError(t, wrapper.Close()) - }() - - result, err := ReplayWriteSet(wrapper, ws) - require.NoError(t, err) - require.Equal(t, 1, result.Blocks) - require.Equal(t, ws.TotalKeys(), result.Keys) - require.Positive(t, result.ApplyDuration) - require.Positive(t, result.CommitDuration) - require.Equal(t, int64(1), wrapper.Version()) - - // The SSTORE'd slot must be readable back through the store key. - changesets, err := ws.BlockChangesets(0) - require.NoError(t, err) - for _, pair := range changesets[0].Changeset.Pairs { - if pair.Key[0] == 0x03 { - value, found, err := wrapper.Read(pair.Key) - require.NoError(t, err) - require.True(t, found) - require.Equal(t, pair.Value, value) - } - } - }) - } -} - -func TestLoadWriteSetValidates(t *testing.T) { - dir := t.TempDir() - path := dir + "/ws.json" - - require.NoError(t, writeFile(path, `{"blocks": [{"writes": [ - {"kind": "storage", "address": "0x1000000000000000000000000000000000000001", - "slot": "0x0000000000000000000000000000000000000000000000000000000000000000", - "value": "0x000000000000000000000000000000000000000000000000000000000000002a"} - ]}]}`)) - ws, err := LoadWriteSet(path) - require.NoError(t, err) - require.Equal(t, 1, ws.TotalKeys()) - - require.NoError(t, writeFile(path, `{"blocks": [{"writes": [{"kind": "bogus"}]}]}`)) - _, err = LoadWriteSet(path) - require.ErrorContains(t, err, "unknown write kind") - - require.NoError(t, writeFile(path, `{"module": "bank", "blocks": [{"writes": []}]}`)) - _, err = LoadWriteSet(path) - require.ErrorContains(t, err, "unsupported module") -} - -func TestValidateRejectsWrongLengthValue(t *testing.T) { - addr := "0x1000000000000000000000000000000000000001" - - // A wrong-length value for a fixed-width kind is rejected up front, so the - // benchmark never feeds divergent data to memiavl (permissive) vs FlatKV - // (which hard-errors on bad lengths deep inside ApplyChangeSets). - for _, tc := range []struct { - name string - entry WriteSetEntry - }{ - {"short nonce", WriteSetEntry{Kind: WriteKindNonce, Address: addr, Value: "0x2a000000"}}, - {"short codehash", WriteSetEntry{Kind: WriteKindCodeHash, Address: addr, Value: "0x2a"}}, - {"short storage", WriteSetEntry{Kind: WriteKindStorage, Address: addr, - Slot: "0x" + hex.EncodeToString(make([]byte, 32)), Value: "0x2a"}}, - } { - t.Run(tc.name, func(t *testing.T) { - ws := &WriteSet{Blocks: []WriteSetBlock{{Writes: []WriteSetEntry{tc.entry}}}} - err := ws.Validate() - require.ErrorContains(t, err, "expected") - _, err = ws.BlockChangesets(0) - require.ErrorContains(t, err, "expected") - }) - } - - // Correctly-sized fixed-width values and unconstrained kinds (code/raw) pass. - ws := &WriteSet{Blocks: []WriteSetBlock{{Writes: []WriteSetEntry{ - {Kind: WriteKindNonce, Address: addr, Value: "0x" + hex.EncodeToString(make([]byte, 8))}, - {Kind: WriteKindCode, Address: addr, Value: "0x602a60005500"}, - }}}} - require.NoError(t, ws.Validate()) -} - -func writeFile(path, contents string) error { - return os.WriteFile(path, []byte(contents), 0o600) -}