mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 06:36:43 +00:00
feat: add logs to track tx and block propagation delay (#1184)
This commit is contained in:
parent
ad6cced99d
commit
c66a003b88
5 changed files with 24 additions and 17 deletions
|
|
@ -794,7 +794,7 @@ 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))
|
log.Trace("Scheduling transaction retrieval", "peer", peer, "len(f.announces[peer])", len(f.announces[peer]), "len(hashes)", len(hashes))
|
||||||
peerAnnounceTxsLenGauge.Update(int64(len(f.announces[peer])))
|
peerAnnounceTxsLenGauge.Update(int64(len(f.announces[peer])))
|
||||||
peerRetrievalTxsLenGauge.Update(int64(len(hashes)))
|
peerRetrievalTxsLenGauge.Update(int64(len(hashes)))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -455,6 +455,8 @@ func (h *handler) BroadcastBlock(block *types.Block, propagate bool) {
|
||||||
hash := block.Hash()
|
hash := block.Hash()
|
||||||
peers := onlyShadowForkPeers(h.shadowForkPeerIDs, h.peers.peersWithoutBlock(hash))
|
peers := onlyShadowForkPeers(h.shadowForkPeerIDs, h.peers.peersWithoutBlock(hash))
|
||||||
|
|
||||||
|
log.Debug("Broadcasting block", "hash", hash.Hex(), "number", block.NumberU64(), "size", block.Size())
|
||||||
|
|
||||||
// If propagation is requested, send to a subset of the peer
|
// If propagation is requested, send to a subset of the peer
|
||||||
if propagate {
|
if propagate {
|
||||||
// Calculate the TD of the block (it's not imported yet, so block.Td is not valid)
|
// Calculate the TD of the block (it's not imported yet, so block.Td is not valid)
|
||||||
|
|
@ -470,7 +472,7 @@ func (h *handler) BroadcastBlock(block *types.Block, propagate bool) {
|
||||||
for _, peer := range transfer {
|
for _, peer := range transfer {
|
||||||
peer.AsyncSendNewBlock(block, td)
|
peer.AsyncSendNewBlock(block, td)
|
||||||
}
|
}
|
||||||
log.Trace("Propagated block", "hash", hash, "recipients", len(transfer), "duration", common.PrettyDuration(time.Since(block.ReceivedAt)))
|
log.Trace("Propagated block", "hash", hash.Hex(), "recipients", len(transfer), "duration", common.PrettyDuration(time.Since(block.ReceivedAt)))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Otherwise if the block is indeed in out own chain, announce it
|
// Otherwise if the block is indeed in out own chain, announce it
|
||||||
|
|
@ -478,7 +480,7 @@ func (h *handler) BroadcastBlock(block *types.Block, propagate bool) {
|
||||||
for _, peer := range peers {
|
for _, peer := range peers {
|
||||||
peer.AsyncSendNewBlockHash(block)
|
peer.AsyncSendNewBlockHash(block)
|
||||||
}
|
}
|
||||||
log.Trace("Announced block", "hash", hash, "recipients", len(peers), "duration", common.PrettyDuration(time.Since(block.ReceivedAt)))
|
log.Trace("Announced block", "hash", hash.Hex(), "recipients", len(peers), "duration", common.PrettyDuration(time.Since(block.ReceivedAt)))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -503,6 +505,7 @@ func (h *handler) BroadcastTransactions(txs types.Transactions) {
|
||||||
if tx.IsL1MessageTx() {
|
if tx.IsL1MessageTx() {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
log.Debug("Broadcasting transaction", "hash", tx.Hash().Hex(), "size", tx.Size())
|
||||||
peers := onlyShadowForkPeers(h.shadowForkPeerIDs, h.peers.peersWithoutTransaction(tx.Hash()))
|
peers := onlyShadowForkPeers(h.shadowForkPeerIDs, h.peers.peersWithoutTransaction(tx.Hash()))
|
||||||
// Send the tx unconditionally to a subset of our peers
|
// Send the tx unconditionally to a subset of our peers
|
||||||
numDirect := int(math.Sqrt(float64(len(peers))))
|
numDirect := int(math.Sqrt(float64(len(peers))))
|
||||||
|
|
@ -518,13 +521,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))
|
log.Trace("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.Trace("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,
|
||||||
|
|
|
||||||
|
|
@ -294,6 +294,7 @@ func handleNewBlock(backend Backend, msg Decoder, peer *Peer) error {
|
||||||
|
|
||||||
// Mark the peer as owning the block
|
// Mark the peer as owning the block
|
||||||
peer.markBlock(ann.Block.Hash())
|
peer.markBlock(ann.Block.Hash())
|
||||||
|
log.Debug("Received new block via gossip", "blockHash", ann.Block.Hash().Hex(), "blockNumber", ann.Block.NumberU64(), "peer", peer.String())
|
||||||
|
|
||||||
return backend.Handle(peer, ann)
|
return backend.Handle(peer, ann)
|
||||||
}
|
}
|
||||||
|
|
@ -362,12 +363,12 @@ 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)
|
log.Trace("Failed to decode `NewPooledTransactionHashesPacket`", "peer", peer.String(), "err", err)
|
||||||
newPooledTxHashesFailMeter.Mark(1)
|
newPooledTxHashesFailMeter.Mark(1)
|
||||||
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))
|
log.Trace("handleNewPooledTransactionHashes", "peer", peer.String(), "len(ann)", len(*ann))
|
||||||
newPooledTxHashesLenGauge.Update(int64(len(*ann)))
|
newPooledTxHashesLenGauge.Update(int64(len(*ann)))
|
||||||
for _, hash := range *ann {
|
for _, hash := range *ann {
|
||||||
peer.markTransaction(hash)
|
peer.markTransaction(hash)
|
||||||
|
|
@ -379,12 +380,15 @@ 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)
|
log.Trace("Failed to decode `GetPooledTransactionsPacket66`", "peer", peer.String(), "err", err)
|
||||||
getPooledTxsFailMeter.Mark(1)
|
getPooledTxsFailMeter.Mark(1)
|
||||||
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))
|
log.Trace("handleGetPooledTransactions", "peer", peer.String(), "RequestId", query.RequestId, "len(query)", len(query.GetPooledTransactionsPacket), "retrieved", len(hashes))
|
||||||
|
for _, hash := range hashes {
|
||||||
|
log.Debug("Received new pooled transaction", "hash", hash.Hex(), "peer", peer.String())
|
||||||
|
}
|
||||||
getPooledTxsQueryLenGauge.Update(int64(len(query.GetPooledTransactionsPacket)))
|
getPooledTxsQueryLenGauge.Update(int64(len(query.GetPooledTransactionsPacket)))
|
||||||
getPooledTxsRetrievedLenGauge.Update(int64(len(hashes)))
|
getPooledTxsRetrievedLenGauge.Update(int64(len(hashes)))
|
||||||
return peer.ReplyPooledTransactionsRLP(query.RequestId, hashes, txs)
|
return peer.ReplyPooledTransactionsRLP(query.RequestId, hashes, txs)
|
||||||
|
|
@ -427,16 +431,16 @@ func handleTransactions(backend Backend, msg Decoder, peer *Peer) error {
|
||||||
var txs TransactionsPacket
|
var txs TransactionsPacket
|
||||||
if err := msg.Decode(&txs); err != nil {
|
if err := msg.Decode(&txs); err != nil {
|
||||||
handleTxsFailMeter.Mark(1)
|
handleTxsFailMeter.Mark(1)
|
||||||
log.Debug("Failed to decode `TransactionsPacket`", "peer", peer.String(), "err", err)
|
log.Trace("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))
|
log.Trace("handleTransactions", "peer", peer.String(), "len(txs)", len(txs))
|
||||||
handleTxsLenGauge.Update(int64(len(txs)))
|
handleTxsLenGauge.Update(int64(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 {
|
||||||
handleTxsNilMeter.Mark(1)
|
handleTxsNilMeter.Mark(1)
|
||||||
log.Debug("handleTransactions: transaction is nil", "peer", peer.String(), "i", i)
|
log.Trace("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())
|
||||||
|
|
@ -453,16 +457,16 @@ func handlePooledTransactions66(backend Backend, msg Decoder, peer *Peer) error
|
||||||
var txs PooledTransactionsPacket66
|
var txs PooledTransactionsPacket66
|
||||||
if err := msg.Decode(&txs); err != nil {
|
if err := msg.Decode(&txs); err != nil {
|
||||||
pooledTxs66FailMeter.Mark(1)
|
pooledTxs66FailMeter.Mark(1)
|
||||||
log.Debug("Failed to decode `PooledTransactionsPacket66`", "peer", peer.String(), "err", err)
|
log.Trace("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)
|
||||||
}
|
}
|
||||||
log.Debug("handlePooledTransactions66", "peer", peer.String(), "len(txs)", len(txs.PooledTransactionsPacket))
|
log.Trace("handlePooledTransactions66", "peer", peer.String(), "len(txs)", len(txs.PooledTransactionsPacket))
|
||||||
pooledTxs66LenGauge.Update(int64(len(txs.PooledTransactionsPacket)))
|
pooledTxs66LenGauge.Update(int64(len(txs.PooledTransactionsPacket)))
|
||||||
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 {
|
||||||
pooledTxs66NilMeter.Mark(1)
|
pooledTxs66NilMeter.Mark(1)
|
||||||
log.Debug("handlePooledTransactions: transaction is nil", "peer", peer.String(), "i", i)
|
log.Trace("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())
|
||||||
|
|
|
||||||
|
|
@ -436,7 +436,7 @@ 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))
|
log.Trace("Requesting transactions", "RequestId", id, "Peer.id", p.id, "count", len(hashes))
|
||||||
peerRequestTxsCntGauge.Update(int64(len(hashes)))
|
peerRequestTxsCntGauge.Update(int64(len(hashes)))
|
||||||
|
|
||||||
requestTracker.Track(p.id, p.version, GetPooledTransactionsMsg, PooledTransactionsMsg, id)
|
requestTracker.Track(p.id, p.version, GetPooledTransactionsMsg, PooledTransactionsMsg, 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 = 8 // Minor version component of the current release
|
VersionMinor = 8 // Minor version component of the current release
|
||||||
VersionPatch = 47 // Patch version component of the current release
|
VersionPatch = 48 // 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