From 9babda7a3f6aa08b33e7197bfeba70ef1c14c65d Mon Sep 17 00:00:00 2001 From: Dmitrij Koniajev Date: Fri, 20 Nov 2020 17:55:48 +0200 Subject: [PATCH] Fixing a potential memory leak problem The underlying Timer is not recovered by the garbage collector until the timer fires. --- engine.go | 12 ++++++++++-- engine_test.go | 6 +++++- manager.go | 5 ++++- 3 files changed, 19 insertions(+), 4 deletions(-) diff --git a/engine.go b/engine.go index 4b672a6..a6fe905 100644 --- a/engine.go +++ b/engine.go @@ -234,13 +234,17 @@ func (e *Engine) startMainLoop() { <-e.txErr <-e.rxErr + timeout := time.Duration(5) * time.Second + idleTimer := time.NewTimer(timeout) + defer idleTimer.Stop() outer: for _, ob := range e.stObservers { for { + idleTimer.Reset(timeout) select { case ob <- e.state: continue outer - case <-time.After(5 * time.Second): + case <-idleTimer.C: log.Printf("Waited 5 seconds for state channel %v\n", ob) } } @@ -310,11 +314,15 @@ func (e *Engine) deliverToObservers(r Reply) { } func (e *Engine) deliverToObserver(c chan<- Reply, r Reply) { + timeout := time.Duration(5) * time.Second + idleTimer := time.NewTimer(timeout) + defer idleTimer.Stop() for { + idleTimer.Reset(timeout) select { case c <- r: return - case <-time.After(time.Duration(5) * time.Second): + case <-idleTimer.C: log.Printf("Waited 5 seconds for reply channel %v\n", c) } } diff --git a/engine_test.go b/engine_test.go index 14e6ace..50b51b0 100644 --- a/engine_test.go +++ b/engine_test.go @@ -25,9 +25,13 @@ func getGatewayURL() string { } func (e *Engine) expect(t *testing.T, seconds int, ch chan Reply, expected []IncomingMessageID) (Reply, error) { + timeout := time.Duration(seconds) * time.Second + idleTimer := time.NewTimer(timeout) + defer idleTimer.Stop() for { + idleTimer.Reset(timeout) select { - case <-time.After(time.Duration(seconds) * time.Second): + case <-idleTimer.C: return nil, errors.New("Timeout waiting") case v := <-ch: if v.code() == 0 { diff --git a/manager.go b/manager.go index ca2a729..3b58812 100644 --- a/manager.go +++ b/manager.go @@ -198,10 +198,13 @@ func (a *AbstractManager) Close() { // require the final result of a Manager (and have no interest in each update). // The Manager is guaranteed to be closed before it returns. func SinkManager(m Manager, timeout time.Duration, updateStop int) (updates int, err error) { + idleTimer := time.NewTimer(timeout) + defer idleTimer.Stop() for { sentClose := false + idleTimer.Reset(timeout) select { - case <-time.After(timeout): + case <-idleTimer.C: m.Close() return updates, fmt.Errorf("SinkManager: no new update in %s", timeout) case _, ok := <-m.Refresh():