From e0468e0442050d6b6b709fab8ac24091fa87322c Mon Sep 17 00:00:00 2001 From: nolash Date: Mon, 20 Mar 2017 02:07:44 +0100 Subject: [PATCH] swarm/network: pssmsg relay through pivotnode --- p2p/testing/protocolsession.go | 3 +- swarm/network/discovery.go | 10 +- swarm/network/discovery_test.go | 2 +- swarm/network/kademlia.go | 14 ++- swarm/network/pss.go | 40 +++++-- swarm/network/pss_test.go | 182 +++++++++++++++++++++++++++----- 6 files changed, 207 insertions(+), 44 deletions(-) diff --git a/p2p/testing/protocolsession.go b/p2p/testing/protocolsession.go index e07075b556..f68bec002e 100644 --- a/p2p/testing/protocolsession.go +++ b/p2p/testing/protocolsession.go @@ -129,7 +129,6 @@ func (self *ProtocolSession) expect(exp Expect) error { glog.V(logger.Detail).Infof("caught timeout") return fmt.Errorf("timout expecting %v sent to peer %v", exp.Msg, exp.Peer) } - // fatal upon encountering first exchange error } // TestExchange tests a series of exchanges againsts the session @@ -137,7 +136,7 @@ func (self *ProtocolSession) TestExchanges(exchanges ...Exchange) error { // launch all triggers of this exchanges for i, e := range exchanges { - errc := make(chan error, 1) + errc := make(chan error) wg := &sync.WaitGroup{} for _, trig := range e.Triggers { err := self.trigger(trig) diff --git a/swarm/network/discovery.go b/swarm/network/discovery.go index cf7c086d3b..39ebe99449 100644 --- a/swarm/network/discovery.go +++ b/swarm/network/discovery.go @@ -38,7 +38,8 @@ 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, MinProxBinSize: 8}) + return self.Send(&SubPeersMsg{ProxLimit: po}) } // new discovery contructor @@ -95,7 +96,7 @@ 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 + //MinProxBinSize uint8 ProxLimit uint8 } @@ -133,6 +134,11 @@ func (p *discPeer) handlePeersMsg(msg interface{}) error { nas = append(nas, addr) p.peers[NodeId(addr).NodeID] = true } + + if len(nas) == 0 { + glog.V(logger.Debug).Infof("whoops, no peers in incoming peersMsg from %v", p) + return nil + } return p.overlay.Register(nas...) } diff --git a/swarm/network/discovery_test.go b/swarm/network/discovery_test.go index 7187596685..e9086e6af2 100644 --- a/swarm/network/discovery_test.go +++ b/swarm/network/discovery_test.go @@ -37,7 +37,7 @@ func TestDiscovery(t *testing.T) { Expects: []p2ptest.Expect{ p2ptest.Expect{ Code: 3, - Msg: &SubPeersMsg{ProxLimit: 0, MinProxBinSize: 8}, + Msg: &SubPeersMsg{ProxLimit: 0}, Peer: s.ProtocolTester.Ids[0], }, }, diff --git a/swarm/network/kademlia.go b/swarm/network/kademlia.go index 2c96d2caf9..94155c30c1 100644 --- a/swarm/network/kademlia.go +++ b/swarm/network/kademlia.go @@ -70,7 +70,7 @@ type KadParams struct { func NewKadParams() *KadParams { return &KadParams{ MaxProxDisplay: 8, - MinProxBinSize: 2, + MinProxBinSize: 4, MinBinSize: 2, MaxBinSize: 4, //RetryInterval: 42000000000, @@ -85,6 +85,7 @@ 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 } // NewKademlia(addr, params) creates a Kademlia table for base address addr @@ -218,10 +219,15 @@ func (self *Kademlia) On(p Peer) { glog.V(logger.Detail).Infof("peer %v nofified", vp) dp := val.(*KadPeer).Peer.(KadDiscovery) dp.NotifyPeer(kp.Peer, uint8(po)) - dp.NotifyProx(uint8(prox)) + if uint8(prox) != self.lastProxLimit { + self.lastProxLimit = uint8(prox) + dp.NotifyProx(uint8(prox)) + } } self.conns.EachNeighbourAsync(kp, 255, 255, f, false) go vp.NotifyProx(uint8(prox)) + + } // Off removes a peer from among live peers @@ -375,9 +381,9 @@ func (self *Kademlia) String() string { var rows []string - rows = append(rows, "=========================================================================") + 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, fmt.Sprintf("population: %d (%d), ProxBinSize: %d, MinBinSize: %d, MaxBinSize: %d", self.conns.Size(), self.peers.Size(), self.MinProxBinSize, self.MinBinSize, self.MaxBinSize)) + 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) peersrows := make([]string, self.MaxProxDisplay) diff --git a/swarm/network/pss.go b/swarm/network/pss.go index 87f506caba..5ce732e7c6 100644 --- a/swarm/network/pss.go +++ b/swarm/network/pss.go @@ -2,25 +2,51 @@ package network import ( "fmt" + "bytes" "github.com/ethereum/go-ethereum/logger" "github.com/ethereum/go-ethereum/logger/glog" ) -type pssPeer struct { - Peer +type Pss struct { + Overlay + LocalAddr []byte + C chan []byte +} + +func NewPss(k Overlay, addr []byte) *Pss { + return &Pss{ + Overlay: k, + LocalAddr: addr, + C: make(chan []byte), + } } type PssMsg struct { - Recipient pssPeer - Payload []byte + To []byte + Data []byte } func (pm *PssMsg) String() string { - return fmt.Sprintf("PssMsg: Recipient: %v", pm.Recipient) + return fmt.Sprintf("PssMsg: Recipient: %v", pm.To) } -func PssMsgHandler(msg interface{}) error { - glog.V(logger.Detail).Infof("Pss Handled!") +func (ps *Pss) HandlePssMsg(msg interface{}) error { + pssmsg := msg.(*PssMsg) + to := pssmsg.To + if bytes.Equal(to, ps.LocalAddr) { + glog.V(logger.Detail).Infof("Pss to us, yay!", to) + ps.C <- pssmsg.Data + return nil + } + + ps.EachLivePeer(to, 255, func(p Peer, po int) bool { + err := p.Send(pssmsg) + if err != nil { + return true + } + return false + }) + return nil } diff --git a/swarm/network/pss_test.go b/swarm/network/pss_test.go index 8c9e3f02e8..ff9e6d9d29 100644 --- a/swarm/network/pss_test.go +++ b/swarm/network/pss_test.go @@ -1,10 +1,12 @@ package network import ( + "bytes" "testing" + "time" - "github.com/ethereum/go-ethereum/logger" - "github.com/ethereum/go-ethereum/logger/glog" + //"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" @@ -14,52 +16,173 @@ import ( type pssTester struct { *p2ptest.ProtocolTester ct *protocols.CodeMap + *Pss } + func TestPssTwoToSelf(t *testing.T) { addr := RandomAddr() pt := newPssTester(t, addr, 2) + payload := []byte("foo42") subpeermsgcode, found := pt.ct.GetCode(&SubPeersMsg{}) if !found { t.Fatalf("peerMsg not defined") } - - /*peermsgcode, found := pt.ct.GetCode(&peersMsg{}) +/* + peersmsgcode, found := pt.ct.GetCode(&peersMsg{}) if !found { - t.Fatalf("peerMsg not defined") - }*/ + 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)) - glog.V(logger.Detail).Infof("Will handshake %v with %v", hs_pivot, hs_sim) <-pt.GetPeer(id).Connc - pt.TestExchanges(bzzHandshakeExchange(hs_pivot, hs_sim, id)...) + err := pt.TestExchanges(bzzHandshakeExchange(hs_pivot, hs_sim, id)...) + if err != nil { + t.Fatalf("Handshake fail: %v", err) + } - pt.TestExchanges( - p2ptest.Exchange{ - Expects: []p2ptest.Expect{ - p2ptest.Expect{ - Code: subpeermsgcode, - Msg: &SubPeersMsg{}, - Peer: id, + err = pt.TestExchanges( + p2ptest.Exchange{ + Expects: []p2ptest.Expect{ + p2ptest.Expect{ + Code: subpeermsgcode, + Msg: &SubPeersMsg{}, + Peer: id, + }, }, + /*Triggers: []p2ptest.Trigger{ + p2ptest.Trigger{ + Code: peersmsgcode, + Msg: &peersMsg{}, + Peer: id, + }, + },*/ }, - },/* - p2ptest.Exchange{ - Expects: []p2ptest.Expect{ - p2ptest.Expect{ - Code: peermsgcode, - Msg: &peersMsg{}, - Peer: id, - }, - }, - },*/ ) - + if err != nil { + t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err) + } } + + err := pt.TestExchanges ( + p2ptest.Exchange{ + Triggers: []p2ptest.Trigger{ + p2ptest.Trigger{ + Code: pssmsgcode, + Msg: &PssMsg{ + To: addr.OverlayAddr(), + Data: payload, + }, + Peer: pt.Ids[0], + }, + }, + }, + ) + 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]) + } +} + + +func TestPssTwoRelaySelf(t *testing.T) { + addr := RandomAddr() + pt := newPssTester(t, addr, 2) + + + 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 + err := pt.TestExchanges(bzzHandshakeExchange(hs_pivot, hs_sim, id)...) + if err != nil { + t.Fatalf("Handshake fail: %v", err) + } + + err = pt.TestExchanges( + p2ptest.Exchange{ + Expects: []p2ptest.Expect{ + p2ptest.Expect{ + Code: subpeermsgcode, + 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 ( + p2ptest.Exchange{ + Expects: []p2ptest.Expect{ + p2ptest.Expect{ + Code: pssmsgcode, + Msg: &PssMsg{ + To: pt.Ids[0].Bytes(), + Data: []byte("foo42"), + }, + Peer: pt.Ids[0], + }, + }, + Triggers: []p2ptest.Trigger{ + p2ptest.Trigger{ + Code: pssmsgcode, + Msg: &PssMsg{ + To: pt.Ids[0].Bytes(), + Data: []byte("foo42"), + }, + Peer: pt.Ids[1], + }, + }, + }, + ) + 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 { @@ -71,13 +194,15 @@ 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 official? + ct.Register(&SubPeersMsg{}) // why is this public? simPipe := adapters.NewSimPipe kp := NewKadParams() + kp.MinProxBinSize = 3 to := NewKademlia(addr.OverlayAddr(), kp) pp := NewHive(NewHiveParams(), to) + ps := NewPss(to, addr.OverlayAddr()) net := simulations.NewNetwork(&simulations.NetworkConfig{}) naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter { na := adapters.NewSimNode(conf.Id, net, simPipe) @@ -86,7 +211,7 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester { net.SetNaf(naf) srv := func(p Peer) error { - p.Register(&PssMsg{}, PssMsgHandler) + p.Register(&PssMsg{}, ps.HandlePssMsg) pp.Add(p) p.DisconnectHook(func(err error) { pp.Remove(p) @@ -103,6 +228,7 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester { return &pssTester{ ProtocolTester: s, ct: ct, + Pss: ps, } }