les: clientpool updates and fixes

This commit is contained in:
Zsolt Felfoldi 2019-11-08 17:21:50 +01:00
parent 57bb045a03
commit 9117084a44
2 changed files with 19 additions and 17 deletions

View file

@ -161,6 +161,14 @@ func (bt *balanceTracker) timeUntil(priority int64) (time.Duration, bool) {
return time.Duration(dt), true return time.Duration(dt), true
} }
// setCapacity updates the capacity value used for priority calculation
func (bt *balanceTracker) setCapacity(capacity uint64) {
bt.lock.Lock()
defer bt.lock.Unlock()
bt.capacity = capacity
}
// getPriority returns the actual priority based on the current balance // getPriority returns the actual priority based on the current balance
func (bt *balanceTracker) getPriority(now mclock.AbsTime) int64 { func (bt *balanceTracker) getPriority(now mclock.AbsTime) int64 {
bt.lock.Lock() bt.lock.Lock()

View file

@ -254,19 +254,7 @@ func (f *clientPool) connect(peer clientPeer, capacity uint64) bool {
if !e.priority { if !e.priority {
capacity = f.freeClientCap capacity = f.freeClientCap
} }
// Ensure the capacity will never lower than the free capacity.
if capacity < f.freeClientCap {
capacity = f.freeClientCap
}
e.capacity = capacity e.capacity = capacity
if e.priority {
f.priorityConnected += capacity
}
// Starts a balance tracker
e.balanceTracker.init(f.clock, capacity)
e.balanceTracker.setBalance(posBalance, negBalance)
e.updatePriceFactors()
// If the number of clients already connected in the clientpool exceeds its // If the number of clients already connected in the clientpool exceeds its
// capacity, evict some clients with lowest priority. // capacity, evict some clients with lowest priority.
@ -310,9 +298,15 @@ func (f *clientPool) connect(peer clientPeer, capacity uint64) bool {
f.connectedQueue.Push(e) f.connectedQueue.Push(e)
f.connectedCap += e.capacity f.connectedCap += e.capacity
// Starts a balance tracker
e.balanceTracker.init(f.clock, capacity)
e.balanceTracker.setBalance(posBalance, negBalance)
e.updatePriceFactors()
// If the current client is a paid client, monitor the status of client, // If the current client is a paid client, monitor the status of client,
// downgrade it to normal client if positive balance is used up. // downgrade it to normal client if positive balance is used up.
if e.priority { if e.priority {
f.priorityConnected += capacity
e.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) }) e.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) })
} }
// If the capacity of client is not the default value(free capacity), notify // If the capacity of client is not the default value(free capacity), notify
@ -495,8 +489,8 @@ func (f *clientPool) setCapacity(c *clientInfo, capacity uint64) error {
oldCapacity := c.capacity oldCapacity := c.capacity
c.capacity = capacity c.capacity = capacity
f.connectedCap += capacity - oldCapacity f.connectedCap += capacity - oldCapacity
f.connectedQueue.Remove(c.queueIndex) c.balanceTracker.setCapacity(capacity)
f.connectedQueue.Push(c) f.connectedQueue.Update(c.queueIndex)
if f.connectedCap > f.capLimit { if f.connectedCap > f.capLimit {
var kickList []*clientInfo var kickList []*clientInfo
kick := true kick := true
@ -516,6 +510,7 @@ func (f *clientPool) setCapacity(c *clientInfo, capacity uint64) error {
} }
} else { } else {
c.capacity = oldCapacity c.capacity = oldCapacity
c.balanceTracker.setCapacity(oldCapacity)
for _, c := range kickList { for _, c := range kickList {
f.connectedCap += c.capacity f.connectedCap += c.capacity
f.connectedQueue.Push(c) f.connectedQueue.Push(c)
@ -525,6 +520,7 @@ func (f *clientPool) setCapacity(c *clientInfo, capacity uint64) error {
} }
totalConnectedGauge.Update(int64(f.connectedCap)) totalConnectedGauge.Update(int64(f.connectedCap))
f.priorityConnected += capacity - oldCapacity f.priorityConnected += capacity - oldCapacity
c.updatePriceFactors()
c.peer.updateCapacity(c.capacity) c.peer.updateCapacity(c.capacity)
return nil return nil
} }
@ -585,9 +581,7 @@ func (f *clientPool) updateBalance(id enode.ID, amount int64, meta string) error
// but we have no idea about the new capacity, need a second // but we have no idea about the new capacity, need a second
// call to udpate it. // call to udpate it.
c.priority = true c.priority = true
if c.priority {
f.priorityConnected += c.capacity f.priorityConnected += c.capacity
}
c.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) }) c.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) })
} }
c.balanceMetaInfo = meta c.balanceMetaInfo = meta