mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les: add waitForPeers to requestDistributor
This commit is contained in:
parent
ba1a974ca4
commit
dfe14f9f3b
1 changed files with 22 additions and 11 deletions
|
|
@ -62,9 +62,10 @@ type distReq struct {
|
||||||
canSend func(distPeer) bool
|
canSend func(distPeer) bool
|
||||||
request func(distPeer) func()
|
request func(distPeer) func()
|
||||||
|
|
||||||
reqOrder uint64
|
reqOrder uint64
|
||||||
sentChn chan distPeer
|
sentChn chan distPeer
|
||||||
element *list.Element
|
element *list.Element
|
||||||
|
waitForPeers mclock.AbsTime
|
||||||
}
|
}
|
||||||
|
|
||||||
// newRequestDistributor creates a new request distributor
|
// newRequestDistributor creates a new request distributor
|
||||||
|
|
@ -106,7 +107,11 @@ func (d *requestDistributor) registerTestPeer(p distPeer) {
|
||||||
|
|
||||||
// distMaxWait is the maximum waiting time after which further necessary waiting
|
// distMaxWait is the maximum waiting time after which further necessary waiting
|
||||||
// times are recalculated based on new feedback from the servers
|
// times are recalculated based on new feedback from the servers
|
||||||
const distMaxWait = time.Millisecond * 10
|
const distMaxWait = time.Millisecond * 50
|
||||||
|
|
||||||
|
// waitForPeers is the time window in which a request does not fail even if it
|
||||||
|
// has no suitable peers to send to at the moment
|
||||||
|
const waitForPeers = time.Second * 3
|
||||||
|
|
||||||
// main event loop
|
// main event loop
|
||||||
func (d *requestDistributor) loop() {
|
func (d *requestDistributor) loop() {
|
||||||
|
|
@ -179,8 +184,6 @@ func (d *requestDistributor) nextRequest() (distPeer, *distReq, time.Duration) {
|
||||||
checkedPeers := make(map[distPeer]struct{})
|
checkedPeers := make(map[distPeer]struct{})
|
||||||
elem := d.reqQueue.Front()
|
elem := d.reqQueue.Front()
|
||||||
var (
|
var (
|
||||||
bestPeer distPeer
|
|
||||||
bestReq *distReq
|
|
||||||
bestWait time.Duration
|
bestWait time.Duration
|
||||||
sel *weightedRandomSelect
|
sel *weightedRandomSelect
|
||||||
)
|
)
|
||||||
|
|
@ -188,9 +191,18 @@ func (d *requestDistributor) nextRequest() (distPeer, *distReq, time.Duration) {
|
||||||
d.peerLock.RLock()
|
d.peerLock.RLock()
|
||||||
defer d.peerLock.RUnlock()
|
defer d.peerLock.RUnlock()
|
||||||
|
|
||||||
for (len(d.peers) > 0 || elem == d.reqQueue.Front()) && elem != nil {
|
peerCount := len(d.peers)
|
||||||
|
for (len(checkedPeers) < peerCount || elem == d.reqQueue.Front()) && elem != nil {
|
||||||
req := elem.Value.(*distReq)
|
req := elem.Value.(*distReq)
|
||||||
canSend := false
|
canSend := false
|
||||||
|
now := d.clock.Now()
|
||||||
|
if req.waitForPeers > now {
|
||||||
|
canSend = true
|
||||||
|
wait := time.Duration(req.waitForPeers - now)
|
||||||
|
if bestWait == 0 || wait < bestWait {
|
||||||
|
bestWait = wait
|
||||||
|
}
|
||||||
|
}
|
||||||
for peer := range d.peers {
|
for peer := range d.peers {
|
||||||
if _, ok := checkedPeers[peer]; !ok && peer.canQueue() && req.canSend(peer) {
|
if _, ok := checkedPeers[peer]; !ok && peer.canQueue() && req.canSend(peer) {
|
||||||
canSend = true
|
canSend = true
|
||||||
|
|
@ -202,9 +214,7 @@ func (d *requestDistributor) nextRequest() (distPeer, *distReq, time.Duration) {
|
||||||
}
|
}
|
||||||
sel.update(selectPeerItem{peer: peer, req: req, weight: int64(bufRemain*1000000) + 1})
|
sel.update(selectPeerItem{peer: peer, req: req, weight: int64(bufRemain*1000000) + 1})
|
||||||
} else {
|
} else {
|
||||||
if bestReq == nil || wait < bestWait {
|
if bestWait == 0 || wait < bestWait {
|
||||||
bestPeer = peer
|
|
||||||
bestReq = req
|
|
||||||
bestWait = wait
|
bestWait = wait
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -223,7 +233,7 @@ func (d *requestDistributor) nextRequest() (distPeer, *distReq, time.Duration) {
|
||||||
c := sel.choose().(selectPeerItem)
|
c := sel.choose().(selectPeerItem)
|
||||||
return c.peer, c.req, 0
|
return c.peer, c.req, 0
|
||||||
}
|
}
|
||||||
return bestPeer, bestReq, bestWait
|
return nil, nil, bestWait
|
||||||
}
|
}
|
||||||
|
|
||||||
// queue adds a request to the distribution queue, returns a channel where the
|
// queue adds a request to the distribution queue, returns a channel where the
|
||||||
|
|
@ -237,6 +247,7 @@ func (d *requestDistributor) queue(r *distReq) chan distPeer {
|
||||||
if r.reqOrder == 0 {
|
if r.reqOrder == 0 {
|
||||||
d.lastReqOrder++
|
d.lastReqOrder++
|
||||||
r.reqOrder = d.lastReqOrder
|
r.reqOrder = d.lastReqOrder
|
||||||
|
r.waitForPeers = d.clock.Now() + mclock.AbsTime(waitForPeers)
|
||||||
}
|
}
|
||||||
|
|
||||||
back := d.reqQueue.Back()
|
back := d.reqQueue.Back()
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue