feat: backport unordered transactions (#23708)
Co-authored-by: Aleksandr Bezobchuk <alexanderbez@users.noreply.github.com> Co-authored-by: yihuang <huang@crypto.com> Co-authored-by: Facundo <facundomedica@gmail.com> Co-authored-by: Facundo Medica <14063057+facundomedica@users.noreply.github.com> Co-authored-by: Alex | Interchain Labs <alex@interchainlabs.io>
This commit is contained in:
co-authored by
Aleksandr Bezobchuk
yihuang
Facundo
Facundo Medica
Alex | Interchain Labs
parent
ff779eca8d
commit
7f7c41e4aa
@@ -4,6 +4,7 @@ import (
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
cmtproto "github.com/cometbft/cometbft/proto/tendermint/types"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -52,6 +53,21 @@ type testTx struct {
|
||||
address sdk.AccAddress
|
||||
// useful for debugging
|
||||
strAddress string
|
||||
unordered bool
|
||||
timeout *time.Time
|
||||
}
|
||||
|
||||
// GetTimeoutTimeStamp implements types.TxWithUnordered.
|
||||
func (tx testTx) GetTimeoutTimeStamp() time.Time {
|
||||
if tx.timeout == nil {
|
||||
return time.Time{}
|
||||
}
|
||||
return *tx.timeout
|
||||
}
|
||||
|
||||
// GetUnordered implements types.TxWithUnordered.
|
||||
func (tx testTx) GetUnordered() bool {
|
||||
return tx.unordered
|
||||
}
|
||||
|
||||
func (tx testTx) GetSigners() ([][]byte, error) { panic("not implemented") }
|
||||
|
||||
@@ -2,6 +2,7 @@ package mempool
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math"
|
||||
"sync"
|
||||
@@ -222,6 +223,16 @@ func (mp *PriorityNonceMempool[C]) Insert(ctx context.Context, tx sdk.Tx) error
|
||||
sender := sig.Signer.String()
|
||||
priority := mp.cfg.TxPriority.GetTxPriority(ctx, tx)
|
||||
nonce := sig.Sequence
|
||||
|
||||
// if it's an unordered tx, we use the timeout timestamp instead of the nonce
|
||||
if unordered, ok := tx.(sdk.TxWithUnordered); ok && unordered.GetUnordered() {
|
||||
timestamp := unordered.GetTimeoutTimeStamp().Unix()
|
||||
if timestamp < 0 {
|
||||
return errors.New("invalid timestamp value")
|
||||
}
|
||||
nonce = uint64(timestamp)
|
||||
}
|
||||
|
||||
key := txMeta[C]{nonce: nonce, priority: priority, sender: sender}
|
||||
|
||||
senderIndex, ok := mp.senderIndices[sender]
|
||||
@@ -458,6 +469,15 @@ func (mp *PriorityNonceMempool[C]) Remove(tx sdk.Tx) error {
|
||||
sender := sig.Signer.String()
|
||||
nonce := sig.Sequence
|
||||
|
||||
// if it's an unordered tx, we use the timeout timestamp instead of the nonce
|
||||
if unordered, ok := tx.(sdk.TxWithUnordered); ok && unordered.GetUnordered() {
|
||||
timestamp := unordered.GetTimeoutTimeStamp().Unix()
|
||||
if timestamp < 0 {
|
||||
return errors.New("invalid timestamp value")
|
||||
}
|
||||
nonce = uint64(timestamp)
|
||||
}
|
||||
|
||||
scoreKey := txMeta[C]{nonce: nonce, sender: sender}
|
||||
score, ok := mp.scores[scoreKey]
|
||||
if !ok {
|
||||
|
||||
@@ -976,3 +976,40 @@ func TestNextSenderTx_TxReplacement(t *testing.T) {
|
||||
iter := mp.Select(ctx, nil)
|
||||
require.Equal(t, txs[3], iter.Tx())
|
||||
}
|
||||
|
||||
func TestPriorityNonceMempool_UnorderedTx(t *testing.T) {
|
||||
ctx := sdk.NewContext(nil, cmtproto.Header{}, false, log.NewNopLogger())
|
||||
accounts := simtypes.RandomAccounts(rand.New(rand.NewSource(0)), 2)
|
||||
sa := accounts[0].Address
|
||||
sb := accounts[1].Address
|
||||
|
||||
mp := mempool.DefaultPriorityMempool()
|
||||
|
||||
now := time.Now()
|
||||
oneHour := now.Add(1 * time.Hour)
|
||||
thirtyMin := now.Add(30 * time.Minute)
|
||||
twoHours := now.Add(2 * time.Hour)
|
||||
fifteenMin := now.Add(15 * time.Minute)
|
||||
|
||||
txs := []testTx{
|
||||
{id: 1, priority: 0, address: sa, timeout: &thirtyMin, unordered: true},
|
||||
{id: 0, priority: 0, address: sa, timeout: &oneHour, unordered: true},
|
||||
{id: 3, priority: 0, address: sb, timeout: &fifteenMin, unordered: true},
|
||||
{id: 2, priority: 0, address: sb, timeout: &twoHours, unordered: true},
|
||||
}
|
||||
|
||||
for _, tx := range txs {
|
||||
c := ctx.WithPriority(tx.priority)
|
||||
require.NoError(t, mp.Insert(c, tx))
|
||||
}
|
||||
|
||||
require.Equal(t, 4, mp.CountTx())
|
||||
|
||||
orderedTxs := fetchTxs(mp.Select(ctx, nil), 100000)
|
||||
require.Equal(t, len(txs), len(orderedTxs))
|
||||
|
||||
// check order
|
||||
for i, tx := range orderedTxs {
|
||||
require.Equal(t, txs[i].id, tx.(testTx).id)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
crand "crypto/rand" // #nosec // crypto/rand is used for seed generation
|
||||
"encoding/binary"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math/rand" // #nosec // math/rand is used for random selection and seeded from crypto/rand
|
||||
"sync"
|
||||
@@ -139,6 +140,15 @@ func (snm *SenderNonceMempool) Insert(_ context.Context, tx sdk.Tx) error {
|
||||
sender := sdk.AccAddress(sig.PubKey.Address()).String()
|
||||
nonce := sig.Sequence
|
||||
|
||||
// if it's an unordered tx, we use the timeout timestamp instead of the nonce
|
||||
if unordered, ok := tx.(sdk.TxWithUnordered); ok && unordered.GetUnordered() {
|
||||
timestamp := unordered.GetTimeoutTimeStamp().Unix()
|
||||
if timestamp < 0 {
|
||||
return errors.New("invalid timestamp value")
|
||||
}
|
||||
nonce = uint64(timestamp)
|
||||
}
|
||||
|
||||
senderTxs, found := snm.senders[sender]
|
||||
if !found {
|
||||
senderTxs = skiplist.New(skiplist.Uint64)
|
||||
@@ -227,6 +237,15 @@ func (snm *SenderNonceMempool) Remove(tx sdk.Tx) error {
|
||||
sender := sdk.AccAddress(sig.PubKey.Address()).String()
|
||||
nonce := sig.Sequence
|
||||
|
||||
// if it's an unordered tx, we use the timeout timestamp instead of the nonce
|
||||
if unordered, ok := tx.(sdk.TxWithUnordered); ok && unordered.GetUnordered() {
|
||||
timestamp := unordered.GetTimeoutTimeStamp().Unix()
|
||||
if timestamp < 0 {
|
||||
return errors.New("invalid timestamp value")
|
||||
}
|
||||
nonce = uint64(timestamp)
|
||||
}
|
||||
|
||||
senderTxs, found := snm.senders[sender]
|
||||
if !found {
|
||||
return ErrTxNotFound
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
cmtproto "github.com/cometbft/cometbft/proto/tendermint/types"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -193,3 +194,67 @@ func (s *MempoolTestSuite) TestTxNotFoundOnSender() {
|
||||
err = mp.Remove(tx)
|
||||
require.Equal(t, mempool.ErrTxNotFound, err)
|
||||
}
|
||||
|
||||
func (s *MempoolTestSuite) TestUnorderedTx() {
|
||||
t := s.T()
|
||||
|
||||
ctx := sdk.NewContext(nil, cmtproto.Header{}, false, log.NewNopLogger())
|
||||
accounts := simtypes.RandomAccounts(rand.New(rand.NewSource(0)), 2)
|
||||
sa := accounts[0].Address
|
||||
sb := accounts[1].Address
|
||||
|
||||
mp := mempool.NewSenderNonceMempool(mempool.SenderNonceMaxTxOpt(5000))
|
||||
|
||||
now := time.Now()
|
||||
oneHour := now.Add(1 * time.Hour)
|
||||
thirtyMin := now.Add(30 * time.Minute)
|
||||
twoHours := now.Add(2 * time.Hour)
|
||||
fifteenMin := now.Add(15 * time.Minute)
|
||||
|
||||
txs := []testTx{
|
||||
{id: 0, address: sa, timeout: &oneHour, unordered: true},
|
||||
{id: 1, address: sa, timeout: &thirtyMin, unordered: true},
|
||||
{id: 2, address: sb, timeout: &twoHours, unordered: true},
|
||||
{id: 3, address: sb, timeout: &fifteenMin, unordered: true},
|
||||
}
|
||||
|
||||
for _, tx := range txs {
|
||||
c := ctx.WithPriority(tx.priority)
|
||||
require.NoError(t, mp.Insert(c, tx))
|
||||
}
|
||||
|
||||
require.Equal(t, 4, mp.CountTx())
|
||||
|
||||
orderedTxs := fetchTxs(mp.Select(ctx, nil), 100000)
|
||||
require.Equal(t, len(txs), len(orderedTxs))
|
||||
|
||||
// Because the sender is selected randomly it can be any of these options
|
||||
acceptableOptions := [][]int{
|
||||
{3, 1, 2, 0},
|
||||
{3, 1, 0, 2},
|
||||
{3, 2, 1, 0},
|
||||
{1, 3, 0, 2},
|
||||
{1, 3, 2, 0},
|
||||
{1, 0, 3, 2},
|
||||
}
|
||||
|
||||
orderedTxsIds := make([]int, len(orderedTxs))
|
||||
for i, tx := range orderedTxs {
|
||||
orderedTxsIds[i] = tx.(testTx).id
|
||||
}
|
||||
|
||||
anyAcceptableOrder := false
|
||||
for _, option := range acceptableOptions {
|
||||
for i, tx := range orderedTxs {
|
||||
if tx.(testTx).id != txs[option[i]].id {
|
||||
break
|
||||
}
|
||||
|
||||
if i == len(orderedTxs)-1 {
|
||||
anyAcceptableOrder = true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
require.True(t, anyAcceptableOrder, "expected any of %v but got %v", acceptableOptions, orderedTxsIds)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user