From 11aeaec42095cfb43d56cfdd4e6341688d4f68f6 Mon Sep 17 00:00:00 2001 From: nolash Date: Fri, 17 Mar 2017 04:30:59 +0100 Subject: [PATCH] swarm/network: WIP kademlia network sim tests --- p2p/protocols/protocol.go | 4 ++-- swarm/network/discovery.go | 3 ++- swarm/network/kademlia.go | 6 +++++- swarm/network/simulations/overlay.go | 30 +++++++++++++++++++++------- 4 files changed, 32 insertions(+), 11 deletions(-) diff --git a/p2p/protocols/protocol.go b/p2p/protocols/protocol.go index cb067dc59c..d555056f90 100644 --- a/p2p/protocols/protocol.go +++ b/p2p/protocols/protocol.go @@ -271,7 +271,7 @@ func (self *Peer) Send(msg interface{}) error { if !found { return errorf(ErrInvalidMsgType, "%v", code) } - glog.V(logger.Detail).Infof("=> %v (%d)", msg, code) + glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg) err := self.m.SendMsg(uint64(code), msg) if err != nil { err = errorf(ErrWrite, "(msg code: %v): %v", code, err) @@ -321,7 +321,7 @@ func (self *Peer) handleIncoming() (interface{}, error) { if err := msg.Decode(val.Interface()); err != nil { return nil, errorf(ErrDecode, "<= %v: %v", msg, err) } - glog.V(logger.Detail).Infof("<= %v %v (%d)", req, typ, msg.Code) + glog.V(logger.Detail).Infof("<= %v FROM %v %v %v", msg, self.ID(), req, typ) // call the registered handler callbacks // a registered callback take the decoded message as argument as an interface diff --git a/swarm/network/discovery.go b/swarm/network/discovery.go index 45bdc662c2..58085c0566 100644 --- a/swarm/network/discovery.go +++ b/swarm/network/discovery.go @@ -26,7 +26,8 @@ func (self *discPeer) NotifyPeer(p Peer, po uint8) error { return nil } resp := &peersMsg{ - Peers: []*peerAddr{p.(*discPeer).Peer.(*bzzPeer).peerAddr}, + //Peers: []*peerAddr{p.(*discPeer).Peer.(*bzzPeer).peerAddr}, + Peers: []*peerAddr{&peerAddr{OAddr: p.OverlayAddr(), UAddr: p.UnderlayAddr()}}, // perhaps the PeerAddr interface is unnecessary generalization } return p.Send(resp) } diff --git a/swarm/network/kademlia.go b/swarm/network/kademlia.go index ddfcb18057..2c96d2caf9 100644 --- a/swarm/network/kademlia.go +++ b/swarm/network/kademlia.go @@ -73,7 +73,8 @@ func NewKadParams() *KadParams { MinProxBinSize: 2, MinBinSize: 2, MaxBinSize: 4, - RetryInterval: 42000000000, + //RetryInterval: 42000000000, + RetryInterval: 420000000, MaxRetries: 42, RetryExponent: 2, } @@ -136,6 +137,7 @@ func (self *KadPeer) String() string { func (self *Kademlia) callable(val pot.PotVal) *KadPeer { kp := val.(*KadPeer) // not callable if peer is live or exceeded maxRetries + glog.V(logger.Detail).Infof(">>>>>>>>>>>>>>>>>>>>>>>>>>> in callable: %T", kp.Peer) if kp.Peer != nil || kp.retries > self.MaxRetries { glog.V(logger.Detail).Infof("peer %v not callable", kp.PeerAddr) return nil @@ -149,6 +151,7 @@ func (self *Kademlia) callable(val pot.PotVal) *KadPeer { } // this is never called concurrently, so safe to increment // peer can be retried again + if retries < kp.retries { glog.V(logger.Detail).Infof("log time needed before retry %v, wait only warrants %v", kp.retries, retries) return nil @@ -313,6 +316,7 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) { }) if p != nil { glog.V(logger.Detail).Infof("candidate prox peer found: %v (%v), %#v", p, ppo, p) + //return p, 0, false return p, 0, false } glog.V(logger.Detail).Infof("no candidate prox peers to connect to (ProxLimit: %v, minProxSize: %v)", proxLimit, self.MinProxBinSize) diff --git a/swarm/network/simulations/overlay.go b/swarm/network/simulations/overlay.go index 94a83e348d..4cc02c7224 100644 --- a/swarm/network/simulations/overlay.go +++ b/swarm/network/simulations/overlay.go @@ -50,8 +50,8 @@ func (self *SimNode) Stop() error { return nil } -func (self *SimNode) RunProtocol(id *adapters.NodeId, rw, rrw p2p.MsgReadWriter, runc chan bool) error { - return self.NodeAdapter.(adapters.ProtocolRunner).RunProtocol(id, rw, rrw, runc) +func (self *SimNode) RunProtocol(id *adapters.NodeId, rw, rrw p2p.MsgReadWriter, peer *adapters.Peer) error { + return self.NodeAdapter.(adapters.ProtocolRunner).RunProtocol(id, rw, rrw, peer) } // NewSimNode creates adapters for nodes in the simulation. @@ -59,13 +59,29 @@ func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapt id := conf.Id na := adapters.NewSimNode(id, self.Network, adapters.NewSimPipe) addr := network.NewPeerAddrFromNodeId(id) - to := network.NewKademlia(addr.OverlayAddr(), nil) // overlay topology driver + kp := network.NewKadParams() + kp.MinProxBinSize = 4 + kp.MinBinSize = 4 + to := network.NewKademlia(addr.OverlayAddr(), kp) // overlay topology driver // to := network.NewTestOverlay(addr.OverlayAddr()) // overlay topology driver - pp := network.NewHive(network.NewHiveParams(), to) // hive + hp := network.NewHiveParams() + hp.CallInterval = 5000 + pp := network.NewHive(hp, to) // hive self.hives = append(self.hives, pp) // remember hive // bzz protocol Run function. messaging through SimPipe + + services := func(p network.Peer) error { + dp := network.NewDiscovery(p, to) + pp.Add(dp) + glog.V(logger.Detail).Infof("kademlia on %v", p) + p.DisconnectHook(func(err error) { + pp.Remove(dp) + }) + return nil + } + ct := network.BzzCodeMap(network.HiveMsgs...) // bzz protocol code map - na.Run = network.Bzz(addr.OverlayAddr(), pp, na, ct, nil).Run + na.Run = network.Bzz(addr.OverlayAddr(), na, ct, services, nil, nil).Run connect := func(s string) error { return self.Connect(id, adapters.NewNodeIdFromHex(s)) } @@ -95,7 +111,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu conf.Backend = true net := NewNetwork(simulations.NewNetwork(conf)) - ids := p2ptest.RandomNodeIds(5) + ids := p2ptest.RandomNodeIds(10) for i, id := range ids { net.NewNode(&simulations.NodeConfig{Id: id}) @@ -134,7 +150,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu // var server func main() { runtime.GOMAXPROCS(runtime.NumCPU()) - glog.SetV(logger.Info) + glog.SetV(logger.Detail) glog.SetToStderr(true) c, quitc := simulations.NewSessionController(nethook)