mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 22:56:43 +00:00
swarm/network: pssmsg relay through pivotnode
This commit is contained in:
parent
5401bf02f1
commit
e0468e0442
6 changed files with 207 additions and 44 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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...)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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],
|
||||
},
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,51 +16,172 @@ 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)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -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,
|
||||
}
|
||||
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue