From b7860c7fc642b21e02e949f3ec14e775eb3aad4f Mon Sep 17 00:00:00 2001 From: zelig Date: Sat, 25 Mar 2017 01:01:40 +0700 Subject: [PATCH] kademlia fixes --- p2p/adapters/inproc.go | 10 +- p2p/protocols/protocol.go | 23 ++- p2p/protocols/protocol_test.go | 5 +- p2p/testing/protocolsession.go | 12 +- pot/address.go | 15 +- pot/pot.go | 18 +- pot/pot_test.go | 98 ++++++++++- swarm/network/discovery.go | 144 ++++++++++------ swarm/network/discovery_test.go | 5 +- swarm/network/hive.go | 119 ++++++------- swarm/network/hive_test.go | 60 +++---- swarm/network/kademlia.go | 166 +++++++++++------- swarm/network/kademlia_test.go | 244 ++++++++++++++++++++++++--- swarm/network/msg | 0 swarm/network/protocol.go | 34 +++- swarm/network/protocol_test.go | 12 ++ swarm/network/pss_test.go | 107 +++++------- swarm/network/simulations/overlay.go | 71 +++++--- 18 files changed, 784 insertions(+), 359 deletions(-) delete mode 100644 swarm/network/msg diff --git a/p2p/adapters/inproc.go b/p2p/adapters/inproc.go index a8edaf8284..7b87310162 100644 --- a/p2p/adapters/inproc.go +++ b/p2p/adapters/inproc.go @@ -109,12 +109,14 @@ func (self *SimNode) setPeer(id *NodeId, m Messenger) *Peer { self.peers = append(self.peers, p) return p } - if self.peers[i] != nil && m != nil { - panic(fmt.Sprintf("pipe for %v already set", id)) - } + // if self.peers[i] != nil && m != nil { + // panic(fmt.Sprintf("pipe for %v already set", id)) + // } // legit reconnect reset disconnection error, p := self.peers[i] p.Messenger = m + p.Connc = make(chan bool) + p.Readyc = make(chan bool) return p } @@ -145,8 +147,6 @@ func (self *SimNode) Connect(rid []byte) error { return fmt.Errorf("node adapter for %v is missing", id) } rw, rrw := p2p.MsgPipe() - // runc := make(chan bool) - // defer close(runc) // // run protocol on remote node with self as peer peer := self.getPeer(id) if peer != nil && peer.Messenger != nil { diff --git a/p2p/protocols/protocol.go b/p2p/protocols/protocol.go index d555056f90..1a3872ae1c 100644 --- a/p2p/protocols/protocol.go +++ b/p2p/protocols/protocol.go @@ -32,6 +32,7 @@ package protocols import ( "fmt" "reflect" + "time" "github.com/ethereum/go-ethereum/logger" "github.com/ethereum/go-ethereum/logger/glog" @@ -196,6 +197,7 @@ type Peer struct { rw p2p.MsgReadWriter // p2p.MsgReadWriter to send messages to and read messages from handlers map[reflect.Type][]func(interface{}) error // message type -> message handler callback(s) map Errc chan error + wErrc chan error // write error channel } // NewPeer returns a new peer @@ -208,6 +210,7 @@ func NewPeer(p *p2p.Peer, ct *CodeMap, m adapters.Messenger) *Peer { m: m, Peer: p, Errc: make(chan error), + wErrc: make(chan error), handlers: make(map[reflect.Type][]func(interface{}) error), } } @@ -272,13 +275,21 @@ func (self *Peer) Send(msg interface{}) error { return errorf(ErrInvalidMsgType, "%v", 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) - self.Drop(err) - return err + go func() { + self.wErrc <- self.m.SendMsg(uint64(code), msg) + }() + var err error + select { + case err = <-self.wErrc: + if err == nil { + return nil + } + case <-time.NewTimer(3000 * time.Millisecond).C: + err = fmt.Errorf("write timeout") } - return nil + err = errorf(ErrWrite, "(msg code: %v): %v", code, err) + self.Drop(err) + return err } func (self *Peer) DisconnectHook(f func(error)) { diff --git a/p2p/protocols/protocol_test.go b/p2p/protocols/protocol_test.go index 9055fb4c9c..a7463622f4 100644 --- a/p2p/protocols/protocol_test.go +++ b/p2p/protocols/protocol_test.go @@ -14,7 +14,7 @@ import ( ) func init() { - glog.SetV(logger.Detail) + glog.SetV(logger.Error) glog.SetToStderr(true) } @@ -69,7 +69,7 @@ func newProtocol(pp *p2ptest.TestPeerPool, wg *sync.WaitGroup) func(adapters.Nod // demonstrates use of peerPool, killing another peer connection as a response to a message peer.Register(&kill{}, func(msg interface{}) error { id := msg.(*kill).C - // pp.Get(id).Drop(fmt.Errorf("killed")) + pp.Get(id).Drop(fmt.Errorf("killed")) glog.V(logger.Detail).Infof("id %v killed", id) return nil }) @@ -342,7 +342,6 @@ func TestMultiplePeersDropSelf(t *testing.T) { } func TestMultiplePeersDropOther(t *testing.T) { - t.Skip("??") runMultiplePeers(t, 1, fmt.Errorf("Message handler error: (msg code 3): dropped"), fmt.Errorf("p2p: read or write on closed message pipe"), diff --git a/p2p/testing/protocolsession.go b/p2p/testing/protocolsession.go index f68bec002e..a57fb7ee40 100644 --- a/p2p/testing/protocolsession.go +++ b/p2p/testing/protocolsession.go @@ -118,7 +118,7 @@ func (self *ProtocolSession) expect(exp Expect) error { t := exp.Timeout if t == time.Duration(0) { - t = 1000 * time.Millisecond + t = 2000 * time.Millisecond } alarm := time.NewTimer(t) select { @@ -126,7 +126,6 @@ func (self *ProtocolSession) expect(exp Expect) error { glog.V(logger.Detail).Infof("expected msg arrives with error %v", err) return err case <-alarm.C: - glog.V(logger.Detail).Infof("caught timeout") return fmt.Errorf("timout expecting %v sent to peer %v", exp.Msg, exp.Peer) } } @@ -207,3 +206,12 @@ func (self *ProtocolSession) TestDisconnected(disconnects ...*Disconnect) error } return nil } + +func (self *ProtocolSession) Stop() { + for _, id := range self.Ids { + p := self.GetPeer(id) + if p != nil && p.Messenger != nil { + p.Close() + } + } +} diff --git a/pot/address.go b/pot/address.go index 6f3019c182..4e4a76d03e 100644 --- a/pot/address.go +++ b/pot/address.go @@ -52,8 +52,12 @@ func (a *Address) UnmarshalJSON(value []byte) error { // the string form of the binary representation of an address (only first 8 bits) func (a Address) Bin() string { + return ToBin(a[:]) +} + +func ToBin(a []byte) string { var bs []string - for _, b := range a[:] { + for _, b := range a { bs = append(bs, fmt.Sprintf("%08b", b)) } return strings.Join(bs, "") @@ -209,7 +213,7 @@ func (a *HashAddress) String() string { return a.Address.Bin() } -func NewHashAddress(s string) *HashAddress { +func NewAddressFromString(s string) []byte { ha := [32]byte{} t := s + string(zerosBin)[:len(zerosBin)-len(s)] @@ -220,8 +224,13 @@ func NewHashAddress(s string) *HashAddress { } binary.BigEndian.PutUint64(ha[i*8:(i+1)*8], uint64(n)) } + return ha[:] +} + +func NewHashAddress(s string) *HashAddress { + ha := NewAddressFromString(s) h := common.Hash{} - copy(h[:], ha[:]) + copy(h[:], ha) return &HashAddress{Address(h)} } diff --git a/pot/pot.go b/pot/pot.go index c2d92b60d3..bb1d4c936e 100644 --- a/pot/pot.go +++ b/pot/pot.go @@ -21,7 +21,6 @@ import ( ) const ( - // keylen = 4 keylen = 256 maxkeylen = 256 ) @@ -398,7 +397,6 @@ func union(t0, t1 *pot) (*pot, int) { if t1 == nil || t1.size == 0 { return t0, 0 } - po, eq := t0.pin.PO(t1.pin, 0) var pin PotVal var bins []*pot var mis []int @@ -410,6 +408,8 @@ func union(t0, t1 *pot) (*pot, int) { var i0, i1 int var common int + po, eq := pin0.PO(pin1, 0) + for { l0 := len(bins0) l1 := len(bins1) @@ -470,14 +470,17 @@ func union(t0, t1 *pot) (*pot, int) { break } + i := i0 + if len(bins0) > i && bins0[i].po == po { + i++ + } var size0 int - for _, n := range bins0[i0:] { + for _, n := range bins0[i:] { size0 += n.size } - np := &pot{ pin: pin0, - bins: bins0[i0:], + bins: bins0[i:], size: size0 + 1, po: po, } @@ -488,11 +491,13 @@ func union(t0, t1 *pot) (*pot, int) { po = maxkeylen + 1 eq = true common-- + } else { bins2 = append(bins2, n0.bins...) pin0 = pin1 pin1 = n0.pin po, eq = pin0.PO(pin1, n0.po) + } bins0 = bins1 bins1 = bins2 @@ -863,6 +868,9 @@ func (t *pot) String() string { } func (t *pot) sstring(indent string) string { + if t == nil { + return "" + } var s string indent += " " s += fmt.Sprintf("%v%v (%v) %v \n", indent, t.pin, t.po, t.size) diff --git a/pot/pot_test.go b/pot/pot_test.go index 7377d17311..42eb3d22e4 100644 --- a/pot/pot_test.go +++ b/pot/pot_test.go @@ -24,6 +24,7 @@ import ( "testing" "time" + "github.com/ethereum/go-ethereum/logger" "github.com/ethereum/go-ethereum/logger/glog" ) @@ -224,7 +225,10 @@ func TestPotSwap(t *testing.T) { return true }) if sum != 2*max { - t.Fatalf("incorrect number of elements. expected %v, got %v", max, sum) + t.Fatalf("incorrect number of elements. expected %v, got %v", 2*max, sum) + } + if sum != n.Size() { + t.Fatalf("incorrect size. expected %v, got %v", sum, n.Size()) } } @@ -288,7 +292,97 @@ func testPotEachNeighbour(n *Pot, val PotVal, expCount int, fs ...func(PotVal, i return err } -func TestPotMerge(t *testing.T) { +const ( + mergeTestCount = 5 + mergeTestChoose = 5 +) + +func TestPotMergeOne(t *testing.T) { + pot1 := NewPot(nil, 0) + pot1.Add(NewTestAddr("10", 0)) + pot1.Add(NewTestAddr("00", 0)) + pot2 := NewPot(nil, 0) + pot2.Add(NewTestAddr("01", 0)) + glog.V(logger.Debug).Infof("\n%v\n%v", pot2, pot1) + pot1.Merge(pot2) + count := 0 + pot1.Each(func(val PotVal, i int) bool { + count++ + return true + }) + if count != 3 { + t.Fatalf("expected count to be 3, got %d\n%v\n%v", count, pot2, pot1) + } +} + +func TestPotMergeCommon(t *testing.T) { + vs := make([]*testBVAddr, mergeTestCount) + for i := 0; i < maxEachNeighbourTests; i++ { + + for i := 0; i < len(vs); i++ { + vs[i] = randomTestBVAddr(keylen, i) + } + max0 := rand.Intn(mergeTestChoose) + 1 + max1 := rand.Intn(mergeTestChoose) + 1 + n0 := NewPot(nil, 0) + n1 := NewPot(nil, 0) + glog.V(3).Infof("round %v: %v - %v", i, max0, max1) + m := make(map[string]bool) + for j := 0; j < max0; { + r := rand.Intn(max0) + v := vs[r] + // v := randomTestBVAddr(keylen, j) + _, found := n0.Add(v) + if !found { + m[v.String()] = false + j++ + } + } + expAdded := 0 + + for j := 0; j < max1; { + r := rand.Intn(max1) + v := vs[r] + _, found := n1.Add(v) + if !found { + j++ + } + _, found = m[v.String()] + if !found { + expAdded++ + m[v.String()] = false + } + } + if i < 6 { + continue + } + expSize := len(m) + glog.V(4).Infof("%v-0: pin: %v, size: %v", i, n0.Pin(), max0) + glog.V(4).Infof("%v-1: pin: %v, size: %v", i, n1.Pin(), max1) + glog.V(4).Infof("%v: merged tree size: %v, newly added: %v", i, expSize, expAdded) + n, common := Union(n0, n1) + added := n1.Size() - common + size := n.Size() + + if expSize != size { + t.Fatalf("%v: incorrect number of elements in merged pot, expected %v, got %v\n%v", i, expSize, size, n) + } + if expAdded != added { + t.Fatalf("%v: incorrect number of added elements in merged pot, expected %v, got %v", i, expAdded, added) + } + if !checkDuplicates(n.pot) { + t.Fatalf("%v: merged pot contains duplicates: \n%v", i, n) + } + for k, _ := range m { + _, found := n.Add(NewTestBVAddr(k, 0)) + if !found { + t.Fatalf("%v: merged pot (size:%v, added: %v) missing element %v\n%v", i, size, added, k, n) + } + } + } +} + +func TestPotMergeScale(t *testing.T) { for i := 0; i < maxEachNeighbourTests; i++ { max0 := rand.Intn(maxEachNeighbour) + 1 max1 := rand.Intn(maxEachNeighbour) + 1 diff --git a/swarm/network/discovery.go b/swarm/network/discovery.go index 39ebe99449..b9b42c6671 100644 --- a/swarm/network/discovery.go +++ b/swarm/network/discovery.go @@ -1,36 +1,61 @@ package network import ( - "bytes" "fmt" "github.com/ethereum/go-ethereum/logger" "github.com/ethereum/go-ethereum/logger/glog" - "github.com/ethereum/go-ethereum/p2p/discover" + // "github.com/ethereum/go-ethereum/p2p/discover" ) // discovery bzz overlay extension doing peer relaying -// can be switched off + +// messages related to peer discovery +var DiscoveryMsgs = []interface{}{ + &getPeersMsg{}, + &peersMsg{}, + &subPeersMsg{}, +} type discPeer struct { Peer - overlay Overlay - proxLimit uint8 - peers map[discover.NodeID]bool - sentPeers bool + overlay Overlay + peers map[string]bool + // peers map[discover.NodeID]bool + proxLimit uint8 // the proximity radius advertised by remote to subscribe to peers + sentPeers bool // set to true when the peer is first notifed of peers close to them +} + +// discovery peer contructor +// registers the handlers for discovery messages +func NewDiscovery(p Peer, o Overlay) *discPeer { + self := &discPeer{ + overlay: o, + Peer: p, + peers: make(map[string]bool), + } + self.seen(self) + + p.Register(&peersMsg{}, self.handlePeersMsg) + p.Register(&getPeersMsg{}, self.handleGetPeersMsg) + p.Register(&subPeersMsg{}, self.handleSubPeersMsg) + + return self } // NotifyPeer notifies the receiver remote end of a peer p or PO po. // callback for overlay driver func (self *discPeer) NotifyPeer(p Peer, po uint8) error { - if po < self.proxLimit || self.peers[p.ID()] || !self.sentPeers { + glog.V(logger.Warn).Infof("peers %v", self.peers) + if po < self.proxLimit || self.seen(p) { return nil } + glog.V(logger.Warn).Infof("notification about %x", p.OverlayAddr()) + resp := &peersMsg{ - //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) + return self.Send(resp) } // NotifyProx sends a subPeers Msg to the receiver notifying them about @@ -38,23 +63,7 @@ func (self *discPeer) NotifyPeer(p Peer, po uint8) error { // or first empty row) // callback for overlay driver func (self *discPeer) NotifyProx(po uint8) error { - //return self.Send(&SubPeersMsg{ProxLimit: po, MinProxBinSize: 8}) - return self.Send(&SubPeersMsg{ProxLimit: po}) -} - -// new discovery contructor -func NewDiscovery(p Peer, o Overlay) *discPeer { - self := &discPeer{ - overlay: o, - Peer: p, - peers: make(map[discover.NodeID]bool), - } - - p.Register(&peersMsg{}, self.handlePeersMsg) - p.Register(&getPeersMsg{}, self.handleGetPeersMsg) - p.Register(&SubPeersMsg{}, self.handleSubPeersMsg) - - return self + return self.Send(&subPeersMsg{ProxLimit: po}) } /* @@ -95,51 +104,54 @@ func (self getPeersMsg) String() string { } // subPeers msg is communicating the depth/sharpness/focus of the overlay table of a peer -type SubPeersMsg struct { - //MinProxBinSize uint8 - ProxLimit uint8 +type subPeersMsg struct { + ProxLimit uint8 } -func (self SubPeersMsg) String() string { +func (self subPeersMsg) String() string { return fmt.Sprintf("%T: request peers > PO%02d. ", self, self.ProxLimit) } func (self *discPeer) handleSubPeersMsg(msg interface{}) error { - spm := msg.(*SubPeersMsg) + spm := msg.(*subPeersMsg) + self.proxLimit = spm.ProxLimit if !self.sentPeers { var peers []*peerAddr self.overlay.EachLivePeer(self.OverlayAddr(), 255, func(p Peer, po int) bool { - if uint8(po) > self.proxLimit { + if uint8(po) < self.proxLimit { return false } - self.peers[p.ID()] = true + self.seen(p) peers = append(peers, &peerAddr{p.OverlayAddr(), p.UnderlayAddr()}) return true }) - self.Send(&peersMsg{Peers: peers}) + glog.V(logger.Warn).Infof("found initial %v peers not farther than %v", len(peers), self.proxLimit) + if len(peers) > 0 { + self.Send(&peersMsg{Peers: peers}) + } } self.sentPeers = true - self.proxLimit = spm.ProxLimit return nil } // handlePeersMsg called by the protocol when receiving peerset (for target address) // list of nodes ([]PeerAddr in peersMsg) is added to the overlay db using the // Register interface method -func (p *discPeer) handlePeersMsg(msg interface{}) error { +func (self *discPeer) handlePeersMsg(msg interface{}) error { // register all addresses var nas []PeerAddr for _, na := range msg.(*peersMsg).Peers { addr := PeerAddr(na) nas = append(nas, addr) - p.peers[NodeId(addr).NodeID] = true + self.seen(addr) } - + if len(nas) == 0 { - glog.V(logger.Debug).Infof("whoops, no peers in incoming peersMsg from %v", p) + glog.V(logger.Debug).Infof("whoops, no peers in incoming peersMsg from %v", self) return nil } - return p.overlay.Register(nas...) + glog.V(logger.Debug).Infof("got peer addresses from %x, %v (%v)", self.OverlayAddr(), nas, len(nas)) + return self.overlay.Register(nas...) } // handleGetPeersMsg is called by the protocol when receiving a @@ -147,29 +159,55 @@ func (p *discPeer) handlePeersMsg(msg interface{}) error { // peers suggestions are retrieved from the overlay topology driver // using the EachLivePeer interface iterator method // peers sent are remembered throughout a session and not sent twice -func (p *discPeer) handleGetPeersMsg(msg interface{}) error { - req := msg.(*getPeersMsg) +func (self *discPeer) handleGetPeersMsg(msg interface{}) error { var peers []*peerAddr - alreadySent := p.peers + req := msg.(*getPeersMsg) i := 0 - p.overlay.EachLivePeer(p.OverlayAddr(), int(req.Order), func(n Peer, po int) bool { + self.overlay.EachLivePeer(self.OverlayAddr(), int(req.Order), func(n Peer, po int) bool { i++ - if bytes.Compare(n.OverlayAddr(), p.OverlayAddr()) != 0 && - // only send peers we have not sent before in this session - !alreadySent[n.ID()] { - alreadySent[n.ID()] = true + // only send peers we have not sent before in this session + if self.seen(n) { peers = append(peers, &peerAddr{n.OverlayAddr(), n.UnderlayAddr()}) } - // return int(req.Order) == po && len(peers) < int(req.Max) return len(peers) < int(req.Max) }) if len(peers) == 0 { - glog.V(logger.Debug).Infof("no peers found for %v", p) + glog.V(logger.Debug).Infof("no peers found for %v", self) return nil } - glog.V(logger.Debug).Infof("%v peers sent to %v", len(peers), p) + glog.V(logger.Debug).Infof("%v peers sent to %v", len(peers), self) resp := &peersMsg{ Peers: peers, } - return p.Send(resp) + return self.Send(resp) +} + +func RequestOrder(k Overlay, order, broadcastSize, maxPeers uint8) { + req := &getPeersMsg{ + Order: uint8(order), + Max: maxPeers, + } + var i uint8 + var err error + k.EachLivePeer(nil, 255, func(n Peer, po int) bool { + glog.V(logger.Detail).Infof("%T sent to %v", req, n.ID()) + err = n.Send(req) + if err == nil { + i++ + if i >= broadcastSize { + return false + } + } + return true + }) + glog.V(logger.Info).Infof("requesting bees of PO%03d from %v/%v (each max %v)", order, i, broadcastSize, maxPeers) +} + +func (self *discPeer) seen(p PeerAddr) bool { + k := NodeId(p).NodeID.String() + if self.peers[k] { + return true + } + self.peers[k] = true + return false } diff --git a/swarm/network/discovery_test.go b/swarm/network/discovery_test.go index e9086e6af2..78ecf05967 100644 --- a/swarm/network/discovery_test.go +++ b/swarm/network/discovery_test.go @@ -16,7 +16,7 @@ import ( func TestDiscovery(t *testing.T) { addr := RandomAddr() to := NewKademlia(addr.OAddr, NewKadParams()) - ct := BzzCodeMap(HiveMsgs...) + ct := BzzCodeMap(DiscoveryMsgs...) services := func(p Peer) error { dp := NewDiscovery(p, to) @@ -29,6 +29,7 @@ func TestDiscovery(t *testing.T) { } s := newBzzBaseTester(t, 1, addr, ct, services) + defer s.Stop() s.runHandshakes() // o := 0 @@ -37,7 +38,7 @@ func TestDiscovery(t *testing.T) { Expects: []p2ptest.Expect{ p2ptest.Expect{ Code: 3, - Msg: &SubPeersMsg{ProxLimit: 0}, + Msg: &subPeersMsg{ProxLimit: 0}, Peer: s.ProtocolTester.Ids[0], }, }, diff --git a/swarm/network/hive.go b/swarm/network/hive.go index 8a43e65be1..df5b758ddb 100644 --- a/swarm/network/hive.go +++ b/swarm/network/hive.go @@ -56,21 +56,14 @@ type Overlay interface { type Hive struct { *HiveParams // settings Overlay // the overlay topology driver - // disc Discovery - - lock sync.Mutex - quit chan bool - toggle chan bool - more chan bool + lock sync.Mutex + quit chan bool + toggle chan bool + more chan bool } -const ( - peersBroadcastSetSize = 2 - maxPeersPerRequest = 5 - callInterval = 1000 -) - type HiveParams struct { + Discovery bool PeersBroadcastSetSize uint8 MaxPeersPerRequest uint8 CallInterval uint @@ -78,9 +71,10 @@ type HiveParams struct { func NewHiveParams() *HiveParams { return &HiveParams{ - PeersBroadcastSetSize: peersBroadcastSetSize, - MaxPeersPerRequest: maxPeersPerRequest, - CallInterval: callInterval, + Discovery: true, + PeersBroadcastSetSize: 2, + MaxPeersPerRequest: 5, + CallInterval: 1000, } } @@ -94,13 +88,6 @@ func NewHive(params *HiveParams, overlay Overlay) *Hive { } } -// messages that hive handles -var HiveMsgs = []interface{}{ - &getPeersMsg{}, - &peersMsg{}, - &SubPeersMsg{}, -} - // Start receives network info only at startup // connectPeer is a function to connect to a peer based on its NodeID or enode URL // these are called on the p2p.Server which runs on the node @@ -108,7 +95,7 @@ var HiveMsgs = []interface{}{ func (self *Hive) Start(connectPeer func(string) error, af func() <-chan time.Time) error { self.toggle = make(chan bool) - self.more = make(chan bool) + self.more = make(chan bool, 1) self.quit = make(chan bool) glog.V(logger.Debug).Infof("hive started") // this loop is doing bootstrapping and maintains a healthy table @@ -133,25 +120,10 @@ func (self *Hive) Start(connectPeer func(string) error, af func() <-chan time.Ti } else { glog.V(logger.Detail).Infof("cannot suggest peers") } + + want = want && self.Discovery if want { - req := &getPeersMsg{ - Order: uint8(order), - Max: self.MaxPeersPerRequest, - } - var i uint8 - var err error - self.EachLivePeer(nil, order, func(n Peer, po int) bool { - glog.V(logger.Detail).Infof("%T sent to %v", req, n.ID()) - err = n.Send(req) - if err == nil { - i++ - if i >= self.PeersBroadcastSetSize { - return false - } - } - return true - }) - glog.V(logger.Info).Infof("requesting bees of PO%03d from %v/%v (each max %v)", order, i, self.PeersBroadcastSetSize, self.MaxPeersPerRequest) + RequestOrder(self.Overlay, uint8(order), self.PeersBroadcastSetSize, self.MaxPeersPerRequest) } select { @@ -197,6 +169,41 @@ func (self *Hive) keepAlive(af func() <-chan time.Time) { } } +// Add is called at the end of a successful protocol handshake +// to register a connected (live) peer +func (self *Hive) Add(p Peer) error { + defer self.wake() + dp := NewDiscovery(p, self.Overlay) + glog.V(logger.Debug).Infof("to add new bee %v", p) + self.On(dp) + self.String() + glog.V(logger.Warn).Infof("%v", self) + return nil +} + +// Remove called after peer is disconnected +func (self *Hive) Remove(p Peer) { + defer self.wake() + glog.V(logger.Debug).Infof("remove bee %v", p) + self.Off(p) +} + +// NodeInfo function is used by the p2p.server RPC interface to display +// protocol specific node information +func (self *Hive) NodeInfo() interface{} { + return interface{}(self.String()) +} + +// PeerInfo function is used by the p2p.server RPC interface to display +// protocol specific information any connected peer referred to by their NodeID +func (self *Hive) PeerInfo(id discover.NodeID) interface{} { + self.lock.Lock() + defer self.lock.Unlock() + addr := NewPeerAddrFromNodeId(adapters.NewNodeId(id[:])) + return interface{}(addr) +} + +// Stop terminates the updateloop func (self *Hive) Stop() { // closing toggle channel quits the updateloop close(self.quit) @@ -212,36 +219,6 @@ func (self *Hive) wake() { } } -// Add is called at the end of a successful protocol handshake -// to register a connected (live) peer -func (self *Hive) Add(p Peer) error { - defer self.wake() - dp := NewDiscovery(p, self.Overlay) - glog.V(logger.Debug).Infof("to add new bee %v", p) - self.On(dp) - glog.V(logger.Warn).Infof("%v", self) - //dp.NotifyProx(0) - return nil -} - -// Remove called after peer is disconnected -func (self *Hive) Remove(p Peer) { - defer self.wake() - glog.V(logger.Debug).Infof("remove bee %v", p) - self.Off(p) -} - -func (self *Hive) NodeInfo() interface{} { - return interface{}(self.String()) -} - -func (self *Hive) PeerInfo(id discover.NodeID) interface{} { - self.lock.Lock() - defer self.lock.Unlock() - addr := NewPeerAddrFromNodeId(adapters.NewNodeId(id[:])) - return interface{}(addr) -} - func HexToBytes(s string) []byte { id := discover.MustHexID(s) return id[:] diff --git a/swarm/network/hive_test.go b/swarm/network/hive_test.go index 7037c022d8..e3f690a702 100644 --- a/swarm/network/hive_test.go +++ b/swarm/network/hive_test.go @@ -5,17 +5,10 @@ import ( "testing" "time" - "github.com/ethereum/go-ethereum/logger" - "github.com/ethereum/go-ethereum/logger/glog" "github.com/ethereum/go-ethereum/p2p/adapters" p2ptest "github.com/ethereum/go-ethereum/p2p/testing" ) -func init() { - glog.SetV(logger.Detail) - glog.SetToStderr(true) -} - type testConnect struct { mu sync.Mutex conns []string @@ -35,12 +28,13 @@ func (self *testConnect) connect(na string) error { return nil } -func TestOverlayRegistration(t *testing.T) { +func newHiveTester(t *testing.T, params *HiveParams) (*bzzTester, *Hive) { // setup addr := RandomAddr() // tested peers peer address - to := NewTestOverlay(addr.OverlayAddr()) // overlay topology driver - pp := NewHive(NewHiveParams(), to) // hive - ct := BzzCodeMap(HiveMsgs...) // bzz protocol code map + to := NewTestOverlay(addr.OverlayAddr()) // overlay topology drive + pp := NewHive(params, to) // hive + + ct := BzzCodeMap(DiscoveryMsgs...) // bzz protocol code map services := func(p Peer) error { pp.Add(p) p.DisconnectHook(func(err error) { @@ -49,40 +43,36 @@ func TestOverlayRegistration(t *testing.T) { return nil } - s := newBzzBaseTester(t, 1, addr, ct, services) + return newBzzBaseTester(t, 1, addr, ct, services), pp +} + +func TestOverlayRegistration(t *testing.T) { + params := NewHiveParams() + params.Discovery = false + s, pp := newHiveTester(t, params) + defer s.Stop() + id := s.Ids[0] raddr := NewPeerAddrFromNodeId(id) s.runHandshakes() // hive should have called the overlay - if to.posMap[string(raddr.OverlayAddr())] == nil { + if pp.Overlay.(*testOverlay).posMap[string(raddr.OverlayAddr())] == nil { t.Fatalf("Overlay#On not called on new peer") } } func TestRegisterAndConnect(t *testing.T) { - // setup - addr := RandomAddr() // tested peers peer address - to := NewTestOverlay(addr.OverlayAddr()) // overlay topology driver - pp := NewHive(NewHiveParams(), to) // hive - ct := BzzCodeMap(HiveMsgs...) // bzz protocol code map - services := func(p Peer) error { - pp.Add(p) - p.DisconnectHook(func(error) { - pp.Remove(p) - }) - return nil - } - - s := newBzzBaseTester(t, 1, addr, ct, services) + params := NewHiveParams() + s, pp := newHiveTester(t, params) + defer s.Stop() id := s.Ids[0] raddr := NewPeerAddrFromNodeId(id) pp.Register(raddr) - glog.V(5).Infof("%v", pp) // start the hive and wait for the connection tc := &testConnect{ @@ -93,16 +83,16 @@ func TestRegisterAndConnect(t *testing.T) { ticker: make(chan time.Time), } pp.Start(tc.connect, tc.ping) + defer pp.Stop() tc.ticker <- time.Now() s.runHandshakes() - if to.posMap[string(raddr.OverlayAddr())] == nil { + if pp.Overlay.(*testOverlay).posMap[string(raddr.OverlayAddr())] == nil { t.Fatalf("Overlay#On not called on new peer") } // retrieve and broadcast - glog.V(6).Infof("check peer requests for %v", id) ord := order(raddr.OverlayAddr()) o := 0 if ord == 0 { @@ -118,14 +108,4 @@ func TestRegisterAndConnect(t *testing.T) { }, }, }) - // s.TestExchanges(p2ptest.Exchange{ - // Label: "SubPeersMsg message outgoing", - // Expects: []p2ptest.Expect{ - // p2ptest.Expect{ - // Code: 3, - // Msg: &SubPeersMsg{ProxLimit: 0, MinProxBinSize: 8}, - // Peer: id, - // }, - // }, - // }) } diff --git a/swarm/network/kademlia.go b/swarm/network/kademlia.go index 94155c30c1..71a5ad8962 100644 --- a/swarm/network/kademlia.go +++ b/swarm/network/kademlia.go @@ -13,9 +13,11 @@ // // You should have received a copy of the GNU Lesser General Public License // along with the go-ethereum library. If not, see . + package network import ( + "bytes" "fmt" "strings" "time" @@ -64,28 +66,30 @@ type KadParams struct { RetryInterval int RetryExponent int MaxRetries int + PruneInterval int } // NewKadParams() returns a params struct with default values func NewKadParams() *KadParams { return &KadParams{ MaxProxDisplay: 8, - MinProxBinSize: 4, + MinProxBinSize: 2, MinBinSize: 2, MaxBinSize: 4, //RetryInterval: 42000000000, - RetryInterval: 420000000, - MaxRetries: 42, - RetryExponent: 2, + RetryInterval: 420000000, + MaxRetries: 42, + RetryExponent: 2, } } // Kademlia is a table of live peers and a db of known peers type Kademlia struct { - addr *pot.HashAddress // immutable baseaddress of the table - *KadParams // Kademlia configuration parameters - conns, peers *pot.Pot // pots container for peers - lastProxLimit uint8 // stores the last calculated proxlimit + addr PeerAddr // immutable baseaddress of the table + // addr *pot.HashAddress // immutable baseaddress of the table + *KadParams // Kademlia configuration parameters + conns, peers *pot.Pot // pots container for peers + lastProxLimit uint8 // stores the last calculated proxlimit } // NewKademlia(addr, params) creates a Kademlia table for base address addr @@ -95,13 +99,46 @@ func NewKademlia(addr []byte, params *KadParams) *Kademlia { if params == nil { params = NewKadParams() } - base := pot.NewHashAddressFromBytes(addr) - return &Kademlia{ - addr: base, + self := &Kademlia{ + addr: &peerAddr{OAddr: addr}, KadParams: params, conns: pot.NewPot(nil, 0), peers: pot.NewPot(nil, 0), } + return self +} + +// Prune implements a forever loop reacting to a ticker time channel given +// as the first argument +// the loop quits if the channel is closed +// it checks each kademlia bin and if the peer count is higher than +// the MaxBinSize parameter it drops the oldest n peers such that +// the bin is reduced to MinBinSize peers thus leaving slots to newly +// connecting peers +func (self *Kademlia) Prune(c <-chan time.Time) { + go func() { + for _ = range c { + glog.V(logger.Debug).Infof("pruning...") + total := 0 + self.peers.EachBin(self.addr, 0, func(po, size int, f func(func(pot.PotVal, int) bool) bool) bool { + extra := size - self.MinBinSize + if size > self.MaxBinSize { + n := 0 + f(func(v pot.PotVal, po int) bool { + p := v.(*KadPeer).Peer + if p != nil { + p.Drop(fmt.Errorf("bucket full")) + } + n++ + return n < extra + }) + total += extra + } + return true + }) + glog.V(logger.Debug).Infof("pruned %v peers", total) + } + }() } // KadPeer represents a Kademlia Peer and extends @@ -109,38 +146,27 @@ func NewKademlia(addr []byte, params *KadParams) *Kademlia { // * Peer interface (id, last seen, drop) // * HashAddress as derived from PeerAddr overlay implement pot.PoVal interface type KadPeer struct { - *pot.HashAddress + // *pot.HashAddress PeerAddr Peer Peer seenAt time.Time retries int } -func (self *KadPeer) PO(val pot.PotVal, pos int) (po int, eq bool) { - kp, ok := val.(*KadPeer) - var ha *pot.HashAddress - if ok { - ha = kp.HashAddress - } else { - ha = val.(*pot.HashAddress) - } - return self.HashAddress.PO(pot.PotVal(ha), pos) -} - func (self *KadPeer) String() string { if self == nil { return "" } // return string(self.OverlayAddr()) - return self.HashAddress.Address.String() + //return self.HashAddress.Address.String() + return fmt.Sprintf("%x", self.OverlayAddr()) } 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) + glog.V(logger.Detail).Infof("peer %v (%T) not callable", kp, kp.Peer) return nil } // calculate the allowed number of retries based on time lapsed since last seen @@ -152,49 +178,61 @@ 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 } kp.retries++ - glog.V(logger.Detail).Infof("peer %v is callable", kp.PeerAddr) + glog.V(logger.Detail).Infof("peer %v is callable", kp) return kp } -// NewKadPeer(na) creates a kademlia peer from a PeerAddr interface +// NewKadPeer creates a kademlia peer from a PeerAddr interface func NewKadPeer(na PeerAddr) *KadPeer { - // o := na.OverlayAddr() - // glog.V(logger.Detail).Infof("newkadpeer from peerAddr overlay address: %x", o[:6]) return &KadPeer{ - HashAddress: pot.NewHashAddressFromBytes(na.OverlayAddr()), - PeerAddr: na, - seenAt: time.Now(), + PeerAddr: na, + seenAt: time.Now(), } } -// Register(nas) enters each PeerAddr as kademlia peers into the +// Register enters each PeerAddr as kademlia peers into the // database of known peers func (self *Kademlia) Register(nas ...PeerAddr) error { + label := fmt.Sprintf("%x", RandomAddr().OverlayAddr()) np := pot.NewPot(nil, 0) for _, na := range nas { + if bytes.Equal(na.OverlayAddr(), self.addr.OverlayAddr()) { + glog.V(logger.Warn).Infof("[%06s] add peers: %x is self.. skipped ", label, self.addr.OverlayAddr()) + continue + } p := NewKadPeer(na) np, _, _ = pot.Add(np, pot.PotVal(p)) } - common := self.peers.Merge(np) - glog.V(logger.Detail).Infof("add peers: %v out of %v new; root %v", np.Size()-common, np.Size(), np.Pin().String()[:6]) + oldpeers := pot.NewPot(nil, 0) + oldpeers.Merge(self.peers) + self.peers.Merge(np) + m := make(map[string]bool) + self.peers.Each(func(val pot.PotVal, i int) bool { + _, found := m[val.String()] + // TODO: remove this check + // glog.V(logger.Debug).Infof("-> %v %v", val, i) + if found { + panic("duplicate found") + } + m[val.String()] = true + return true + }) return nil } // On(p) inserts the peer as a kademlia peer into the live peers func (self *Kademlia) On(p Peer) { - pp := NewKadPeer(p) - // var pp *KadPeer - kp := pp - self.conns.Swap(pp, func(v pot.PotVal) pot.PotVal { + kp := NewKadPeer(p) + self.conns.Swap(kp, func(v pot.PotVal) pot.PotVal { if v == nil { - self.peers.Swap(pp, func(v pot.PotVal) pot.PotVal { + self.peers.Swap(kp, func(v pot.PotVal) pot.PotVal { if v != nil { kp = v.(*KadPeer) } @@ -211,23 +249,21 @@ func (self *Kademlia) On(p Peer) { vp, ok := kp.Peer.(KadDiscovery) if !ok { - glog.V(logger.Detail).Infof("not discovery peer") + glog.V(logger.Detail).Infof("not discovery peer %T", kp) return } - // vp.NotifyProx(uint8(prox)) + go vp.NotifyProx(uint8(prox)) f := func(val pot.PotVal, po int) { - glog.V(logger.Detail).Infof("peer %v nofified", vp) dp := val.(*KadPeer).Peer.(KadDiscovery) + glog.V(logger.Debug).Infof("peer %v notified of %v (%v)", dp, kp, po) dp.NotifyPeer(kp.Peer, uint8(po)) if uint8(prox) != self.lastProxLimit { self.lastProxLimit = uint8(prox) dp.NotifyProx(uint8(prox)) } + glog.V(logger.Debug).Infof("peer notified") } - self.conns.EachNeighbourAsync(kp, 255, 255, f, false) - go vp.NotifyProx(uint8(prox)) - - + self.conns.EachNeighbourAsync(kp, 1024, 255, f, false) } // Off removes a peer from among live peers @@ -246,6 +282,10 @@ func (self *Kademlia) Off(p Peer) { kp.seenAt = time.Now() } +type ByteAddr struct { + key []byte +} + // EachLivePeer(base, po, f) is an iterator applying f to each live peer // that has proximity order po or less as measured from the base // if base is nil, kademlia base address is used @@ -254,7 +294,7 @@ func (self *Kademlia) EachLivePeer(base []byte, o int, f func(Peer, int) bool) { if base == nil { p = pot.PotVal(self.addr) } else { - p = pot.NewHashAddressFromBytes(base) + p = pot.PotVal(&peerAddr{OAddr: base}) } self.conns.EachNeighbour(p, func(val pot.PotVal, po int) bool { if po > o { @@ -314,15 +354,14 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) { self.peers.EachNeighbour(self.addr, func(val pot.PotVal, po int) bool { r := self.callable(val) if r == nil { - glog.V(logger.Detail).Infof("candidate peer not callable: %#v", r) - return po > proxLimit + return po >= proxLimit } p = r + ppo = po return false }) if p != nil { - glog.V(logger.Detail).Infof("candidate prox peer found: %v (%v), %#v", p, ppo, p) - //return p, 0, false + glog.V(logger.Detail).Infof("candidate prox peer found: %v (%v), %v", p, ppo, p) return p, 0, false } glog.V(logger.Detail).Infof("no candidate prox peers to connect to (ProxLimit: %v, minProxSize: %v)", proxLimit, self.MinProxBinSize) @@ -343,7 +382,7 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) { return size > 0 && po < proxLimit }) // all buckets are full - // minsize == self.BucketSize + // minsize == self.MinBinSize if len(bpo) == 0 { return nil, 0, false } @@ -381,8 +420,8 @@ func (self *Kademlia) String() string { var rows []string - rows = append(rows, "============================================================================") - rows = append(rows, fmt.Sprintf("%v KΛÐΞMLIΛ hive: queen's address: %v", time.Now().UTC().Format(time.UnixDate), self.addr.Address.String()[:6])) + rows = append(rows, "=========================================================================") + rows = append(rows, fmt.Sprintf("%v KΛÐΞMLIΛ hive: queen's address: %v", time.Now().UTC().Format(time.UnixDate), fmt.Sprintf("%x", self.addr.OverlayAddr()[:3]))) rows = append(rows, fmt.Sprintf("population: %d (%d), MinProxBinSize: %d, MinBinSize: %d, MaxBinSize: %d", self.conns.Size(), self.peers.Size(), self.MinProxBinSize, self.MinBinSize, self.MaxBinSize)) liverows := make([]string, self.MaxProxDisplay) @@ -393,6 +432,9 @@ func (self *Kademlia) String() string { rest := self.conns.Size() self.conns.EachBin(self.addr, 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool { var rowlen int + if po >= self.MaxProxDisplay { + po = self.MaxProxDisplay - 1 + } row := []string{fmt.Sprintf("%2d", size)} rest -= size f(func(val pot.PotVal, vpo int) bool { @@ -414,6 +456,12 @@ func (self *Kademlia) String() string { self.peers.EachBin(self.addr, 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool { var rowlen int + if po >= self.MaxProxDisplay { + po = self.MaxProxDisplay - 1 + } + if size < 0 { + panic("wtf") + } row := []string{fmt.Sprintf("%2d", size)} f(func(val pot.PotVal, vpo int) bool { kp := val.(*KadPeer) @@ -424,6 +472,10 @@ func (self *Kademlia) String() string { rowlen++ return rowlen < 4 }) + // glog.V(logger.Debug).Infof("po: %v, peerrows length: %v, maxProxDisplay: %v", po, len(peersrows), self.MaxProxDisplay) + // if po < self.MaxProxDisplay { + // peersrows[po] = strings.Join(row, " ") + // } peersrows[po] = strings.Join(row, " ") return true }) diff --git a/swarm/network/kademlia_test.go b/swarm/network/kademlia_test.go index a830f2902c..a1a853ddd0 100644 --- a/swarm/network/kademlia_test.go +++ b/swarm/network/kademlia_test.go @@ -17,11 +17,12 @@ package network import ( "fmt" + "sync" "testing" "time" // "github.com/ethereum/go-ethereum/logger" - // "github.com/ethereum/go-ethereum/logger/glog" + "github.com/ethereum/go-ethereum/logger/glog" "github.com/ethereum/go-ethereum/pot" ) @@ -30,25 +31,61 @@ func testKadPeerAddr(s string) *peerAddr { return &peerAddr{OAddr: a, UAddr: a} } -func testKadPeer(s string) *bzzPeer { - return &bzzPeer{peerAddr: testKadPeerAddr(s)} +type testDropPeer struct { + Peer + dropc chan error } -func testStr(a PeerAddr) string { - s, _ := a.(*KadPeer) - if s == nil { - return "" - } - // return s.String() - // return fmt.Sprintf("%06x", s.OverlayAddr()) - return pot.NewHashAddressFromBytes(a.(*KadPeer).PeerAddr.OverlayAddr()).Bin()[:6] - // return a.(*KadPeer).String()[:6] //.HashAddress.String()[:6] - // return a.(*KadPeer).String() //[:6] //.HashAddress.String()[:6] - // return "wtf" +type testDiscPeer struct { + *testDropPeer + lock *sync.Mutex + notifications map[string]uint8 +} + +type testPeerNotification struct { + rec string + addr string + po uint8 +} + +type testProxNotification struct { + rec string + po uint8 +} + +type dropError struct { + error + addr string +} + +func (self *testDropPeer) Drop(err error) { + err2 := &dropError{err, overlayStr(self)} + self.dropc <- err2 +} + +func (self *testDiscPeer) NotifyProx(po uint8) error { + key := overlayStr(self) + self.lock.Lock() + defer self.lock.Unlock() + self.notifications[key] = po + return nil +} + +func (self *testDiscPeer) NotifyPeer(p Peer, po uint8) error { + key := overlayStr(self) + key += overlayStr(p) + self.lock.Lock() + defer self.lock.Unlock() + self.notifications[key] = po + return nil } type testKademlia struct { *Kademlia + Discovery bool + dropc chan error + lock *sync.Mutex + notifications map[string]uint8 } func newTestKademlia(b string) *testKademlia { @@ -56,20 +93,63 @@ func newTestKademlia(b string) *testKademlia { params.MinBinSize = 1 params.MinProxBinSize = 2 base := pot.NewHashAddress(b).Bytes() - return &testKademlia{NewKademlia(base, params)} + return &testKademlia{ + NewKademlia(base, params), + false, + make(chan error), + &sync.Mutex{}, + make(map[string]uint8), + } +} + +func (k *testKademlia) newTestKadPeer(s string) Peer { + dp := &testDropPeer{&bzzPeer{peerAddr: testKadPeerAddr(s)}, k.dropc} + if k.Discovery { + return Peer(&testDiscPeer{dp, k.lock, k.notifications}) + } + return Peer(dp) +} + +func overlayStr(a PeerAddr) string { + glog.V(6).Infof("PeerAddr: %v (%T)", a, a) + // if a == (*KadPeer)(nil) || a == (*testDiscPeer)(nil) || a == (*bzzPeer)(nil) || a == nil { + // return "" + // } + // var p Peer + // s, ok := a.(*KadPeer) + // if ok { + // p = s.Peer + // } else { + // p = a.(*testDiscPeer).Peer + // } + // glog.V(6).Infof("PeerAddr: %v (%T)", p, p) + // if p == (Peer)(nil) || p == (*testDiscPeer)(nil) || p == (*bzzPeer)(nil) { + // return "" + // } + // return pot.NewHashAddressFromBytes(p.OverlayAddr()).Bin()[:6] + if a == nil { + return "" + } + k, ok := a.(*KadPeer) + if ok && k.Peer != nil { + return pot.ToBin(a.(*KadPeer).Peer.OverlayAddr())[:6] + } + return pot.ToBin(a.OverlayAddr())[:6] } func (k *testKademlia) On(ons ...string) *testKademlia { for _, s := range ons { - k.Kademlia.On(testKadPeer(s)) + p := k.newTestKadPeer(s) + k.Kademlia.On(p) } return k } func (k *testKademlia) Off(offs ...string) *testKademlia { for _, s := range offs { - k.Kademlia.Off(testKadPeer(s)) + k.Kademlia.Off(k.newTestKadPeer(s)) } + return k } @@ -84,8 +164,8 @@ func (k *testKademlia) Register(regs ...string) *testKademlia { func testSuggestPeer(t *testing.T, k *testKademlia, expAddr string, expPo int, expWant bool) error { addr, o, want := k.SuggestPeer() - if testStr(addr) != expAddr { - return fmt.Errorf("incorrect peer address suggested. expected %v, got %v", expAddr, testStr(addr)) + if overlayStr(addr) != expAddr { + return fmt.Errorf("incorrect peer address suggested. expected %v, got %v", expAddr, overlayStr(addr)) } if o != expPo { return fmt.Errorf("incorrect prox order suggested. expected %v, got %v", expPo, o) @@ -99,6 +179,8 @@ func testSuggestPeer(t *testing.T, k *testKademlia, expAddr string, expPo int, e func TestSuggestPeerFindPeers(t *testing.T) { // 2 row gap, unsaturated proxbin, no callables -> want PO 0 k := newTestKademlia("000000").On("001000") + // k.MinProxBinSize = 2 + // k.MinBinSize = 2 err := testSuggestPeer(t, k, "", 0, true) if err != nil { t.Fatal(err.Error()) @@ -272,11 +354,133 @@ func TestSuggestPeerRetries(t *testing.T) { } +func TestPruning(t *testing.T) { + k := newTestKademlia("000000") + k.On("100000", "110000", "101000", "100100", "100010") + k.On("010000", "011000", "010100", "010010", "010001") + k.On("001000", "001100", "001010", "001001") + k.MaxBinSize = 4 + k.MinBinSize = 3 + prune := make(chan time.Time) + defer close(prune) + k.Prune((<-chan time.Time)(prune)) + prune <- time.Now() + errc := make(chan error) + timeout := time.NewTimer(1000 * time.Millisecond) + n := 0 + dropped := make(map[string]error) + go func() { + for e := range k.dropc { + err := e.(*dropError) + dropped[err.addr] = err.error + n++ + if n == 4 { + break + } + } + close(errc) + }() + select { + case <-errc: + case <-timeout.C: + t.Fatalf("timeout waiting for 4 peers to be dropped") + } + // TODO: this is now based on just taking the first 2 peers + // in order of connecting + expDropped := []string{ + "101000", + "110000", + "010100", + "011000", + } + for _, addr := range expDropped { + err := dropped[addr] + if err == nil { + t.Fatalf("expected peer %v to be dropped", addr) + } + if err.Error() != "bucket full" { + t.Fatalf("incorrect error. expected %v, got %v", "bucket full", err) + } + } +} + func TestKademliaHiveString(t *testing.T) { k := newTestKademlia("000000").On("010000", "001000").Register("100000", "100001") h := k.String() - expH := "\n=========================================================================\nMon Feb 27 12:10:28 UTC 2017 KΛÐΞMLIΛ hive: queen's address: 000000\npopulation: 2 (4), ProxBinSize: 2, MinBinSize: 1, MaxBinSize: 4\n============ PROX LIMIT: 0 ==========================================\n000 0 | 2 840000 800000\n001 1 400000 | 1 400000\n002 1 200000 | 1 200000\n003 0 | 0\n004 0 | 0\n005 0 | 0\n006 0 | 0\n007 0 | 0\n=========================================================================" + expH := "\n=========================================================================\nMon Feb 27 12:10:28 UTC 2017 KΛÐΞMLIΛ hive: queen's address: 000000\npopulation: 2 (4), MinProxBinSize: 2, MinBinSize: 1, MaxBinSize: 4\n============ PROX LIMIT: 0 ==========================================\n000 0 | 2 840000 800000\n001 1 400000 | 1 400000\n002 1 200000 | 1 200000\n003 0 | 0\n004 0 | 0\n005 0 | 0\n006 0 | 0\n007 0 | 0\n=========================================================================" if expH[100:] != h[100:] { t.Fatalf("incorrect hive output. expected %v, got %v", expH, h) } } + +func (self *testKademlia) checkNotifications(npeers []*testPeerNotification, nprox []*testProxNotification) error { + for _, pn := range npeers { + key := pn.rec + pn.addr + po, found := self.notifications[key] + if !found || pn.po != po { + return fmt.Errorf("%v, expected to have notified %v about peer %v (%v)", key, pn.rec, pn.addr, pn.po) + } + delete(self.notifications, key) + } + for _, pn := range nprox { + key := pn.rec + po, found := self.notifications[key] + if !found || pn.po != po { + return fmt.Errorf("expected to have notified %v about new prox limit %v", pn.rec, pn.po) + } + delete(self.notifications, key) + } + if len(self.notifications) > 0 { + return fmt.Errorf("%v unexpected notifications", len(self.notifications)) + } + return nil +} + +func TestNotifications(t *testing.T) { + k := newTestKademlia("000000") + k.Discovery = true + k.MinProxBinSize = 3 + k.On("010000", "001000") + time.Sleep(100 * time.Millisecond) + err := k.checkNotifications( + []*testPeerNotification{ + &testPeerNotification{"010000", "001000", 1}, + }, + []*testProxNotification{ + &testProxNotification{"001000", 0}, + &testProxNotification{"010000", 0}, + }, + ) + if err != nil { + t.Fatal(err.Error()) + } + k = k.On("100000") + time.Sleep(100 * time.Millisecond) + + k.checkNotifications( + []*testPeerNotification{ + &testPeerNotification{"010000", "100000", 0}, + &testPeerNotification{"001000", "100000", 0}, + }, + []*testProxNotification{ + &testProxNotification{"100000", 0}, + }, + ) + + k = k.On("010001") + time.Sleep(100 * time.Millisecond) + + k.checkNotifications( + []*testPeerNotification{ + &testPeerNotification{"010000", "010001", 5}, + &testPeerNotification{"001000", "010001", 1}, + &testPeerNotification{"100000", "010001", 0}, + }, + []*testProxNotification{ + &testProxNotification{"100000", 0}, + &testProxNotification{"010000", 0}, + &testProxNotification{"010001", 0}, + &testProxNotification{"001000", 0}, + }, + ) +} diff --git a/swarm/network/msg b/swarm/network/msg deleted file mode 100644 index e69de29bb2..0000000000 diff --git a/swarm/network/protocol.go b/swarm/network/protocol.go index dbecfc1e44..83c6d05dd6 100644 --- a/swarm/network/protocol.go +++ b/swarm/network/protocol.go @@ -27,6 +27,7 @@ import ( "github.com/ethereum/go-ethereum/p2p/adapters" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/protocols" + "github.com/ethereum/go-ethereum/pot" ) const ( @@ -54,12 +55,14 @@ func (self *bzzPeer) LastActive() time.Time { type PeerAddr interface { OverlayAddr() []byte UnderlayAddr() []byte + PO(pot.PotVal, int) (int, bool) + String() string } // the Peer interface that peerPool needs type Peer interface { PeerAddr - String() string // pretty printable the Node + // String() string // pretty printable the Node ID() discover.NodeID // the key that uniquely identifies the Node for the peerPool Send(interface{}) error // can send messages Drop(error) // disconnect this peer @@ -139,6 +142,35 @@ func (self *peerAddr) UnderlayAddr() []byte { return self.UAddr } +func (self *peerAddr) PO(val pot.PotVal, pos int) (int, bool) { + kp := val.(PeerAddr) + one := kp.OverlayAddr() + other := self.OAddr + for i := pos / 8; i < len(one); i++ { + if one[i] == other[i] { + continue + } + oxo := one[i] ^ other[i] + start := 0 + if i == pos/8 { + start = pos % 8 + } + for j := start; j < 8; j++ { + if (uint8(oxo)>>uint8(7-j))&0x01 != 0 { + return i*8 + j, false + } + } + } + return len(one) * 8, true + // var ha *pot.HashAddress + // var left, right string + // if ok { + // ha = kp.HashAddress + // } else { + // ha = val.(*pot.HashAddress) + // } +} + func (self *peerAddr) String() string { return fmt.Sprintf("%x <%x>", self.OAddr, self.UAddr) } diff --git a/swarm/network/protocol_test.go b/swarm/network/protocol_test.go index e26146cb39..b09db764e7 100644 --- a/swarm/network/protocol_test.go +++ b/swarm/network/protocol_test.go @@ -125,6 +125,8 @@ func TestBzzHandshakeNetworkIdMismatch(t *testing.T) { pp := p2ptest.NewTestPeerPool() addr := RandomAddr() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() + id := s.Ids[0] s.testHandshake( correctBzzHandshake(addr), @@ -137,6 +139,8 @@ func TestBzzHandshakeVersionMismatch(t *testing.T) { pp := p2ptest.NewTestPeerPool() addr := RandomAddr() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() + id := s.Ids[0] s.testHandshake( correctBzzHandshake(addr), @@ -149,6 +153,8 @@ func TestBzzHandshakeSuccess(t *testing.T) { pp := p2ptest.NewTestPeerPool() addr := RandomAddr() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() + id := s.Ids[0] s.testHandshake( correctBzzHandshake(addr), @@ -160,6 +166,7 @@ func TestBzzPeerPoolAdd(t *testing.T) { pp := p2ptest.NewTestPeerPool() addr := RandomAddr() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() id := s.Ids[0] glog.V(logger.Detail).Infof("handshake with %v", id) @@ -174,6 +181,8 @@ func TestBzzPeerPoolRemove(t *testing.T) { addr := RandomAddr() pp := p2ptest.NewTestPeerPool() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() + s.runHandshakes() id := s.Ids[0] @@ -188,6 +197,8 @@ func TestBzzPeerPoolBothAddRemove(t *testing.T) { addr := RandomAddr() pp := p2ptest.NewTestPeerPool() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() + s.runHandshakes() id := s.Ids[0] @@ -206,6 +217,7 @@ func TestBzzPeerPoolNotAdd(t *testing.T) { addr := RandomAddr() pp := p2ptest.NewTestPeerPool() s := newBzzTester(t, 1, addr, pp, nil, nil) + defer s.Stop() id := s.Ids[0] s.testHandshake(correctBzzHandshake(addr), &bzzHandshake{0, 321, NewPeerAddrFromNodeId(id)}, &p2ptest.Disconnect{Peer: id, Error: fmt.Errorf("network id mismatch 321 (!= 322)")}) diff --git a/swarm/network/pss_test.go b/swarm/network/pss_test.go index ff9e6d9d29..e2cd307de8 100644 --- a/swarm/network/pss_test.go +++ b/swarm/network/pss_test.go @@ -8,40 +8,34 @@ import ( //"github.com/ethereum/go-ethereum/logger" //"github.com/ethereum/go-ethereum/logger/glog" "github.com/ethereum/go-ethereum/p2p/adapters" - "github.com/ethereum/go-ethereum/p2p/simulations" "github.com/ethereum/go-ethereum/p2p/protocols" + "github.com/ethereum/go-ethereum/p2p/simulations" p2ptest "github.com/ethereum/go-ethereum/p2p/testing" ) type pssTester struct { *p2ptest.ProtocolTester - ct *protocols.CodeMap + ct *protocols.CodeMap *Pss } - func TestPssTwoToSelf(t *testing.T) { addr := RandomAddr() pt := newPssTester(t, addr, 2) + defer pt.Stop() payload := []byte("foo42") - - subpeermsgcode, found := pt.ct.GetCode(&SubPeersMsg{}) + + subpeermsgcode, found := pt.ct.GetCode(&subPeersMsg{}) if !found { t.Fatalf("peerMsg not defined") } -/* - peersmsgcode, found := pt.ct.GetCode(&peersMsg{}) - if !found { - t.Fatalf("PssMsg not defined") - } -*/ pssmsgcode, found := pt.ct.GetCode(&PssMsg{}) if !found { t.Fatalf("PssMsg not defined") } - + hs_pivot := correctBzzHandshake(addr) - + for _, id := range pt.Ids { hs_sim := correctBzzHandshake(NewPeerAddrFromNodeId(id)) <-pt.GetPeer(id).Connc @@ -49,37 +43,30 @@ func TestPssTwoToSelf(t *testing.T) { if err != nil { t.Fatalf("Handshake fail: %v", err) } - + err = pt.TestExchanges( p2ptest.Exchange{ Expects: []p2ptest.Expect{ p2ptest.Expect{ Code: subpeermsgcode, - Msg: &SubPeersMsg{}, + Msg: &subPeersMsg{}, Peer: id, }, }, - /*Triggers: []p2ptest.Trigger{ - p2ptest.Trigger{ - Code: peersmsgcode, - Msg: &peersMsg{}, - Peer: id, - }, - },*/ }, ) if err != nil { t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err) } } - - err := pt.TestExchanges ( + + err := pt.TestExchanges( p2ptest.Exchange{ Triggers: []p2ptest.Trigger{ p2ptest.Trigger{ Code: pssmsgcode, - Msg: &PssMsg{ - To: addr.OverlayAddr(), + Msg: &PssMsg{ + To: addr.OverlayAddr(), Data: payload, }, Peer: pt.Ids[0], @@ -89,42 +76,36 @@ func TestPssTwoToSelf(t *testing.T) { ) if err != nil { t.Fatalf("PssMsg sending %v to %v (pivot) fail: %v", pt.Ids[0], addr.OverlayAddr(), err) - } - + } + alarm := time.NewTimer(1000 * time.Millisecond) select { - case data := <-pt.C: - if !bytes.Equal(data, payload) { - t.Fatalf("Data transfer failed, expected: %v, got: %v", payload, data) - } - case <-alarm.C: - t.Fatalf("Pivot receive of PssMsg from %v timeout", pt.Ids[0]) + case data := <-pt.C: + if !bytes.Equal(data, payload) { + t.Fatalf("Data transfer failed, expected: %v, got: %v", payload, data) + } + case <-alarm.C: + t.Fatalf("Pivot receive of PssMsg from %v timeout", pt.Ids[0]) } } - func TestPssTwoRelaySelf(t *testing.T) { + // <-(chan bool)(nil) addr := RandomAddr() pt := newPssTester(t, addr, 2) - - - subpeermsgcode, found := pt.ct.GetCode(&SubPeersMsg{}) + defer pt.Stop() + + subpeermsgcode, found := pt.ct.GetCode(&subPeersMsg{}) if !found { t.Fatalf("peerMsg not defined") } -/* - peersmsgcode, found := pt.ct.GetCode(&peersMsg{}) - if !found { - t.Fatalf("PssMsg not defined") - } -*/ pssmsgcode, found := pt.ct.GetCode(&PssMsg{}) if !found { t.Fatalf("PssMsg not defined") } - + hs_pivot := correctBzzHandshake(addr) - + for _, id := range pt.Ids { hs_sim := correctBzzHandshake(NewPeerAddrFromNodeId(id)) <-pt.GetPeer(id).Connc @@ -132,37 +113,30 @@ func TestPssTwoRelaySelf(t *testing.T) { if err != nil { t.Fatalf("Handshake fail: %v", err) } - + err = pt.TestExchanges( p2ptest.Exchange{ Expects: []p2ptest.Expect{ p2ptest.Expect{ Code: subpeermsgcode, - Msg: &SubPeersMsg{}, + Msg: &subPeersMsg{}, Peer: id, }, }, - /*Triggers: []p2ptest.Trigger{ - p2ptest.Trigger{ - Code: peersmsgcode, - Msg: &peersMsg{}, - Peer: id, - }, - },*/ }, ) if err != nil { t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err) } } - - err := pt.TestExchanges ( + + err := pt.TestExchanges( p2ptest.Exchange{ Expects: []p2ptest.Expect{ p2ptest.Expect{ Code: pssmsgcode, - Msg: &PssMsg{ - To: pt.Ids[0].Bytes(), + Msg: &PssMsg{ + To: pt.Ids[0].Bytes(), Data: []byte("foo42"), }, Peer: pt.Ids[0], @@ -171,8 +145,8 @@ func TestPssTwoRelaySelf(t *testing.T) { Triggers: []p2ptest.Trigger{ p2ptest.Trigger{ Code: pssmsgcode, - Msg: &PssMsg{ - To: pt.Ids[0].Bytes(), + Msg: &PssMsg{ + To: pt.Ids[0].Bytes(), Data: []byte("foo42"), }, Peer: pt.Ids[1], @@ -182,7 +156,7 @@ func TestPssTwoRelaySelf(t *testing.T) { ) if err != nil { t.Fatalf("PssMsg routing from %v to %v fail: %v", pt.Ids[0], pt.Ids[1], err) - } + } } func newPssTester(t *testing.T, addr *peerAddr, n int) *pssTester { @@ -194,9 +168,8 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester { ct.Register(&PssMsg{}) ct.Register(&peersMsg{}) ct.Register(&getPeersMsg{}) - ct.Register(&SubPeersMsg{}) // why is this public? - - + ct.Register(&subPeersMsg{}) // why is this public? + simPipe := adapters.NewSimPipe kp := NewKadParams() kp.MinProxBinSize = 3 @@ -227,8 +200,8 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester { return &pssTester{ ProtocolTester: s, - ct: ct, - Pss: ps, + ct: ct, + Pss: ps, } } diff --git a/swarm/network/simulations/overlay.go b/swarm/network/simulations/overlay.go index 41c19ebefd..84daa5e2b2 100644 --- a/swarm/network/simulations/overlay.go +++ b/swarm/network/simulations/overlay.go @@ -7,6 +7,7 @@ package main import ( "math/rand" + "reflect" "runtime" "time" @@ -15,7 +16,6 @@ import ( "github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p/adapters" "github.com/ethereum/go-ethereum/p2p/simulations" - //p2ptest "github.com/ethereum/go-ethereum/p2p/testing" "github.com/ethereum/go-ethereum/swarm/network" ) @@ -60,27 +60,33 @@ func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapt na := adapters.NewSimNode(id, self.Network, adapters.NewSimPipe) addr := network.NewPeerAddrFromNodeId(id) kp := network.NewKadParams() - kp.MinProxBinSize = 4 - kp.MinBinSize = 4 + + kp.MinProxBinSize = 2 + kp.MaxBinSize = 3 + kp.MinBinSize = 1 + kp.MaxRetries = 1000 + kp.RetryExponent = 2 + kp.RetryInterval = 1000000 + to := network.NewKademlia(addr.OverlayAddr(), kp) // overlay topology driver // to := network.NewTestOverlay(addr.OverlayAddr()) // overlay topology driver hp := network.NewHiveParams() hp.CallInterval = 5000 - pp := network.NewHive(hp, to) // hive - self.hives = append(self.hives, pp) // remember hive + 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) + glog.V(logger.Detail).Infof("kademlia on %v", dp) p.DisconnectHook(func(err error) { pp.Remove(dp) }) return nil } - - ct := network.BzzCodeMap(network.HiveMsgs...) // bzz protocol code map + + ct := network.BzzCodeMap(network.DiscoveryMsgs...) // bzz protocol code map na.Run = network.Bzz(addr.OverlayAddr(), na, ct, services, nil, nil).Run connect := func(s string) error { return self.Connect(id, adapters.NewNodeIdFromHex(s)) @@ -94,7 +100,6 @@ func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapt func NewNetwork(network *simulations.Network) *Network { n := &Network{ - // hives: Network: network, } n.SetNaf(n.NewSimNode) @@ -104,7 +109,7 @@ func NewNetwork(network *simulations.Network) *Network { func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simulations.ResourceController) { conf.DefaultMockerConfig = simulations.DefaultMockerConfig() conf.DefaultMockerConfig.SwitchonRate = 100 - // conf.DefaultMockerConfig.NodesTarget = 15 + conf.DefaultMockerConfig.NodesTarget = 15 conf.DefaultMockerConfig.NewConnCount = 1 conf.DefaultMockerConfig.DegreeTarget = 0 conf.Id = "0" @@ -112,7 +117,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu net := NewNetwork(simulations.NewNetwork(conf)) //ids := p2ptest.RandomNodeIds(10) - ids := adapters.RandomNodeIds(10) + ids := adapters.RandomNodeIds(3) for i, id := range ids { net.NewNode(&simulations.NodeConfig{Id: id}) @@ -129,23 +134,45 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu } go func() { for _, id := range ids { + n := rand.Intn(1000) + time.Sleep(time.Duration(n) * time.Millisecond) net.NewNode(&simulations.NodeConfig{Id: id}) net.Start(id) glog.V(logger.Debug).Infof("node %v starting up", id) - n := rand.Intn(1000) - time.Sleep(time.Duration(n) * time.Millisecond) + // time.Sleep(1000 * time.Millisecond) // net.Stop(id) } }() - // go func() { - // for _, id := range ids { - // net.Stop(id) - // glog.V(logger.Debug).Infof("node %v shutting down", id) - // n := rand.Intn(500) - // time.Sleep(time.Duration(n) * time.Millisecond) + // for i, id := range ids { + // n := 3000 + i*1000 + // go func() { + // for { + // // n := rand.Intn(5000) + // // n := 3000 + // time.Sleep(time.Duration(n) * time.Millisecond) + // glog.V(logger.Debug).Infof("node %v shutting down", id) + // net.Stop(id) + // // n = rand.Intn(5000) + // n = 2000 + // time.Sleep(time.Duration(n) * time.Millisecond) + // glog.V(logger.Debug).Infof("node %v starting up", id) + // net.Start(id) + // n = 5000 + // } + // }() // } - // }() - return net, nil + nodes := simulations.NewResourceContoller( + &simulations.ResourceHandlers{ + Retrieve: &simulations.ResourceHandler{ + Handle: func(msg interface{}, parent *simulations.ResourceController) (interface{}, error) { + id := msg.(string) + pp := net.GetNode(adapters.NewNodeIdFromHex(id)).Adapter().(*SimNode).hive + return pp.String(), nil + }, + Type: reflect.TypeOf([]string{}), // this is input not output param structure + }, + }) + return net, nodes } // var server