eth/catalyst: in zero-period dev mode, if transactions or blocks are available, commit as many blocks as possible until we have no more transactions/withdrawals to include

This commit is contained in:
Jared Wasinger 2024-07-01 13:23:48 -07:00
parent 47d5b4fc21
commit 6862283d19
2 changed files with 45 additions and 24 deletions

View file

@ -56,6 +56,10 @@ func (w *withdrawalQueue) add(withdrawal *types.Withdrawal) error {
return errors.New("withdrawal queue full") return errors.New("withdrawal queue full")
} }
w.queue = append(w.queue, withdrawal) w.queue = append(w.queue, withdrawal)
select {
case w.pending <- struct{}{}:
default:
}
return nil return nil
} }
@ -75,6 +79,12 @@ func (w *withdrawalQueue) popFront(count int) {
defer w.mu.Unlock() defer w.mu.Unlock()
w.queue = w.queue[count:] w.queue = w.queue[count:]
if len(w.queue) > 0 {
select {
case w.pending <- struct{}{}:
default:
}
}
} }
type SimulatedBeacon struct { type SimulatedBeacon struct {
@ -148,7 +158,7 @@ func (c *SimulatedBeacon) Stop() error {
// sealBlock initiates payload building for a new block and creates a new block // sealBlock initiates payload building for a new block and creates a new block
// with the completed payload. // with the completed payload.
func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp uint64) error { func (c *SimulatedBeacon) sealBlock(sealEmpty bool, timestamp uint64) (bool, error) {
if timestamp <= c.lastBlockTime { if timestamp <= c.lastBlockTime {
timestamp = c.lastBlockTime + 1 timestamp = c.lastBlockTime + 1
} }
@ -162,6 +172,7 @@ func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp u
c.setCurrentState(header.Hash(), *finalizedHash) c.setCurrentState(header.Hash(), *finalizedHash)
} }
withdrawals := c.withdrawals.gatherPending(maxWithdrawalCount)
var random [32]byte var random [32]byte
rand.Read(random[:]) rand.Read(random[:])
fcResponse, err := c.engineAPI.forkchoiceUpdated(c.curForkchoiceState, &engine.PayloadAttributes{ fcResponse, err := c.engineAPI.forkchoiceUpdated(c.curForkchoiceState, &engine.PayloadAttributes{
@ -172,23 +183,26 @@ func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp u
BeaconRoot: &common.Hash{}, BeaconRoot: &common.Hash{},
}, engine.PayloadV3, true) }, engine.PayloadV3, true)
if err != nil { if err != nil {
return err return false, err
} }
if fcResponse == engine.STATUS_SYNCING { if fcResponse == engine.STATUS_SYNCING {
return errors.New("chain rewind prevented invocation of payload creation") return false, errors.New("chain rewind prevented invocation of payload creation")
} }
envelope, err := c.engineAPI.getPayload(*fcResponse.PayloadID, true) envelope, err := c.engineAPI.getPayload(*fcResponse.PayloadID, true)
if err != nil { if err != nil {
return err return false, err
} }
payload := envelope.ExecutionPayload payload := envelope.ExecutionPayload
if !sealEmpty && len(payload.Transactions) == 0 && len(payload.Withdrawals) == 0 {
return false, nil
}
var finalizedHash common.Hash var finalizedHash common.Hash
if payload.Number%devEpochLength == 0 { if payload.Number%devEpochLength == 0 {
finalizedHash = payload.BlockHash finalizedHash = payload.BlockHash
} else { } else {
if fh := c.finalizedBlockHash(payload.Number); fh == nil { if fh := c.finalizedBlockHash(payload.Number); fh == nil {
return errors.New("chain rewind interrupted calculation of finalized block hash") return false, errors.New("chain rewind interrupted calculation of finalized block hash")
} else { } else {
finalizedHash = *fh finalizedHash = *fh
} }
@ -201,7 +215,7 @@ func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp u
for _, commit := range envelope.BlobsBundle.Commitments { for _, commit := range envelope.BlobsBundle.Commitments {
var c kzg4844.Commitment var c kzg4844.Commitment
if len(commit) != len(c) { if len(commit) != len(c) {
return errors.New("invalid commitment length") return false, errors.New("invalid commitment length")
} }
copy(c[:], commit) copy(c[:], commit)
blobHashes = append(blobHashes, kzg4844.CalcBlobHashV1(hasher, &c)) blobHashes = append(blobHashes, kzg4844.CalcBlobHashV1(hasher, &c))
@ -209,16 +223,17 @@ func (c *SimulatedBeacon) sealBlock(withdrawals []*types.Withdrawal, timestamp u
} }
// Mark the payload as canon // Mark the payload as canon
if _, err = c.engineAPI.NewPayloadV3(*payload, blobHashes, &common.Hash{}); err != nil { if _, err = c.engineAPI.NewPayloadV3(*payload, blobHashes, &common.Hash{}); err != nil {
return err return false, err
} }
c.setCurrentState(payload.BlockHash, finalizedHash) c.setCurrentState(payload.BlockHash, finalizedHash)
// Mark the block containing the payload as canonical // Mark the block containing the payload as canonical
if _, err = c.engineAPI.ForkchoiceUpdatedV2(c.curForkchoiceState, nil); err != nil { if _, err = c.engineAPI.ForkchoiceUpdatedV2(c.curForkchoiceState, nil); err != nil {
return err return false, err
} }
c.lastBlockTime = payload.Timestamp c.lastBlockTime = payload.Timestamp
return nil c.withdrawals.popFront(len(withdrawals))
return true, nil
} }
// loop runs the block production loop for non-zero period configuration // loop runs the block production loop for non-zero period configuration
@ -229,8 +244,7 @@ func (c *SimulatedBeacon) loop() {
case <-c.shutdownCh: case <-c.shutdownCh:
return return
case <-timer.C: case <-timer.C:
withdrawals := c.withdrawals.gatherPending(10) if _, err := c.sealBlock(true, uint64(time.Now().Unix())); err != nil {
if err := c.sealBlock(withdrawals, uint64(time.Now().Unix())); err != nil {
log.Warn("Error performing sealing work", "err", err) log.Warn("Error performing sealing work", "err", err)
} else { } else {
timer.Reset(time.Second * time.Duration(c.period)) timer.Reset(time.Second * time.Duration(c.period))
@ -266,13 +280,26 @@ func (c *SimulatedBeacon) setCurrentState(headHash, finalizedHash common.Hash) {
// Commit seals a block on demand. // Commit seals a block on demand.
func (c *SimulatedBeacon) Commit() common.Hash { func (c *SimulatedBeacon) Commit() common.Hash {
withdrawals := c.withdrawals.gatherPending(10) if _, err := c.sealBlock(true, uint64(time.Now().Unix())); err != nil {
if err := c.sealBlock(withdrawals, uint64(time.Now().Unix())); err != nil {
log.Warn("Error performing sealing work", "err", err) log.Warn("Error performing sealing work", "err", err)
} }
return c.eth.BlockChain().CurrentBlock().Hash() return c.eth.BlockChain().CurrentBlock().Hash()
} }
// commitUntilEmpty seals blocks until there are now transactions or withdrawals
// left to include
func (c *SimulatedBeacon) commitUntilEmpty() {
for {
committed, err := c.sealBlock(false, uint64(time.Now().Unix()))
if err != nil {
log.Error("failed to seal block", "err", err)
}
if !committed {
return
}
}
}
// Rollback un-sends previously added transactions. // Rollback un-sends previously added transactions.
func (c *SimulatedBeacon) Rollback() { func (c *SimulatedBeacon) Rollback() {
// Flush all transactions from the transaction pools // Flush all transactions from the transaction pools
@ -307,8 +334,8 @@ func (c *SimulatedBeacon) AdjustTime(adjustment time.Duration) error {
if parent == nil { if parent == nil {
return errors.New("parent not found") return errors.New("parent not found")
} }
withdrawals := c.withdrawals.gatherPending(10) _, err := c.sealBlock(true, parent.Time+uint64(adjustment))
return c.sealBlock(withdrawals, parent.Time+uint64(adjustment)) return err
} }
func RegisterSimulatedBeaconAPIs(stack *node.Node, sim *SimulatedBeacon) { func RegisterSimulatedBeaconAPIs(stack *node.Node, sim *SimulatedBeacon) {

View file

@ -18,12 +18,9 @@ package catalyst
import ( import (
"context" "context"
"time"
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core" "github.com/ethereum/go-ethereum/core"
"github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/log"
) )
type api struct { type api struct {
@ -41,13 +38,10 @@ func (a *api) loop() {
select { select {
case <-a.sim.shutdownCh: case <-a.sim.shutdownCh:
return return
case w := <-a.sim.withdrawals.pending: case <-a.sim.withdrawals.pending:
withdrawals := append(a.sim.withdrawals.gatherPending(9), w) a.sim.commitUntilEmpty()
if err := a.sim.sealBlock(withdrawals, uint64(time.Now().Unix())); err != nil {
log.Warn("Error performing sealing work", "err", err)
}
case <-newTxs: case <-newTxs:
a.sim.Commit() a.sim.commitUntilEmpty()
} }
} }
} }