feat(BB): Better proposal, better docs, better logging, papa johns (#172)
This commit is contained in:
+31
-13
@@ -45,13 +45,21 @@ func (h *ProposalHandler) PrepareProposalHandler() sdk.PrepareProposalHandler {
|
||||
}
|
||||
}()
|
||||
|
||||
proposal := h.prepareLanesHandler(ctx, blockbuster.NewProposal(req.MaxTxBytes))
|
||||
|
||||
resp = abci.ResponsePrepareProposal{
|
||||
Txs: proposal.Txs,
|
||||
proposal, err := h.prepareLanesHandler(ctx, blockbuster.NewProposal(req.MaxTxBytes))
|
||||
if err != nil {
|
||||
h.logger.Error("failed to prepare proposal", "err", err)
|
||||
return abci.ResponsePrepareProposal{Txs: make([][]byte, 0)}
|
||||
}
|
||||
|
||||
return
|
||||
h.logger.Info(
|
||||
"prepared proposal",
|
||||
"num_txs", proposal.GetNumTxs(),
|
||||
"total_tx_bytes", proposal.GetTotalTxBytes(),
|
||||
)
|
||||
|
||||
return abci.ResponsePrepareProposal{
|
||||
Txs: proposal.GetProposal(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -69,8 +77,14 @@ func (h *ProposalHandler) ProcessProposalHandler() sdk.ProcessProposalHandler {
|
||||
}
|
||||
}()
|
||||
|
||||
txs := req.Txs
|
||||
if len(txs) == 0 {
|
||||
h.logger.Info("accepted empty proposal")
|
||||
return abci.ResponseProcessProposal{Status: abci.ResponseProcessProposal_ACCEPT}
|
||||
}
|
||||
|
||||
// Decode the transactions from the proposal.
|
||||
decodedTxs, err := utils.GetDecodedTxs(h.txDecoder, req.Txs)
|
||||
decodedTxs, err := utils.GetDecodedTxs(h.txDecoder, txs)
|
||||
if err != nil {
|
||||
h.logger.Error("failed to decode transactions", "err", err)
|
||||
return abci.ResponseProcessProposal{Status: abci.ResponseProcessProposal_REJECT}
|
||||
@@ -82,6 +96,8 @@ func (h *ProposalHandler) ProcessProposalHandler() sdk.ProcessProposalHandler {
|
||||
return abci.ResponseProcessProposal{Status: abci.ResponseProcessProposal_REJECT}
|
||||
}
|
||||
|
||||
h.logger.Info("validated proposal", "num_txs", len(txs))
|
||||
|
||||
return abci.ResponseProcessProposal{Status: abci.ResponseProcessProposal_ACCEPT}
|
||||
}
|
||||
}
|
||||
@@ -102,7 +118,7 @@ func ChainPrepareLanes(chain ...blockbuster.Lane) blockbuster.PrepareLanesHandle
|
||||
chain = append(chain, terminator.Terminator{})
|
||||
}
|
||||
|
||||
return func(ctx sdk.Context, partialProposal *blockbuster.Proposal) (finalProposal *blockbuster.Proposal) {
|
||||
return func(ctx sdk.Context, partialProposal blockbuster.BlockProposal) (finalProposal blockbuster.BlockProposal, err error) {
|
||||
lane := chain[0]
|
||||
lane.Logger().Info("preparing lane", "lane", lane.Name())
|
||||
|
||||
@@ -110,8 +126,8 @@ func ChainPrepareLanes(chain ...blockbuster.Lane) blockbuster.PrepareLanesHandle
|
||||
cacheCtx, write := ctx.CacheContext()
|
||||
|
||||
defer func() {
|
||||
if err := recover(); err != nil {
|
||||
lane.Logger().Error("failed to prepare lane", "lane", lane.Name(), "err", err)
|
||||
if rec := recover(); rec != nil || err != nil {
|
||||
lane.Logger().Error("failed to prepare lane", "lane", lane.Name(), "err", err, "recover_error", rec)
|
||||
|
||||
lanesRemaining := len(chain)
|
||||
switch {
|
||||
@@ -119,18 +135,19 @@ func ChainPrepareLanes(chain ...blockbuster.Lane) blockbuster.PrepareLanesHandle
|
||||
// If there are only two lanes remaining, then the first lane in the chain
|
||||
// is the lane that failed to prepare the partial proposal and the second lane in the
|
||||
// chain is the terminator lane. We return the proposal as is.
|
||||
finalProposal = partialProposal
|
||||
finalProposal, err = partialProposal, nil
|
||||
default:
|
||||
// If there are more than two lanes remaining, then the first lane in the chain
|
||||
// is the lane that failed to prepare the proposal but the second lane in the
|
||||
// chain is not the terminator lane so there could potentially be more transactions
|
||||
// added to the proposal
|
||||
maxTxBytesForLane := utils.GetMaxTxBytesForLane(
|
||||
partialProposal,
|
||||
partialProposal.GetMaxTxBytes(),
|
||||
partialProposal.GetTotalTxBytes(),
|
||||
chain[1].GetMaxBlockSpace(),
|
||||
)
|
||||
|
||||
finalProposal = chain[1].PrepareLane(
|
||||
finalProposal, err = chain[1].PrepareLane(
|
||||
ctx,
|
||||
partialProposal,
|
||||
maxTxBytesForLane,
|
||||
@@ -146,7 +163,8 @@ func ChainPrepareLanes(chain ...blockbuster.Lane) blockbuster.PrepareLanesHandle
|
||||
|
||||
// Get the maximum number of bytes that can be included in the proposal for this lane.
|
||||
maxTxBytesForLane := utils.GetMaxTxBytesForLane(
|
||||
partialProposal,
|
||||
partialProposal.GetMaxTxBytes(),
|
||||
partialProposal.GetTotalTxBytes(),
|
||||
lane.GetMaxBlockSpace(),
|
||||
)
|
||||
|
||||
|
||||
@@ -26,8 +26,7 @@ import (
|
||||
|
||||
type ABCITestSuite struct {
|
||||
suite.Suite
|
||||
logger log.Logger
|
||||
ctx sdk.Context
|
||||
ctx sdk.Context
|
||||
|
||||
// Define basic tx configuration
|
||||
encodingConfig testutils.EncodingConfig
|
||||
@@ -333,8 +332,9 @@ func (suite *ABCITestSuite) TestPrepareProposal() {
|
||||
suite.Require().NoError(err)
|
||||
freeSize := int64(len(freeTxBytes))
|
||||
|
||||
maxTxBytes = tobSize + freeSize - 1
|
||||
maxTxBytes = tobSize*2 + freeSize - 1
|
||||
suite.tobConfig.MaxBlockSpace = sdk.ZeroDec()
|
||||
suite.freeConfig.MaxBlockSpace = sdk.MustNewDecFromStr("0.1")
|
||||
|
||||
txs = []sdk.Tx{freeTx}
|
||||
auctionTxs = []sdk.Tx{bidTx}
|
||||
@@ -388,7 +388,7 @@ func (suite *ABCITestSuite) TestPrepareProposal() {
|
||||
suite.Require().NoError(err)
|
||||
normalSize := int64(len(normalTxBytes))
|
||||
|
||||
maxTxBytes = tobSize + freeSize + normalSize + 1
|
||||
maxTxBytes = tobSize*2 + freeSize + normalSize + 1
|
||||
|
||||
// Tob can take up as much space as it wants
|
||||
suite.tobConfig.MaxBlockSpace = sdk.ZeroDec()
|
||||
@@ -693,7 +693,7 @@ func (suite *ABCITestSuite) TestPrepareProposal() {
|
||||
}
|
||||
|
||||
// Create a new proposal handler
|
||||
suite.proposalHandler = abci.NewProposalHandler(suite.logger, suite.encodingConfig.TxConfig.TxDecoder(), suite.mempool)
|
||||
suite.proposalHandler = abci.NewProposalHandler(log.NewNopLogger(), suite.encodingConfig.TxConfig.TxDecoder(), suite.mempool)
|
||||
handler := suite.proposalHandler.PrepareProposalHandler()
|
||||
res := handler(suite.ctx, abcitypes.RequestPrepareProposal{
|
||||
MaxTxBytes: maxTxBytes,
|
||||
|
||||
+2
-43
@@ -1,8 +1,6 @@
|
||||
package blockbuster
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
|
||||
"github.com/cometbft/cometbft/libs/log"
|
||||
@@ -11,25 +9,10 @@ import (
|
||||
)
|
||||
|
||||
type (
|
||||
// Proposal defines a block proposal type.
|
||||
Proposal struct {
|
||||
// Txs is the list of transactions in the proposal.
|
||||
Txs [][]byte
|
||||
|
||||
// Cache is a cache of the selected transactions in the proposal.
|
||||
Cache map[string]struct{}
|
||||
|
||||
// TotalTxBytes is the total number of bytes currently included in the proposal.
|
||||
TotalTxBytes int64
|
||||
|
||||
// MaxTxBytes is the maximum number of bytes that can be included in the proposal.
|
||||
MaxTxBytes int64
|
||||
}
|
||||
|
||||
// PrepareLanesHandler wraps all of the lanes Prepare function into a single chained
|
||||
// function. You can think of it like an AnteHandler, but for preparing proposals in the
|
||||
// context of lanes instead of modules.
|
||||
PrepareLanesHandler func(ctx sdk.Context, proposal *Proposal) *Proposal
|
||||
PrepareLanesHandler func(ctx sdk.Context, proposal BlockProposal) (BlockProposal, error)
|
||||
|
||||
// ProcessLanesHandler wraps all of the lanes Process functions into a single chained
|
||||
// function. You can think of it like an AnteHandler, but for processing proposals in the
|
||||
@@ -78,7 +61,7 @@ type (
|
||||
// included in the proposal for the given lane, the partial proposal, and a function
|
||||
// to call the next lane in the chain. The next lane in the chain will be called with
|
||||
// the updated proposal and context.
|
||||
PrepareLane(ctx sdk.Context, proposal *Proposal, maxTxBytes int64, next PrepareLanesHandler) *Proposal
|
||||
PrepareLane(ctx sdk.Context, proposal BlockProposal, maxTxBytes int64, next PrepareLanesHandler) (BlockProposal, error)
|
||||
|
||||
// ProcessLaneBasic validates that transactions belonging to this lane are not misplaced
|
||||
// in the block proposal.
|
||||
@@ -131,27 +114,3 @@ func (c *BaseLaneConfig) ValidateBasic() error {
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// NewProposal returns a new empty proposal.
|
||||
func NewProposal(maxTxBytes int64) *Proposal {
|
||||
return &Proposal{
|
||||
Txs: make([][]byte, 0),
|
||||
Cache: make(map[string]struct{}),
|
||||
MaxTxBytes: maxTxBytes,
|
||||
}
|
||||
}
|
||||
|
||||
// UpdateProposal updates the proposal with the given transactions and total size.
|
||||
func (p *Proposal) UpdateProposal(txs [][]byte, totalSize int64) *Proposal {
|
||||
p.TotalTxBytes += totalSize
|
||||
p.Txs = append(p.Txs, txs...)
|
||||
|
||||
for _, tx := range txs {
|
||||
txHash := sha256.Sum256(tx)
|
||||
txHashStr := hex.EncodeToString(txHash[:])
|
||||
|
||||
p.Cache[txHashStr] = struct{}{}
|
||||
}
|
||||
|
||||
return p
|
||||
}
|
||||
|
||||
@@ -14,13 +14,12 @@ import (
|
||||
// will return an empty partial proposal if no valid bids are found.
|
||||
func (l *TOBLane) PrepareLane(
|
||||
ctx sdk.Context,
|
||||
proposal *blockbuster.Proposal,
|
||||
proposal blockbuster.BlockProposal,
|
||||
maxTxBytes int64,
|
||||
next blockbuster.PrepareLanesHandler,
|
||||
) *blockbuster.Proposal {
|
||||
) (blockbuster.BlockProposal, error) {
|
||||
// Define all of the info we need to select transactions for the partial proposal.
|
||||
var (
|
||||
totalSize int64
|
||||
txs [][]byte
|
||||
txsToRemove = make(map[sdk.Tx]struct{}, 0)
|
||||
)
|
||||
@@ -33,14 +32,14 @@ selectBidTxLoop:
|
||||
cacheCtx, write := ctx.CacheContext()
|
||||
tmpBidTx := bidTxIterator.Tx()
|
||||
|
||||
bidTxBz, txHash, err := utils.GetTxHashStr(l.Cfg.TxEncoder, tmpBidTx)
|
||||
bidTxBz, _, err := utils.GetTxHashStr(l.Cfg.TxEncoder, tmpBidTx)
|
||||
if err != nil {
|
||||
txsToRemove[tmpBidTx] = struct{}{}
|
||||
continue
|
||||
continue selectBidTxLoop
|
||||
}
|
||||
|
||||
// if the transaction is already in the (partial) block proposal, we skip it.
|
||||
if _, ok := proposal.Cache[txHash]; ok {
|
||||
if proposal.Contains(bidTxBz) {
|
||||
continue selectBidTxLoop
|
||||
}
|
||||
|
||||
@@ -71,14 +70,14 @@ selectBidTxLoop:
|
||||
continue selectBidTxLoop
|
||||
}
|
||||
|
||||
sdkTxBz, hash, err := utils.GetTxHashStr(l.Cfg.TxEncoder, sdkTx)
|
||||
sdkTxBz, _, err := utils.GetTxHashStr(l.Cfg.TxEncoder, sdkTx)
|
||||
if err != nil {
|
||||
txsToRemove[tmpBidTx] = struct{}{}
|
||||
continue selectBidTxLoop
|
||||
}
|
||||
|
||||
// if the transaction is already in the (partial) block proposal, we skip it.
|
||||
if _, ok := proposal.Cache[hash]; ok {
|
||||
if proposal.Contains(sdkTxBz) {
|
||||
continue selectBidTxLoop
|
||||
}
|
||||
|
||||
@@ -93,7 +92,6 @@ selectBidTxLoop:
|
||||
// update the total size selected thus far.
|
||||
txs = append(txs, bidTxBz)
|
||||
txs = append(txs, bundledTxBz...)
|
||||
totalSize = bidTxSize
|
||||
|
||||
// Write the cache context to the original context when we know we have a
|
||||
// valid top of block bundle.
|
||||
@@ -111,12 +109,15 @@ selectBidTxLoop:
|
||||
|
||||
// Remove all transactions that were invalid during the creation of the partial proposal.
|
||||
if err := utils.RemoveTxsFromLane(txsToRemove, l.Mempool); err != nil {
|
||||
l.Cfg.Logger.Error("failed to remove txs from mempool", "lane", l.Name(), "err", err)
|
||||
return proposal
|
||||
return proposal, err
|
||||
}
|
||||
|
||||
// Update the proposal with the selected transactions.
|
||||
proposal.UpdateProposal(txs, totalSize)
|
||||
// Update the proposal with the selected transactions. This will only return an error
|
||||
// if the invarient checks are not passed. In the case when this errors, the original proposal
|
||||
// will be returned (without the selected transactions from this lane).
|
||||
if err := proposal.UpdateProposal(l, txs); err != nil {
|
||||
return proposal, err
|
||||
}
|
||||
|
||||
return next(ctx, proposal)
|
||||
}
|
||||
|
||||
@@ -8,13 +8,15 @@ import (
|
||||
"github.com/skip-mev/pob/blockbuster/utils"
|
||||
)
|
||||
|
||||
// PrepareLane will prepare a partial proposal for the base lane.
|
||||
// PrepareLane will prepare a partial proposal for the default lane. It will select and include
|
||||
// all valid transactions in the mempool that are not already in the partial proposal.
|
||||
// The default lane orders transactions by the sdk.Context priority.
|
||||
func (l *DefaultLane) PrepareLane(
|
||||
ctx sdk.Context,
|
||||
proposal *blockbuster.Proposal,
|
||||
proposal blockbuster.BlockProposal,
|
||||
maxTxBytes int64,
|
||||
next blockbuster.PrepareLanesHandler,
|
||||
) *blockbuster.Proposal {
|
||||
) (blockbuster.BlockProposal, error) {
|
||||
// Define all of the info we need to select transactions for the partial proposal.
|
||||
var (
|
||||
totalSize int64
|
||||
@@ -27,14 +29,14 @@ func (l *DefaultLane) PrepareLane(
|
||||
for iterator := l.Mempool.Select(ctx, nil); iterator != nil; iterator = iterator.Next() {
|
||||
tx := iterator.Tx()
|
||||
|
||||
txBytes, hash, err := utils.GetTxHashStr(l.Cfg.TxEncoder, tx)
|
||||
txBytes, _, err := utils.GetTxHashStr(l.Cfg.TxEncoder, tx)
|
||||
if err != nil {
|
||||
txsToRemove[tx] = struct{}{}
|
||||
continue
|
||||
}
|
||||
|
||||
// if the transaction is already in the (partial) block proposal, we skip it.
|
||||
if _, ok := proposal.Cache[hash]; ok {
|
||||
if proposal.Contains(txBytes) {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -56,16 +58,22 @@ func (l *DefaultLane) PrepareLane(
|
||||
|
||||
// Remove all transactions that were invalid during the creation of the partial proposal.
|
||||
if err := utils.RemoveTxsFromLane(txsToRemove, l.Mempool); err != nil {
|
||||
l.Cfg.Logger.Error("failed to remove txs from mempool", "lane", l.Name(), "err", err)
|
||||
return proposal
|
||||
return proposal, err
|
||||
}
|
||||
|
||||
proposal.UpdateProposal(txs, totalSize)
|
||||
// Update the partial proposal with the selected transactions. If the proposal is unable to
|
||||
// be updated, we return an error. The proposal will only be modified if it passes all
|
||||
// of the invarient checks.
|
||||
if err := proposal.UpdateProposal(l, txs); err != nil {
|
||||
return proposal, err
|
||||
}
|
||||
|
||||
return next(ctx, proposal)
|
||||
}
|
||||
|
||||
// ProcessLane verifies the default lane's portion of a block proposal.
|
||||
// ProcessLane verifies the default lane's portion of a block proposal. Since the default lane's
|
||||
// ProcessLaneBasic function ensures that all of the default transactions are in the correct order,
|
||||
// we only need to verify the contiguous set of transactions that match to the default lane.
|
||||
func (l *DefaultLane) ProcessLane(ctx sdk.Context, txs []sdk.Tx, next blockbuster.ProcessLanesHandler) (sdk.Context, error) {
|
||||
for index, tx := range txs {
|
||||
if l.Match(tx) {
|
||||
@@ -81,8 +89,9 @@ func (l *DefaultLane) ProcessLane(ctx sdk.Context, txs []sdk.Tx, next blockbuste
|
||||
return ctx, nil
|
||||
}
|
||||
|
||||
// ProcessLaneBasic does basic validation on the block proposal to ensure that
|
||||
// transactions that belong to this lane are not misplaced in the block proposal.
|
||||
// transactions that belong to this lane are not misplaced in the block proposal i.e.
|
||||
// the proposal only contains contiguous transactions that belong to this lane - there
|
||||
// can be no interleaving of transactions from other lanes.
|
||||
func (l *DefaultLane) ProcessLaneBasic(txs []sdk.Tx) error {
|
||||
seenOtherLaneTx := false
|
||||
lastSeenIndex := 0
|
||||
|
||||
@@ -13,8 +13,12 @@ const (
|
||||
|
||||
var _ blockbuster.Lane = (*DefaultLane)(nil)
|
||||
|
||||
// DefaultLane defines a default lane implementation. It contains a priority-nonce
|
||||
// index along with core lane functionality.
|
||||
// DefaultLane defines a default lane implementation. The default lane orders
|
||||
// transactions by the sdk.Context priority. The default lane will accept any
|
||||
// transaction that is not a part of the lane's IgnoreList. By default, the IgnoreList
|
||||
// is empty and the default lane will accept any transaction. The default lane on its
|
||||
// own implements the same functionality as the pre v0.47.0 tendermint mempool and proposal
|
||||
// handlers.
|
||||
type DefaultLane struct {
|
||||
// Mempool defines the mempool for the lane.
|
||||
Mempool
|
||||
|
||||
@@ -40,6 +40,8 @@ type (
|
||||
}
|
||||
)
|
||||
|
||||
// NewDefaultMempool returns a new default mempool instance. The default mempool
|
||||
// orders transactions by the sdk.Context priority.
|
||||
func NewDefaultMempool(txEncoder sdk.TxEncoder) *DefaultMempool {
|
||||
return &DefaultMempool{
|
||||
index: blockbuster.NewPriorityMempool(
|
||||
|
||||
@@ -35,8 +35,8 @@ type Terminator struct{}
|
||||
var _ blockbuster.Lane = (*Terminator)(nil)
|
||||
|
||||
// PrepareLane is a no-op
|
||||
func (t Terminator) PrepareLane(_ sdk.Context, proposal *blockbuster.Proposal, _ int64, _ blockbuster.PrepareLanesHandler) *blockbuster.Proposal {
|
||||
return proposal
|
||||
func (t Terminator) PrepareLane(_ sdk.Context, proposal blockbuster.BlockProposal, _ int64, _ blockbuster.PrepareLanesHandler) (blockbuster.BlockProposal, error) {
|
||||
return proposal, nil
|
||||
}
|
||||
|
||||
// ProcessLane is a no-op
|
||||
|
||||
+44
-2
@@ -38,13 +38,28 @@ type (
|
||||
}
|
||||
)
|
||||
|
||||
// NewMempool returns a new Blockbuster mempool. The blockbuster mempool is
|
||||
// comprised of a registry of lanes. Each lane is responsible for selecting
|
||||
// transactions according to its own selection logic. The lanes are ordered
|
||||
// according to their priority. The first lane in the registry has the highest
|
||||
// priority. Proposals are verified according to the order of the lanes in the
|
||||
// registry. Basic mempool API, such as insertion, removal, and contains, are
|
||||
// delegated to the first lane that matches the transaction. Each transaction
|
||||
// should only belong in one lane.
|
||||
func NewMempool(lanes ...Lane) *BBMempool {
|
||||
return &BBMempool{
|
||||
mempool := &BBMempool{
|
||||
registry: lanes,
|
||||
}
|
||||
|
||||
if err := mempool.ValidateBasic(); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
return mempool
|
||||
}
|
||||
|
||||
// CountTx returns the total number of transactions in the mempool.
|
||||
// CountTx returns the total number of transactions in the mempool. This will
|
||||
// be the sum of the number of transactions in each lane.
|
||||
func (m *BBMempool) CountTx() int {
|
||||
var total int
|
||||
for _, lane := range m.registry {
|
||||
@@ -125,6 +140,33 @@ func (m *BBMempool) Registry() []Lane {
|
||||
return m.registry
|
||||
}
|
||||
|
||||
// ValidateBasic validates the mempools configuration.
|
||||
func (m *BBMempool) ValidateBasic() error {
|
||||
sum := sdk.ZeroDec()
|
||||
seenZeroMaxBlockSpace := false
|
||||
|
||||
for _, lane := range m.registry {
|
||||
maxBlockSpace := lane.GetMaxBlockSpace()
|
||||
if maxBlockSpace.IsZero() {
|
||||
seenZeroMaxBlockSpace = true
|
||||
}
|
||||
|
||||
sum = sum.Add(lane.GetMaxBlockSpace())
|
||||
}
|
||||
|
||||
switch {
|
||||
// Ensure that the sum of the lane max block space percentages is less than
|
||||
// or equal to 1.
|
||||
case sum.GT(sdk.OneDec()):
|
||||
return fmt.Errorf("sum of lane max block space percentages must be less than or equal to 1, got %s", sum)
|
||||
// Ensure that there is no unused block space.
|
||||
case sum.LT(sdk.OneDec()) && !seenZeroMaxBlockSpace:
|
||||
return fmt.Errorf("sum of total block space percentages will be less than 1")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetLane returns the lane with the given name.
|
||||
func (m *BBMempool) GetLane(name string) (Lane, error) {
|
||||
for _, lane := range m.registry {
|
||||
|
||||
@@ -0,0 +1,197 @@
|
||||
package blockbuster
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
|
||||
"github.com/cometbft/cometbft/libs/log"
|
||||
sdk "github.com/cosmos/cosmos-sdk/types"
|
||||
"github.com/skip-mev/pob/blockbuster/utils"
|
||||
)
|
||||
|
||||
var _ BlockProposal = (*Proposal)(nil)
|
||||
|
||||
type (
|
||||
// LaneProposal defines the interface/APIs that are required for the proposal to interact
|
||||
// with a lane.
|
||||
LaneProposal interface {
|
||||
// Logger returns the lane's logger.
|
||||
Logger() log.Logger
|
||||
|
||||
// GetMaxBlockSpace returns the maximum block space for the lane as a relative percentage.
|
||||
GetMaxBlockSpace() sdk.Dec
|
||||
|
||||
// Name returns the name of the lane.
|
||||
Name() string
|
||||
}
|
||||
|
||||
// BlockProposal is the interface/APIs that are required for proposal creation + interacting with
|
||||
// and updating proposals. BlockProposals are iteratively updated as each lane prepares its
|
||||
// partial proposal. Each lane must call UpdateProposal with its partial proposal in PrepareLane. BlockProposals
|
||||
// can also include vote extensions, which are included at the top of the proposal.
|
||||
BlockProposal interface {
|
||||
// UpdateProposal updates the proposal with the given transactions. There are a
|
||||
// few invarients that are checked:
|
||||
// 1. The total size of the proposal must be less than the maximum number of bytes allowed.
|
||||
// 2. The total size of the partial proposal must be less than the maximum number of bytes allowed for
|
||||
// the lane.
|
||||
UpdateProposal(lane LaneProposal, partialProposalTxs [][]byte) error
|
||||
|
||||
// GetMaxTxBytes returns the maximum number of bytes that can be included in the proposal.
|
||||
GetMaxTxBytes() int64
|
||||
|
||||
// GetTotalTxBytes returns the total number of bytes currently included in the proposal.
|
||||
GetTotalTxBytes() int64
|
||||
|
||||
// GetTxs returns the transactions in the proposal.
|
||||
GetTxs() [][]byte
|
||||
|
||||
// GetNumTxs returns the number of transactions in the proposal.
|
||||
GetNumTxs() int
|
||||
|
||||
// Contains returns true if the proposal contains the given transaction.
|
||||
Contains(tx []byte) bool
|
||||
|
||||
// AddVoteExtension adds a vote extension to the proposal.
|
||||
AddVoteExtension(voteExtension []byte)
|
||||
|
||||
// GetVoteExtensions returns the vote extensions in the proposal.
|
||||
GetVoteExtensions() [][]byte
|
||||
|
||||
// GetProposal returns all of the transactions in the proposal along with the vote extensions
|
||||
// at the top of the proposal.
|
||||
GetProposal() [][]byte
|
||||
}
|
||||
|
||||
// Proposal defines a block proposal type.
|
||||
Proposal struct {
|
||||
// txs is the list of transactions in the proposal.
|
||||
txs [][]byte
|
||||
|
||||
// voteExtensions is the list of vote extensions in the proposal.
|
||||
voteExtensions [][]byte
|
||||
|
||||
// cache is a cache of the selected transactions in the proposal.
|
||||
cache map[string]struct{}
|
||||
|
||||
// totalTxBytes is the total number of bytes currently included in the proposal.
|
||||
totalTxBytes int64
|
||||
|
||||
// maxTxBytes is the maximum number of bytes that can be included in the proposal.
|
||||
maxTxBytes int64
|
||||
}
|
||||
)
|
||||
|
||||
// NewProposal returns a new empty proposal.
|
||||
func NewProposal(maxTxBytes int64) *Proposal {
|
||||
return &Proposal{
|
||||
txs: make([][]byte, 0),
|
||||
voteExtensions: make([][]byte, 0),
|
||||
cache: make(map[string]struct{}),
|
||||
maxTxBytes: maxTxBytes,
|
||||
}
|
||||
}
|
||||
|
||||
// UpdateProposal updates the proposal with the given transactions and total size. There are a
|
||||
// few invarients that are checked:
|
||||
// 1. The total size of the proposal must be less than the maximum number of bytes allowed.
|
||||
// 2. The total size of the partial proposal must be less than the maximum number of bytes allowed for
|
||||
// the lane.
|
||||
func (p *Proposal) UpdateProposal(lane LaneProposal, partialProposalTxs [][]byte) error {
|
||||
if len(partialProposalTxs) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
partialProposalSize := int64(0)
|
||||
for _, tx := range partialProposalTxs {
|
||||
partialProposalSize += int64(len(tx))
|
||||
}
|
||||
|
||||
// Invarient check: Ensure that the lane did not prepare a partial proposal that is too large.
|
||||
maxTxBytesForLane := utils.GetMaxTxBytesForLane(p.GetMaxTxBytes(), p.GetTotalTxBytes(), lane.GetMaxBlockSpace())
|
||||
if partialProposalSize > maxTxBytesForLane {
|
||||
return fmt.Errorf(
|
||||
"%s lane prepared a partial proposal that is too large: %d > %d",
|
||||
lane.Name(),
|
||||
partialProposalSize,
|
||||
maxTxBytesForLane,
|
||||
)
|
||||
}
|
||||
|
||||
// Invarient check: Ensure that the lane did not prepare a block proposal that is too large.
|
||||
updatedSize := p.totalTxBytes + partialProposalSize
|
||||
if updatedSize > p.maxTxBytes {
|
||||
return fmt.Errorf(
|
||||
"lane %s prepared a block proposal that is too large: %d > %d",
|
||||
lane.Name(),
|
||||
p.totalTxBytes,
|
||||
p.maxTxBytes,
|
||||
)
|
||||
}
|
||||
p.totalTxBytes = updatedSize
|
||||
|
||||
lane.Logger().Info(
|
||||
"adding transactions to proposal",
|
||||
"lane", lane.Name(),
|
||||
"num_txs", len(partialProposalTxs),
|
||||
"total_tx_bytes", partialProposalSize,
|
||||
"cumulative_size", updatedSize,
|
||||
)
|
||||
|
||||
p.txs = append(p.txs, partialProposalTxs...)
|
||||
|
||||
for _, tx := range partialProposalTxs {
|
||||
txHash := sha256.Sum256(tx)
|
||||
txHashStr := hex.EncodeToString(txHash[:])
|
||||
|
||||
p.cache[txHashStr] = struct{}{}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetProposal returns all of the transactions in the proposal along with the vote extensions
|
||||
// at the top of the proposal.
|
||||
func (p *Proposal) GetProposal() [][]byte {
|
||||
return append(p.voteExtensions, p.txs...)
|
||||
}
|
||||
|
||||
// AddVoteExtension adds a vote extension to the proposal.
|
||||
func (p *Proposal) AddVoteExtension(voteExtension []byte) {
|
||||
p.voteExtensions = append(p.voteExtensions, voteExtension)
|
||||
}
|
||||
|
||||
// GetVoteExtensions returns the vote extensions in the proposal.
|
||||
func (p *Proposal) GetVoteExtensions() [][]byte {
|
||||
return p.voteExtensions
|
||||
}
|
||||
|
||||
// GetMaxTxBytes returns the maximum number of bytes that can be included in the proposal.
|
||||
func (p *Proposal) GetMaxTxBytes() int64 {
|
||||
return p.maxTxBytes
|
||||
}
|
||||
|
||||
// GetTotalTxBytes returns the total number of bytes currently included in the proposal.
|
||||
func (p *Proposal) GetTotalTxBytes() int64 {
|
||||
return p.totalTxBytes
|
||||
}
|
||||
|
||||
// GetTxs returns the transactions in the proposal.
|
||||
func (p *Proposal) GetTxs() [][]byte {
|
||||
return p.txs
|
||||
}
|
||||
|
||||
// GetNumTxs returns the number of transactions in the proposal.
|
||||
func (p *Proposal) GetNumTxs() int {
|
||||
return len(p.txs)
|
||||
}
|
||||
|
||||
// Contains returns true if the proposal contains the given transaction.
|
||||
func (p *Proposal) Contains(tx []byte) bool {
|
||||
txHash := sha256.Sum256(tx)
|
||||
txHashStr := hex.EncodeToString(txHash[:])
|
||||
|
||||
_, ok := p.cache[txHashStr]
|
||||
return ok
|
||||
}
|
||||
@@ -7,7 +7,6 @@ import (
|
||||
|
||||
sdk "github.com/cosmos/cosmos-sdk/types"
|
||||
sdkmempool "github.com/cosmos/cosmos-sdk/types/mempool"
|
||||
"github.com/skip-mev/pob/blockbuster"
|
||||
)
|
||||
|
||||
// GetTxHashStr returns the hex-encoded hash of the transaction alongside the
|
||||
@@ -52,13 +51,13 @@ func RemoveTxsFromLane(txs map[sdk.Tx]struct{}, mempool sdkmempool.Mempool) erro
|
||||
|
||||
// GetMaxTxBytesForLane returns the maximum number of bytes that can be included in the proposal
|
||||
// for the given lane.
|
||||
func GetMaxTxBytesForLane(proposal *blockbuster.Proposal, ratio sdk.Dec) int64 {
|
||||
func GetMaxTxBytesForLane(maxTxBytes, totalTxBytes int64, ratio sdk.Dec) int64 {
|
||||
// In the case where the ratio is zero, we return the max tx bytes remaining. Note, the only
|
||||
// lane that should have a ratio of zero is the default lane. This means the default lane
|
||||
// will have no limit on the number of transactions it can include in a block and is only
|
||||
// limited by the maxTxBytes included in the PrepareProposalRequest.
|
||||
if ratio.IsZero() {
|
||||
remainder := proposal.MaxTxBytes - proposal.TotalTxBytes
|
||||
remainder := maxTxBytes - totalTxBytes
|
||||
if remainder < 0 {
|
||||
return 0
|
||||
}
|
||||
@@ -67,5 +66,5 @@ func GetMaxTxBytesForLane(proposal *blockbuster.Proposal, ratio sdk.Dec) int64 {
|
||||
}
|
||||
|
||||
// Otherwise, we calculate the max tx bytes for the lane based on the ratio.
|
||||
return ratio.MulInt64(proposal.MaxTxBytes).TruncateInt().Int64()
|
||||
return ratio.MulInt64(maxTxBytes).TruncateInt().Int64()
|
||||
}
|
||||
|
||||
@@ -1,79 +1,66 @@
|
||||
package utils
|
||||
package utils_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
sdk "github.com/cosmos/cosmos-sdk/types"
|
||||
"github.com/skip-mev/pob/blockbuster"
|
||||
"github.com/skip-mev/pob/blockbuster/utils"
|
||||
)
|
||||
|
||||
func TestGetMaxTxBytesForLane(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
proposal *blockbuster.Proposal
|
||||
ratio sdk.Dec
|
||||
expected int64
|
||||
name string
|
||||
maxTxBytes int64
|
||||
totalTxBytes int64
|
||||
ratio sdk.Dec
|
||||
expected int64
|
||||
}{
|
||||
{
|
||||
"ratio is zero",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 100,
|
||||
TotalTxBytes: 50,
|
||||
},
|
||||
100,
|
||||
50,
|
||||
sdk.ZeroDec(),
|
||||
50,
|
||||
},
|
||||
{
|
||||
"ratio is zero",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 100,
|
||||
TotalTxBytes: 100,
|
||||
},
|
||||
100,
|
||||
100,
|
||||
sdk.ZeroDec(),
|
||||
0,
|
||||
},
|
||||
{
|
||||
"ratio is zero",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 100,
|
||||
TotalTxBytes: 150,
|
||||
},
|
||||
100,
|
||||
150,
|
||||
sdk.ZeroDec(),
|
||||
0,
|
||||
},
|
||||
{
|
||||
"ratio is 1",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 100,
|
||||
TotalTxBytes: 50,
|
||||
},
|
||||
100,
|
||||
50,
|
||||
sdk.OneDec(),
|
||||
100,
|
||||
},
|
||||
{
|
||||
"ratio is 10%",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 100,
|
||||
TotalTxBytes: 50,
|
||||
},
|
||||
100,
|
||||
50,
|
||||
sdk.MustNewDecFromStr("0.1"),
|
||||
10,
|
||||
},
|
||||
{
|
||||
"ratio is 25%",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 100,
|
||||
TotalTxBytes: 50,
|
||||
},
|
||||
100,
|
||||
50,
|
||||
sdk.MustNewDecFromStr("0.25"),
|
||||
25,
|
||||
},
|
||||
{
|
||||
"ratio is 50%",
|
||||
&blockbuster.Proposal{
|
||||
MaxTxBytes: 101,
|
||||
TotalTxBytes: 50,
|
||||
},
|
||||
101,
|
||||
50,
|
||||
sdk.MustNewDecFromStr("0.5"),
|
||||
50,
|
||||
},
|
||||
@@ -81,7 +68,7 @@ func TestGetMaxTxBytesForLane(t *testing.T) {
|
||||
|
||||
for _, tc := range testCases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
actual := GetMaxTxBytesForLane(tc.proposal, tc.ratio)
|
||||
actual := utils.GetMaxTxBytesForLane(tc.maxTxBytes, tc.totalTxBytes, tc.ratio)
|
||||
if actual != tc.expected {
|
||||
t.Errorf("expected %d, got %d", tc.expected, actual)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user