mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 10:22:23 +00:00
logging logic in tx pool when 1) nonce gap occurs and 2) txns originate from different accounts
This commit is contained in:
parent
5b6e854eb3
commit
5df21a1e6f
4 changed files with 20 additions and 1 deletions
|
|
@ -186,6 +186,7 @@ func (m *txSortedMap) Ready(start uint64) types.Transactions {
|
||||||
ready = append(ready, m.items[next])
|
ready = append(ready, m.items[next])
|
||||||
delete(m.items, next)
|
delete(m.items, next)
|
||||||
heap.Pop(m.index)
|
heap.Pop(m.index)
|
||||||
|
log.Info("=====> Txn List adding txn with sequentially valid nonce that is greater than current pool nonce", "txn", m.items[next])
|
||||||
}
|
}
|
||||||
m.cache = nil
|
m.cache = nil
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -661,14 +661,17 @@ func (pool *TxPool) add(tx *types.Transaction, local bool) (bool, error) {
|
||||||
// If the transaction is replacing an already pending one, do directly
|
// If the transaction is replacing an already pending one, do directly
|
||||||
from, _ := types.Sender(pool.signer, tx) // already validated
|
from, _ := types.Sender(pool.signer, tx) // already validated
|
||||||
if list := pool.pending[from]; list != nil && list.Overlaps(tx) {
|
if list := pool.pending[from]; list != nil && list.Overlaps(tx) {
|
||||||
|
log.Info("======> Pool detected txn nonce already in txn list", "nonce", tx.Nonce())
|
||||||
// Nonce already pending, check if required price bump is met
|
// Nonce already pending, check if required price bump is met
|
||||||
inserted, old := list.Add(tx, pool.config.PriceBump)
|
inserted, old := list.Add(tx, pool.config.PriceBump)
|
||||||
if !inserted {
|
if !inserted {
|
||||||
|
log.Info("======> Older txn with shared nonce has higher gas price, accepted over new txn", "oldHash", old.Hash(), "oldGasPrice", old.GasPrice(), "txnGasPrice", tx.GasPrice())
|
||||||
pendingDiscardCounter.Inc(1)
|
pendingDiscardCounter.Inc(1)
|
||||||
return false, ErrReplaceUnderpriced
|
return false, ErrReplaceUnderpriced
|
||||||
}
|
}
|
||||||
// New transaction is better, replace old one
|
// New transaction is better, replace old one
|
||||||
if old != nil {
|
if old != nil {
|
||||||
|
log.Info("======> New txn with shared nonce has higher gas price, overwriting old txn", "oldHash", old.Hash())
|
||||||
pool.all.Remove(old.Hash())
|
pool.all.Remove(old.Hash())
|
||||||
pool.priced.Removed()
|
pool.priced.Removed()
|
||||||
pendingReplaceCounter.Inc(1)
|
pendingReplaceCounter.Inc(1)
|
||||||
|
|
@ -750,9 +753,11 @@ func (pool *TxPool) promoteTx(addr common.Address, hash common.Hash, tx *types.T
|
||||||
// Try to insert the transaction into the pending queue
|
// Try to insert the transaction into the pending queue
|
||||||
if pool.pending[addr] == nil {
|
if pool.pending[addr] == nil {
|
||||||
pool.pending[addr] = newTxList(true)
|
pool.pending[addr] = newTxList(true)
|
||||||
|
log.Info("======> Pool intitializing new txn list in pending queue under address", "address", addr)
|
||||||
}
|
}
|
||||||
list := pool.pending[addr]
|
list := pool.pending[addr]
|
||||||
|
|
||||||
|
log.Info("=====> Pool attempting to promote txn with hash", "hash", hash)
|
||||||
inserted, old := list.Add(tx, pool.config.PriceBump)
|
inserted, old := list.Add(tx, pool.config.PriceBump)
|
||||||
if !inserted {
|
if !inserted {
|
||||||
// An older transaction was better, discard this
|
// An older transaction was better, discard this
|
||||||
|
|
@ -777,6 +782,7 @@ func (pool *TxPool) promoteTx(addr common.Address, hash common.Hash, tx *types.T
|
||||||
// Set the potentially new pending nonce and notify any subsystems of the new tx
|
// Set the potentially new pending nonce and notify any subsystems of the new tx
|
||||||
pool.beats[addr] = time.Now()
|
pool.beats[addr] = time.Now()
|
||||||
pool.pendingState.SetNonce(addr, tx.Nonce()+1)
|
pool.pendingState.SetNonce(addr, tx.Nonce()+1)
|
||||||
|
log.Info("======> Pool promoted next txn for execution, incrementing nonce at address", "addr", addr, "nonce", pool.pendingState.GetNonce(addr))
|
||||||
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
@ -971,7 +977,9 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
}
|
}
|
||||||
// Gather all executable transactions and promote them
|
// Gather all executable transactions and promote them
|
||||||
log.Info("Pool promoting txns ready for execution")
|
log.Info("Pool promoting txns ready for execution")
|
||||||
|
log.Info("=====> Pool calling list.Ready() to grab all txns with nonce greater than current pool state under this address", "address", addr)
|
||||||
for _, tx := range list.Ready(pool.pendingState.GetNonce(addr)) {
|
for _, tx := range list.Ready(pool.pendingState.GetNonce(addr)) {
|
||||||
|
log.Info("=====> Pool retrieved txns with nonce greater than currentState.nonce, sorted in sequentially increasing order", "list", list.Ready(pool.pendingState.GetNonce(addr), "current nonce", pool.pendingState.GetNonce(addr)))
|
||||||
hash := tx.Hash()
|
hash := tx.Hash()
|
||||||
if pool.promoteTx(addr, hash, tx) {
|
if pool.promoteTx(addr, hash, tx) {
|
||||||
log.Trace("Promoting queued transaction", "hash", hash)
|
log.Trace("Promoting queued transaction", "hash", hash)
|
||||||
|
|
@ -996,6 +1004,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
}
|
}
|
||||||
// Notify subsystem for new promoted transactions.
|
// Notify subsystem for new promoted transactions.
|
||||||
if len(promoted) > 0 {
|
if len(promoted) > 0 {
|
||||||
|
log.Info("======> New promoted txns", "promoted", promoted)
|
||||||
go pool.txFeed.Send(NewTxsEvent{promoted})
|
go pool.txFeed.Send(NewTxsEvent{promoted})
|
||||||
}
|
}
|
||||||
// If the pending limit is overflown, start equalizing allowances
|
// If the pending limit is overflown, start equalizing allowances
|
||||||
|
|
@ -1004,6 +1013,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
pending += uint64(list.Len())
|
pending += uint64(list.Len())
|
||||||
}
|
}
|
||||||
if pending > pool.config.GlobalSlots {
|
if pending > pool.config.GlobalSlots {
|
||||||
|
log.Info("======> Pending exceeds pool size limit", "pending", pending, "limit", pool.config.GlobalSlots)
|
||||||
pendingBeforeCap := pending
|
pendingBeforeCap := pending
|
||||||
// Assemble a spam order to penalize large transactors first
|
// Assemble a spam order to penalize large transactors first
|
||||||
spammers := prque.New(nil)
|
spammers := prque.New(nil)
|
||||||
|
|
@ -1015,6 +1025,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
}
|
}
|
||||||
// Gradually drop transactions from offenders
|
// Gradually drop transactions from offenders
|
||||||
offenders := []common.Address{}
|
offenders := []common.Address{}
|
||||||
|
log.Info("======> Found accounts with too many transactions, dropping their transactions", "spammers", spammers)
|
||||||
for pending > pool.config.GlobalSlots && !spammers.Empty() {
|
for pending > pool.config.GlobalSlots && !spammers.Empty() {
|
||||||
// Retrieve the next offender if not local address
|
// Retrieve the next offender if not local address
|
||||||
offender, _ := spammers.Pop()
|
offender, _ := spammers.Pop()
|
||||||
|
|
@ -1031,6 +1042,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
list := pool.pending[offenders[i]]
|
list := pool.pending[offenders[i]]
|
||||||
for _, tx := range list.Cap(list.Len() - 1) {
|
for _, tx := range list.Cap(list.Len() - 1) {
|
||||||
// Drop the transaction from the global pools too
|
// Drop the transaction from the global pools too
|
||||||
|
log.Info("=======> Dropping transaction from pool", "hash", tx.Hash(), "offender", offenders[i])
|
||||||
hash := tx.Hash()
|
hash := tx.Hash()
|
||||||
pool.all.Remove(hash)
|
pool.all.Remove(hash)
|
||||||
pool.priced.Removed()
|
pool.priced.Removed()
|
||||||
|
|
@ -1038,6 +1050,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
// Update the account nonce to the dropped transaction
|
// Update the account nonce to the dropped transaction
|
||||||
if nonce := tx.Nonce(); pool.pendingState.GetNonce(offenders[i]) > nonce {
|
if nonce := tx.Nonce(); pool.pendingState.GetNonce(offenders[i]) > nonce {
|
||||||
pool.pendingState.SetNonce(offenders[i], nonce)
|
pool.pendingState.SetNonce(offenders[i], nonce)
|
||||||
|
log.Info("=======> Setting offender's nonce to the pool's nonce state", "nonce", nonce)
|
||||||
}
|
}
|
||||||
log.Trace("Removed fairness-exceeding pending transaction", "hash", hash)
|
log.Trace("Removed fairness-exceeding pending transaction", "hash", hash)
|
||||||
}
|
}
|
||||||
|
|
@ -1060,6 +1073,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) {
|
||||||
// Update the account nonce to the dropped transaction
|
// Update the account nonce to the dropped transaction
|
||||||
if nonce := tx.Nonce(); pool.pendingState.GetNonce(addr) > nonce {
|
if nonce := tx.Nonce(); pool.pendingState.GetNonce(addr) > nonce {
|
||||||
pool.pendingState.SetNonce(addr, nonce)
|
pool.pendingState.SetNonce(addr, nonce)
|
||||||
|
log.Info("=======> Setting offender's nonce to the pool's nonce state", "nonce", nonce)
|
||||||
}
|
}
|
||||||
log.Trace("Removed fairness-exceeding pending transaction", "hash", hash)
|
log.Trace("Removed fairness-exceeding pending transaction", "hash", hash)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -26,6 +26,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/common/hexutil"
|
"github.com/ethereum/go-ethereum/common/hexutil"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -335,7 +336,9 @@ type TransactionsByPriceAndNonce struct {
|
||||||
func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transactions) *TransactionsByPriceAndNonce {
|
func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transactions) *TransactionsByPriceAndNonce {
|
||||||
// Initialize a price based heap with the head transactions
|
// Initialize a price based heap with the head transactions
|
||||||
heads := make(TxByPrice, 0, len(txs))
|
heads := make(TxByPrice, 0, len(txs))
|
||||||
|
log.Info("======> Sorting the txns from pool by initializing a heap and iteratively deleting from heap, causing next best txn price to bubble up")
|
||||||
for from, accTxs := range txs {
|
for from, accTxs := range txs {
|
||||||
|
log.Info("======> Txn with next best price/nonce", "nextBestTxn", accTxs[0])
|
||||||
heads = append(heads, accTxs[0])
|
heads = append(heads, accTxs[0])
|
||||||
// Ensure the sender address is from the signer
|
// Ensure the sender address is from the signer
|
||||||
acc, _ := Sender(signer, accTxs[0])
|
acc, _ := Sender(signer, accTxs[0])
|
||||||
|
|
@ -345,6 +348,7 @@ func NewTransactionsByPriceAndNonce(signer Signer, txs map[common.Address]Transa
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
heap.Init(&heads)
|
heap.Init(&heads)
|
||||||
|
log.Info("======> Returning sorted txn heap", "heap", heads)
|
||||||
|
|
||||||
// Assemble and return the transaction set
|
// Assemble and return the transaction set
|
||||||
return &TransactionsByPriceAndNonce{
|
return &TransactionsByPriceAndNonce{
|
||||||
|
|
|
||||||
|
|
@ -918,7 +918,7 @@ func (w *worker) commitNewWork(interrupt *int32, noempty bool, timestamp int64)
|
||||||
|
|
||||||
// Fill the block with all available pending transactions.
|
// Fill the block with all available pending transactions.
|
||||||
pending, err := w.eth.TxPool().Pending()
|
pending, err := w.eth.TxPool().Pending()
|
||||||
log.Info("Worker found new pending txns from pool", "num_txns", len(pending), "txns", pending, "location", whereami.WhereAmI())
|
log.Info("Worker found new pending txns from pending queue set by tx pool", "num_txns", len(pending), "txns", pending, "location", whereami.WhereAmI())
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Failed to fetch pending transactions", "err", err)
|
log.Error("Failed to fetch pending transactions", "err", err)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue