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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,9 @@ This document outlines major changes between releases.
## [Unreleased]

New features:
* `PrepareRequestExtensionEnablingHeight` configuration parameter and
`NewPrepareRequestExtended` callback to attach full transaction list to
`PrepareRequest` instead of hashes starting from the given height (#160)

Behaviour changes:

Expand Down
43 changes: 38 additions & 5 deletions config.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,9 @@ type Config[H Hash] struct {
// AntiMEVExtensionEnablingHeight denotes the height starting from which dBFT
// Anti-MEV extensions should be enabled. -1 means no extension is enabled.
AntiMEVExtensionEnablingHeight int64
// PrepareRequestExtensionEnablingHeight denotes the height starting from which
// an extended PrepareRequest format should be enabled. -1 means no extension is enabled.
PrepareRequestExtensionEnablingHeight int64
// GetKeyPair returns an index of the node in the list of validators
// together with it's key pair.
GetKeyPair func([]PublicKey) (int, PrivateKey, PublicKey)
Expand All @@ -53,6 +56,8 @@ type Config[H Hash] struct {
StopTxFlow func()
// GetTx returns a transaction from memory pool.
GetTx func(h H) Transaction[H]
// GetTxData returns arbitrary verified data associated with the transaction.
GetTxData func(h H) any
// GetVerified returns a slice of verified transactions
// to be proposed in a new block.
GetVerified func() []Transaction[H]
Expand Down Expand Up @@ -81,8 +86,12 @@ type Config[H Hash] struct {
GetValidators func(...Transaction[H]) []PublicKey
// NewConsensusPayload is a constructor for payload.ConsensusPayload.
NewConsensusPayload func(*Context[H], MessageType, any) ConsensusPayload[H]
// NewPrepareRequest is a constructor for payload.PrepareRequest.
// NewPrepareRequest is a constructor for payload.PrepareRequest that
// builds a request carrying transaction hashes only.
NewPrepareRequest func(ts uint64, nonce uint64, transactionHashes []H) PrepareRequest[H]
// NewPrepareRequestExtended is a constructor for payload.PrepareRequest that
// builds a request carrying full transaction list.
NewPrepareRequestExtended func(ts uint64, nonce uint64, transactions []Transaction[H]) PrepareRequest[H]
// NewPrepareResponse is a constructor for payload.PrepareResponse.
NewPrepareResponse func(preparationHash H) PrepareResponse[H]
// NewChangeView is a constructor for payload.ChangeView.
Expand Down Expand Up @@ -137,9 +146,10 @@ func defaultConfig[H Hash]() *Config[H] {
VerifyPrepareResponse: func(ConsensusPayload[H]) error { return nil },
VerifyCommit: func(ConsensusPayload[H]) error { return nil },

AntiMEVExtensionEnablingHeight: -1,
VerifyPreBlock: func(PreBlock[H]) bool { return true },
VerifyPreCommit: func(ConsensusPayload[H]) error { return nil },
AntiMEVExtensionEnablingHeight: -1,
PrepareRequestExtensionEnablingHeight: -1,
VerifyPreBlock: func(PreBlock[H]) bool { return true },
VerifyPreCommit: func(ConsensusPayload[H]) error { return nil },
}
}

Expand Down Expand Up @@ -204,6 +214,15 @@ func checkConfig[H Hash](cfg *Config[H]) error {
return errors.New("NewPreCommit is set, but AntiMEVExtensionEnablingHeight is not specified")
}
}
if cfg.PrepareRequestExtensionEnablingHeight >= 0 {
if cfg.NewPrepareRequestExtended == nil {
return errors.New("NewPrepareRequestExtended is nil")
}
} else {
if cfg.NewPrepareRequestExtended != nil {
return errors.New("NewPrepareRequestExtended is set, but PrepareRequestExtensionEnablingHeight is not specified")
}
}
if (cfg.MaxTimePerBlock == nil) != (cfg.SubscribeForTxs == nil) {
return errors.New("MaxTimePerBlock and SubscribeForTxs should be specified/not specified at the same time")
}
Expand Down Expand Up @@ -253,6 +272,13 @@ func WithAntiMEVExtensionEnablingHeight[H Hash](h int64) func(config *Config[H])
}
}

// WithPrepareRequestExtensionEnablingHeight sets PrepareRequestExtensionEnablingHeight.
func WithPrepareRequestExtensionEnablingHeight[H Hash](h int64) func(config *Config[H]) {
return func(cfg *Config[H]) {
cfg.PrepareRequestExtensionEnablingHeight = h
}
}

// WithTimestampIncrement sets TimestampIncrement.
func WithTimestampIncrement[H Hash](u uint64) func(config *Config[H]) {
return func(cfg *Config[H]) {
Expand Down Expand Up @@ -388,12 +414,19 @@ func WithNewConsensusPayload[H Hash](f func(ctx *Context[H], typ MessageType, ms
}

// WithNewPrepareRequest sets NewPrepareRequest.
func WithNewPrepareRequest[H Hash](f func(ts uint64, nonce uint64, transactionsHashes []H) PrepareRequest[H]) func(config *Config[H]) {
func WithNewPrepareRequest[H Hash](f func(ts uint64, nonce uint64, transactionHashes []H) PrepareRequest[H]) func(config *Config[H]) {
Comment thread
AnnaShaleva marked this conversation as resolved.
return func(cfg *Config[H]) {
cfg.NewPrepareRequest = f
}
}

// WithNewPrepareRequestExtended sets NewPrepareRequestExtended.
func WithNewPrepareRequestExtended[H Hash](f func(ts uint64, nonce uint64, transactions []Transaction[H]) PrepareRequest[H]) func(config *Config[H]) {
return func(cfg *Config[H]) {
cfg.NewPrepareRequestExtended = f
}
}

// WithNewPrepareResponse sets NewPrepareResponse.
func WithNewPrepareResponse[H Hash](f func(preparationHash H) PrepareResponse[H]) func(config *Config[H]) {
return func(cfg *Config[H]) {
Expand Down
70 changes: 48 additions & 22 deletions context.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,12 +56,24 @@ type Context[H Hash] struct {
Timestamp uint64
Nonce uint64
// TransactionHashes is a slice of hashes of proposed transactions in the current block.
// It's used when PrepareRequestExtensionEnabled is false.
TransactionHashes []H
// MissingTransactions is a slice of hashes containing missing transactions for the current block.
// MissingTransactions is a slice of hashes containing missing transactions (in case of
// disabled PrepareRequestExtension) or hashes of those transactions that require
// additional data to be fetched prior to the proposal verification (in case of enabled
// PrepareRequestExtension).
MissingTransactions []H
// Transactions is a map containing actual transactions for the current block.
// It's used when PrepareRequestExtensionEnabled is false.
Transactions map[H]Transaction[H]

// TransactionList is a full ordered list of transactions proposed in the current block.
// It's used when PrepareRequestExtensionEnabled is true.
TransactionList []Transaction[H]
Comment thread
AnnaShaleva marked this conversation as resolved.
// PrepareRequestExtensionEnabled tells whether full transaction list is used for the
// PrepareRequest construction. Refreshed once per dBFT reset.
PrepareRequestExtensionEnabled bool

// PreparationPayloads stores consensus Prepare* payloads for the current epoch.
PreparationPayloads []ConsensusPayload[H]
// PreCommitPayloads stores consensus PreCommit payloads sent through all epochs
Expand Down Expand Up @@ -287,6 +299,7 @@ func (c *Context[H]) reset(view byte, ts uint64) {
}
c.PreparationPayloads = emptyReusableSlice(c.PreparationPayloads, n)

c.TransactionList = nil
if c.Transactions == nil { // Init.
c.Transactions = make(map[H]Transaction[H])
} else { // Regular use.
Expand All @@ -302,6 +315,7 @@ func (c *Context[H]) reset(view byte, ts uint64) {
if c.MyIndex >= 0 {
c.LastSeenMessage[c.MyIndex] = &HeightView{c.BlockIndex, c.ViewNumber}
}
c.PrepareRequestExtensionEnabled = c.isPrepareRequestExtensionEnabled()
}

func emptyReusableSlice[E any](s []E, n int) []E {
Expand All @@ -325,12 +339,15 @@ func (c *Context[H]) Fill(force bool) bool {
_, _ = rand.Read(b)

c.Nonce = binary.LittleEndian.Uint64(b)
c.TransactionHashes = make([]H, len(txx))

for i := range txx {
h := txx[i].Hash()
c.TransactionHashes[i] = h
c.Transactions[h] = txx[i]
if c.PrepareRequestExtensionEnabled {
c.TransactionList = txx
} else {
c.TransactionHashes = make([]H, len(txx))
for i := range txx {
h := txx[i].Hash()
c.TransactionHashes[i] = h
c.Transactions[h] = txx[i]
}
}

c.Timestamp = c.lastBlockTimestamp + c.Config.TimestampIncrement
Expand All @@ -353,18 +370,12 @@ func (c *Context[H]) CreateBlock() Block[H] {
return nil
}

txx := make([]Transaction[H], len(c.TransactionHashes))

for i, h := range c.TransactionHashes {
txx[i] = c.Transactions[h]
}

// Anti-MEV extension properly sets PreBlock transactions once during PreBlock
// construction and then never updates these transactions in the dBFT context.
// Thus, user must not reuse txx if anti-MEV extension is enabled. However,
// we don't skip a call to Block.SetTransactions since it may be used as a
// signal to the user's code to finalize the block.
c.block.SetTransactions(txx)
c.block.SetTransactions(c.collectTransactions())
}

return c.block
Expand All @@ -377,24 +388,39 @@ func (c *Context[H]) CreatePreBlock() PreBlock[H] {
return nil
}

txx := make([]Transaction[H], len(c.TransactionHashes))

for i, h := range c.TransactionHashes {
txx[i] = c.Transactions[h]
}

c.preBlock.SetTransactions(txx)
c.preBlock.SetTransactions(c.collectTransactions())
}

return c.preBlock
}

// collectTransactions returns the ordered list of transactions to be proposed in
// the resulting block. If PrepareRequestExtension is enabled, they are taken
// directly from TransactionList; otherwise they are reconstructed, preserving
// TransactionHashes order, from the already-collected Transactions map.
func (c *Context[H]) collectTransactions() []Transaction[H] {
if c.PrepareRequestExtensionEnabled {
return c.TransactionList
}
txx := make([]Transaction[H], len(c.TransactionHashes))
for i, h := range c.TransactionHashes {
txx[i] = c.Transactions[h]
}
return txx
}

// isAntiMEVExtensionEnabled returns whether Anti-MEV dBFT extension is enabled
// at the currently processing block height.
func (c *Context[H]) isAntiMEVExtensionEnabled() bool {
return c.Config.AntiMEVExtensionEnablingHeight >= 0 && uint32(c.Config.AntiMEVExtensionEnablingHeight) <= c.BlockIndex
}

// isPrepareRequestExtensionEnabled returns whether PrepareRequest dBFT extension is enabled
// at the currently processing block height.
func (c *Context[H]) isPrepareRequestExtensionEnabled() bool {
return c.Config.PrepareRequestExtensionEnablingHeight >= 0 && uint32(c.Config.PrepareRequestExtensionEnablingHeight) <= c.BlockIndex
}

// MakeHeader returns half-filled block for the current epoch.
// All hashable fields will be filled.
func (c *Context[H]) MakeHeader() Block[H] {
Expand Down Expand Up @@ -432,7 +458,7 @@ func (c *Context[H]) MakePreHeader() PreBlock[H] {
// hasAllTransactions returns true iff all transactions were received
// for the proposed block.
func (c *Context[H]) hasAllTransactions() bool {
return len(c.TransactionHashes) == len(c.Transactions)
return len(c.MissingTransactions) == 0
}

func (c *Context[H]) subscribeForTransactions() {
Expand Down
44 changes: 33 additions & 11 deletions dbft.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,14 @@ func New[H Hash](options ...func(config *Config[H])) (*DBFT[H], error) {
return d, nil
}

// addTransaction adds a missing transaction to the context (in case of disabled
// PrepareRequestExtension) and advances the state machine if all transactions
// (or dependent data in case of enabled PrepareRequestExtension) are collected
// and it's possible to build a valid block.
func (d *DBFT[H]) addTransaction(tx Transaction[H]) {
d.Transactions[tx.Hash()] = tx
Comment thread
AnnaShaleva marked this conversation as resolved.
if !d.PrepareRequestExtensionEnabled {
d.Transactions[tx.Hash()] = tx
}
if d.hasAllTransactions() {
if d.IsPrimary() || d.Context.WatchOnly() {
return
Expand Down Expand Up @@ -159,8 +165,9 @@ func (d *DBFT[H]) initializeConsensus(view byte, ts uint64) {
d.changeTimer(timeout)
}

// OnTransaction notifies service about receiving new transaction from the
// proposed list of transactions.
// OnTransaction notifies service about receiving new transaction (in case of
// disabled PrepareRequestExtension) or transaction data (in case of enabled
// PrepareRequestExtension) from the proposed list of transactions.
func (d *DBFT[H]) OnTransaction(tx Transaction[H]) {
// d.Logger.Debug("OnTransaction",
// zap.Bool("backup", d.IsBackup()),
Expand All @@ -178,13 +185,8 @@ func (d *DBFT[H]) OnTransaction(tx Transaction[H]) {
if i < 0 {
return
}
d.addTransaction(tx)
// `addTransaction` checks for responses and commits. If this was the last transaction
// Context could be initialized on a new height, clearing this field.
if len(d.MissingTransactions) == 0 {
return
}
d.MissingTransactions = slices.Delete(d.MissingTransactions, i, i+1)
d.addTransaction(tx)
}

// OnTimeout advances state machine as if timeout was fired.
Expand Down Expand Up @@ -349,9 +351,16 @@ func (d *DBFT[H]) onPrepareRequest(msg ConsensusPayload[H]) {

d.Timestamp = p.Timestamp()
d.Nonce = p.Nonce()
d.TransactionHashes = p.TransactionHashes()
var txCount int
if d.PrepareRequestExtensionEnabled {
d.TransactionList = p.Transactions()
txCount = len(d.TransactionList)
} else {
d.TransactionHashes = p.TransactionHashes()
txCount = len(d.TransactionHashes)
}

d.Logger.Info("received PrepareRequest", zap.Uint16("validator", msg.ValidatorIndex()), zap.Int("tx", len(d.TransactionHashes)))
d.Logger.Info("received PrepareRequest", zap.Uint16("validator", msg.ValidatorIndex()), zap.Int("tx", txCount))
d.processMissingTx()
d.updateExistingPayloads(msg)
d.PreparationPayloads[msg.ValidatorIndex()] = msg
Expand All @@ -365,6 +374,19 @@ func (d *DBFT[H]) onPrepareRequest(msg ConsensusPayload[H]) {
}

func (d *DBFT[H]) processMissingTx() {
if d.PrepareRequestExtensionEnabled {
for _, tx := range d.TransactionList {
if tx.HasData() {
continue
}
h := tx.Hash()
if data := d.GetTxData(tx.Hash()); data == nil {
d.MissingTransactions = append(d.MissingTransactions, h)
}
}

return
}
for _, h := range d.TransactionHashes {
if _, ok := d.Transactions[h]; ok {
continue
Expand Down
Loading
Loading