From 425397526910176f3a185d7150275dd4c8f08483 Mon Sep 17 00:00:00 2001 From: Roberto Bayardo Date: Sun, 7 Jul 2024 13:12:39 -0700 Subject: [PATCH] - fetch transactions from a peer in the order they were announced to minimize nonce-gaps (which cause blob txs to be rejected) - don't wait on fetching blob transactions after announcement is received, since they are not broadcast --- cmd/devp2p/internal/ethtest/suite.go | 11 ++- eth/fetcher/tx_fetcher.go | 122 ++++++++++++--------------- eth/fetcher/tx_fetcher_test.go | 83 ++++++++++++------ 3 files changed, 122 insertions(+), 94 deletions(-) diff --git a/cmd/devp2p/internal/ethtest/suite.go b/cmd/devp2p/internal/ethtest/suite.go index b5cc27a2b5..876861f9d8 100644 --- a/cmd/devp2p/internal/ethtest/suite.go +++ b/cmd/devp2p/internal/ethtest/suite.go @@ -849,7 +849,16 @@ func (s *Suite) TestBlobViolations(t *utesting.T) { if code, _, err := conn.Read(); err != nil { t.Fatalf("expected disconnect on blob violation, got err: %v", err) } else if code != discMsg { - t.Fatalf("expected disconnect on blob violation, got msg code: %d", code) + if code == 24 { + // sometimes we'll get a blob transaction hashes announcement before the disconnect + // because blob transactions are scheduled to be fetched right away. + if code, _, err = conn.Read(); err != nil { + t.Fatalf("expected disconnect on blob violation, got err on second read: %v", err) + } + } + if code != discMsg { + t.Fatalf("expected disconnect on blob violation, got msg code: %d", code) + } } conn.Close() } diff --git a/eth/fetcher/tx_fetcher.go b/eth/fetcher/tx_fetcher.go index a113155009..49f906f910 100644 --- a/eth/fetcher/tx_fetcher.go +++ b/eth/fetcher/tx_fetcher.go @@ -17,7 +17,6 @@ package fetcher import ( - "bytes" "errors" "fmt" "math" @@ -35,7 +34,7 @@ import ( ) const ( - // maxTxAnnounces is the maximum number of unique transaction a peer + // maxTxAnnounces is the maximum number of unique transactions a peer // can announce in a short time. maxTxAnnounces = 4096 @@ -114,14 +113,17 @@ var errTerminated = errors.New("terminated") type txAnnounce struct { origin string // Identifier of the peer originating the notification hashes []common.Hash // Batch of transaction hashes being announced - metas []*txMetadata // Batch of metadata associated with the hashes + metas []txMetadata // Batch of metadatas associated with the hashes } -// txMetadata is a set of extra data transmitted along the announcement for better -// fetch scheduling. +// txMetadata provides the extra data transmitted along with the announcement +// for better fetch scheduling ('kind' & 'size'), plus an extra field +// ('arrival') to keep track of its order of arrival. 'size==0' can be used to +// test for 0 pre-eth/68 announcements. In this case, kind will also be 0. type txMetadata struct { - kind byte // Transaction consensus type - size uint32 // Transaction size in bytes + kind byte // Transaction consensus type + size uint32 // Transaction size in bytes, or 0 if the announcement didn't include metadata + arrival uint64 // Value that can be used to sort announcements by order of arrival } // txRequest represents an in-flight transaction retrieval request destined to @@ -159,7 +161,7 @@ type txDrop struct { // The invariants of the fetcher are: // - Each tracked transaction (hash) must only be present in one of the // three stages. This ensures that the fetcher operates akin to a finite -// state automata and there's do data leak. +// state automata and there's no data leak. // - Each peer that announced transactions may be scheduled retrievals, but // only ever one concurrently. This ensures we can immediately know what is // missing from a reply and reschedule it. @@ -173,14 +175,14 @@ type TxFetcher struct { // Stage 1: Waiting lists for newly discovered transactions that might be // broadcast without needing explicit request/reply round trips. - waitlist map[common.Hash]map[string]struct{} // Transactions waiting for an potential broadcast - waittime map[common.Hash]mclock.AbsTime // Timestamps when transactions were added to the waitlist - waitslots map[string]map[common.Hash]*txMetadata // Waiting announcements grouped by peer (DoS protection) + waitlist map[common.Hash]map[string]struct{} // Transactions waiting for an potential broadcast + waittime map[common.Hash]mclock.AbsTime // Timestamps when transactions were added to the waitlist + waitslots map[string]map[common.Hash]txMetadata // Waiting announcements grouped by peer (DoS protection) // Stage 2: Queue of transactions that waiting to be allocated to some peer // to be retrieved directly. - announces map[string]map[common.Hash]*txMetadata // Set of announced transactions, grouped by origin peer - announced map[common.Hash]map[string]struct{} // Set of download locations, grouped by transaction hash + announces map[string]map[common.Hash]txMetadata // Set of announced transactions, grouped by origin peer + announced map[common.Hash]map[string]struct{} // Set of download locations, grouped by transaction hash // Stage 3: Set of transactions currently being retrieved, some which may be // fulfilled and some rescheduled. Note, this step shares 'announces' from the @@ -218,8 +220,8 @@ func NewTxFetcherForTests( quit: make(chan struct{}), waitlist: make(map[common.Hash]map[string]struct{}), waittime: make(map[common.Hash]mclock.AbsTime), - waitslots: make(map[string]map[common.Hash]*txMetadata), - announces: make(map[string]map[common.Hash]*txMetadata), + waitslots: make(map[string]map[common.Hash]txMetadata), + announces: make(map[string]map[common.Hash]txMetadata), announced: make(map[common.Hash]map[string]struct{}), fetching: make(map[common.Hash]string), requests: make(map[string]*txRequest), @@ -247,7 +249,7 @@ func (f *TxFetcher) Notify(peer string, types []byte, sizes []uint32, hashes []c // loop, so anything caught here is time saved internally. var ( unknownHashes = make([]common.Hash, 0, len(hashes)) - unknownMetas = make([]*txMetadata, 0, len(hashes)) + unknownMetas = make([]txMetadata, 0, len(hashes)) duplicate int64 underpriced int64 @@ -264,7 +266,7 @@ func (f *TxFetcher) Notify(peer string, types []byte, sizes []uint32, hashes []c // Transaction metadata has been available since eth68, and all // legacy eth protocols (prior to eth68) have been deprecated. // Therefore, metadata is always expected in the announcement. - unknownMetas = append(unknownMetas, &txMetadata{kind: types[i], size: sizes[i]}) + unknownMetas = append(unknownMetas, txMetadata{kind: types[i], size: sizes[i]}) } } txAnnounceKnownMeter.Mark(duplicate) @@ -445,7 +447,7 @@ func (f *TxFetcher) loop() { if announces := f.announces[ann.origin]; announces != nil { announces[hash] = ann.metas[i] } else { - f.announces[ann.origin] = map[common.Hash]*txMetadata{hash: ann.metas[i]} + f.announces[ann.origin] = map[common.Hash]txMetadata{hash: ann.metas[i]} } continue } @@ -458,7 +460,7 @@ func (f *TxFetcher) loop() { if announces := f.announces[ann.origin]; announces != nil { announces[hash] = ann.metas[i] } else { - f.announces[ann.origin] = map[common.Hash]*txMetadata{hash: ann.metas[i]} + f.announces[ann.origin] = map[common.Hash]txMetadata{hash: ann.metas[i]} } continue } @@ -477,18 +479,26 @@ func (f *TxFetcher) loop() { if waitslots := f.waitslots[ann.origin]; waitslots != nil { waitslots[hash] = ann.metas[i] } else { - f.waitslots[ann.origin] = map[common.Hash]*txMetadata{hash: ann.metas[i]} + f.waitslots[ann.origin] = map[common.Hash]txMetadata{hash: ann.metas[i]} } continue } // Transaction unknown to the fetcher, insert it into the waiting list f.waitlist[hash] = map[string]struct{}{ann.origin: {}} - f.waittime[hash] = f.clock.Now() + if ann.metas[i].kind == types.BlobTxType { + // blob transactions are never broadcast, so to force them + // to be fetched immediately we pretend they arrived + // earlier. + f.waittime[hash] = f.clock.Now() - mclock.AbsTime(txArriveTimeout) + idleWait = true // may need to reschedule fetcher due to "time travel" + } else { + f.waittime[hash] = f.clock.Now() + } if waitslots := f.waitslots[ann.origin]; waitslots != nil { waitslots[hash] = ann.metas[i] } else { - f.waitslots[ann.origin] = map[common.Hash]*txMetadata{hash: ann.metas[i]} + f.waitslots[ann.origin] = map[common.Hash]txMetadata{hash: ann.metas[i]} } } // If a new item was added to the waitlist, schedule it into the fetcher @@ -516,7 +526,7 @@ func (f *TxFetcher) loop() { if announces := f.announces[peer]; announces != nil { announces[hash] = f.waitslots[peer][hash] } else { - f.announces[peer] = map[common.Hash]*txMetadata{hash: f.waitslots[peer][hash]} + f.announces[peer] = map[common.Hash]txMetadata{hash: f.waitslots[peer][hash]} } delete(f.waitslots[peer], hash) if len(f.waitslots[peer]) == 0 { @@ -590,7 +600,7 @@ func (f *TxFetcher) loop() { for i, hash := range delivery.hashes { if _, ok := f.waitlist[hash]; ok { for peer, txset := range f.waitslots { - if meta := txset[hash]; meta != nil { + if meta, ok := txset[hash]; ok && meta.size != 0 { if delivery.metas[i].kind != meta.kind { log.Warn("Announced transaction type mismatch", "peer", peer, "tx", hash, "type", delivery.metas[i].kind, "ann", meta.kind) f.dropPeer(peer) @@ -616,7 +626,7 @@ func (f *TxFetcher) loop() { delete(f.waittime, hash) } else { for peer, txset := range f.announces { - if meta := txset[hash]; meta != nil { + if meta, ok := txset[hash]; ok && meta.size != 0 { if delivery.metas[i].kind != meta.kind { log.Warn("Announced transaction type mismatch", "peer", peer, "tx", hash, "type", delivery.metas[i].kind, "ann", meta.kind) f.dropPeer(peer) @@ -873,7 +883,7 @@ func (f *TxFetcher) scheduleFetches(timer *mclock.Timer, timeout chan struct{}, hashes = make([]common.Hash, 0, maxTxRetrievals) bytes uint64 ) - f.forEachAnnounce(f.announces[peer], func(hash common.Hash, meta *txMetadata) bool { + f.forEachAnnounce(f.announces[peer], func(hash common.Hash, meta txMetadata) bool { // If the transaction is already fetching, skip to the next one if _, ok := f.fetching[hash]; ok { return true @@ -938,28 +948,25 @@ func (f *TxFetcher) forEachPeer(peers map[string]struct{}, do func(peer string)) } } -// forEachAnnounce does a range loop over a map of announcements in production, -// but during testing it does a deterministic sorted random to allow reproducing -// issues. -func (f *TxFetcher) forEachAnnounce(announces map[common.Hash]*txMetadata, do func(hash common.Hash, meta *txMetadata) bool) { - // If we're running production, use whatever Go's map gives us - if f.rand == nil { - for hash, meta := range announces { - if !do(hash, meta) { - return - } - } - return +// forEachAnnounce loops over the given announcements in arrival order, invoking +// the do function for each until it returns false. We enforce an arrival +// ordering to minimize the chances of mempool nonce-gaps, which result in blob +// transactions being rejected by the mempool. +func (f *TxFetcher) forEachAnnounce(announces map[common.Hash]txMetadata, do func(hash common.Hash, meta txMetadata) bool) { + type announcement struct { + hash common.Hash + meta txMetadata } - // We're running the test suite, make iteration deterministic - list := make([]common.Hash, 0, len(announces)) - for hash := range announces { - list = append(list, hash) + // process announcements by their arrival order + list := make([]announcement, 0, len(announces)) + for hash, metadata := range announces { + list = append(list, announcement{hash: hash, meta: metadata}) } - sortHashes(list) - rotateHashes(list, f.rand.Intn(len(list))) - for _, hash := range list { - if !do(hash, announces[hash]) { + sort.Slice(list, func(i, j int) bool { + return list[i].meta.arrival < list[j].meta.arrival + }) + for i := range list { + if !do(list[i].hash, list[i].meta) { return } } @@ -975,26 +982,3 @@ func rotateStrings(slice []string, n int) { slice[i] = orig[(i+n)%len(orig)] } } - -// sortHashes sorts a slice of hashes. This method is only used in tests in order -// to simulate random map iteration but keep it deterministic. -func sortHashes(slice []common.Hash) { - for i := 0; i < len(slice); i++ { - for j := i + 1; j < len(slice); j++ { - if bytes.Compare(slice[i][:], slice[j][:]) > 0 { - slice[i], slice[j] = slice[j], slice[i] - } - } - } -} - -// rotateHashes rotates the contents of a slice by n steps. This method is only -// used in tests to simulate random map iteration but keep it deterministic. -func rotateHashes(slice []common.Hash, n int) { - orig := make([]common.Hash, len(slice)) - copy(orig, slice) - - for i := 0; i < len(orig); i++ { - slice[i] = orig[(i+n)%len(orig)] - } -} diff --git a/eth/fetcher/tx_fetcher_test.go b/eth/fetcher/tx_fetcher_test.go index 0b47646669..26b177ad41 100644 --- a/eth/fetcher/tx_fetcher_test.go +++ b/eth/fetcher/tx_fetcher_test.go @@ -179,6 +179,38 @@ func TestTransactionFetcherWaiting(t *testing.T) { }, }), isScheduled{tracking: nil, fetching: nil}, + // Announce a non-conflicting blob tx, which should immediately go + // to fetching after a trivial wait. + doTxNotify{peer: "D", hashes: []common.Hash{{0x0b}}, types: []byte{types.BlobTxType}, sizes: []uint32{1000}}, + doWait{time: 0, step: true}, + isWaiting(map[string][]announce{ + "A": { + {common.Hash{0x01}, types.LegacyTxType, 111}, + {common.Hash{0x02}, types.LegacyTxType, 222}, + {common.Hash{0x03}, types.LegacyTxType, 333}, + {common.Hash{0x05}, types.LegacyTxType, 555}, + }, + "B": { + {common.Hash{0x03}, types.LegacyTxType, 333}, + {common.Hash{0x04}, types.LegacyTxType, 444}, + }, + "C": { + {common.Hash{0x01}, types.LegacyTxType, 111}, + {common.Hash{0x04}, types.LegacyTxType, 444}, + }, + "D": { + {common.Hash{0x01}, types.LegacyTxType, 999}, + {common.Hash{0x02}, types.BlobTxType, 222}, + }, + }), + isScheduled{ + tracking: map[string][]announce{ + "D": {{common.Hash{0x0B}, types.BlobTxType, 1000}}, + }, + fetching: map[string][]common.Hash{ + "D": {{0x0B}}, + }, + }, // Wait for the arrival timeout which should move all expired items // from the wait list to the scheduler @@ -203,19 +235,20 @@ func TestTransactionFetcherWaiting(t *testing.T) { "D": { {common.Hash{0x01}, types.LegacyTxType, 999}, {common.Hash{0x02}, types.BlobTxType, 222}, + {common.Hash{0x0B}, types.BlobTxType, 1000}, }, }, fetching: map[string][]common.Hash{ // Depends on deterministic test randomizer - "A": {{0x03}, {0x05}}, - "C": {{0x01}, {0x04}}, - "D": {{0x02}}, + "A": {{0x01}, {0x02}, {0x03}, {0x05}}, + "B": {{0x04}}, + "D": {{0x0B}}, }, }, // Queue up a non-fetchable transaction and then trigger it with a new // peer (weird case to test 1 line in the fetcher) - doTxNotify{peer: "C", hashes: []common.Hash{{0x06}, {0x07}}, types: []byte{types.LegacyTxType, types.LegacyTxType}, sizes: []uint32{666, 777}}, + doTxNotify{peer: "B", hashes: []common.Hash{{0x06}, {0x07}}, types: []byte{types.LegacyTxType, types.LegacyTxType}, sizes: []uint32{666, 777}}, isWaiting(map[string][]announce{ - "C": { + "B": { {common.Hash{0x06}, types.LegacyTxType, 666}, {common.Hash{0x07}, types.LegacyTxType, 777}, }, @@ -232,22 +265,23 @@ func TestTransactionFetcherWaiting(t *testing.T) { "B": { {common.Hash{0x03}, types.LegacyTxType, 333}, {common.Hash{0x04}, types.LegacyTxType, 444}, + {common.Hash{0x06}, types.LegacyTxType, 666}, + {common.Hash{0x07}, types.LegacyTxType, 777}, }, "C": { {common.Hash{0x01}, types.LegacyTxType, 111}, {common.Hash{0x04}, types.LegacyTxType, 444}, - {common.Hash{0x06}, types.LegacyTxType, 666}, - {common.Hash{0x07}, types.LegacyTxType, 777}, }, "D": { {common.Hash{0x01}, types.LegacyTxType, 999}, {common.Hash{0x02}, types.BlobTxType, 222}, + {common.Hash{0x0B}, types.BlobTxType, 1000}, }, }, fetching: map[string][]common.Hash{ - "A": {{0x03}, {0x05}}, - "C": {{0x01}, {0x04}}, - "D": {{0x02}}, + "A": {{0x01}, {0x02}, {0x03}, {0x05}}, + "B": {{0x04}}, + "D": {{0x0B}}, }, }, doTxNotify{peer: "E", hashes: []common.Hash{{0x06}, {0x07}}, types: []byte{types.LegacyTxType, types.LegacyTxType}, sizes: []uint32{666, 777}}, @@ -262,16 +296,17 @@ func TestTransactionFetcherWaiting(t *testing.T) { "B": { {common.Hash{0x03}, types.LegacyTxType, 333}, {common.Hash{0x04}, types.LegacyTxType, 444}, + {common.Hash{0x06}, types.LegacyTxType, 666}, + {common.Hash{0x07}, types.LegacyTxType, 777}, }, "C": { {common.Hash{0x01}, types.LegacyTxType, 111}, {common.Hash{0x04}, types.LegacyTxType, 444}, - {common.Hash{0x06}, types.LegacyTxType, 666}, - {common.Hash{0x07}, types.LegacyTxType, 777}, }, "D": { {common.Hash{0x01}, types.LegacyTxType, 999}, {common.Hash{0x02}, types.BlobTxType, 222}, + {common.Hash{0x0B}, types.BlobTxType, 1000}, }, "E": { {common.Hash{0x06}, types.LegacyTxType, 666}, @@ -279,9 +314,9 @@ func TestTransactionFetcherWaiting(t *testing.T) { }, }, fetching: map[string][]common.Hash{ - "A": {{0x03}, {0x05}}, - "C": {{0x01}, {0x04}}, - "D": {{0x02}}, + "A": {{0x01}, {0x02}, {0x03}, {0x05}}, + "B": {{0x04}}, + "D": {{0x0B}}, "E": {{0x06}, {0x07}}, }, }, @@ -701,7 +736,7 @@ func TestTransactionFetcherMissingRescheduling(t *testing.T) { }, // Deliver the middle transaction requested, the one before which // should be dropped and the one after re-requested. - doTxEnqueue{peer: "A", txs: []*types.Transaction{testTxs[0]}, direct: true}, // This depends on the deterministic random + doTxEnqueue{peer: "A", txs: []*types.Transaction{testTxs[1]}, direct: true}, isScheduled{ tracking: map[string][]announce{ "A": { @@ -1070,7 +1105,7 @@ func TestTransactionFetcherRateLimiting(t *testing.T) { "A": announces, }, fetching: map[string][]common.Hash{ - "A": hashes[1643 : 1643+maxTxRetrievals], + "A": hashes[:maxTxRetrievals], }, }, }, @@ -1130,9 +1165,9 @@ func TestTransactionFetcherBandwidthLimiting(t *testing.T) { }, }, fetching: map[string][]common.Hash{ - "A": {{0x02}, {0x03}, {0x04}}, - "B": {{0x06}}, - "C": {{0x08}}, + "A": {{0x01}, {0x02}, {0x03}}, + "B": {{0x05}}, + "C": {{0x07}}, }, }, }, @@ -1209,8 +1244,8 @@ func TestTransactionFetcherDoSProtection(t *testing.T) { "B": announceB[:maxTxAnnounces/2-1], }, fetching: map[string][]common.Hash{ - "A": hashesA[1643 : 1643+maxTxRetrievals], - "B": append(append([]common.Hash{}, hashesB[maxTxAnnounces/2-3:maxTxAnnounces/2-1]...), hashesB[:maxTxRetrievals-2]...), + "A": hashesA[:maxTxRetrievals], + "B": hashesB[:maxTxRetrievals], }, }, // Ensure that adding even one more hash results in dropping the hash @@ -1227,8 +1262,8 @@ func TestTransactionFetcherDoSProtection(t *testing.T) { "B": announceB[:maxTxAnnounces/2-1], }, fetching: map[string][]common.Hash{ - "A": hashesA[1643 : 1643+maxTxRetrievals], - "B": append(append([]common.Hash{}, hashesB[maxTxAnnounces/2-3:maxTxAnnounces/2-1]...), hashesB[:maxTxRetrievals-2]...), + "A": hashesA[:maxTxRetrievals], + "B": hashesB[:maxTxRetrievals], }, }, },