core: send tx events synchronously

This commit is contained in:
Felix Lange 2018-05-17 13:09:48 +02:00
parent 7b4f6c04dc
commit 70c1f89caf

View file

@ -192,7 +192,6 @@ type TxPool struct {
chainHeadSub event.Subscription chainHeadSub event.Subscription
signer types.Signer signer types.Signer
mu sync.RWMutex mu sync.RWMutex
promotedTxCh chan TxPreEvent
currentState *state.StateDB // Current state in the blockchain head currentState *state.StateDB // Current state in the blockchain head
pendingState *state.ManagedState // Pending state tracking virtual nonces pendingState *state.ManagedState // Pending state tracking virtual nonces
@ -220,17 +219,16 @@ func NewTxPool(config TxPoolConfig, chainconfig *params.ChainConfig, chain block
// Create the transaction pool with its initial settings // Create the transaction pool with its initial settings
pool := &TxPool{ pool := &TxPool{
config: config, config: config,
chainconfig: chainconfig, chainconfig: chainconfig,
chain: chain, chain: chain,
signer: types.NewEIP155Signer(chainconfig.ChainId), signer: types.NewEIP155Signer(chainconfig.ChainId),
pending: make(map[common.Address]*txList), pending: make(map[common.Address]*txList),
queue: make(map[common.Address]*txList), queue: make(map[common.Address]*txList),
beats: make(map[common.Address]time.Time), beats: make(map[common.Address]time.Time),
all: make(map[common.Hash]*types.Transaction), all: make(map[common.Hash]*types.Transaction),
chainHeadCh: make(chan ChainHeadEvent, chainHeadChanSize), chainHeadCh: make(chan ChainHeadEvent, chainHeadChanSize),
gasPrice: new(big.Int).SetUint64(config.PriceLimit), gasPrice: new(big.Int).SetUint64(config.PriceLimit),
promotedTxCh: make(chan TxPreEvent, config.GlobalSlots),
} }
pool.locals = newAccountSet(pool.signer) pool.locals = newAccountSet(pool.signer)
pool.priced = newTxPricedList(&pool.all) pool.priced = newTxPricedList(&pool.all)
@ -335,10 +333,6 @@ func (pool *TxPool) loop() {
} }
pool.mu.Unlock() pool.mu.Unlock()
} }
// Handle promoted transaction
case txPreEvent := <-pool.promotedTxCh:
pool.txFeed.Send(txPreEvent)
} }
} }
} }
@ -658,10 +652,8 @@ func (pool *TxPool) add(tx *types.Transaction, local bool) (bool, error) {
log.Trace("Pooled new executable transaction", "hash", hash, "from", from, "to", tx.To()) log.Trace("Pooled new executable transaction", "hash", hash, "from", from, "to", tx.To())
// We've directly injected a replacement transaction, notify subsystems // We've directly injected a replacement transaction, notify subsystems.
go func(tx *types.Transaction) { pool.txFeed.Send(TxPreEvent{tx})
pool.promotedTxCh <- TxPreEvent{tx}
}(tx)
return old != nil, nil return old != nil, nil
} }
@ -755,9 +747,7 @@ func (pool *TxPool) promoteTx(addr common.Address, hash common.Hash, tx *types.T
pool.beats[addr] = time.Now() pool.beats[addr] = time.Now()
pool.pendingState.SetNonce(addr, tx.Nonce()+1) pool.pendingState.SetNonce(addr, tx.Nonce()+1)
go func(tx *types.Transaction) { pool.txFeed.Send(TxPreEvent{tx})
pool.promotedTxCh <- TxPreEvent{tx}
}(tx)
} }
// AddLocal enqueues a single transaction into the pool if it is valid, marking // AddLocal enqueues a single transaction into the pool if it is valid, marking