mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-17 09:23:48 +00:00
Extract thread safe functions to manipulate Streamer.peers map
This commit is contained in:
parent
3836dd0da5
commit
83b6cc4280
1 changed files with 27 additions and 11 deletions
|
|
@ -306,6 +306,30 @@ func (self *Streamer) Retrieve(chunk *storage.Chunk) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (self *Streamer) getPeer(peerId discover.NodeID) *StreamerPeer {
|
||||||
|
if self.peers == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
self.peersLock.RLock()
|
||||||
|
defer self.peersLock.RUnlock()
|
||||||
|
return self.peers[peerId]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (self *Streamer) setPeer(peer *StreamerPeer) {
|
||||||
|
if self.peers == nil {
|
||||||
|
self.peers = make(map[discover.NodeID]*StreamerPeer)
|
||||||
|
}
|
||||||
|
self.peersLock.Lock()
|
||||||
|
self.peers[peer.ID()] = peer
|
||||||
|
self.peersLock.Unlock()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (self *Streamer) deletePeer(peer *StreamerPeer) {
|
||||||
|
self.peersLock.Lock()
|
||||||
|
delete(self.peers, peer.ID())
|
||||||
|
self.peersLock.Unlock()
|
||||||
|
}
|
||||||
|
|
||||||
func (self *StreamerPeer) handleChunkDeliveryMsg(req *ChunkDeliveryMsg) error {
|
func (self *StreamerPeer) handleChunkDeliveryMsg(req *ChunkDeliveryMsg) error {
|
||||||
chunk, err := self.dbAccess.get(req.Key)
|
chunk, err := self.dbAccess.get(req.Key)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -426,9 +450,7 @@ func (self *Streamer) Subscribe(peerId discover.NodeID, s string, t []byte, from
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
self.peersLock.RLock()
|
peer := self.getPeer(peerId)
|
||||||
peer := self.peers[peerId]
|
|
||||||
self.peersLock.RUnlock()
|
|
||||||
if peer == nil {
|
if peer == nil {
|
||||||
return fmt.Errorf("peer not found %v", peerId)
|
return fmt.Errorf("peer not found %v", peerId)
|
||||||
}
|
}
|
||||||
|
|
@ -630,15 +652,9 @@ func (s *Streamer) Run(p *bzzPeer) error {
|
||||||
// })
|
// })
|
||||||
// subscribe to request handling ; only with non-light nodes
|
// subscribe to request handling ; only with non-light nodes
|
||||||
|
|
||||||
s.peersLock.Lock()
|
s.setPeer(sp)
|
||||||
s.peers[sp.ID()] = sp
|
|
||||||
s.peersLock.Unlock()
|
|
||||||
|
|
||||||
defer func() {
|
defer s.deletePeer(sp)
|
||||||
s.peersLock.Lock()
|
|
||||||
delete(s.peers, sp.ID())
|
|
||||||
s.peersLock.Unlock()
|
|
||||||
}()
|
|
||||||
|
|
||||||
s.Subscribe(sp.ID(), retrieveRequestStream, nil, 0, 0, Top, true)
|
s.Subscribe(sp.ID(), retrieveRequestStream, nil, 0, 0, Top, true)
|
||||||
defer close(sp.quit)
|
defer close(sp.quit)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue