eth/fetcher: separate hasTx from chainTx and refactor test code

Signed-off-by: Csaba Kiraly <csaba.kiraly@gmail.com>
This commit is contained in:
Csaba Kiraly 2025-12-12 09:49:05 +01:00
parent f5ff0e837d
commit 259705a9e8
No known key found for this signature in database
GPG key ID: 0FE274EE8C95166E
3 changed files with 79 additions and 270 deletions

View file

@ -320,7 +320,6 @@ func (f *TxFetcher) Enqueue(peer string, txs []*types.Transaction, direct bool)
otherreject int64
)
batch := txs[i:end]
for j, err := range f.addTxs(batch) {
// Track the transaction hash if the price is too low for us.
// Avoid re-request this transaction when we receive another

View file

@ -83,6 +83,17 @@ type txFetcherTest struct {
steps []interface{}
}
func initDefaultTxFetcher() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
}
// Tests that transaction announcements with associated metadata are added to a
// waitlist, and none of them are scheduled for retrieval until the wait expires.
//
@ -91,14 +102,7 @@ type txFetcherTest struct {
// with all the useless extra fields.
func TestTransactionFetcherWaiting(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Initial announcement to get something into the waitlist
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}, {0x02}}, types: []byte{types.LegacyTxType, types.LegacyTxType}, sizes: []uint32{111, 222}},
@ -293,14 +297,7 @@ func TestTransactionFetcherWaiting(t *testing.T) {
// already scheduled.
func TestTransactionFetcherSkipWaiting(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{
@ -383,14 +380,7 @@ func TestTransactionFetcherSkipWaiting(t *testing.T) {
// and subsequent announces block or get allotted to someone else.
func TestTransactionFetcherSingletonRequesting(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}, {0x02}}, types: []byte{types.LegacyTxType, types.LegacyTxType}, sizes: []uint32{111, 222}},
@ -489,15 +479,12 @@ func TestTransactionFetcherFailedRescheduling(t *testing.T) {
proceed := make(chan struct{})
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(origin string, hashes []common.Hash) error {
<-proceed
return errors.New("peer disconnected")
},
nil,
)
f := initDefaultTxFetcher()
f.fetchTxs = func(origin string, hashes []common.Hash) error {
<-proceed
return errors.New("peer disconnected")
}
return f
},
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
@ -572,16 +559,7 @@ func TestTransactionFetcherFailedRescheduling(t *testing.T) {
// are cleaned up.
func TestTransactionFetcherCleanup(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{peer: "A", hashes: []common.Hash{testTxsHashes[0]}, types: []byte{testTxs[0].Type()}, sizes: []uint32{uint32(testTxs[0].Size())}},
@ -616,16 +594,7 @@ func TestTransactionFetcherCleanup(t *testing.T) {
// this was a bug)).
func TestTransactionFetcherCleanupEmpty(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{peer: "A", hashes: []common.Hash{testTxsHashes[0]}, types: []byte{testTxs[0].Type()}, sizes: []uint32{uint32(testTxs[0].Size())}},
@ -659,16 +628,7 @@ func TestTransactionFetcherCleanupEmpty(t *testing.T) {
// different peer, or self if they are after the cutoff point.
func TestTransactionFetcherMissingRescheduling(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{peer: "A",
@ -720,16 +680,7 @@ func TestTransactionFetcherMissingRescheduling(t *testing.T) {
// delivered, the peer gets properly cleaned out from the internal state.
func TestTransactionFetcherMissingCleanup(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{peer: "A",
@ -769,16 +720,7 @@ func TestTransactionFetcherMissingCleanup(t *testing.T) {
// Tests that transaction broadcasts properly clean up announcements.
func TestTransactionFetcherBroadcasts(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Set up three transactions to be in different stats, waiting, queued and fetching
doTxNotify{peer: "A", hashes: []common.Hash{testTxsHashes[0]}, types: []byte{testTxs[0].Type()}, sizes: []uint32{uint32(testTxs[0].Size())}},
@ -825,14 +767,7 @@ func TestTransactionFetcherBroadcasts(t *testing.T) {
// Tests that the waiting list timers properly reset and reschedule.
func TestTransactionFetcherWaitTimerResets(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}}, types: []byte{types.LegacyTxType}, sizes: []uint32{111}},
isWaiting(map[string][]announce{
@ -895,16 +830,7 @@ func TestTransactionFetcherWaitTimerResets(t *testing.T) {
// out and be re-scheduled for someone else.
func TestTransactionFetcherTimeoutRescheduling(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Push an initial announcement through to the scheduled stage
doTxNotify{
@ -973,14 +899,7 @@ func TestTransactionFetcherTimeoutRescheduling(t *testing.T) {
// Tests that the fetching timeout timers properly reset and reschedule.
func TestTransactionFetcherTimeoutTimerResets(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}}, types: []byte{types.LegacyTxType}, sizes: []uint32{111}},
doWait{time: txArriveTimeout, step: true},
@ -1051,14 +970,7 @@ func TestTransactionFetcherRateLimiting(t *testing.T) {
})
}
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Announce all the transactions, wait a bit and ensure only a small
// percentage gets requested
@ -1081,14 +993,7 @@ func TestTransactionFetcherRateLimiting(t *testing.T) {
// be requested at a time, to keep the responses below a reasonable level.
func TestTransactionFetcherBandwidthLimiting(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Announce mid size transactions from A to verify that multiple
// ones can be piled into a single request.
@ -1198,14 +1103,7 @@ func TestTransactionFetcherDoSProtection(t *testing.T) {
})
}
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Announce half of the transaction and wait for them to be scheduled
doTxNotify{peer: "A", hashes: hashesA[:maxTxAnnounces/2], types: typesA[:maxTxAnnounces/2], sizes: sizesA[:maxTxAnnounces/2]},
@ -1266,24 +1164,21 @@ func TestTransactionFetcherDoSProtection(t *testing.T) {
func TestTransactionFetcherUnderpricedDedup(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
errs := make([]error, len(txs))
for i := 0; i < len(errs); i++ {
if i%3 == 0 {
errs[i] = txpool.ErrUnderpriced
} else if i%3 == 1 {
errs[i] = txpool.ErrReplaceUnderpriced
} else {
errs[i] = txpool.ErrTxGasPriceTooLow
}
f := initDefaultTxFetcher()
f.addTxs = func(txs []*types.Transaction) []error {
errs := make([]error, len(txs))
for i := 0; i < len(errs); i++ {
if i%3 == 0 {
errs[i] = txpool.ErrUnderpriced
} else if i%3 == 1 {
errs[i] = txpool.ErrReplaceUnderpriced
} else {
errs[i] = txpool.ErrTxGasPriceTooLow
}
return errs
},
func(string, []common.Hash) error { return nil },
nil,
)
}
return errs
}
return f
},
steps: []interface{}{
// Deliver a transaction through the fetcher, but reject as underpriced
@ -1367,18 +1262,15 @@ func TestTransactionFetcherUnderpricedDoSProtection(t *testing.T) {
}
testTransactionFetcher(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
errs := make([]error, len(txs))
for i := 0; i < len(errs); i++ {
errs[i] = txpool.ErrUnderpriced
}
return errs
},
func(string, []common.Hash) error { return nil },
nil,
)
f := initDefaultTxFetcher()
f.addTxs = func(txs []*types.Transaction) []error {
errs := make([]error, len(txs))
for i := 0; i < len(errs); i++ {
errs[i] = txpool.ErrUnderpriced
}
return errs
}
return f
},
steps: append(steps, []interface{}{
// The preparation of the test has already been done in `steps`, add the last check
@ -1398,16 +1290,7 @@ func TestTransactionFetcherUnderpricedDoSProtection(t *testing.T) {
// Tests that unexpected deliveries don't corrupt the internal state.
func TestTransactionFetcherOutOfBoundDeliveries(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Deliver something out of the blue
isWaiting(nil),
@ -1457,16 +1340,7 @@ func TestTransactionFetcherOutOfBoundDeliveries(t *testing.T) {
// live or dangling stages.
func TestTransactionFetcherDrop(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Set up a few hashes into various stages
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}}, types: []byte{types.LegacyTxType}, sizes: []uint32{111}},
@ -1531,16 +1405,7 @@ func TestTransactionFetcherDrop(t *testing.T) {
// available peer.
func TestTransactionFetcherDropRescheduling(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Set up a few hashes into various stages
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}}, types: []byte{types.LegacyTxType}, sizes: []uint32{111}},
@ -1578,14 +1443,9 @@ func TestInvalidAnnounceMetadata(t *testing.T) {
drop := make(chan string, 2)
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
func(peer string) { drop <- peer },
)
f := initDefaultTxFetcher()
f.dropPeer = func(peer string) { drop <- peer }
return f
},
steps: []interface{}{
// Initial announcement to get something into the waitlist
@ -1660,16 +1520,7 @@ func TestInvalidAnnounceMetadata(t *testing.T) {
// announced one.
func TestTransactionFetcherFuzzCrash01(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Get a transaction into fetching mode and make it dangling with a broadcast
doTxNotify{peer: "A", hashes: []common.Hash{testTxsHashes[0]}, types: []byte{testTxs[0].Type()}, sizes: []uint32{uint32(testTxs[0].Size())}},
@ -1688,16 +1539,7 @@ func TestTransactionFetcherFuzzCrash01(t *testing.T) {
// concurrently announced one.
func TestTransactionFetcherFuzzCrash02(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Get a transaction into fetching mode and make it dangling with a broadcast
doTxNotify{peer: "A", hashes: []common.Hash{testTxsHashes[0]}, types: []byte{testTxs[0].Type()}, sizes: []uint32{uint32(testTxs[0].Size())}},
@ -1718,16 +1560,7 @@ func TestTransactionFetcherFuzzCrash02(t *testing.T) {
// with a concurrent notify.
func TestTransactionFetcherFuzzCrash03(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Get a transaction into fetching mode and make it dangling with a broadcast
doTxNotify{
@ -1758,17 +1591,12 @@ func TestTransactionFetcherFuzzCrash04(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error {
<-proceed
return errors.New("peer disconnected")
},
nil,
)
f := initDefaultTxFetcher()
f.fetchTxs = func(string, []common.Hash) error {
<-proceed
return errors.New("peer disconnected")
}
return f
},
steps: []interface{}{
// Get a transaction into fetching mode and make it dangling with a broadcast
@ -1792,14 +1620,7 @@ func TestTransactionFetcherFuzzCrash04(t *testing.T) {
// once they are announced in the network.
func TestBlobTransactionAnnounce(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
nil,
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
// Initial announcement to get something into the waitlist
doTxNotify{peer: "A", hashes: []common.Hash{{0x01}, {0x02}}, types: []byte{types.LegacyTxType, types.LegacyTxType}, sizes: []uint32{111, 222}},
@ -1860,16 +1681,7 @@ func TestBlobTransactionAnnounce(t *testing.T) {
func TestTransactionFetcherDropAlternates(t *testing.T) {
testTransactionFetcherParallel(t, txFetcherTest{
init: func() *TxFetcher {
return NewTxFetcher(
func(common.Hash, byte) error { return nil },
func(txs []*types.Transaction) []error {
return make([]error, len(txs))
},
func(string, []common.Hash) error { return nil },
nil,
)
},
init: initDefaultTxFetcher,
steps: []interface{}{
doTxNotify{peer: "A", hashes: []common.Hash{testTxsHashes[0]}, types: []byte{testTxs[0].Type()}, sizes: []uint32{uint32(testTxs[0].Size())}},
doWait{time: txArriveTimeout, step: true},

View file

@ -180,26 +180,24 @@ func newHandler(config *handlerConfig) (*handler, error) {
return h.txpool.Add(txs, false)
}
hasTx := func(hash common.Hash) bool {
if h.txpool.Has(hash) {
return true
}
// check on chain as well (no need to check limbo separately, as chain checks limbo too)
if h.chain.HasCanonicalTransaction(hash, false) {
return true
}
// tx not found
return false
return h.txpool.Has(hash)
}
chainTx := func(hash common.Hash) bool {
// check on chain (no need to check limbo separately, as chain checks limbo too)
return h.chain.HasCanonicalTransaction(hash, false)
}
validateMeta := func(tx common.Hash, kind byte) error {
if hasTx(tx) {
return txpool.ErrAlreadyKnown
}
if chainTx(tx) {
return txpool.ErrAlreadyKnown
}
if !h.txpool.FilterType(kind) {
return types.ErrTxTypeNotSupported
}
return nil
}
h.txFetcher = fetcher.NewTxFetcher(validateMeta, addTxs, fetchTx, h.removePeer)
return h, nil
}