From 77ccd372f7d40905f65b217bd108b3b256af4c67 Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Thu, 30 May 2019 01:04:03 +0200 Subject: [PATCH] les/flowcontrol: regularly update total capacity --- les/flowcontrol/manager.go | 23 +++++++++++++++++++++++ les/server.go | 1 + 2 files changed, 24 insertions(+) diff --git a/les/flowcontrol/manager.go b/les/flowcontrol/manager.go index 7fa16728c8..68f1a47c97 100644 --- a/les/flowcontrol/manager.go +++ b/les/flowcontrol/manager.go @@ -61,6 +61,7 @@ type ClientManager struct { clock mclock.Clock lock sync.Mutex enabledCh chan struct{} + stop chan chan struct{} curve PieceWiseLinear sumRecharge, totalRecharge, totalConnected uint64 @@ -109,13 +110,35 @@ func NewClientManager(curve PieceWiseLinear, clock mclock.Clock) *ClientManager clock: clock, rcQueue: prque.New(func(a interface{}, i int) { a.(*ClientNode).queueIndex = i }), capLastUpdate: clock.Now(), + stop: make(chan chan struct{}), } if curve != nil { cm.SetRechargeCurve(curve) } + go func() { + // regularly recalculate and update total capacity + for { + select { + case <-time.After(time.Minute): + cm.lock.Lock() + cm.updateTotalCapacity(cm.clock.Now(), true) + cm.lock.Unlock() + case stop := <-cm.stop: + close(stop) + return + } + } + }() return cm } +// Stop stops the client manager +func (cm *ClientManager) Stop() { + stop := make(chan struct{}) + cm.stop <- stop + <-stop +} + // SetRechargeCurve updates the recharge curve func (cm *ClientManager) SetRechargeCurve(curve PieceWiseLinear) { cm.lock.Lock() diff --git a/les/server.go b/les/server.go index 18e497777e..836fa0d552 100644 --- a/les/server.go +++ b/les/server.go @@ -290,6 +290,7 @@ func (s *LesServer) SetBloomBitsIndexer(bloomIndexer *core.ChainIndexer) { // Stop stops the LES service func (s *LesServer) Stop() { + s.fcManager.Stop() s.chtIndexer.Close() // bloom trie indexer is closed by parent bloombits indexer go func() {