core: dedup txpool announced events, discard stales

This commit is contained in:
Péter Szilágyi 2019-06-20 10:57:06 +03:00
parent 9df266354e
commit 0dbe8bdc02
No known key found for this signature in database
GPG key ID: E9AE538CEDF8293D

View file

@ -915,16 +915,20 @@ func (pool *TxPool) scheduleReorgLoop() {
launchNextRun bool launchNextRun bool
reset *txpoolResetRequest reset *txpoolResetRequest
dirtyAccounts *accountSet dirtyAccounts *accountSet
queuedEvents []*types.Transaction queuedEvents = make(map[common.Address]*txSortedMap)
) )
for { for {
// Launch next run if needed. // Launch next background reorg if needed
if curDone == nil && launchNextRun { if curDone == nil && launchNextRun {
// Run the background reorg and announcements
go pool.runReorg(nextDone, reset, dirtyAccounts, queuedEvents) go pool.runReorg(nextDone, reset, dirtyAccounts, queuedEvents)
curDone = nextDone
nextDone = make(chan struct{}) // Prepare everything for the next round of reorg
curDone, nextDone = nextDone, make(chan struct{})
launchNextRun = false launchNextRun = false
reset, dirtyAccounts, queuedEvents = nil, nil, nil
reset, dirtyAccounts = nil, nil
queuedEvents = make(map[common.Address]*txSortedMap)
} }
select { select {
@ -951,7 +955,11 @@ func (pool *TxPool) scheduleReorgLoop() {
case tx := <-pool.queueTxEventCh: case tx := <-pool.queueTxEventCh:
// Queue up the event, but don't schedule a reorg. It's up to the caller to // Queue up the event, but don't schedule a reorg. It's up to the caller to
// request one later if they want the events sent. // request one later if they want the events sent.
queuedEvents = append(queuedEvents, tx) addr, _ := types.Sender(pool.signer, tx)
if _, ok := queuedEvents[addr]; !ok {
queuedEvents[addr] = newTxSortedMap()
}
queuedEvents[addr].Put(tx)
case <-curDone: case <-curDone:
curDone = nil curDone = nil
@ -968,30 +976,48 @@ func (pool *TxPool) scheduleReorgLoop() {
} }
// runReorg runs reset and promoteExecutables on behalf of scheduleReorgLoop. // runReorg runs reset and promoteExecutables on behalf of scheduleReorgLoop.
func (pool *TxPool) runReorg(done chan struct{}, reset *txpoolResetRequest, dirtyAccounts *accountSet, events []*types.Transaction) { func (pool *TxPool) runReorg(done chan struct{}, reset *txpoolResetRequest, dirtyAccounts *accountSet, events map[common.Address]*txSortedMap) {
defer close(done) defer close(done)
var promoteAddrs []common.Address var promoteAddrs []common.Address
if dirtyAccounts != nil { if dirtyAccounts != nil {
promoteAddrs = dirtyAccounts.flatten() promoteAddrs = dirtyAccounts.flatten()
} }
pool.mu.Lock() pool.mu.Lock()
if reset != nil { if reset != nil {
// Reset from the old head to the new, rescheduling any reorged transactions
pool.reset(reset.oldHead, reset.newHead) pool.reset(reset.oldHead, reset.newHead)
// Reset needs promote for all addresses.
// Nonces were reset, discard any events that became stale
for addr := range events {
events[addr].Forward(pool.pendingState.GetNonce(addr))
if events[addr].Len() == 0 {
delete(events, addr)
}
}
// Reset needs promote for all addresses
promoteAddrs = promoteAddrs[:0] promoteAddrs = promoteAddrs[:0]
for addr := range pool.queue { for addr := range pool.queue {
promoteAddrs = append(promoteAddrs, addr) promoteAddrs = append(promoteAddrs, addr)
} }
} }
promoted := pool.promoteExecutables(promoteAddrs) promoted := pool.promoteExecutables(promoteAddrs)
events = append(events, promoted...) for _, tx := range promoted {
addr, _ := types.Sender(pool.signer, tx)
if _, ok := events[addr]; !ok {
events[addr] = newTxSortedMap()
}
events[addr].Put(tx)
}
pool.mu.Unlock() pool.mu.Unlock()
// Notify subsystems for newly added transactions. // Notify subsystems for newly added transactions
if len(events) > 0 { if len(events) > 0 {
pool.txFeed.Send(NewTxsEvent{events}) var txs []*types.Transaction
for _, set := range events {
txs = append(txs, set.Flatten()...)
}
pool.txFeed.Send(NewTxsEvent{txs})
} }
} }