les: separate peer into clientPeer and serverPeer (#19991)
* les: separate peer into clientPeer and serverPeer * les: address comments
This commit is contained in:
@ -27,14 +27,13 @@ import (
|
||||
|
||||
type ltrInfo struct {
|
||||
tx *types.Transaction
|
||||
sentTo map[*peer]struct{}
|
||||
sentTo map[*serverPeer]struct{}
|
||||
}
|
||||
|
||||
type lesTxRelay struct {
|
||||
txSent map[common.Hash]*ltrInfo
|
||||
txPending map[common.Hash]struct{}
|
||||
ps *peerSet
|
||||
peerList []*peer
|
||||
peerList []*serverPeer
|
||||
peerStartPos int
|
||||
lock sync.RWMutex
|
||||
stop chan struct{}
|
||||
@ -42,15 +41,14 @@ type lesTxRelay struct {
|
||||
retriever *retrieveManager
|
||||
}
|
||||
|
||||
func newLesTxRelay(ps *peerSet, retriever *retrieveManager) *lesTxRelay {
|
||||
func newLesTxRelay(ps *serverPeerSet, retriever *retrieveManager) *lesTxRelay {
|
||||
r := &lesTxRelay{
|
||||
txSent: make(map[common.Hash]*ltrInfo),
|
||||
txPending: make(map[common.Hash]struct{}),
|
||||
ps: ps,
|
||||
retriever: retriever,
|
||||
stop: make(chan struct{}),
|
||||
}
|
||||
ps.notify(r)
|
||||
ps.subscribe(r)
|
||||
return r
|
||||
}
|
||||
|
||||
@ -58,24 +56,34 @@ func (ltrx *lesTxRelay) Stop() {
|
||||
close(ltrx.stop)
|
||||
}
|
||||
|
||||
func (ltrx *lesTxRelay) registerPeer(p *peer) {
|
||||
func (ltrx *lesTxRelay) registerPeer(p *serverPeer) {
|
||||
ltrx.lock.Lock()
|
||||
defer ltrx.lock.Unlock()
|
||||
|
||||
ltrx.peerList = ltrx.ps.AllPeers()
|
||||
// Short circuit if the peer is announce only.
|
||||
if p.onlyAnnounce {
|
||||
return
|
||||
}
|
||||
ltrx.peerList = append(ltrx.peerList, p)
|
||||
}
|
||||
|
||||
func (ltrx *lesTxRelay) unregisterPeer(p *peer) {
|
||||
func (ltrx *lesTxRelay) unregisterPeer(p *serverPeer) {
|
||||
ltrx.lock.Lock()
|
||||
defer ltrx.lock.Unlock()
|
||||
|
||||
ltrx.peerList = ltrx.ps.AllPeers()
|
||||
for i, peer := range ltrx.peerList {
|
||||
if peer == p {
|
||||
// Remove from the peer list
|
||||
ltrx.peerList = append(ltrx.peerList[:i], ltrx.peerList[i+1:]...)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// send sends a list of transactions to at most a given number of peers at
|
||||
// once, never resending any particular transaction to the same peer twice
|
||||
func (ltrx *lesTxRelay) send(txs types.Transactions, count int) {
|
||||
sendTo := make(map[*peer]types.Transactions)
|
||||
sendTo := make(map[*serverPeer]types.Transactions)
|
||||
|
||||
ltrx.peerStartPos++ // rotate the starting position of the peer list
|
||||
if ltrx.peerStartPos >= len(ltrx.peerList) {
|
||||
@ -88,7 +96,7 @@ func (ltrx *lesTxRelay) send(txs types.Transactions, count int) {
|
||||
if !ok {
|
||||
ltr = <rInfo{
|
||||
tx: tx,
|
||||
sentTo: make(map[*peer]struct{}),
|
||||
sentTo: make(map[*serverPeer]struct{}),
|
||||
}
|
||||
ltrx.txSent[hash] = ltr
|
||||
ltrx.txPending[hash] = struct{}{}
|
||||
@ -126,17 +134,17 @@ func (ltrx *lesTxRelay) send(txs types.Transactions, count int) {
|
||||
reqID := genReqID()
|
||||
rq := &distReq{
|
||||
getCost: func(dp distPeer) uint64 {
|
||||
peer := dp.(*peer)
|
||||
return peer.GetTxRelayCost(len(ll), len(enc))
|
||||
peer := dp.(*serverPeer)
|
||||
return peer.getTxRelayCost(len(ll), len(enc))
|
||||
},
|
||||
canSend: func(dp distPeer) bool {
|
||||
return !dp.(*peer).onlyAnnounce && dp.(*peer) == pp
|
||||
return !dp.(*serverPeer).onlyAnnounce && dp.(*serverPeer) == pp
|
||||
},
|
||||
request: func(dp distPeer) func() {
|
||||
peer := dp.(*peer)
|
||||
cost := peer.GetTxRelayCost(len(ll), len(enc))
|
||||
peer := dp.(*serverPeer)
|
||||
cost := peer.getTxRelayCost(len(ll), len(enc))
|
||||
peer.fcServer.QueuedRequest(reqID, cost)
|
||||
return func() { peer.SendTxs(reqID, cost, enc) }
|
||||
return func() { peer.sendTxs(reqID, enc) }
|
||||
},
|
||||
}
|
||||
go ltrx.retriever.retrieve(context.Background(), reqID, rq, func(p distPeer, msg *Msg) error { return nil }, ltrx.stop)
|
||||
|
Reference in New Issue
Block a user