core: split commit and flush in prep for witnesses

This commit is contained in:
Péter Szilágyi 2024-06-19 15:23:10 +03:00
parent c11aac249d
commit 56c016b7cf
3 changed files with 54 additions and 36 deletions

View file

@ -1472,8 +1472,8 @@ func (bc *BlockChain) writeBlockWithState(block *types.Block, receipts []*types.
if err := blockBatch.Write(); err != nil { if err := blockBatch.Write(); err != nil {
log.Crit("Failed to write block into disk", "err", err) log.Crit("Failed to write block into disk", "err", err)
} }
// Commit all cached state changes into underlying memory database. // Flush all cached state changes into underlying memory database.
root, err := statedb.Commit(block.NumberU64(), bc.chainConfig.IsEIP158(block.Number())) root, err := statedb.Flush(block.NumberU64())
if err != nil { if err != nil {
return err return err
} }
@ -1929,6 +1929,13 @@ func (bc *BlockChain) processBlock(block *types.Block, statedb *state.StateDB, s
return nil, err return nil, err
} }
vtime := time.Since(vstart) vtime := time.Since(vstart)
// Commit the tries to in-memory buffers to finish populating any witnesses
if err = statedb.CommitWithoutFlush(bc.chainConfig.IsEIP158(block.Number())); err != nil {
return nil, err
}
// TODO(karalabe): Do the stateless cross validation here when merging it
proctime := time.Since(start) // processing + validation proctime := time.Since(start) // processing + validation
// Update the metrics touched during block processing and validation // Update the metrics touched during block processing and validation

View file

@ -146,6 +146,9 @@ type StateDB struct {
validRevisions []revision validRevisions []revision
nextRevisionId int nextRevisionId int
// Temporary storage area for comitted but not yet flushed state updates
commit *stateUpdate
// Measurements gathered during execution for debugging purposes // Measurements gathered during execution for debugging purposes
AccountReads time.Duration AccountReads time.Duration
AccountHashes time.Duration AccountHashes time.Duration
@ -1099,12 +1102,17 @@ func (s *StateDB) GetTrie() Trie {
return s.trie return s.trie
} }
// commit gathers the state mutations accumulated along with the associated // CommitWithoutFlush gathers the state mutations accumulated along with the
// trie changes, resetting all internal flags with the new state as the base. // associated trie changes, resetting all internal flags with the new state
func (s *StateDB) commit(deleteEmptyObjects bool) (*stateUpdate, error) { // as the base.
//
// The mutations will be stored in the statedb but will not be flushed to disk.
// This can be useful in cases where additional validation (stateless witness
// checks across clients) it be be performed on the state before storing it.
func (s *StateDB) CommitWithoutFlush(deleteEmptyObjects bool) error {
// Short circuit in case any database failure occurred earlier. // Short circuit in case any database failure occurred earlier.
if s.dbErr != nil { if s.dbErr != nil {
return nil, fmt.Errorf("commit aborted due to earlier error: %v", s.dbErr) return fmt.Errorf("commit aborted due to earlier error: %v", s.dbErr)
} }
// Finalize any pending changes and merge everything into the tries // Finalize any pending changes and merge everything into the tries
s.IntermediateRoot(deleteEmptyObjects) s.IntermediateRoot(deleteEmptyObjects)
@ -1153,11 +1161,11 @@ func (s *StateDB) commit(deleteEmptyObjects bool) (*stateUpdate, error) {
// during subsequent resurrection can be combined correctly. // during subsequent resurrection can be combined correctly.
deletes, delNodes, err := s.handleDestruction() deletes, delNodes, err := s.handleDestruction()
if err != nil { if err != nil {
return nil, err return err
} }
for _, set := range delNodes { for _, set := range delNodes {
if err := merge(set); err != nil { if err := merge(set); err != nil {
return nil, err return err
} }
} }
// Handle all state updates afterwards, concurrently to one another to shave // Handle all state updates afterwards, concurrently to one another to shave
@ -1202,7 +1210,7 @@ func (s *StateDB) commit(deleteEmptyObjects bool) (*stateUpdate, error) {
// Write any contract code associated with the state object // Write any contract code associated with the state object
obj := s.stateObjects[addr] obj := s.stateObjects[addr]
if obj == nil { if obj == nil {
return nil, errors.New("missing state object") return errors.New("missing state object")
} }
// Run the storage updates concurrently to one another // Run the storage updates concurrently to one another
workers.Go(func() error { workers.Go(func() error {
@ -1223,7 +1231,7 @@ func (s *StateDB) commit(deleteEmptyObjects bool) (*stateUpdate, error) {
} }
// Wait for everything to finish and update the metrics // Wait for everything to finish and update the metrics
if err := workers.Wait(); err != nil { if err := workers.Wait(); err != nil {
return nil, err return err
} }
accountUpdatedMeter.Mark(int64(s.AccountUpdated)) accountUpdatedMeter.Mark(int64(s.AccountUpdated))
storageUpdatedMeter.Mark(s.StorageUpdated.Load()) storageUpdatedMeter.Mark(s.StorageUpdated.Load())
@ -1243,55 +1251,59 @@ func (s *StateDB) commit(deleteEmptyObjects bool) (*stateUpdate, error) {
origin := s.originalRoot origin := s.originalRoot
s.originalRoot = root s.originalRoot = root
return newStateUpdate(origin, root, deletes, updates, nodes), nil
s.commit = newStateUpdate(origin, root, deletes, updates, nodes)
return nil
} }
// commitAndFlush is a wrapper of commit which also commits the state mutations // Flush operates on a comitted-but-not-flushed state database to iterate and
// to the configured data stores. // push the state diffs into the database. Unless you need to do something in
func (s *StateDB) commitAndFlush(block uint64, deleteEmptyObjects bool) (*stateUpdate, error) { // between commit and flush, just use statedb.Commit directly.
ret, err := s.commit(deleteEmptyObjects) func (s *StateDB) Flush(block uint64) (common.Hash, error) {
if err != nil { // Sanity check that commit was called before flush
return nil, err if s.commit == nil {
log.Error("Statedb not committed before flush")
return common.Hash{}, errors.New("statedb uncomitted")
} }
// Commit dirty contract code if any exists // Commit dirty contract code if any exists
if db := s.db.DiskDB(); db != nil && len(ret.codes) > 0 { if db := s.db.DiskDB(); db != nil && len(s.commit.codes) > 0 {
batch := db.NewBatch() batch := db.NewBatch()
for _, code := range ret.codes { for _, code := range s.commit.codes {
rawdb.WriteCode(batch, code.hash, code.blob) rawdb.WriteCode(batch, code.hash, code.blob)
} }
if err := batch.Write(); err != nil { if err := batch.Write(); err != nil {
return nil, err return common.Hash{}, err
} }
} }
if !ret.empty() { if !s.commit.empty() {
// If snapshotting is enabled, update the snapshot tree with this new version // If snapshotting is enabled, update the snapshot tree with this new version
if s.snap != nil { if s.snap != nil {
s.snap = nil s.snap = nil
start := time.Now() start := time.Now()
if err := s.snaps.Update(ret.root, ret.originRoot, ret.destructs, ret.accounts, ret.storages); err != nil { if err := s.snaps.Update(s.commit.root, s.commit.originRoot, s.commit.destructs, s.commit.accounts, s.commit.storages); err != nil {
log.Warn("Failed to update snapshot tree", "from", ret.originRoot, "to", ret.root, "err", err) log.Warn("Failed to update snapshot tree", "from", s.commit.originRoot, "to", s.commit.root, "err", err)
} }
// Keep 128 diff layers in the memory, persistent layer is 129th. // Keep 128 diff layers in the memory, persistent layer is 129th.
// - head layer is paired with HEAD state // - head layer is paired with HEAD state
// - head-1 layer is paired with HEAD-1 state // - head-1 layer is paired with HEAD-1 state
// - head-127 layer(bottom-most diff layer) is paired with HEAD-127 state // - head-127 layer(bottom-most diff layer) is paired with HEAD-127 state
if err := s.snaps.Cap(ret.root, TriesInMemory); err != nil { if err := s.snaps.Cap(s.commit.root, TriesInMemory); err != nil {
log.Warn("Failed to cap snapshot tree", "root", ret.root, "layers", TriesInMemory, "err", err) log.Warn("Failed to cap snapshot tree", "root", s.commit.root, "layers", TriesInMemory, "err", err)
} }
s.SnapshotCommits += time.Since(start) s.SnapshotCommits += time.Since(start)
} }
// If trie database is enabled, commit the state update as a new layer // If trie database is enabled, commit the state update as a new layer
if db := s.db.TrieDB(); db != nil { if db := s.db.TrieDB(); db != nil {
start := time.Now() start := time.Now()
set := triestate.New(ret.accountsOrigin, ret.storagesOrigin) set := triestate.New(s.commit.accountsOrigin, s.commit.storagesOrigin)
if err := db.Update(ret.root, ret.originRoot, block, ret.nodes, set); err != nil { if err := db.Update(s.commit.root, s.commit.originRoot, block, s.commit.nodes, set); err != nil {
return nil, err return common.Hash{}, err
} }
s.TrieDBCommits += time.Since(start) s.TrieDBCommits += time.Since(start)
} }
} }
return ret, err return s.commit.root, nil
} }
// Commit writes the state mutations into the configured data stores. // Commit writes the state mutations into the configured data stores.
@ -1304,11 +1316,10 @@ func (s *StateDB) commitAndFlush(block uint64, deleteEmptyObjects bool) (*stateU
// The associated block number of the state transition is also provided // The associated block number of the state transition is also provided
// for more chain context. // for more chain context.
func (s *StateDB) Commit(block uint64, deleteEmptyObjects bool) (common.Hash, error) { func (s *StateDB) Commit(block uint64, deleteEmptyObjects bool) (common.Hash, error) {
ret, err := s.commitAndFlush(block, deleteEmptyObjects) if err := s.CommitWithoutFlush(deleteEmptyObjects); err != nil {
if err != nil {
return common.Hash{}, err return common.Hash{}, err
} }
return ret.root, nil return s.Flush(block)
} }
// Prepare handles the preparatory steps for executing a state transition with. // Prepare handles the preparatory steps for executing a state transition with.

View file

@ -236,15 +236,15 @@ func (test *stateTest) run() bool {
} else { } else {
state.IntermediateRoot(true) // call intermediateRoot at the transaction boundary state.IntermediateRoot(true) // call intermediateRoot at the transaction boundary
} }
ret, err := state.commitAndFlush(0, true) // call commit at the block boundary _, err = state.Commit(0, true) // call commit at the block boundary
if err != nil { if err != nil {
panic(err) panic(err)
} }
if ret.empty() { if state.commit.empty() {
return true return true
} }
copyUpdate(ret) copyUpdate(state.commit)
roots = append(roots, ret.root) roots = append(roots, state.commit.root)
} }
for i := 0; i < len(test.actions); i++ { for i := 0; i < len(test.actions); i++ {
root := types.EmptyRootHash root := types.EmptyRootHash