swarm/network/kademlia+hive: unwanted peers (due to full kad bucket) are now properly dropped with ErrUnwanted

This commit is contained in:
zelig 2016-06-16 11:26:00 +01:00
parent f72d3b26a2
commit babb3a7526
3 changed files with 25 additions and 14 deletions

View file

@ -198,9 +198,18 @@ func (self *Hive) Stop() error {
} }
// called at the end of a successful protocol handshake // called at the end of a successful protocol handshake
func (self *Hive) addPeer(p *peer) { func (self *Hive) addPeer(p *peer) error {
defer func() {
select {
case self.more <- true:
default:
}
}()
glog.V(logger.Detail).Infof("[BZZ] KΛÐΞMLIΛ hive: hi new bee %v", p) glog.V(logger.Detail).Infof("[BZZ] KΛÐΞMLIΛ hive: hi new bee %v", p)
self.kad.On(p, loadSync) err := self.kad.On(p, loadSync)
if err != nil {
return err
}
// self lookup (can be encoded as nil/zero key since peers addr known) + no id () // self lookup (can be encoded as nil/zero key since peers addr known) + no id ()
// the most common way of saying hi in bzz is initiation of gossip // the most common way of saying hi in bzz is initiation of gossip
// let me know about anyone new from my hood , here is the storageradius // let me know about anyone new from my hood , here is the storageradius
@ -208,10 +217,8 @@ func (self *Hive) addPeer(p *peer) {
// we do not record as request or forward it, just reply with peers // we do not record as request or forward it, just reply with peers
p.retrieve(&retrieveRequestMsgData{}) p.retrieve(&retrieveRequestMsgData{})
glog.V(logger.Detail).Infof("[BZZ] KΛÐΞMLIΛ hive: 'whatsup wheresdaparty' sent to %v", p) glog.V(logger.Detail).Infof("[BZZ] KΛÐΞMLIΛ hive: 'whatsup wheresdaparty' sent to %v", p)
select {
case self.more <- true: return nil
default:
}
} }
// called after peer disconnected // called after peer disconnected

View file

@ -136,13 +136,14 @@ func (self *Kademlia) On(node Node, cb func(*NodeRecord, Node) error) (err error
} }
if replaced != nil { if replaced != nil {
glog.V(logger.Debug).Infof("[KΛÐ]: node %v replaced by %v ", replaced, node) glog.V(logger.Debug).Infof("[KΛÐ]: node %v replaced by %v ", replaced, node)
return replaced.Drop()
return nil
} }
// new node added // new node added
glog.V(logger.Info).Infof("[KΛÐ]: add node %v to table", node) glog.V(logger.Info).Infof("[KΛÐ]: add node %v to table", node)
self.count++ self.count++
self.setProxLimit(index, false) self.setProxLimit(index, false)
return return nil
} }
// is the entrypoint called when a node is taken offline // is the entrypoint called when a node is taken offline
@ -161,7 +162,8 @@ func (self *Kademlia) Off(node Node, cb func(*NodeRecord, Node)) (err error) {
} }
if !found { if !found {
return // gracefully return without error if peer already offline
return nil
} }
glog.V(logger.Info).Infof("[KΛÐ]: remove node %v from table", node) glog.V(logger.Info).Infof("[KΛÐ]: remove node %v from table", node)
@ -215,8 +217,6 @@ func (self *Kademlia) setProxLimit(r int, off bool) {
self.proxLimit++ self.proxLimit++
glog.V(logger.Detail).Infof("[KΛÐ]: proxbin contraction (size: %v, limit: %v, bin: %v, off: %v)", self.proxSize, self.proxLimit, r, off) glog.V(logger.Detail).Infof("[KΛÐ]: proxbin contraction (size: %v, limit: %v, bin: %v, off: %v)", self.proxSize, self.proxLimit, r, off)
} }
// glog.V(logger.Detail).Infof("%v", self)
} }
/* /*

View file

@ -49,6 +49,7 @@ const (
ErrExtraStatusMsg ErrExtraStatusMsg
ErrSwap ErrSwap
ErrSync ErrSync
ErrUnwanted
) )
var errorToString = map[int]string{ var errorToString = map[int]string{
@ -61,6 +62,7 @@ var errorToString = map[int]string{
ErrExtraStatusMsg: "Extra status message", ErrExtraStatusMsg: "Extra status message",
ErrSwap: "SWAP error", ErrSwap: "SWAP error",
ErrSync: "Sync error", ErrSync: "Sync error",
ErrUnwanted: "Unwanted peer",
} }
// bzz represents the swarm wire protocol // bzz represents the swarm wire protocol
@ -254,7 +256,7 @@ func (self *bzz) handle() error {
return self.protoError(ErrDecode, "<- %v: %v", msg, err) return self.protoError(ErrDecode, "<- %v: %v", msg, err)
} }
req.from = &peer{bzz: self} req.from = &peer{bzz: self}
glog.V(logger.Debug).Infof("[BZZ] <- peer addresses: %v", req) glog.V(logger.Detail).Infof("[BZZ] <- peer addresses: %v", req)
self.hive.HandlePeersMsg(&req, &peer{bzz: self}) self.hive.HandlePeersMsg(&req, &peer{bzz: self})
case syncRequestMsg: case syncRequestMsg:
@ -366,7 +368,10 @@ func (self *bzz) handleStatus() (err error) {
} }
glog.V(logger.Info).Infof("[BZZ] Peer %08x is [bzz] capable (%d/%d)", self.remoteAddr.Addr[:4], status.Version, status.NetworkId) glog.V(logger.Info).Infof("[BZZ] Peer %08x is [bzz] capable (%d/%d)", self.remoteAddr.Addr[:4], status.Version, status.NetworkId)
self.hive.addPeer(&peer{bzz: self}) err = self.hive.addPeer(&peer{bzz: self})
if err != nil {
return self.protoError(ErrUnwanted, "%v", err)
}
// hive sets syncstate so sync should start after node added // hive sets syncstate so sync should start after node added
glog.V(logger.Info).Infof("[BZZ] syncronisation request sent with %v", self.syncState) glog.V(logger.Info).Infof("[BZZ] syncronisation request sent with %v", self.syncState)
@ -516,7 +521,6 @@ func (self *bzz) send(msg uint64, data interface{}) error {
if self.hive.blockWrite { if self.hive.blockWrite {
return fmt.Errorf("network write blocked") return fmt.Errorf("network write blocked")
} }
// self.messages = append(self.messages, "")
glog.V(logger.Detail).Infof("[BZZ] -> %v: %v (%T) to %v", msg, data, data, self) glog.V(logger.Detail).Infof("[BZZ] -> %v: %v (%T) to %v", msg, data, data, self)
err := p2p.Send(self.rw, msg, data) err := p2p.Send(self.rw, msg, data)
if err != nil { if err != nil {