mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 15:16:43 +00:00
feat: add more logs to investigate tx broadcast issue (#1005)
* 1. broadcast + announce * 2. broadcast * 3. broadcast * 5. announce * 6. announce * 7. announce * 8. announce * 9. announce * 10. announce * minor * chore: auto version bump [bot] --------- Co-authored-by: HAOYUatHZ <HAOYUatHZ@users.noreply.github.com>
This commit is contained in:
parent
c252ee1ee0
commit
77a565943f
6 changed files with 27 additions and 1 deletions
|
|
@ -790,6 +790,9 @@ func (f *TxFetcher) scheduleFetches(timer *mclock.Timer, timeout chan struct{},
|
||||||
}
|
}
|
||||||
return true // continue in the for-each
|
return true // continue in the for-each
|
||||||
})
|
})
|
||||||
|
|
||||||
|
log.Debug("Scheduling transaction retrieval", "peer", peer, "len(f.announces[peer])", len(f.announces[peer]), "len(hashes)", len(hashes))
|
||||||
|
|
||||||
// If any hashes were allocated, request them from the peer
|
// If any hashes were allocated, request them from the peer
|
||||||
if len(hashes) > 0 {
|
if len(hashes) > 0 {
|
||||||
f.requests[peer] = &txRequest{hashes: hashes, time: f.clock.Now()}
|
f.requests[peer] = &txRequest{hashes: hashes, time: f.clock.Now()}
|
||||||
|
|
|
||||||
|
|
@ -512,11 +512,13 @@ func (h *handler) BroadcastTransactions(txs types.Transactions) {
|
||||||
directPeers++
|
directPeers++
|
||||||
directCount += len(hashes)
|
directCount += len(hashes)
|
||||||
peer.AsyncSendTransactions(hashes)
|
peer.AsyncSendTransactions(hashes)
|
||||||
|
log.Debug("Transactions being broadcasted to", "peer", peer.String(), "len", len(hashes))
|
||||||
}
|
}
|
||||||
for peer, hashes := range annos {
|
for peer, hashes := range annos {
|
||||||
annoPeers++
|
annoPeers++
|
||||||
annoCount += len(hashes)
|
annoCount += len(hashes)
|
||||||
peer.AsyncSendPooledTransactionHashes(hashes)
|
peer.AsyncSendPooledTransactionHashes(hashes)
|
||||||
|
log.Debug("Transactions being announced to", "peer", peer.String(), "len", len(hashes))
|
||||||
}
|
}
|
||||||
log.Debug("Transaction broadcast", "txs", len(txs),
|
log.Debug("Transaction broadcast", "txs", len(txs),
|
||||||
"announce packs", annoPeers, "announced hashes", annoCount,
|
"announce packs", annoPeers, "announced hashes", annoCount,
|
||||||
|
|
|
||||||
|
|
@ -21,6 +21,7 @@ import (
|
||||||
|
|
||||||
"github.com/scroll-tech/go-ethereum/common"
|
"github.com/scroll-tech/go-ethereum/common"
|
||||||
"github.com/scroll-tech/go-ethereum/core/types"
|
"github.com/scroll-tech/go-ethereum/core/types"
|
||||||
|
"github.com/scroll-tech/go-ethereum/log"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
|
|
@ -92,10 +93,13 @@ func (p *Peer) broadcastTransactions() {
|
||||||
if len(txs) > 0 {
|
if len(txs) > 0 {
|
||||||
done = make(chan struct{})
|
done = make(chan struct{})
|
||||||
go func() {
|
go func() {
|
||||||
|
log.Debug("Sending transactions", "count", len(txs))
|
||||||
if err := p.SendTransactions(txs); err != nil {
|
if err := p.SendTransactions(txs); err != nil {
|
||||||
|
log.Debug("Sending transactions", "count", len(txs), "err", err)
|
||||||
fail <- err
|
fail <- err
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
log.Debug("Sent transactions", "count", len(txs))
|
||||||
close(done)
|
close(done)
|
||||||
p.Log().Trace("Sent transactions", "count", len(txs))
|
p.Log().Trace("Sent transactions", "count", len(txs))
|
||||||
}()
|
}()
|
||||||
|
|
@ -110,6 +114,7 @@ func (p *Peer) broadcastTransactions() {
|
||||||
}
|
}
|
||||||
// New batch of transactions to be broadcast, queue them (with cap)
|
// New batch of transactions to be broadcast, queue them (with cap)
|
||||||
queue = append(queue, hashes...)
|
queue = append(queue, hashes...)
|
||||||
|
log.Debug("Queue size in broadcastTransactions", "len(hashes)", len(hashes), "len(queue)", len(queue), "maxQueuedTxs", maxQueuedTxs)
|
||||||
if len(queue) > maxQueuedTxs {
|
if len(queue) > maxQueuedTxs {
|
||||||
// Fancy copy and resize to ensure buffer doesn't grow indefinitely
|
// Fancy copy and resize to ensure buffer doesn't grow indefinitely
|
||||||
queue = queue[:copy(queue, queue[len(queue)-maxQueuedTxs:])]
|
queue = queue[:copy(queue, queue[len(queue)-maxQueuedTxs:])]
|
||||||
|
|
@ -159,10 +164,13 @@ func (p *Peer) announceTransactions() {
|
||||||
if len(pending) > 0 {
|
if len(pending) > 0 {
|
||||||
done = make(chan struct{})
|
done = make(chan struct{})
|
||||||
go func() {
|
go func() {
|
||||||
|
log.Debug("Sending transaction announcements", "count", len(pending))
|
||||||
if err := p.sendPooledTransactionHashes(pending); err != nil {
|
if err := p.sendPooledTransactionHashes(pending); err != nil {
|
||||||
|
log.Debug("Sending transaction announcements", "count", len(pending), "err", err)
|
||||||
fail <- err
|
fail <- err
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
log.Debug("Sent transaction announcements", "count", len(pending))
|
||||||
close(done)
|
close(done)
|
||||||
p.Log().Trace("Sent transaction announcements", "count", len(pending))
|
p.Log().Trace("Sent transaction announcements", "count", len(pending))
|
||||||
}()
|
}()
|
||||||
|
|
@ -177,6 +185,7 @@ func (p *Peer) announceTransactions() {
|
||||||
}
|
}
|
||||||
// New batch of transactions to be broadcast, queue them (with cap)
|
// New batch of transactions to be broadcast, queue them (with cap)
|
||||||
queue = append(queue, hashes...)
|
queue = append(queue, hashes...)
|
||||||
|
log.Debug("Queue size in announceTransactions", "len(hashes)", len(hashes), "len(queue)", len(queue), "maxQueuedTxAnns", maxQueuedTxAnns)
|
||||||
if len(queue) > maxQueuedTxAnns {
|
if len(queue) > maxQueuedTxAnns {
|
||||||
// Fancy copy and resize to ensure buffer doesn't grow indefinitely
|
// Fancy copy and resize to ensure buffer doesn't grow indefinitely
|
||||||
queue = queue[:copy(queue, queue[len(queue)-maxQueuedTxAnns:])]
|
queue = queue[:copy(queue, queue[len(queue)-maxQueuedTxAnns:])]
|
||||||
|
|
|
||||||
|
|
@ -323,9 +323,11 @@ func handleNewPooledTransactionHashes(backend Backend, msg Decoder, peer *Peer)
|
||||||
}
|
}
|
||||||
ann := new(NewPooledTransactionHashesPacket)
|
ann := new(NewPooledTransactionHashesPacket)
|
||||||
if err := msg.Decode(ann); err != nil {
|
if err := msg.Decode(ann); err != nil {
|
||||||
|
log.Debug("Failed to decode `NewPooledTransactionHashesPacket`", "peer", peer.String(), "err", err)
|
||||||
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
||||||
}
|
}
|
||||||
// Schedule all the unknown hashes for retrieval
|
// Schedule all the unknown hashes for retrieval
|
||||||
|
log.Debug("handleNewPooledTransactionHashes", "peer", peer.String(), "len(ann)", len(*ann))
|
||||||
for _, hash := range *ann {
|
for _, hash := range *ann {
|
||||||
peer.markTransaction(hash)
|
peer.markTransaction(hash)
|
||||||
}
|
}
|
||||||
|
|
@ -336,9 +338,11 @@ func handleGetPooledTransactions66(backend Backend, msg Decoder, peer *Peer) err
|
||||||
// Decode the pooled transactions retrieval message
|
// Decode the pooled transactions retrieval message
|
||||||
var query GetPooledTransactionsPacket66
|
var query GetPooledTransactionsPacket66
|
||||||
if err := msg.Decode(&query); err != nil {
|
if err := msg.Decode(&query); err != nil {
|
||||||
|
log.Debug("Failed to decode `GetPooledTransactionsPacket66`", "peer", peer.String(), "err", err)
|
||||||
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
||||||
}
|
}
|
||||||
hashes, txs := answerGetPooledTransactions(backend, query.GetPooledTransactionsPacket, peer)
|
hashes, txs := answerGetPooledTransactions(backend, query.GetPooledTransactionsPacket, peer)
|
||||||
|
log.Debug("handleGetPooledTransactions", "peer", peer.String(), "RequestId", query.RequestId, "len(query)", len(query.GetPooledTransactionsPacket), "retrieved", len(hashes))
|
||||||
return peer.ReplyPooledTransactionsRLP(query.RequestId, hashes, txs)
|
return peer.ReplyPooledTransactionsRLP(query.RequestId, hashes, txs)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -378,11 +382,14 @@ func handleTransactions(backend Backend, msg Decoder, peer *Peer) error {
|
||||||
// Transactions can be processed, parse all of them and deliver to the pool
|
// Transactions can be processed, parse all of them and deliver to the pool
|
||||||
var txs TransactionsPacket
|
var txs TransactionsPacket
|
||||||
if err := msg.Decode(&txs); err != nil {
|
if err := msg.Decode(&txs); err != nil {
|
||||||
|
log.Debug("Failed to decode `TransactionsPacket`", "peer", peer.String(), "err", err)
|
||||||
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
||||||
}
|
}
|
||||||
|
log.Debug("handleTransactions", "peer", peer.String(), "len(txs)", len(txs))
|
||||||
for i, tx := range txs {
|
for i, tx := range txs {
|
||||||
// Validate and mark the remote transaction
|
// Validate and mark the remote transaction
|
||||||
if tx == nil {
|
if tx == nil {
|
||||||
|
log.Debug("handleTransactions: transaction is nil", "peer", peer.String(), "i", i)
|
||||||
return fmt.Errorf("%w: transaction %d is nil", errDecode, i)
|
return fmt.Errorf("%w: transaction %d is nil", errDecode, i)
|
||||||
}
|
}
|
||||||
peer.markTransaction(tx.Hash())
|
peer.markTransaction(tx.Hash())
|
||||||
|
|
@ -398,11 +405,13 @@ func handlePooledTransactions66(backend Backend, msg Decoder, peer *Peer) error
|
||||||
// Transactions can be processed, parse all of them and deliver to the pool
|
// Transactions can be processed, parse all of them and deliver to the pool
|
||||||
var txs PooledTransactionsPacket66
|
var txs PooledTransactionsPacket66
|
||||||
if err := msg.Decode(&txs); err != nil {
|
if err := msg.Decode(&txs); err != nil {
|
||||||
|
log.Debug("Failed to decode `PooledTransactionsPacket66`", "peer", peer.String(), "err", err)
|
||||||
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
||||||
}
|
}
|
||||||
for i, tx := range txs.PooledTransactionsPacket {
|
for i, tx := range txs.PooledTransactionsPacket {
|
||||||
// Validate and mark the remote transaction
|
// Validate and mark the remote transaction
|
||||||
if tx == nil {
|
if tx == nil {
|
||||||
|
log.Debug("handlePooledTransactions: transaction is nil", "peer", peer.String(), "i", i)
|
||||||
return fmt.Errorf("%w: transaction %d is nil", errDecode, i)
|
return fmt.Errorf("%w: transaction %d is nil", errDecode, i)
|
||||||
}
|
}
|
||||||
peer.markTransaction(tx.Hash())
|
peer.markTransaction(tx.Hash())
|
||||||
|
|
|
||||||
|
|
@ -25,6 +25,7 @@ import (
|
||||||
|
|
||||||
"github.com/scroll-tech/go-ethereum/common"
|
"github.com/scroll-tech/go-ethereum/common"
|
||||||
"github.com/scroll-tech/go-ethereum/core/types"
|
"github.com/scroll-tech/go-ethereum/core/types"
|
||||||
|
"github.com/scroll-tech/go-ethereum/log"
|
||||||
"github.com/scroll-tech/go-ethereum/p2p"
|
"github.com/scroll-tech/go-ethereum/p2p"
|
||||||
"github.com/scroll-tech/go-ethereum/rlp"
|
"github.com/scroll-tech/go-ethereum/rlp"
|
||||||
)
|
)
|
||||||
|
|
@ -419,6 +420,8 @@ func (p *Peer) RequestTxs(hashes []common.Hash) error {
|
||||||
p.Log().Debug("Fetching batch of transactions", "count", len(hashes))
|
p.Log().Debug("Fetching batch of transactions", "count", len(hashes))
|
||||||
id := rand.Uint64()
|
id := rand.Uint64()
|
||||||
|
|
||||||
|
log.Debug("Requesting transactions", "RequestId", id, "Peer.id", p.id, "count", len(hashes))
|
||||||
|
|
||||||
requestTracker.Track(p.id, p.version, GetPooledTransactionsMsg, PooledTransactionsMsg, id)
|
requestTracker.Track(p.id, p.version, GetPooledTransactionsMsg, PooledTransactionsMsg, id)
|
||||||
return p2p.Send(p.rw, GetPooledTransactionsMsg, &GetPooledTransactionsPacket66{
|
return p2p.Send(p.rw, GetPooledTransactionsMsg, &GetPooledTransactionsPacket66{
|
||||||
RequestId: id,
|
RequestId: id,
|
||||||
|
|
|
||||||
|
|
@ -24,7 +24,7 @@ import (
|
||||||
const (
|
const (
|
||||||
VersionMajor = 5 // Major version component of the current release
|
VersionMajor = 5 // Major version component of the current release
|
||||||
VersionMinor = 7 // Minor version component of the current release
|
VersionMinor = 7 // Minor version component of the current release
|
||||||
VersionPatch = 1 // Patch version component of the current release
|
VersionPatch = 2 // Patch version component of the current release
|
||||||
VersionMeta = "mainnet" // Version metadata to append to the version string
|
VersionMeta = "mainnet" // Version metadata to append to the version string
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue