mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les: address comments
This commit is contained in:
parent
a6fe625580
commit
77ae572cfe
5 changed files with 22 additions and 18 deletions
|
|
@ -22,6 +22,7 @@ import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/eth/downloader"
|
"github.com/ethereum/go-ethereum/eth/downloader"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
|
|
@ -115,10 +116,10 @@ func (h *clientHandler) handle(p *peer) error {
|
||||||
}
|
}
|
||||||
serverConnectionGauge.Update(int64(h.backend.peers.Len()))
|
serverConnectionGauge.Update(int64(h.backend.peers.Len()))
|
||||||
|
|
||||||
connectedAt := time.Now()
|
connectedAt := mclock.Now()
|
||||||
defer func() {
|
defer func() {
|
||||||
h.backend.peers.Unregister(p.id)
|
h.backend.peers.Unregister(p.id)
|
||||||
connectionTimer.UpdateSince(connectedAt)
|
connectionTimer.Update(time.Duration(mclock.Now() - connectedAt))
|
||||||
serverConnectionGauge.Update(int64(h.backend.peers.Len()))
|
serverConnectionGauge.Update(int64(h.backend.peers.Len()))
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -72,7 +72,7 @@ var (
|
||||||
|
|
||||||
requestServedMeter = metrics.NewRegisteredMeter("les/server/req/avgServedTime", nil)
|
requestServedMeter = metrics.NewRegisteredMeter("les/server/req/avgServedTime", nil)
|
||||||
requestServedTimer = metrics.NewRegisteredTimer("les/server/req/servedTime", nil)
|
requestServedTimer = metrics.NewRegisteredTimer("les/server/req/servedTime", nil)
|
||||||
requestEstimatedMeter = metrics.NewRegisteredMeter("les/server/req/argEstimatedTime", nil)
|
requestEstimatedMeter = metrics.NewRegisteredMeter("les/server/req/avgEstimatedTime", nil)
|
||||||
requestEstimatedTimer = metrics.NewRegisteredTimer("les/server/req/estimatedTime", nil)
|
requestEstimatedTimer = metrics.NewRegisteredTimer("les/server/req/estimatedTime", nil)
|
||||||
relativeCostHistogram = metrics.NewRegisteredHistogram("les/server/req/relative", nil, metrics.NewExpDecaySample(1028, 0.015))
|
relativeCostHistogram = metrics.NewRegisteredHistogram("les/server/req/relative", nil, metrics.NewExpDecaySample(1028, 0.015))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ import (
|
||||||
"context"
|
"context"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
"github.com/ethereum/go-ethereum/core"
|
"github.com/ethereum/go-ethereum/core"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
|
|
@ -121,11 +122,11 @@ func (odr *LesOdr) Retrieve(ctx context.Context, req light.OdrRequest) (err erro
|
||||||
return func() { lreq.Request(reqID, p) }
|
return func() { lreq.Request(reqID, p) }
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
sent := time.Now()
|
sent := mclock.Now()
|
||||||
if err = odr.retriever.retrieve(ctx, reqID, rq, func(p distPeer, msg *Msg) error { return lreq.Validate(odr.db, msg) }, odr.stop); err == nil {
|
if err = odr.retriever.retrieve(ctx, reqID, rq, func(p distPeer, msg *Msg) error { return lreq.Validate(odr.db, msg) }, odr.stop); err == nil {
|
||||||
// retrieved from network, store in db
|
// retrieved from network, store in db
|
||||||
req.StoreResult(odr.db)
|
req.StoreResult(odr.db)
|
||||||
requestRTT.UpdateSince(sent)
|
requestRTT.Update(time.Duration(mclock.Now() - sent))
|
||||||
} else {
|
} else {
|
||||||
log.Debug("Failed to retrieve data from network", "err", err)
|
log.Debug("Failed to retrieve data from network", "err", err)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -221,7 +221,7 @@ func (s *LesServer) capacityManagement() {
|
||||||
var (
|
var (
|
||||||
busy bool
|
busy bool
|
||||||
freePeers uint64
|
freePeers uint64
|
||||||
blockProcess time.Time
|
blockProcess mclock.AbsTime
|
||||||
)
|
)
|
||||||
updateRecharge := func() {
|
updateRecharge := func() {
|
||||||
if busy {
|
if busy {
|
||||||
|
|
@ -238,9 +238,9 @@ func (s *LesServer) capacityManagement() {
|
||||||
select {
|
select {
|
||||||
case busy = <-processCh:
|
case busy = <-processCh:
|
||||||
if busy {
|
if busy {
|
||||||
blockProcess = time.Now()
|
blockProcess = mclock.Now()
|
||||||
} else {
|
} else {
|
||||||
blockProcessingTimer.UpdateSince(blockProcess)
|
blockProcessingTimer.Update(time.Duration(mclock.Now() - blockProcess))
|
||||||
}
|
}
|
||||||
updateRecharge()
|
updateRecharge()
|
||||||
case totalRecharge = <-totalRechargeCh:
|
case totalRecharge = <-totalRechargeCh:
|
||||||
|
|
|
||||||
|
|
@ -137,12 +137,12 @@ func (h *serverHandler) handle(p *peer) error {
|
||||||
p.balanceTracker.init(&mclock.System{}, 1)
|
p.balanceTracker.init(&mclock.System{}, 1)
|
||||||
}
|
}
|
||||||
|
|
||||||
connectedAt := time.Now()
|
connectedAt := mclock.Now()
|
||||||
defer func() {
|
defer func() {
|
||||||
p.balanceTracker = nil
|
p.balanceTracker = nil
|
||||||
h.server.peers.Unregister(p.id)
|
h.server.peers.Unregister(p.id)
|
||||||
clientConnectionGauge.Update(int64(h.server.peers.Len()))
|
clientConnectionGauge.Update(int64(h.server.peers.Len()))
|
||||||
connectionTimer.UpdateSince(connectedAt)
|
connectionTimer.Update(time.Duration(mclock.Now() - connectedAt))
|
||||||
}()
|
}()
|
||||||
|
|
||||||
// Spawn a main loop to handle all incoming messages.
|
// Spawn a main loop to handle all incoming messages.
|
||||||
|
|
@ -180,8 +180,9 @@ func (h *serverHandler) handleMsg(p *peer) error {
|
||||||
var (
|
var (
|
||||||
maxCost uint64
|
maxCost uint64
|
||||||
task *servingTask
|
task *servingTask
|
||||||
respId = p.responseID()
|
|
||||||
)
|
)
|
||||||
|
p.responseCount++
|
||||||
|
responseCount := p.responseCount
|
||||||
// accept returns an indicator whether the request can be served.
|
// accept returns an indicator whether the request can be served.
|
||||||
// If so, deduct the max cost from the flow control buffer.
|
// If so, deduct the max cost from the flow control buffer.
|
||||||
accept := func(reqID, reqCnt, maxCnt uint64) bool {
|
accept := func(reqID, reqCnt, maxCnt uint64) bool {
|
||||||
|
|
@ -193,7 +194,7 @@ func (h *serverHandler) handleMsg(p *peer) error {
|
||||||
}
|
}
|
||||||
// Prepaid max cost units before request been serving.
|
// Prepaid max cost units before request been serving.
|
||||||
maxCost = p.fcCosts.getMaxCost(msg.Code, reqCnt)
|
maxCost = p.fcCosts.getMaxCost(msg.Code, reqCnt)
|
||||||
accepted, bufShort, priority := p.fcClient.AcceptRequest(reqID, respId, maxCost)
|
accepted, bufShort, priority := p.fcClient.AcceptRequest(reqID, responseCount, maxCost)
|
||||||
if !accepted {
|
if !accepted {
|
||||||
p.freezeClient()
|
p.freezeClient()
|
||||||
p.Log().Error("Request came too early", "remaining", common.PrettyDuration(time.Duration(bufShort*1000000/p.fcParams.MinRecharge)))
|
p.Log().Error("Request came too early", "remaining", common.PrettyDuration(time.Duration(bufShort*1000000/p.fcParams.MinRecharge)))
|
||||||
|
|
@ -212,7 +213,7 @@ func (h *serverHandler) handleMsg(p *peer) error {
|
||||||
if task.start() {
|
if task.start() {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
p.fcClient.RequestProcessed(reqID, respId, maxCost, inSizeCost)
|
p.fcClient.RequestProcessed(reqID, responseCount, maxCost, inSizeCost)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
// sendResponse sends back the response and updates the flow control statistic.
|
// sendResponse sends back the response and updates the flow control statistic.
|
||||||
|
|
@ -223,7 +224,7 @@ func (h *serverHandler) handleMsg(p *peer) error {
|
||||||
// Short circuit if the client is already frozen.
|
// Short circuit if the client is already frozen.
|
||||||
if p.isFrozen() {
|
if p.isFrozen() {
|
||||||
realCost := h.server.costTracker.realCost(servingTime, msg.Size, 0)
|
realCost := h.server.costTracker.realCost(servingTime, msg.Size, 0)
|
||||||
p.fcClient.RequestProcessed(reqID, respId, maxCost, realCost)
|
p.fcClient.RequestProcessed(reqID, responseCount, maxCost, realCost)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Positive correction buffer value with real cost.
|
// Positive correction buffer value with real cost.
|
||||||
|
|
@ -231,12 +232,13 @@ func (h *serverHandler) handleMsg(p *peer) error {
|
||||||
if reply != nil {
|
if reply != nil {
|
||||||
replySize = reply.size()
|
replySize = reply.size()
|
||||||
}
|
}
|
||||||
realCost := h.server.costTracker.realCost(servingTime, msg.Size, replySize)
|
var realCost uint64
|
||||||
// Assign a fake cost for testing purpose.
|
|
||||||
if h.server.costTracker.testing {
|
if h.server.costTracker.testing {
|
||||||
realCost = maxCost
|
realCost = maxCost // Assign a fake cost for testing purpose
|
||||||
|
} else {
|
||||||
|
realCost = h.server.costTracker.realCost(servingTime, msg.Size, replySize)
|
||||||
}
|
}
|
||||||
bv := p.fcClient.RequestProcessed(reqID, respId, maxCost, realCost)
|
bv := p.fcClient.RequestProcessed(reqID, responseCount, maxCost, realCost)
|
||||||
if amount != 0 {
|
if amount != 0 {
|
||||||
// Feed cost tracker request serving statistic.
|
// Feed cost tracker request serving statistic.
|
||||||
h.server.costTracker.updateStats(msg.Code, amount, servingTime, realCost)
|
h.server.costTracker.updateStats(msg.Code, amount, servingTime, realCost)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue