les/vflux/server: fix metrics (#22946)
* les/vflux/server: fix metrics * les/vflux/server: fix metrics
This commit is contained in:
parent
53b1420ede
commit
088bc34194
@ -143,8 +143,10 @@ func NewClientPool(balanceDb ethdb.KeyValueStore, minCap uint64, connectedBias t
|
|||||||
if oldState.HasAll(cp.setup.activeFlag) && oldState.HasNone(cp.setup.activeFlag) {
|
if oldState.HasAll(cp.setup.activeFlag) && oldState.HasNone(cp.setup.activeFlag) {
|
||||||
clientDeactivatedMeter.Mark(1)
|
clientDeactivatedMeter.Mark(1)
|
||||||
}
|
}
|
||||||
_, connected := cp.Active()
|
activeCount, activeCap := cp.Active()
|
||||||
totalConnectedGauge.Update(int64(connected))
|
totalActiveCountGauge.Update(int64(activeCount))
|
||||||
|
totalActiveCapacityGauge.Update(int64(activeCap))
|
||||||
|
totalInactiveCountGauge.Update(int64(cp.Inactive()))
|
||||||
})
|
})
|
||||||
return cp
|
return cp
|
||||||
}
|
}
|
||||||
|
@ -21,7 +21,9 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
totalConnectedGauge = metrics.NewRegisteredGauge("vflux/server/totalConnected", nil)
|
totalActiveCapacityGauge = metrics.NewRegisteredGauge("vflux/server/active/capacity", nil)
|
||||||
|
totalActiveCountGauge = metrics.NewRegisteredGauge("vflux/server/active/count", nil)
|
||||||
|
totalInactiveCountGauge = metrics.NewRegisteredGauge("vflux/server/inactive/count", nil)
|
||||||
|
|
||||||
clientConnectedMeter = metrics.NewRegisteredMeter("vflux/server/clientEvent/connected", nil)
|
clientConnectedMeter = metrics.NewRegisteredMeter("vflux/server/clientEvent/connected", nil)
|
||||||
clientActivatedMeter = metrics.NewRegisteredMeter("vflux/server/clientEvent/activated", nil)
|
clientActivatedMeter = metrics.NewRegisteredMeter("vflux/server/clientEvent/activated", nil)
|
||||||
|
@ -128,7 +128,7 @@ func newPriorityPool(ns *nodestate.NodeStateMachine, setup *serverSetup, clock m
|
|||||||
} else {
|
} else {
|
||||||
ns.SetStateSub(node, nodestate.Flags{}, pp.setup.activeFlag.Or(pp.setup.inactiveFlag), 0)
|
ns.SetStateSub(node, nodestate.Flags{}, pp.setup.activeFlag.Or(pp.setup.inactiveFlag), 0)
|
||||||
if n, _ := pp.ns.GetField(node, pp.setup.queueField).(*ppNodeInfo); n != nil {
|
if n, _ := pp.ns.GetField(node, pp.setup.queueField).(*ppNodeInfo); n != nil {
|
||||||
pp.disconnectedNode(n)
|
pp.disconnectNode(n)
|
||||||
}
|
}
|
||||||
ns.SetFieldSub(node, pp.setup.capacityField, nil)
|
ns.SetFieldSub(node, pp.setup.capacityField, nil)
|
||||||
ns.SetFieldSub(node, pp.setup.queueField, nil)
|
ns.SetFieldSub(node, pp.setup.queueField, nil)
|
||||||
@ -137,10 +137,10 @@ func newPriorityPool(ns *nodestate.NodeStateMachine, setup *serverSetup, clock m
|
|||||||
ns.SubscribeState(pp.setup.activeFlag.Or(pp.setup.inactiveFlag), func(node *enode.Node, oldState, newState nodestate.Flags) {
|
ns.SubscribeState(pp.setup.activeFlag.Or(pp.setup.inactiveFlag), func(node *enode.Node, oldState, newState nodestate.Flags) {
|
||||||
if c, _ := pp.ns.GetField(node, pp.setup.queueField).(*ppNodeInfo); c != nil {
|
if c, _ := pp.ns.GetField(node, pp.setup.queueField).(*ppNodeInfo); c != nil {
|
||||||
if oldState.IsEmpty() {
|
if oldState.IsEmpty() {
|
||||||
pp.connectedNode(c)
|
pp.connectNode(c)
|
||||||
}
|
}
|
||||||
if newState.IsEmpty() {
|
if newState.IsEmpty() {
|
||||||
pp.disconnectedNode(c)
|
pp.disconnectNode(c)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
@ -233,6 +233,14 @@ func (pp *priorityPool) Active() (uint64, uint64) {
|
|||||||
return pp.activeCount, pp.activeCap
|
return pp.activeCount, pp.activeCap
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Inactive returns the number of currently inactive nodes
|
||||||
|
func (pp *priorityPool) Inactive() int {
|
||||||
|
pp.lock.Lock()
|
||||||
|
defer pp.lock.Unlock()
|
||||||
|
|
||||||
|
return pp.inactiveQueue.Size()
|
||||||
|
}
|
||||||
|
|
||||||
// Limits returns the maximum allowed number and total capacity of active nodes
|
// Limits returns the maximum allowed number and total capacity of active nodes
|
||||||
func (pp *priorityPool) Limits() (uint64, uint64) {
|
func (pp *priorityPool) Limits() (uint64, uint64) {
|
||||||
pp.lock.Lock()
|
pp.lock.Lock()
|
||||||
@ -285,9 +293,9 @@ func (pp *priorityPool) inactivePriority(p *ppNodeInfo) int64 {
|
|||||||
return p.nodePriority.priority(pp.minCap)
|
return p.nodePriority.priority(pp.minCap)
|
||||||
}
|
}
|
||||||
|
|
||||||
// connectedNode is called when a new node has been added to the pool (inactiveFlag set)
|
// connectNode is called when a new node has been added to the pool (inactiveFlag set)
|
||||||
// Note: this function should run inside a NodeStateMachine operation
|
// Note: this function should run inside a NodeStateMachine operation
|
||||||
func (pp *priorityPool) connectedNode(c *ppNodeInfo) {
|
func (pp *priorityPool) connectNode(c *ppNodeInfo) {
|
||||||
pp.lock.Lock()
|
pp.lock.Lock()
|
||||||
pp.activeQueue.Refresh()
|
pp.activeQueue.Refresh()
|
||||||
if c.connected {
|
if c.connected {
|
||||||
@ -301,10 +309,10 @@ func (pp *priorityPool) connectedNode(c *ppNodeInfo) {
|
|||||||
pp.updateFlags(updates)
|
pp.updateFlags(updates)
|
||||||
}
|
}
|
||||||
|
|
||||||
// disconnectedNode is called when a node has been removed from the pool (both inactiveFlag
|
// disconnectNode is called when a node has been removed from the pool (both inactiveFlag
|
||||||
// and activeFlag reset)
|
// and activeFlag reset)
|
||||||
// Note: this function should run inside a NodeStateMachine operation
|
// Note: this function should run inside a NodeStateMachine operation
|
||||||
func (pp *priorityPool) disconnectedNode(c *ppNodeInfo) {
|
func (pp *priorityPool) disconnectNode(c *ppNodeInfo) {
|
||||||
pp.lock.Lock()
|
pp.lock.Lock()
|
||||||
pp.activeQueue.Refresh()
|
pp.activeQueue.Refresh()
|
||||||
if !c.connected {
|
if !c.connected {
|
||||||
|
Loading…
Reference in New Issue
Block a user