diff --git a/swarm/network/hive.go b/swarm/network/hive.go index 9f21440724..dcf666d0b9 100644 --- a/swarm/network/hive.go +++ b/swarm/network/hive.go @@ -198,9 +198,18 @@ func (self *Hive) Stop() error { } // 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) - 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 () // 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 @@ -208,10 +217,8 @@ func (self *Hive) addPeer(p *peer) { // we do not record as request or forward it, just reply with peers p.retrieve(&retrieveRequestMsgData{}) glog.V(logger.Detail).Infof("[BZZ] KΛÐΞMLIΛ hive: 'whatsup wheresdaparty' sent to %v", p) - select { - case self.more <- true: - default: - } + + return nil } // called after peer disconnected diff --git a/swarm/network/kademlia/kademlia.go b/swarm/network/kademlia/kademlia.go index 50fb6b6f2c..e635c042c2 100644 --- a/swarm/network/kademlia/kademlia.go +++ b/swarm/network/kademlia/kademlia.go @@ -136,13 +136,14 @@ func (self *Kademlia) On(node Node, cb func(*NodeRecord, Node) error) (err error } if replaced != nil { glog.V(logger.Debug).Infof("[KΛÐ]: node %v replaced by %v ", replaced, node) - return + replaced.Drop() + return nil } // new node added glog.V(logger.Info).Infof("[KΛÐ]: add node %v to table", node) self.count++ self.setProxLimit(index, false) - return + return nil } // 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 { - return + // gracefully return without error if peer already offline + return nil } 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++ 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) - } /* diff --git a/swarm/network/protocol.go b/swarm/network/protocol.go index 6f27f68b26..e4e247d458 100644 --- a/swarm/network/protocol.go +++ b/swarm/network/protocol.go @@ -49,6 +49,7 @@ const ( ErrExtraStatusMsg ErrSwap ErrSync + ErrUnwanted ) var errorToString = map[int]string{ @@ -61,6 +62,7 @@ var errorToString = map[int]string{ ErrExtraStatusMsg: "Extra status message", ErrSwap: "SWAP error", ErrSync: "Sync error", + ErrUnwanted: "Unwanted peer", } // bzz represents the swarm wire protocol @@ -254,7 +256,7 @@ func (self *bzz) handle() error { return self.protoError(ErrDecode, "<- %v: %v", msg, err) } 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}) 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) - 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 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 { 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) err := p2p.Send(self.rw, msg, data) if err != nil {