Merge pull request #21300 from rjl493456442/txpool-fix-queued-evictions
core: fix queued transaction eviction
This commit is contained in:
		
						commit
						6793ffa12b
					
				| @ -98,6 +98,7 @@ var ( | |||||||
| 	queuedReplaceMeter   = metrics.NewRegisteredMeter("txpool/queued/replace", nil) | 	queuedReplaceMeter   = metrics.NewRegisteredMeter("txpool/queued/replace", nil) | ||||||
| 	queuedRateLimitMeter = metrics.NewRegisteredMeter("txpool/queued/ratelimit", nil) // Dropped due to rate limiting
 | 	queuedRateLimitMeter = metrics.NewRegisteredMeter("txpool/queued/ratelimit", nil) // Dropped due to rate limiting
 | ||||||
| 	queuedNofundsMeter   = metrics.NewRegisteredMeter("txpool/queued/nofunds", nil)   // Dropped due to out-of-funds
 | 	queuedNofundsMeter   = metrics.NewRegisteredMeter("txpool/queued/nofunds", nil)   // Dropped due to out-of-funds
 | ||||||
|  | 	queuedEvictionMeter  = metrics.NewRegisteredMeter("txpool/queued/eviction", nil)  // Dropped due to lifetime
 | ||||||
| 
 | 
 | ||||||
| 	// General tx metrics
 | 	// General tx metrics
 | ||||||
| 	knownTxMeter       = metrics.NewRegisteredMeter("txpool/known", nil) | 	knownTxMeter       = metrics.NewRegisteredMeter("txpool/known", nil) | ||||||
| @ -362,9 +363,11 @@ func (pool *TxPool) loop() { | |||||||
| 				} | 				} | ||||||
| 				// Any non-locals old enough should be removed
 | 				// Any non-locals old enough should be removed
 | ||||||
| 				if time.Since(pool.beats[addr]) > pool.config.Lifetime { | 				if time.Since(pool.beats[addr]) > pool.config.Lifetime { | ||||||
| 					for _, tx := range pool.queue[addr].Flatten() { | 					list := pool.queue[addr].Flatten() | ||||||
|  | 					for _, tx := range list { | ||||||
| 						pool.removeTx(tx.Hash(), true) | 						pool.removeTx(tx.Hash(), true) | ||||||
| 					} | 					} | ||||||
|  | 					queuedEvictionMeter.Mark(int64(len(list))) | ||||||
| 				} | 				} | ||||||
| 			} | 			} | ||||||
| 			pool.mu.Unlock() | 			pool.mu.Unlock() | ||||||
| @ -614,6 +617,9 @@ func (pool *TxPool) add(tx *types.Transaction, local bool) (replaced bool, err e | |||||||
| 		pool.journalTx(from, tx) | 		pool.journalTx(from, tx) | ||||||
| 		pool.queueTxEvent(tx) | 		pool.queueTxEvent(tx) | ||||||
| 		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()) | ||||||
|  | 
 | ||||||
|  | 		// Successful promotion, bump the heartbeat
 | ||||||
|  | 		pool.beats[from] = time.Now() | ||||||
| 		return old != nil, nil | 		return old != nil, nil | ||||||
| 	} | 	} | ||||||
| 	// New transaction isn't replacing a pending one, push into queue
 | 	// New transaction isn't replacing a pending one, push into queue
 | ||||||
| @ -665,6 +671,10 @@ func (pool *TxPool) enqueueTx(hash common.Hash, tx *types.Transaction) (bool, er | |||||||
| 		pool.all.Add(tx) | 		pool.all.Add(tx) | ||||||
| 		pool.priced.Put(tx) | 		pool.priced.Put(tx) | ||||||
| 	} | 	} | ||||||
|  | 	// If we never record the heartbeat, do it right now.
 | ||||||
|  | 	if _, exist := pool.beats[from]; !exist { | ||||||
|  | 		pool.beats[from] = time.Now() | ||||||
|  | 	} | ||||||
| 	return old != nil, nil | 	return old != nil, nil | ||||||
| } | } | ||||||
| 
 | 
 | ||||||
| @ -696,7 +706,6 @@ func (pool *TxPool) promoteTx(addr common.Address, hash common.Hash, tx *types.T | |||||||
| 		// An older transaction was better, discard this
 | 		// An older transaction was better, discard this
 | ||||||
| 		pool.all.Remove(hash) | 		pool.all.Remove(hash) | ||||||
| 		pool.priced.Removed(1) | 		pool.priced.Removed(1) | ||||||
| 
 |  | ||||||
| 		pendingDiscardMeter.Mark(1) | 		pendingDiscardMeter.Mark(1) | ||||||
| 		return false | 		return false | ||||||
| 	} | 	} | ||||||
| @ -704,7 +713,6 @@ func (pool *TxPool) promoteTx(addr common.Address, hash common.Hash, tx *types.T | |||||||
| 	if old != nil { | 	if old != nil { | ||||||
| 		pool.all.Remove(old.Hash()) | 		pool.all.Remove(old.Hash()) | ||||||
| 		pool.priced.Removed(1) | 		pool.priced.Removed(1) | ||||||
| 
 |  | ||||||
| 		pendingReplaceMeter.Mark(1) | 		pendingReplaceMeter.Mark(1) | ||||||
| 	} else { | 	} else { | ||||||
| 		// Nothing was replaced, bump the pending counter
 | 		// Nothing was replaced, bump the pending counter
 | ||||||
| @ -716,9 +724,10 @@ func (pool *TxPool) promoteTx(addr common.Address, hash common.Hash, tx *types.T | |||||||
| 		pool.priced.Put(tx) | 		pool.priced.Put(tx) | ||||||
| 	} | 	} | ||||||
| 	// 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.pendingNonces.set(addr, tx.Nonce()+1) | 	pool.pendingNonces.set(addr, tx.Nonce()+1) | ||||||
| 
 | 
 | ||||||
|  | 	// Successful promotion, bump the heartbeat
 | ||||||
|  | 	pool.beats[addr] = time.Now() | ||||||
| 	return true | 	return true | ||||||
| } | } | ||||||
| 
 | 
 | ||||||
| @ -891,7 +900,6 @@ func (pool *TxPool) removeTx(hash common.Hash, outofbound bool) { | |||||||
| 			// If no more pending transactions are left, remove the list
 | 			// If no more pending transactions are left, remove the list
 | ||||||
| 			if pending.Empty() { | 			if pending.Empty() { | ||||||
| 				delete(pool.pending, addr) | 				delete(pool.pending, addr) | ||||||
| 				delete(pool.beats, addr) |  | ||||||
| 			} | 			} | ||||||
| 			// Postpone any invalidated transactions
 | 			// Postpone any invalidated transactions
 | ||||||
| 			for _, tx := range invalids { | 			for _, tx := range invalids { | ||||||
| @ -912,6 +920,7 @@ func (pool *TxPool) removeTx(hash common.Hash, outofbound bool) { | |||||||
| 		} | 		} | ||||||
| 		if future.Empty() { | 		if future.Empty() { | ||||||
| 			delete(pool.queue, addr) | 			delete(pool.queue, addr) | ||||||
|  | 			delete(pool.beats, addr) | ||||||
| 		} | 		} | ||||||
| 	} | 	} | ||||||
| } | } | ||||||
| @ -1229,6 +1238,7 @@ func (pool *TxPool) promoteExecutables(accounts []common.Address) []*types.Trans | |||||||
| 		// Delete the entire queue entry if it became empty.
 | 		// Delete the entire queue entry if it became empty.
 | ||||||
| 		if list.Empty() { | 		if list.Empty() { | ||||||
| 			delete(pool.queue, addr) | 			delete(pool.queue, addr) | ||||||
|  | 			delete(pool.beats, addr) | ||||||
| 		} | 		} | ||||||
| 	} | 	} | ||||||
| 	return promoted | 	return promoted | ||||||
| @ -1410,10 +1420,9 @@ func (pool *TxPool) demoteUnexecutables() { | |||||||
| 			} | 			} | ||||||
| 			pendingGauge.Dec(int64(len(gapped))) | 			pendingGauge.Dec(int64(len(gapped))) | ||||||
| 		} | 		} | ||||||
| 		// Delete the entire queue entry if it became empty.
 | 		// Delete the entire pending entry if it became empty.
 | ||||||
| 		if list.Empty() { | 		if list.Empty() { | ||||||
| 			delete(pool.pending, addr) | 			delete(pool.pending, addr) | ||||||
| 			delete(pool.beats, addr) |  | ||||||
| 		} | 		} | ||||||
| 	} | 	} | ||||||
| } | } | ||||||
|  | |||||||
| @ -109,6 +109,7 @@ func validateTxPoolInternals(pool *TxPool) error { | |||||||
| 	if priced := pool.priced.items.Len() - pool.priced.stales; priced != pending+queued { | 	if priced := pool.priced.items.Len() - pool.priced.stales; priced != pending+queued { | ||||||
| 		return fmt.Errorf("total priced transaction count %d != %d pending + %d queued", priced, pending, queued) | 		return fmt.Errorf("total priced transaction count %d != %d pending + %d queued", priced, pending, queued) | ||||||
| 	} | 	} | ||||||
|  | 
 | ||||||
| 	// Ensure the next nonce to assign is the correct one
 | 	// Ensure the next nonce to assign is the correct one
 | ||||||
| 	for addr, txs := range pool.pending { | 	for addr, txs := range pool.pending { | ||||||
| 		// Find the last transaction
 | 		// Find the last transaction
 | ||||||
| @ -868,7 +869,7 @@ func TestTransactionQueueTimeLimitingNoLocals(t *testing.T) { | |||||||
| func testTransactionQueueTimeLimiting(t *testing.T, nolocals bool) { | func testTransactionQueueTimeLimiting(t *testing.T, nolocals bool) { | ||||||
| 	// Reduce the eviction interval to a testable amount
 | 	// Reduce the eviction interval to a testable amount
 | ||||||
| 	defer func(old time.Duration) { evictionInterval = old }(evictionInterval) | 	defer func(old time.Duration) { evictionInterval = old }(evictionInterval) | ||||||
| 	evictionInterval = time.Second | 	evictionInterval = time.Millisecond * 100 | ||||||
| 
 | 
 | ||||||
| 	// Create the pool to test the non-expiration enforcement
 | 	// Create the pool to test the non-expiration enforcement
 | ||||||
| 	statedb, _ := state.New(common.Hash{}, state.NewDatabase(rawdb.NewMemoryDatabase()), nil) | 	statedb, _ := state.New(common.Hash{}, state.NewDatabase(rawdb.NewMemoryDatabase()), nil) | ||||||
| @ -905,6 +906,22 @@ func testTransactionQueueTimeLimiting(t *testing.T, nolocals bool) { | |||||||
| 	if err := validateTxPoolInternals(pool); err != nil { | 	if err := validateTxPoolInternals(pool); err != nil { | ||||||
| 		t.Fatalf("pool internal state corrupted: %v", err) | 		t.Fatalf("pool internal state corrupted: %v", err) | ||||||
| 	} | 	} | ||||||
|  | 
 | ||||||
|  | 	// Allow the eviction interval to run
 | ||||||
|  | 	time.Sleep(2 * evictionInterval) | ||||||
|  | 
 | ||||||
|  | 	// Transactions should not be evicted from the queue yet since lifetime duration has not passed
 | ||||||
|  | 	pending, queued = pool.Stats() | ||||||
|  | 	if pending != 0 { | ||||||
|  | 		t.Fatalf("pending transactions mismatched: have %d, want %d", pending, 0) | ||||||
|  | 	} | ||||||
|  | 	if queued != 2 { | ||||||
|  | 		t.Fatalf("queued transactions mismatched: have %d, want %d", queued, 2) | ||||||
|  | 	} | ||||||
|  | 	if err := validateTxPoolInternals(pool); err != nil { | ||||||
|  | 		t.Fatalf("pool internal state corrupted: %v", err) | ||||||
|  | 	} | ||||||
|  | 
 | ||||||
| 	// Wait a bit for eviction to run and clean up any leftovers, and ensure only the local remains
 | 	// Wait a bit for eviction to run and clean up any leftovers, and ensure only the local remains
 | ||||||
| 	time.Sleep(2 * config.Lifetime) | 	time.Sleep(2 * config.Lifetime) | ||||||
| 
 | 
 | ||||||
| @ -924,6 +941,72 @@ func testTransactionQueueTimeLimiting(t *testing.T, nolocals bool) { | |||||||
| 	if err := validateTxPoolInternals(pool); err != nil { | 	if err := validateTxPoolInternals(pool); err != nil { | ||||||
| 		t.Fatalf("pool internal state corrupted: %v", err) | 		t.Fatalf("pool internal state corrupted: %v", err) | ||||||
| 	} | 	} | ||||||
|  | 
 | ||||||
|  | 	// remove current transactions and increase nonce to prepare for a reset and cleanup
 | ||||||
|  | 	statedb.SetNonce(crypto.PubkeyToAddress(remote.PublicKey), 2) | ||||||
|  | 	statedb.SetNonce(crypto.PubkeyToAddress(local.PublicKey), 2) | ||||||
|  | 	<-pool.requestReset(nil, nil) | ||||||
|  | 
 | ||||||
|  | 	// make sure queue, pending are cleared
 | ||||||
|  | 	pending, queued = pool.Stats() | ||||||
|  | 	if pending != 0 { | ||||||
|  | 		t.Fatalf("pending transactions mismatched: have %d, want %d", pending, 0) | ||||||
|  | 	} | ||||||
|  | 	if queued != 0 { | ||||||
|  | 		t.Fatalf("queued transactions mismatched: have %d, want %d", queued, 0) | ||||||
|  | 	} | ||||||
|  | 	if err := validateTxPoolInternals(pool); err != nil { | ||||||
|  | 		t.Fatalf("pool internal state corrupted: %v", err) | ||||||
|  | 	} | ||||||
|  | 
 | ||||||
|  | 	// Queue gapped transactions
 | ||||||
|  | 	if err := pool.AddLocal(pricedTransaction(4, 100000, big.NewInt(1), local)); err != nil { | ||||||
|  | 		t.Fatalf("failed to add remote transaction: %v", err) | ||||||
|  | 	} | ||||||
|  | 	if err := pool.addRemoteSync(pricedTransaction(4, 100000, big.NewInt(1), remote)); err != nil { | ||||||
|  | 		t.Fatalf("failed to add remote transaction: %v", err) | ||||||
|  | 	} | ||||||
|  | 	time.Sleep(5 * evictionInterval) // A half lifetime pass
 | ||||||
|  | 
 | ||||||
|  | 	// Queue executable transactions, the life cycle should be restarted.
 | ||||||
|  | 	if err := pool.AddLocal(pricedTransaction(2, 100000, big.NewInt(1), local)); err != nil { | ||||||
|  | 		t.Fatalf("failed to add remote transaction: %v", err) | ||||||
|  | 	} | ||||||
|  | 	if err := pool.addRemoteSync(pricedTransaction(2, 100000, big.NewInt(1), remote)); err != nil { | ||||||
|  | 		t.Fatalf("failed to add remote transaction: %v", err) | ||||||
|  | 	} | ||||||
|  | 	time.Sleep(6 * evictionInterval) | ||||||
|  | 
 | ||||||
|  | 	// All gapped transactions shouldn't be kicked out
 | ||||||
|  | 	pending, queued = pool.Stats() | ||||||
|  | 	if pending != 2 { | ||||||
|  | 		t.Fatalf("pending transactions mismatched: have %d, want %d", pending, 2) | ||||||
|  | 	} | ||||||
|  | 	if queued != 2 { | ||||||
|  | 		t.Fatalf("queued transactions mismatched: have %d, want %d", queued, 3) | ||||||
|  | 	} | ||||||
|  | 	if err := validateTxPoolInternals(pool); err != nil { | ||||||
|  | 		t.Fatalf("pool internal state corrupted: %v", err) | ||||||
|  | 	} | ||||||
|  | 
 | ||||||
|  | 	// The whole life time pass after last promotion, kick out stale transactions
 | ||||||
|  | 	time.Sleep(2 * config.Lifetime) | ||||||
|  | 	pending, queued = pool.Stats() | ||||||
|  | 	if pending != 2 { | ||||||
|  | 		t.Fatalf("pending transactions mismatched: have %d, want %d", pending, 2) | ||||||
|  | 	} | ||||||
|  | 	if nolocals { | ||||||
|  | 		if queued != 0 { | ||||||
|  | 			t.Fatalf("queued transactions mismatched: have %d, want %d", queued, 0) | ||||||
|  | 		} | ||||||
|  | 	} else { | ||||||
|  | 		if queued != 1 { | ||||||
|  | 			t.Fatalf("queued transactions mismatched: have %d, want %d", queued, 1) | ||||||
|  | 		} | ||||||
|  | 	} | ||||||
|  | 	if err := validateTxPoolInternals(pool); err != nil { | ||||||
|  | 		t.Fatalf("pool internal state corrupted: %v", err) | ||||||
|  | 	} | ||||||
| } | } | ||||||
| 
 | 
 | ||||||
| // Tests that even if the transaction count belonging to a single account goes
 | // Tests that even if the transaction count belonging to a single account goes
 | ||||||
|  | |||||||
		Loading…
	
		Reference in New Issue
	
	Block a user