mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-17 01:13:45 +00:00
swarm/network: Streamer.Retrieve allows peer list to skip
This commit is contained in:
parent
b2f0655c57
commit
791cac2403
1 changed files with 17 additions and 5 deletions
|
|
@ -273,8 +273,8 @@ func (self *StreamerPeer) handleRetrieveRequestMsg(req *RetrieveRequestMsg) erro
|
||||||
streamer := s.OutgoingStreamer.(*RetrieveRequestStreamer)
|
streamer := s.OutgoingStreamer.(*RetrieveRequestStreamer)
|
||||||
if chunk.ReqC != nil {
|
if chunk.ReqC != nil {
|
||||||
if created {
|
if created {
|
||||||
if err := self.streamer.Retrieve(chunk); err != nil {
|
if err := self.streamer.Retrieve(chunk, self.ID()); err != nil {
|
||||||
return err
|
return nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
go func() {
|
go func() {
|
||||||
|
|
@ -299,16 +299,27 @@ func (self *StreamerPeer) handleRetrieveRequestMsg(req *RetrieveRequestMsg) erro
|
||||||
}
|
}
|
||||||
|
|
||||||
// Retrieve sends a chunk retrieve request to
|
// Retrieve sends a chunk retrieve request to
|
||||||
func (self *Streamer) Retrieve(chunk *storage.Chunk) error {
|
func (self *Streamer) Retrieve(chunk *storage.Chunk, peersToSkip ...discover.NodeID) error {
|
||||||
|
var success bool
|
||||||
self.overlay.EachConn(chunk.Key[:], 255, func(p OverlayConn, po int, nn bool) bool {
|
self.overlay.EachConn(chunk.Key[:], 255, func(p OverlayConn, po int, nn bool) bool {
|
||||||
sp := p.(*StreamerPeer)
|
spId := p.(Peer).ID()
|
||||||
|
for _, p := range peersToSkip {
|
||||||
|
if p == spId {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
sp := self.getPeer(spId)
|
||||||
// TODO: skip light nodes that do not accept retrieve requests
|
// TODO: skip light nodes that do not accept retrieve requests
|
||||||
sp.SendPriority(&RetrieveRequestMsg{
|
sp.SendPriority(&RetrieveRequestMsg{
|
||||||
Key: chunk.Key[:],
|
Key: chunk.Key[:],
|
||||||
}, Top)
|
}, Top)
|
||||||
|
success = true
|
||||||
return false
|
return false
|
||||||
})
|
})
|
||||||
return nil
|
if success {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return errors.New("no peer found")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *Streamer) getPeer(peerId discover.NodeID) *StreamerPeer {
|
func (self *Streamer) getPeer(peerId discover.NodeID) *StreamerPeer {
|
||||||
|
|
@ -391,6 +402,7 @@ func (self *StreamerPeer) setOutgoingStreamer(s string, o OutgoingStreamer, prio
|
||||||
os := &outgoingStreamer{
|
os := &outgoingStreamer{
|
||||||
OutgoingStreamer: o,
|
OutgoingStreamer: o,
|
||||||
priority: priority,
|
priority: priority,
|
||||||
|
stream: s,
|
||||||
}
|
}
|
||||||
self.outgoing[s] = os
|
self.outgoing[s] = os
|
||||||
return os, nil
|
return os, nil
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue