From 89f5aa610460f84162fd31d4a7f697c1e5c82515 Mon Sep 17 00:00:00 2001 From: songzhibin97 <718428482@qq.com> Date: Wed, 13 Mar 2024 11:00:45 +0800 Subject: [PATCH] fix: use time.After in loop, possible memory leaks --- core/bloombits/matcher.go | 6 +++++- eth/downloader/beaconsync.go | 7 ++++++- eth/downloader/downloader.go | 30 ++++++++++++++++++++++-------- ethstats/ethstats.go | 6 +++++- p2p/simulations/adapters/exec.go | 6 +++++- p2p/simulations/mocker.go | 7 ++++++- p2p/simulations/network.go | 5 ++++- 7 files changed, 53 insertions(+), 14 deletions(-) diff --git a/core/bloombits/matcher.go b/core/bloombits/matcher.go index 6a4cfb23db..486581fe23 100644 --- a/core/bloombits/matcher.go +++ b/core/bloombits/matcher.go @@ -596,6 +596,9 @@ func (s *MatcherSession) deliverSections(bit uint, sections []uint64, bitsets [] // of the session, any request in-flight need to be responded to! Empty responses // are fine though in that case. func (s *MatcherSession) Multiplex(batch int, wait time.Duration, mux chan chan *Retrieval) { + waitTimer := time.NewTimer(wait) + defer waitTimer.Stop() + for { // Allocate a new bloom bit index to retrieve data for, stopping when done bit, ok := s.allocateRetrieval() @@ -604,6 +607,7 @@ func (s *MatcherSession) Multiplex(batch int, wait time.Duration, mux chan chan } // Bit allocated, throttle a bit if we're below our batch limit if s.pendingSections(bit) < batch { + waitTimer.Reset(wait) select { case <-s.quit: // Session terminating, we can't meaningfully service, abort @@ -611,7 +615,7 @@ func (s *MatcherSession) Multiplex(batch int, wait time.Duration, mux chan chan s.deliverSections(bit, []uint64{}, [][]byte{}) return - case <-time.After(wait): + case <-waitTimer.C: // Throttling up, fetch whatever is available } } diff --git a/eth/downloader/beaconsync.go b/eth/downloader/beaconsync.go index d3f75c8527..ee6039f5ba 100644 --- a/eth/downloader/beaconsync.go +++ b/eth/downloader/beaconsync.go @@ -289,6 +289,10 @@ func (d *Downloader) fetchBeaconHeaders(from uint64) error { localHeaders = d.readHeaderRange(tail, int(count)) log.Warn("Retrieved beacon headers from local", "from", from, "count", count) } + + fsHeaderContCheckTimer := time.NewTimer(fsHeaderContCheck) + defer fsHeaderContCheckTimer.Stop() + for { // Some beacon headers might have appeared since the last cycle, make // sure we're always syncing to all available ones @@ -381,8 +385,9 @@ func (d *Downloader) fetchBeaconHeaders(from uint64) error { } // State sync still going, wait a bit for new headers and retry log.Trace("Pivot not yet committed, waiting...") + fsHeaderContCheckTimer.Reset(fsHeaderContCheck) select { - case <-time.After(fsHeaderContCheck): + case <-fsHeaderContCheckTimer.C: case <-d.cancelCh: return errCanceled } diff --git a/eth/downloader/downloader.go b/eth/downloader/downloader.go index 6e7c5dcf02..bb2c835b39 100644 --- a/eth/downloader/downloader.go +++ b/eth/downloader/downloader.go @@ -1020,6 +1020,10 @@ func (d *Downloader) fetchHeaders(p *peerConnection, from uint64, head uint64) e ancestor = from mode = d.getMode() ) + + fsHeaderContCheckTimer := time.NewTimer(fsHeaderContCheck) + defer fsHeaderContCheckTimer.Stop() + for { // Pull the next batch of headers, it either: // - Pivot check to see if the chain moved too far @@ -1124,8 +1128,9 @@ func (d *Downloader) fetchHeaders(p *peerConnection, from uint64, head uint64) e // Don't abort header fetches while the pivot is downloading if !d.committed.Load() && pivot <= from { p.log.Debug("No headers, waiting for pivot commit") + fsHeaderContCheckTimer.Reset(fsHeaderContCheck) select { - case <-time.After(fsHeaderContCheck): + case <-fsHeaderContCheckTimer.C: continue case <-d.cancelCh: return errCanceled @@ -1194,9 +1199,10 @@ func (d *Downloader) fetchHeaders(p *peerConnection, from uint64, head uint64) e // sleep a bit and retry. Take care with headers already consumed during // skeleton filling if len(headers) == 0 && !progressed { + fsHeaderContCheckTimer.Reset(fsHeaderContCheck) p.log.Trace("All headers delayed, waiting") select { - case <-time.After(fsHeaderContCheck): + case <-fsHeaderContCheckTimer.C: continue case <-d.cancelCh: return errCanceled @@ -1274,9 +1280,13 @@ func (d *Downloader) fetchReceipts(from uint64, beaconMode bool) error { // queue until the stream ends or a failure occurs. func (d *Downloader) processHeaders(origin uint64, td, ttd *big.Int, beaconMode bool) error { var ( - mode = d.getMode() - gotHeaders = false // Wait for batches of headers to process + mode = d.getMode() + gotHeaders = false // Wait for batches of headers to process + secondTimer = time.NewTimer(time.Second) ) + + defer secondTimer.Stop() + for { select { case <-d.cancelCh: @@ -1397,10 +1407,11 @@ func (d *Downloader) processHeaders(origin uint64, td, ttd *big.Int, beaconMode if mode == FullSync || mode == SnapSync { // If we've reached the allowed number of pending headers, stall a bit for d.queue.PendingBodies() >= maxQueuedHeaders || d.queue.PendingReceipts() >= maxQueuedHeaders { + secondTimer.Reset(time.Second) select { case <-d.cancelCh: return errCanceled - case <-time.After(time.Second): + case <-secondTimer.C: } } // Otherwise insert the headers for content retrieval @@ -1565,9 +1576,11 @@ func (d *Downloader) processSnapSyncContent() error { // Note, there's no issue with memory piling up since after 64 blocks the // pivot will forcefully move so these accumulators will be dropped. var ( - oldPivot *fetchResult // Locked in pivot block, might change eventually - oldTail []*fetchResult // Downloaded content after the pivot + oldPivot *fetchResult // Locked in pivot block, might change eventually + oldTail []*fetchResult // Downloaded content after the pivot + secondTimer = time.NewTimer(time.Second) ) + defer secondTimer.Stop() for { // Wait for the next batch of downloaded data to be available. If we have // not yet reached the pivot point, wait blockingly as there's no need to @@ -1650,6 +1663,7 @@ func (d *Downloader) processSnapSyncContent() error { oldPivot = P } // Wait for completion, occasionally checking for pivot staleness + secondTimer.Reset(time.Second) select { case <-sync.done: if sync.err != nil { @@ -1660,7 +1674,7 @@ func (d *Downloader) processSnapSyncContent() error { } oldPivot = nil - case <-time.After(time.Second): + case <-secondTimer.C: oldTail = afterP continue } diff --git a/ethstats/ethstats.go b/ethstats/ethstats.go index 6e71666ec1..0c334df748 100644 --- a/ethstats/ethstats.go +++ b/ethstats/ethstats.go @@ -544,10 +544,14 @@ func (s *Service) reportLatency(conn *connWrapper) error { return err } // Wait for the pong request to arrive back + + timer := time.NewTimer(5 * time.Second) + defer timer.Stop() + select { case <-s.pongCh: // Pong delivered, report the latency - case <-time.After(5 * time.Second): + case <-timer.C: // Ping timeout, abort return errors.New("ping timed out") } diff --git a/p2p/simulations/adapters/exec.go b/p2p/simulations/adapters/exec.go index 17e0f75d5a..10cfd12d71 100644 --- a/p2p/simulations/adapters/exec.go +++ b/p2p/simulations/adapters/exec.go @@ -303,10 +303,14 @@ func (n *ExecNode) Stop() error { go func() { waitErr <- n.Cmd.Wait() }() + + timer := time.NewTimer(5 * time.Second) + defer timer.Stop() + select { case err := <-waitErr: return err - case <-time.After(5 * time.Second): + case <-timer.C: return n.Cmd.Process.Kill() } } diff --git a/p2p/simulations/mocker.go b/p2p/simulations/mocker.go index 0dc04e65f9..1388ef503f 100644 --- a/p2p/simulations/mocker.go +++ b/p2p/simulations/mocker.go @@ -67,6 +67,10 @@ func startStop(net *Network, quit chan struct{}, nodeCount int) { } tick := time.NewTicker(10 * time.Second) defer tick.Stop() + + timer := time.NewTimer(3 * time.Second) + defer timer.Stop() + for { select { case <-quit: @@ -80,11 +84,12 @@ func startStop(net *Network, quit chan struct{}, nodeCount int) { return } + timer.Reset(3 * time.Second) select { case <-quit: log.Info("Terminating simulation loop") return - case <-time.After(3 * time.Second): + case <-timer.C: } log.Debug("starting node", "id", id) diff --git a/p2p/simulations/network.go b/p2p/simulations/network.go index 4735e5cfa6..9c6230d832 100644 --- a/p2p/simulations/network.go +++ b/p2p/simulations/network.go @@ -1028,11 +1028,14 @@ func (net *Network) Load(snap *Snapshot) error { } } + snapshotLoadTimeoutTimer := time.NewTimer(snapshotLoadTimeout) + defer snapshotLoadTimeoutTimer.Stop() + select { // Wait until all connections from the snapshot are established. case <-allConnected: // Make sure that we do not wait forever. - case <-time.After(snapshotLoadTimeout): + case <-snapshotLoadTimeoutTimer.C: return errors.New("snapshot connections not established") } return nil