kademlia fixes

This commit is contained in:
zelig 2017-03-25 01:01:40 +07:00 committed by Lewis Marshall
parent e0468e0442
commit b7860c7fc6
18 changed files with 784 additions and 359 deletions

View file

@ -109,12 +109,14 @@ func (self *SimNode) setPeer(id *NodeId, m Messenger) *Peer {
self.peers = append(self.peers, p) self.peers = append(self.peers, p)
return p return p
} }
if self.peers[i] != nil && m != nil { // if self.peers[i] != nil && m != nil {
panic(fmt.Sprintf("pipe for %v already set", id)) // panic(fmt.Sprintf("pipe for %v already set", id))
} // }
// legit reconnect reset disconnection error, // legit reconnect reset disconnection error,
p := self.peers[i] p := self.peers[i]
p.Messenger = m p.Messenger = m
p.Connc = make(chan bool)
p.Readyc = make(chan bool)
return p return p
} }
@ -145,8 +147,6 @@ func (self *SimNode) Connect(rid []byte) error {
return fmt.Errorf("node adapter for %v is missing", id) return fmt.Errorf("node adapter for %v is missing", id)
} }
rw, rrw := p2p.MsgPipe() rw, rrw := p2p.MsgPipe()
// runc := make(chan bool)
// defer close(runc)
// // run protocol on remote node with self as peer // // run protocol on remote node with self as peer
peer := self.getPeer(id) peer := self.getPeer(id)
if peer != nil && peer.Messenger != nil { if peer != nil && peer.Messenger != nil {

View file

@ -32,6 +32,7 @@ package protocols
import ( import (
"fmt" "fmt"
"reflect" "reflect"
"time"
"github.com/ethereum/go-ethereum/logger" "github.com/ethereum/go-ethereum/logger"
"github.com/ethereum/go-ethereum/logger/glog" "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 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 handlers map[reflect.Type][]func(interface{}) error // message type -> message handler callback(s) map
Errc chan error Errc chan error
wErrc chan error // write error channel
} }
// NewPeer returns a new peer // NewPeer returns a new peer
@ -208,6 +210,7 @@ func NewPeer(p *p2p.Peer, ct *CodeMap, m adapters.Messenger) *Peer {
m: m, m: m,
Peer: p, Peer: p,
Errc: make(chan error), Errc: make(chan error),
wErrc: make(chan error),
handlers: make(map[reflect.Type][]func(interface{}) 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) return errorf(ErrInvalidMsgType, "%v", code)
} }
glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg) glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg)
err := self.m.SendMsg(uint64(code), msg) go func() {
if err != nil { self.wErrc <- self.m.SendMsg(uint64(code), msg)
err = errorf(ErrWrite, "(msg code: %v): %v", code, err) }()
self.Drop(err) var err error
return err 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)) { func (self *Peer) DisconnectHook(f func(error)) {

View file

@ -14,7 +14,7 @@ import (
) )
func init() { func init() {
glog.SetV(logger.Detail) glog.SetV(logger.Error)
glog.SetToStderr(true) 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 // demonstrates use of peerPool, killing another peer connection as a response to a message
peer.Register(&kill{}, func(msg interface{}) error { peer.Register(&kill{}, func(msg interface{}) error {
id := msg.(*kill).C 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) glog.V(logger.Detail).Infof("id %v killed", id)
return nil return nil
}) })
@ -342,7 +342,6 @@ func TestMultiplePeersDropSelf(t *testing.T) {
} }
func TestMultiplePeersDropOther(t *testing.T) { func TestMultiplePeersDropOther(t *testing.T) {
t.Skip("??")
runMultiplePeers(t, 1, runMultiplePeers(t, 1,
fmt.Errorf("Message handler error: (msg code 3): dropped"), fmt.Errorf("Message handler error: (msg code 3): dropped"),
fmt.Errorf("p2p: read or write on closed message pipe"), fmt.Errorf("p2p: read or write on closed message pipe"),

View file

@ -118,7 +118,7 @@ func (self *ProtocolSession) expect(exp Expect) error {
t := exp.Timeout t := exp.Timeout
if t == time.Duration(0) { if t == time.Duration(0) {
t = 1000 * time.Millisecond t = 2000 * time.Millisecond
} }
alarm := time.NewTimer(t) alarm := time.NewTimer(t)
select { select {
@ -126,7 +126,6 @@ func (self *ProtocolSession) expect(exp Expect) error {
glog.V(logger.Detail).Infof("expected msg arrives with error %v", err) glog.V(logger.Detail).Infof("expected msg arrives with error %v", err)
return err return err
case <-alarm.C: case <-alarm.C:
glog.V(logger.Detail).Infof("caught timeout")
return fmt.Errorf("timout expecting %v sent to peer %v", exp.Msg, exp.Peer) 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 return nil
} }
func (self *ProtocolSession) Stop() {
for _, id := range self.Ids {
p := self.GetPeer(id)
if p != nil && p.Messenger != nil {
p.Close()
}
}
}

View file

@ -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) // the string form of the binary representation of an address (only first 8 bits)
func (a Address) Bin() string { func (a Address) Bin() string {
return ToBin(a[:])
}
func ToBin(a []byte) string {
var bs []string var bs []string
for _, b := range a[:] { for _, b := range a {
bs = append(bs, fmt.Sprintf("%08b", b)) bs = append(bs, fmt.Sprintf("%08b", b))
} }
return strings.Join(bs, "") return strings.Join(bs, "")
@ -209,7 +213,7 @@ func (a *HashAddress) String() string {
return a.Address.Bin() return a.Address.Bin()
} }
func NewHashAddress(s string) *HashAddress { func NewAddressFromString(s string) []byte {
ha := [32]byte{} ha := [32]byte{}
t := s + string(zerosBin)[:len(zerosBin)-len(s)] 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)) binary.BigEndian.PutUint64(ha[i*8:(i+1)*8], uint64(n))
} }
return ha[:]
}
func NewHashAddress(s string) *HashAddress {
ha := NewAddressFromString(s)
h := common.Hash{} h := common.Hash{}
copy(h[:], ha[:]) copy(h[:], ha)
return &HashAddress{Address(h)} return &HashAddress{Address(h)}
} }

View file

@ -21,7 +21,6 @@ import (
) )
const ( const (
// keylen = 4
keylen = 256 keylen = 256
maxkeylen = 256 maxkeylen = 256
) )
@ -398,7 +397,6 @@ func union(t0, t1 *pot) (*pot, int) {
if t1 == nil || t1.size == 0 { if t1 == nil || t1.size == 0 {
return t0, 0 return t0, 0
} }
po, eq := t0.pin.PO(t1.pin, 0)
var pin PotVal var pin PotVal
var bins []*pot var bins []*pot
var mis []int var mis []int
@ -410,6 +408,8 @@ func union(t0, t1 *pot) (*pot, int) {
var i0, i1 int var i0, i1 int
var common int var common int
po, eq := pin0.PO(pin1, 0)
for { for {
l0 := len(bins0) l0 := len(bins0)
l1 := len(bins1) l1 := len(bins1)
@ -470,14 +470,17 @@ func union(t0, t1 *pot) (*pot, int) {
break break
} }
i := i0
if len(bins0) > i && bins0[i].po == po {
i++
}
var size0 int var size0 int
for _, n := range bins0[i0:] { for _, n := range bins0[i:] {
size0 += n.size size0 += n.size
} }
np := &pot{ np := &pot{
pin: pin0, pin: pin0,
bins: bins0[i0:], bins: bins0[i:],
size: size0 + 1, size: size0 + 1,
po: po, po: po,
} }
@ -488,11 +491,13 @@ func union(t0, t1 *pot) (*pot, int) {
po = maxkeylen + 1 po = maxkeylen + 1
eq = true eq = true
common-- common--
} else { } else {
bins2 = append(bins2, n0.bins...) bins2 = append(bins2, n0.bins...)
pin0 = pin1 pin0 = pin1
pin1 = n0.pin pin1 = n0.pin
po, eq = pin0.PO(pin1, n0.po) po, eq = pin0.PO(pin1, n0.po)
} }
bins0 = bins1 bins0 = bins1
bins1 = bins2 bins1 = bins2
@ -863,6 +868,9 @@ func (t *pot) String() string {
} }
func (t *pot) sstring(indent string) string { func (t *pot) sstring(indent string) string {
if t == nil {
return "<nil>"
}
var s string var s string
indent += " " indent += " "
s += fmt.Sprintf("%v%v (%v) %v \n", indent, t.pin, t.po, t.size) s += fmt.Sprintf("%v%v (%v) %v \n", indent, t.pin, t.po, t.size)

View file

@ -24,6 +24,7 @@ import (
"testing" "testing"
"time" "time"
"github.com/ethereum/go-ethereum/logger"
"github.com/ethereum/go-ethereum/logger/glog" "github.com/ethereum/go-ethereum/logger/glog"
) )
@ -224,7 +225,10 @@ func TestPotSwap(t *testing.T) {
return true return true
}) })
if sum != 2*max { 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 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++ { for i := 0; i < maxEachNeighbourTests; i++ {
max0 := rand.Intn(maxEachNeighbour) + 1 max0 := rand.Intn(maxEachNeighbour) + 1
max1 := rand.Intn(maxEachNeighbour) + 1 max1 := rand.Intn(maxEachNeighbour) + 1

View file

@ -1,36 +1,61 @@
package network package network
import ( import (
"bytes"
"fmt" "fmt"
"github.com/ethereum/go-ethereum/logger" "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/p2p/discover" // "github.com/ethereum/go-ethereum/p2p/discover"
) )
// discovery bzz overlay extension doing peer relaying // 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 { type discPeer struct {
Peer Peer
overlay Overlay overlay Overlay
proxLimit uint8 peers map[string]bool
peers map[discover.NodeID]bool // peers map[discover.NodeID]bool
sentPeers 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. // NotifyPeer notifies the receiver remote end of a peer p or PO po.
// callback for overlay driver // callback for overlay driver
func (self *discPeer) NotifyPeer(p Peer, po uint8) error { 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 return nil
} }
glog.V(logger.Warn).Infof("notification about %x", p.OverlayAddr())
resp := &peersMsg{ 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 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 // 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) // or first empty row)
// callback for overlay driver // callback for overlay driver
func (self *discPeer) NotifyProx(po uint8) error { func (self *discPeer) NotifyProx(po uint8) error {
//return self.Send(&SubPeersMsg{ProxLimit: po, MinProxBinSize: 8}) return self.Send(&subPeersMsg{ProxLimit: po})
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
} }
/* /*
@ -95,51 +104,54 @@ func (self getPeersMsg) String() string {
} }
// subPeers msg is communicating the depth/sharpness/focus of the overlay table of a peer // subPeers msg is communicating the depth/sharpness/focus of the overlay table of a peer
type SubPeersMsg struct { type subPeersMsg struct {
//MinProxBinSize uint8 ProxLimit uint8
ProxLimit uint8
} }
func (self SubPeersMsg) String() string { func (self subPeersMsg) String() string {
return fmt.Sprintf("%T: request peers > PO%02d. ", self, self.ProxLimit) return fmt.Sprintf("%T: request peers > PO%02d. ", self, self.ProxLimit)
} }
func (self *discPeer) handleSubPeersMsg(msg interface{}) error { func (self *discPeer) handleSubPeersMsg(msg interface{}) error {
spm := msg.(*SubPeersMsg) spm := msg.(*subPeersMsg)
self.proxLimit = spm.ProxLimit
if !self.sentPeers { if !self.sentPeers {
var peers []*peerAddr var peers []*peerAddr
self.overlay.EachLivePeer(self.OverlayAddr(), 255, func(p Peer, po int) bool { self.overlay.EachLivePeer(self.OverlayAddr(), 255, func(p Peer, po int) bool {
if uint8(po) > self.proxLimit { if uint8(po) < self.proxLimit {
return false return false
} }
self.peers[p.ID()] = true self.seen(p)
peers = append(peers, &peerAddr{p.OverlayAddr(), p.UnderlayAddr()}) peers = append(peers, &peerAddr{p.OverlayAddr(), p.UnderlayAddr()})
return true 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.sentPeers = true
self.proxLimit = spm.ProxLimit
return nil return nil
} }
// handlePeersMsg called by the protocol when receiving peerset (for target address) // 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 // list of nodes ([]PeerAddr in peersMsg) is added to the overlay db using the
// Register interface method // Register interface method
func (p *discPeer) handlePeersMsg(msg interface{}) error { func (self *discPeer) handlePeersMsg(msg interface{}) error {
// register all addresses // register all addresses
var nas []PeerAddr var nas []PeerAddr
for _, na := range msg.(*peersMsg).Peers { for _, na := range msg.(*peersMsg).Peers {
addr := PeerAddr(na) addr := PeerAddr(na)
nas = append(nas, addr) nas = append(nas, addr)
p.peers[NodeId(addr).NodeID] = true self.seen(addr)
} }
if len(nas) == 0 { 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 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 // 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 // peers suggestions are retrieved from the overlay topology driver
// using the EachLivePeer interface iterator method // using the EachLivePeer interface iterator method
// peers sent are remembered throughout a session and not sent twice // peers sent are remembered throughout a session and not sent twice
func (p *discPeer) handleGetPeersMsg(msg interface{}) error { func (self *discPeer) handleGetPeersMsg(msg interface{}) error {
req := msg.(*getPeersMsg)
var peers []*peerAddr var peers []*peerAddr
alreadySent := p.peers req := msg.(*getPeersMsg)
i := 0 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++ i++
if bytes.Compare(n.OverlayAddr(), p.OverlayAddr()) != 0 && // only send peers we have not sent before in this session
// only send peers we have not sent before in this session if self.seen(n) {
!alreadySent[n.ID()] {
alreadySent[n.ID()] = true
peers = append(peers, &peerAddr{n.OverlayAddr(), n.UnderlayAddr()}) peers = append(peers, &peerAddr{n.OverlayAddr(), n.UnderlayAddr()})
} }
// return int(req.Order) == po && len(peers) < int(req.Max)
return len(peers) < int(req.Max) return len(peers) < int(req.Max)
}) })
if len(peers) == 0 { 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 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{ resp := &peersMsg{
Peers: peers, 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
} }

View file

@ -16,7 +16,7 @@ import (
func TestDiscovery(t *testing.T) { func TestDiscovery(t *testing.T) {
addr := RandomAddr() addr := RandomAddr()
to := NewKademlia(addr.OAddr, NewKadParams()) to := NewKademlia(addr.OAddr, NewKadParams())
ct := BzzCodeMap(HiveMsgs...) ct := BzzCodeMap(DiscoveryMsgs...)
services := func(p Peer) error { services := func(p Peer) error {
dp := NewDiscovery(p, to) dp := NewDiscovery(p, to)
@ -29,6 +29,7 @@ func TestDiscovery(t *testing.T) {
} }
s := newBzzBaseTester(t, 1, addr, ct, services) s := newBzzBaseTester(t, 1, addr, ct, services)
defer s.Stop()
s.runHandshakes() s.runHandshakes()
// o := 0 // o := 0
@ -37,7 +38,7 @@ func TestDiscovery(t *testing.T) {
Expects: []p2ptest.Expect{ Expects: []p2ptest.Expect{
p2ptest.Expect{ p2ptest.Expect{
Code: 3, Code: 3,
Msg: &SubPeersMsg{ProxLimit: 0}, Msg: &subPeersMsg{ProxLimit: 0},
Peer: s.ProtocolTester.Ids[0], Peer: s.ProtocolTester.Ids[0],
}, },
}, },

View file

@ -56,21 +56,14 @@ type Overlay interface {
type Hive struct { type Hive struct {
*HiveParams // settings *HiveParams // settings
Overlay // the overlay topology driver Overlay // the overlay topology driver
// disc Discovery lock sync.Mutex
quit chan bool
lock sync.Mutex toggle chan bool
quit chan bool more chan bool
toggle chan bool
more chan bool
} }
const (
peersBroadcastSetSize = 2
maxPeersPerRequest = 5
callInterval = 1000
)
type HiveParams struct { type HiveParams struct {
Discovery bool
PeersBroadcastSetSize uint8 PeersBroadcastSetSize uint8
MaxPeersPerRequest uint8 MaxPeersPerRequest uint8
CallInterval uint CallInterval uint
@ -78,9 +71,10 @@ type HiveParams struct {
func NewHiveParams() *HiveParams { func NewHiveParams() *HiveParams {
return &HiveParams{ return &HiveParams{
PeersBroadcastSetSize: peersBroadcastSetSize, Discovery: true,
MaxPeersPerRequest: maxPeersPerRequest, PeersBroadcastSetSize: 2,
CallInterval: callInterval, 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 // Start receives network info only at startup
// connectPeer is a function to connect to a peer based on its NodeID or enode URL // 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 // 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 { func (self *Hive) Start(connectPeer func(string) error, af func() <-chan time.Time) error {
self.toggle = make(chan bool) self.toggle = make(chan bool)
self.more = make(chan bool) self.more = make(chan bool, 1)
self.quit = make(chan bool) self.quit = make(chan bool)
glog.V(logger.Debug).Infof("hive started") glog.V(logger.Debug).Infof("hive started")
// this loop is doing bootstrapping and maintains a healthy table // 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 { } else {
glog.V(logger.Detail).Infof("cannot suggest peers") glog.V(logger.Detail).Infof("cannot suggest peers")
} }
want = want && self.Discovery
if want { if want {
req := &getPeersMsg{ RequestOrder(self.Overlay, uint8(order), self.PeersBroadcastSetSize, self.MaxPeersPerRequest)
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)
} }
select { 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() { func (self *Hive) Stop() {
// closing toggle channel quits the updateloop // closing toggle channel quits the updateloop
close(self.quit) 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 { func HexToBytes(s string) []byte {
id := discover.MustHexID(s) id := discover.MustHexID(s)
return id[:] return id[:]

View file

@ -5,17 +5,10 @@ import (
"testing" "testing"
"time" "time"
"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/adapters"
p2ptest "github.com/ethereum/go-ethereum/p2p/testing" p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
) )
func init() {
glog.SetV(logger.Detail)
glog.SetToStderr(true)
}
type testConnect struct { type testConnect struct {
mu sync.Mutex mu sync.Mutex
conns []string conns []string
@ -35,12 +28,13 @@ func (self *testConnect) connect(na string) error {
return nil return nil
} }
func TestOverlayRegistration(t *testing.T) { func newHiveTester(t *testing.T, params *HiveParams) (*bzzTester, *Hive) {
// setup // setup
addr := RandomAddr() // tested peers peer address addr := RandomAddr() // tested peers peer address
to := NewTestOverlay(addr.OverlayAddr()) // overlay topology driver to := NewTestOverlay(addr.OverlayAddr()) // overlay topology drive
pp := NewHive(NewHiveParams(), to) // hive pp := NewHive(params, to) // hive
ct := BzzCodeMap(HiveMsgs...) // bzz protocol code map
ct := BzzCodeMap(DiscoveryMsgs...) // bzz protocol code map
services := func(p Peer) error { services := func(p Peer) error {
pp.Add(p) pp.Add(p)
p.DisconnectHook(func(err error) { p.DisconnectHook(func(err error) {
@ -49,40 +43,36 @@ func TestOverlayRegistration(t *testing.T) {
return nil 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] id := s.Ids[0]
raddr := NewPeerAddrFromNodeId(id) raddr := NewPeerAddrFromNodeId(id)
s.runHandshakes() s.runHandshakes()
// hive should have called the overlay // 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") t.Fatalf("Overlay#On not called on new peer")
} }
} }
func TestRegisterAndConnect(t *testing.T) { func TestRegisterAndConnect(t *testing.T) {
// setup params := NewHiveParams()
addr := RandomAddr() // tested peers peer address s, pp := newHiveTester(t, params)
to := NewTestOverlay(addr.OverlayAddr()) // overlay topology driver defer s.Stop()
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)
id := s.Ids[0] id := s.Ids[0]
raddr := NewPeerAddrFromNodeId(id) raddr := NewPeerAddrFromNodeId(id)
pp.Register(raddr) pp.Register(raddr)
glog.V(5).Infof("%v", pp)
// start the hive and wait for the connection // start the hive and wait for the connection
tc := &testConnect{ tc := &testConnect{
@ -93,16 +83,16 @@ func TestRegisterAndConnect(t *testing.T) {
ticker: make(chan time.Time), ticker: make(chan time.Time),
} }
pp.Start(tc.connect, tc.ping) pp.Start(tc.connect, tc.ping)
defer pp.Stop()
tc.ticker <- time.Now() tc.ticker <- time.Now()
s.runHandshakes() 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") t.Fatalf("Overlay#On not called on new peer")
} }
// retrieve and broadcast // retrieve and broadcast
glog.V(6).Infof("check peer requests for %v", id)
ord := order(raddr.OverlayAddr()) ord := order(raddr.OverlayAddr())
o := 0 o := 0
if ord == 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,
// },
// },
// })
} }

View file

@ -13,9 +13,11 @@
// //
// You should have received a copy of the GNU Lesser General Public License // You should have received a copy of the GNU Lesser General Public License
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>. // along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
package network package network
import ( import (
"bytes"
"fmt" "fmt"
"strings" "strings"
"time" "time"
@ -64,28 +66,30 @@ type KadParams struct {
RetryInterval int RetryInterval int
RetryExponent int RetryExponent int
MaxRetries int MaxRetries int
PruneInterval int
} }
// NewKadParams() returns a params struct with default values // NewKadParams() returns a params struct with default values
func NewKadParams() *KadParams { func NewKadParams() *KadParams {
return &KadParams{ return &KadParams{
MaxProxDisplay: 8, MaxProxDisplay: 8,
MinProxBinSize: 4, MinProxBinSize: 2,
MinBinSize: 2, MinBinSize: 2,
MaxBinSize: 4, MaxBinSize: 4,
//RetryInterval: 42000000000, //RetryInterval: 42000000000,
RetryInterval: 420000000, RetryInterval: 420000000,
MaxRetries: 42, MaxRetries: 42,
RetryExponent: 2, RetryExponent: 2,
} }
} }
// Kademlia is a table of live peers and a db of known peers // Kademlia is a table of live peers and a db of known peers
type Kademlia struct { type Kademlia struct {
addr *pot.HashAddress // immutable baseaddress of the table addr PeerAddr // immutable baseaddress of the table
*KadParams // Kademlia configuration parameters // addr *pot.HashAddress // immutable baseaddress of the table
conns, peers *pot.Pot // pots container for peers *KadParams // Kademlia configuration parameters
lastProxLimit uint8 // stores the last calculated proxlimit 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 // 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 { if params == nil {
params = NewKadParams() params = NewKadParams()
} }
base := pot.NewHashAddressFromBytes(addr) self := &Kademlia{
return &Kademlia{ addr: &peerAddr{OAddr: addr},
addr: base,
KadParams: params, KadParams: params,
conns: pot.NewPot(nil, 0), conns: pot.NewPot(nil, 0),
peers: 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 // KadPeer represents a Kademlia Peer and extends
@ -109,38 +146,27 @@ func NewKademlia(addr []byte, params *KadParams) *Kademlia {
// * Peer interface (id, last seen, drop) // * Peer interface (id, last seen, drop)
// * HashAddress as derived from PeerAddr overlay implement pot.PoVal interface // * HashAddress as derived from PeerAddr overlay implement pot.PoVal interface
type KadPeer struct { type KadPeer struct {
*pot.HashAddress // *pot.HashAddress
PeerAddr PeerAddr
Peer Peer Peer Peer
seenAt time.Time seenAt time.Time
retries int 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 { func (self *KadPeer) String() string {
if self == nil { if self == nil {
return "<nil>" return "<nil>"
} }
// return string(self.OverlayAddr()) // 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 { func (self *Kademlia) callable(val pot.PotVal) *KadPeer {
kp := val.(*KadPeer) kp := val.(*KadPeer)
// not callable if peer is live or exceeded maxRetries // 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 { 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 return nil
} }
// calculate the allowed number of retries based on time lapsed since last seen // 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 // this is never called concurrently, so safe to increment
// peer can be retried again // peer can be retried again
if retries < kp.retries { if retries < kp.retries {
glog.V(logger.Detail).Infof("log time needed before retry %v, wait only warrants %v", kp.retries, retries) glog.V(logger.Detail).Infof("log time needed before retry %v, wait only warrants %v", kp.retries, retries)
return nil return nil
} }
kp.retries++ 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 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 { func NewKadPeer(na PeerAddr) *KadPeer {
// o := na.OverlayAddr()
// glog.V(logger.Detail).Infof("newkadpeer from peerAddr overlay address: %x", o[:6])
return &KadPeer{ return &KadPeer{
HashAddress: pot.NewHashAddressFromBytes(na.OverlayAddr()), PeerAddr: na,
PeerAddr: na, seenAt: time.Now(),
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 // database of known peers
func (self *Kademlia) Register(nas ...PeerAddr) error { func (self *Kademlia) Register(nas ...PeerAddr) error {
label := fmt.Sprintf("%x", RandomAddr().OverlayAddr())
np := pot.NewPot(nil, 0) np := pot.NewPot(nil, 0)
for _, na := range nas { 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) p := NewKadPeer(na)
np, _, _ = pot.Add(np, pot.PotVal(p)) np, _, _ = pot.Add(np, pot.PotVal(p))
} }
common := self.peers.Merge(np) oldpeers := pot.NewPot(nil, 0)
glog.V(logger.Detail).Infof("add peers: %v out of %v new; root %v", np.Size()-common, np.Size(), np.Pin().String()[:6]) 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 return nil
} }
// On(p) inserts the peer as a kademlia peer into the live peers // On(p) inserts the peer as a kademlia peer into the live peers
func (self *Kademlia) On(p Peer) { func (self *Kademlia) On(p Peer) {
pp := NewKadPeer(p) kp := NewKadPeer(p)
// var pp *KadPeer self.conns.Swap(kp, func(v pot.PotVal) pot.PotVal {
kp := pp
self.conns.Swap(pp, func(v pot.PotVal) pot.PotVal {
if v == nil { 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 { if v != nil {
kp = v.(*KadPeer) kp = v.(*KadPeer)
} }
@ -211,23 +249,21 @@ func (self *Kademlia) On(p Peer) {
vp, ok := kp.Peer.(KadDiscovery) vp, ok := kp.Peer.(KadDiscovery)
if !ok { if !ok {
glog.V(logger.Detail).Infof("not discovery peer") glog.V(logger.Detail).Infof("not discovery peer %T", kp)
return return
} }
// vp.NotifyProx(uint8(prox)) go vp.NotifyProx(uint8(prox))
f := func(val pot.PotVal, po int) { f := func(val pot.PotVal, po int) {
glog.V(logger.Detail).Infof("peer %v nofified", vp)
dp := val.(*KadPeer).Peer.(KadDiscovery) 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)) dp.NotifyPeer(kp.Peer, uint8(po))
if uint8(prox) != self.lastProxLimit { if uint8(prox) != self.lastProxLimit {
self.lastProxLimit = uint8(prox) self.lastProxLimit = uint8(prox)
dp.NotifyProx(uint8(prox)) dp.NotifyProx(uint8(prox))
} }
glog.V(logger.Debug).Infof("peer notified")
} }
self.conns.EachNeighbourAsync(kp, 255, 255, f, false) self.conns.EachNeighbourAsync(kp, 1024, 255, f, false)
go vp.NotifyProx(uint8(prox))
} }
// Off removes a peer from among live peers // Off removes a peer from among live peers
@ -246,6 +282,10 @@ func (self *Kademlia) Off(p Peer) {
kp.seenAt = time.Now() kp.seenAt = time.Now()
} }
type ByteAddr struct {
key []byte
}
// EachLivePeer(base, po, f) is an iterator applying f to each live peer // 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 // that has proximity order po or less as measured from the base
// if base is nil, kademlia base address is used // 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 { if base == nil {
p = pot.PotVal(self.addr) p = pot.PotVal(self.addr)
} else { } else {
p = pot.NewHashAddressFromBytes(base) p = pot.PotVal(&peerAddr{OAddr: base})
} }
self.conns.EachNeighbour(p, func(val pot.PotVal, po int) bool { self.conns.EachNeighbour(p, func(val pot.PotVal, po int) bool {
if po > o { 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 { self.peers.EachNeighbour(self.addr, func(val pot.PotVal, po int) bool {
r := self.callable(val) r := self.callable(val)
if r == nil { if r == nil {
glog.V(logger.Detail).Infof("candidate peer not callable: %#v", r) return po >= proxLimit
return po > proxLimit
} }
p = r p = r
ppo = po
return false return false
}) })
if p != nil { if p != nil {
glog.V(logger.Detail).Infof("candidate prox peer found: %v (%v), %#v", p, ppo, p) glog.V(logger.Detail).Infof("candidate prox peer found: %v (%v), %v", p, ppo, p)
//return p, 0, false
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) 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 return size > 0 && po < proxLimit
}) })
// all buckets are full // all buckets are full
// minsize == self.BucketSize // minsize == self.MinBinSize
if len(bpo) == 0 { if len(bpo) == 0 {
return nil, 0, false return nil, 0, false
} }
@ -381,8 +420,8 @@ func (self *Kademlia) String() string {
var rows []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("%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)) 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) liverows := make([]string, self.MaxProxDisplay)
@ -393,6 +432,9 @@ func (self *Kademlia) String() string {
rest := self.conns.Size() 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 { self.conns.EachBin(self.addr, 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
var rowlen int var rowlen int
if po >= self.MaxProxDisplay {
po = self.MaxProxDisplay - 1
}
row := []string{fmt.Sprintf("%2d", size)} row := []string{fmt.Sprintf("%2d", size)}
rest -= size rest -= size
f(func(val pot.PotVal, vpo int) bool { 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 { self.peers.EachBin(self.addr, 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
var rowlen int var rowlen int
if po >= self.MaxProxDisplay {
po = self.MaxProxDisplay - 1
}
if size < 0 {
panic("wtf")
}
row := []string{fmt.Sprintf("%2d", size)} row := []string{fmt.Sprintf("%2d", size)}
f(func(val pot.PotVal, vpo int) bool { f(func(val pot.PotVal, vpo int) bool {
kp := val.(*KadPeer) kp := val.(*KadPeer)
@ -424,6 +472,10 @@ func (self *Kademlia) String() string {
rowlen++ rowlen++
return rowlen < 4 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, " ") peersrows[po] = strings.Join(row, " ")
return true return true
}) })

View file

@ -17,11 +17,12 @@ package network
import ( import (
"fmt" "fmt"
"sync"
"testing" "testing"
"time" "time"
// "github.com/ethereum/go-ethereum/logger" // "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" "github.com/ethereum/go-ethereum/pot"
) )
@ -30,25 +31,61 @@ func testKadPeerAddr(s string) *peerAddr {
return &peerAddr{OAddr: a, UAddr: a} return &peerAddr{OAddr: a, UAddr: a}
} }
func testKadPeer(s string) *bzzPeer { type testDropPeer struct {
return &bzzPeer{peerAddr: testKadPeerAddr(s)} Peer
dropc chan error
} }
func testStr(a PeerAddr) string { type testDiscPeer struct {
s, _ := a.(*KadPeer) *testDropPeer
if s == nil { lock *sync.Mutex
return "<nil>" notifications map[string]uint8
} }
// return s.String()
// return fmt.Sprintf("%06x", s.OverlayAddr()) type testPeerNotification struct {
return pot.NewHashAddressFromBytes(a.(*KadPeer).PeerAddr.OverlayAddr()).Bin()[:6] rec string
// return a.(*KadPeer).String()[:6] //.HashAddress.String()[:6] addr string
// return a.(*KadPeer).String() //[:6] //.HashAddress.String()[:6] po uint8
// return "wtf" }
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 { type testKademlia struct {
*Kademlia *Kademlia
Discovery bool
dropc chan error
lock *sync.Mutex
notifications map[string]uint8
} }
func newTestKademlia(b string) *testKademlia { func newTestKademlia(b string) *testKademlia {
@ -56,20 +93,63 @@ func newTestKademlia(b string) *testKademlia {
params.MinBinSize = 1 params.MinBinSize = 1
params.MinProxBinSize = 2 params.MinProxBinSize = 2
base := pot.NewHashAddress(b).Bytes() 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 "<nil>"
// }
// 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 "<nil>"
// }
// return pot.NewHashAddressFromBytes(p.OverlayAddr()).Bin()[:6]
if a == nil {
return "<nil>"
}
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 { func (k *testKademlia) On(ons ...string) *testKademlia {
for _, s := range ons { for _, s := range ons {
k.Kademlia.On(testKadPeer(s)) p := k.newTestKadPeer(s)
k.Kademlia.On(p)
} }
return k return k
} }
func (k *testKademlia) Off(offs ...string) *testKademlia { func (k *testKademlia) Off(offs ...string) *testKademlia {
for _, s := range offs { for _, s := range offs {
k.Kademlia.Off(testKadPeer(s)) k.Kademlia.Off(k.newTestKadPeer(s))
} }
return k 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 { func testSuggestPeer(t *testing.T, k *testKademlia, expAddr string, expPo int, expWant bool) error {
addr, o, want := k.SuggestPeer() addr, o, want := k.SuggestPeer()
if testStr(addr) != expAddr { if overlayStr(addr) != expAddr {
return fmt.Errorf("incorrect peer address suggested. expected %v, got %v", expAddr, testStr(addr)) return fmt.Errorf("incorrect peer address suggested. expected %v, got %v", expAddr, overlayStr(addr))
} }
if o != expPo { if o != expPo {
return fmt.Errorf("incorrect prox order suggested. expected %v, got %v", expPo, o) 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) { func TestSuggestPeerFindPeers(t *testing.T) {
// 2 row gap, unsaturated proxbin, no callables -> want PO 0 // 2 row gap, unsaturated proxbin, no callables -> want PO 0
k := newTestKademlia("000000").On("001000") k := newTestKademlia("000000").On("001000")
// k.MinProxBinSize = 2
// k.MinBinSize = 2
err := testSuggestPeer(t, k, "<nil>", 0, true) err := testSuggestPeer(t, k, "<nil>", 0, true)
if err != nil { if err != nil {
t.Fatal(err.Error()) 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) { func TestKademliaHiveString(t *testing.T) {
k := newTestKademlia("000000").On("010000", "001000").Register("100000", "100001") k := newTestKademlia("000000").On("010000", "001000").Register("100000", "100001")
h := k.String() 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:] { if expH[100:] != h[100:] {
t.Fatalf("incorrect hive output. expected %v, got %v", expH, h) 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},
},
)
}

View file

View file

@ -27,6 +27,7 @@ import (
"github.com/ethereum/go-ethereum/p2p/adapters" "github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/protocols" "github.com/ethereum/go-ethereum/p2p/protocols"
"github.com/ethereum/go-ethereum/pot"
) )
const ( const (
@ -54,12 +55,14 @@ func (self *bzzPeer) LastActive() time.Time {
type PeerAddr interface { type PeerAddr interface {
OverlayAddr() []byte OverlayAddr() []byte
UnderlayAddr() []byte UnderlayAddr() []byte
PO(pot.PotVal, int) (int, bool)
String() string
} }
// the Peer interface that peerPool needs // the Peer interface that peerPool needs
type Peer interface { type Peer interface {
PeerAddr 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 ID() discover.NodeID // the key that uniquely identifies the Node for the peerPool
Send(interface{}) error // can send messages Send(interface{}) error // can send messages
Drop(error) // disconnect this peer Drop(error) // disconnect this peer
@ -139,6 +142,35 @@ func (self *peerAddr) UnderlayAddr() []byte {
return self.UAddr 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 { func (self *peerAddr) String() string {
return fmt.Sprintf("%x <%x>", self.OAddr, self.UAddr) return fmt.Sprintf("%x <%x>", self.OAddr, self.UAddr)
} }

View file

@ -125,6 +125,8 @@ func TestBzzHandshakeNetworkIdMismatch(t *testing.T) {
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
addr := RandomAddr() addr := RandomAddr()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
id := s.Ids[0] id := s.Ids[0]
s.testHandshake( s.testHandshake(
correctBzzHandshake(addr), correctBzzHandshake(addr),
@ -137,6 +139,8 @@ func TestBzzHandshakeVersionMismatch(t *testing.T) {
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
addr := RandomAddr() addr := RandomAddr()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
id := s.Ids[0] id := s.Ids[0]
s.testHandshake( s.testHandshake(
correctBzzHandshake(addr), correctBzzHandshake(addr),
@ -149,6 +153,8 @@ func TestBzzHandshakeSuccess(t *testing.T) {
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
addr := RandomAddr() addr := RandomAddr()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
id := s.Ids[0] id := s.Ids[0]
s.testHandshake( s.testHandshake(
correctBzzHandshake(addr), correctBzzHandshake(addr),
@ -160,6 +166,7 @@ func TestBzzPeerPoolAdd(t *testing.T) {
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
addr := RandomAddr() addr := RandomAddr()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
id := s.Ids[0] id := s.Ids[0]
glog.V(logger.Detail).Infof("handshake with %v", id) glog.V(logger.Detail).Infof("handshake with %v", id)
@ -174,6 +181,8 @@ func TestBzzPeerPoolRemove(t *testing.T) {
addr := RandomAddr() addr := RandomAddr()
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
s.runHandshakes() s.runHandshakes()
id := s.Ids[0] id := s.Ids[0]
@ -188,6 +197,8 @@ func TestBzzPeerPoolBothAddRemove(t *testing.T) {
addr := RandomAddr() addr := RandomAddr()
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
s.runHandshakes() s.runHandshakes()
id := s.Ids[0] id := s.Ids[0]
@ -206,6 +217,7 @@ func TestBzzPeerPoolNotAdd(t *testing.T) {
addr := RandomAddr() addr := RandomAddr()
pp := p2ptest.NewTestPeerPool() pp := p2ptest.NewTestPeerPool()
s := newBzzTester(t, 1, addr, pp, nil, nil) s := newBzzTester(t, 1, addr, pp, nil, nil)
defer s.Stop()
id := s.Ids[0] 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)")}) s.testHandshake(correctBzzHandshake(addr), &bzzHandshake{0, 321, NewPeerAddrFromNodeId(id)}, &p2ptest.Disconnect{Peer: id, Error: fmt.Errorf("network id mismatch 321 (!= 322)")})

View file

@ -8,40 +8,34 @@ import (
//"github.com/ethereum/go-ethereum/logger" //"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/p2p/adapters" "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/protocols"
"github.com/ethereum/go-ethereum/p2p/simulations"
p2ptest "github.com/ethereum/go-ethereum/p2p/testing" p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
) )
type pssTester struct { type pssTester struct {
*p2ptest.ProtocolTester *p2ptest.ProtocolTester
ct *protocols.CodeMap ct *protocols.CodeMap
*Pss *Pss
} }
func TestPssTwoToSelf(t *testing.T) { func TestPssTwoToSelf(t *testing.T) {
addr := RandomAddr() addr := RandomAddr()
pt := newPssTester(t, addr, 2) pt := newPssTester(t, addr, 2)
defer pt.Stop()
payload := []byte("foo42") payload := []byte("foo42")
subpeermsgcode, found := pt.ct.GetCode(&SubPeersMsg{}) subpeermsgcode, found := pt.ct.GetCode(&subPeersMsg{})
if !found { if !found {
t.Fatalf("peerMsg not defined") t.Fatalf("peerMsg not defined")
} }
/*
peersmsgcode, found := pt.ct.GetCode(&peersMsg{})
if !found {
t.Fatalf("PssMsg not defined")
}
*/
pssmsgcode, found := pt.ct.GetCode(&PssMsg{}) pssmsgcode, found := pt.ct.GetCode(&PssMsg{})
if !found { if !found {
t.Fatalf("PssMsg not defined") t.Fatalf("PssMsg not defined")
} }
hs_pivot := correctBzzHandshake(addr) hs_pivot := correctBzzHandshake(addr)
for _, id := range pt.Ids { for _, id := range pt.Ids {
hs_sim := correctBzzHandshake(NewPeerAddrFromNodeId(id)) hs_sim := correctBzzHandshake(NewPeerAddrFromNodeId(id))
<-pt.GetPeer(id).Connc <-pt.GetPeer(id).Connc
@ -49,37 +43,30 @@ func TestPssTwoToSelf(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("Handshake fail: %v", err) t.Fatalf("Handshake fail: %v", err)
} }
err = pt.TestExchanges( err = pt.TestExchanges(
p2ptest.Exchange{ p2ptest.Exchange{
Expects: []p2ptest.Expect{ Expects: []p2ptest.Expect{
p2ptest.Expect{ p2ptest.Expect{
Code: subpeermsgcode, Code: subpeermsgcode,
Msg: &SubPeersMsg{}, Msg: &subPeersMsg{},
Peer: id, Peer: id,
}, },
}, },
/*Triggers: []p2ptest.Trigger{
p2ptest.Trigger{
Code: peersmsgcode,
Msg: &peersMsg{},
Peer: id,
},
},*/
}, },
) )
if err != nil { if err != nil {
t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err) t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err)
} }
} }
err := pt.TestExchanges ( err := pt.TestExchanges(
p2ptest.Exchange{ p2ptest.Exchange{
Triggers: []p2ptest.Trigger{ Triggers: []p2ptest.Trigger{
p2ptest.Trigger{ p2ptest.Trigger{
Code: pssmsgcode, Code: pssmsgcode,
Msg: &PssMsg{ Msg: &PssMsg{
To: addr.OverlayAddr(), To: addr.OverlayAddr(),
Data: payload, Data: payload,
}, },
Peer: pt.Ids[0], Peer: pt.Ids[0],
@ -89,42 +76,36 @@ func TestPssTwoToSelf(t *testing.T) {
) )
if err != nil { if err != nil {
t.Fatalf("PssMsg sending %v to %v (pivot) fail: %v", pt.Ids[0], addr.OverlayAddr(), err) t.Fatalf("PssMsg sending %v to %v (pivot) fail: %v", pt.Ids[0], addr.OverlayAddr(), err)
} }
alarm := time.NewTimer(1000 * time.Millisecond) alarm := time.NewTimer(1000 * time.Millisecond)
select { select {
case data := <-pt.C: case data := <-pt.C:
if !bytes.Equal(data, payload) { if !bytes.Equal(data, payload) {
t.Fatalf("Data transfer failed, expected: %v, got: %v", payload, data) t.Fatalf("Data transfer failed, expected: %v, got: %v", payload, data)
} }
case <-alarm.C: case <-alarm.C:
t.Fatalf("Pivot receive of PssMsg from %v timeout", pt.Ids[0]) t.Fatalf("Pivot receive of PssMsg from %v timeout", pt.Ids[0])
} }
} }
func TestPssTwoRelaySelf(t *testing.T) { func TestPssTwoRelaySelf(t *testing.T) {
// <-(chan bool)(nil)
addr := RandomAddr() addr := RandomAddr()
pt := newPssTester(t, addr, 2) pt := newPssTester(t, addr, 2)
defer pt.Stop()
subpeermsgcode, found := pt.ct.GetCode(&SubPeersMsg{}) subpeermsgcode, found := pt.ct.GetCode(&subPeersMsg{})
if !found { if !found {
t.Fatalf("peerMsg not defined") t.Fatalf("peerMsg not defined")
} }
/*
peersmsgcode, found := pt.ct.GetCode(&peersMsg{})
if !found {
t.Fatalf("PssMsg not defined")
}
*/
pssmsgcode, found := pt.ct.GetCode(&PssMsg{}) pssmsgcode, found := pt.ct.GetCode(&PssMsg{})
if !found { if !found {
t.Fatalf("PssMsg not defined") t.Fatalf("PssMsg not defined")
} }
hs_pivot := correctBzzHandshake(addr) hs_pivot := correctBzzHandshake(addr)
for _, id := range pt.Ids { for _, id := range pt.Ids {
hs_sim := correctBzzHandshake(NewPeerAddrFromNodeId(id)) hs_sim := correctBzzHandshake(NewPeerAddrFromNodeId(id))
<-pt.GetPeer(id).Connc <-pt.GetPeer(id).Connc
@ -132,37 +113,30 @@ func TestPssTwoRelaySelf(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("Handshake fail: %v", err) t.Fatalf("Handshake fail: %v", err)
} }
err = pt.TestExchanges( err = pt.TestExchanges(
p2ptest.Exchange{ p2ptest.Exchange{
Expects: []p2ptest.Expect{ Expects: []p2ptest.Expect{
p2ptest.Expect{ p2ptest.Expect{
Code: subpeermsgcode, Code: subpeermsgcode,
Msg: &SubPeersMsg{}, Msg: &subPeersMsg{},
Peer: id, Peer: id,
}, },
}, },
/*Triggers: []p2ptest.Trigger{
p2ptest.Trigger{
Code: peersmsgcode,
Msg: &peersMsg{},
Peer: id,
},
},*/
}, },
) )
if err != nil { if err != nil {
t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err) t.Fatalf("Subpeersmsg to peer %v fail: %v", id, err)
} }
} }
err := pt.TestExchanges ( err := pt.TestExchanges(
p2ptest.Exchange{ p2ptest.Exchange{
Expects: []p2ptest.Expect{ Expects: []p2ptest.Expect{
p2ptest.Expect{ p2ptest.Expect{
Code: pssmsgcode, Code: pssmsgcode,
Msg: &PssMsg{ Msg: &PssMsg{
To: pt.Ids[0].Bytes(), To: pt.Ids[0].Bytes(),
Data: []byte("foo42"), Data: []byte("foo42"),
}, },
Peer: pt.Ids[0], Peer: pt.Ids[0],
@ -171,8 +145,8 @@ func TestPssTwoRelaySelf(t *testing.T) {
Triggers: []p2ptest.Trigger{ Triggers: []p2ptest.Trigger{
p2ptest.Trigger{ p2ptest.Trigger{
Code: pssmsgcode, Code: pssmsgcode,
Msg: &PssMsg{ Msg: &PssMsg{
To: pt.Ids[0].Bytes(), To: pt.Ids[0].Bytes(),
Data: []byte("foo42"), Data: []byte("foo42"),
}, },
Peer: pt.Ids[1], Peer: pt.Ids[1],
@ -182,7 +156,7 @@ func TestPssTwoRelaySelf(t *testing.T) {
) )
if err != nil { if err != nil {
t.Fatalf("PssMsg routing from %v to %v fail: %v", pt.Ids[0], pt.Ids[1], err) 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 { 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(&PssMsg{})
ct.Register(&peersMsg{}) ct.Register(&peersMsg{})
ct.Register(&getPeersMsg{}) ct.Register(&getPeersMsg{})
ct.Register(&SubPeersMsg{}) // why is this public? ct.Register(&subPeersMsg{}) // why is this public?
simPipe := adapters.NewSimPipe simPipe := adapters.NewSimPipe
kp := NewKadParams() kp := NewKadParams()
kp.MinProxBinSize = 3 kp.MinProxBinSize = 3
@ -227,8 +200,8 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester {
return &pssTester{ return &pssTester{
ProtocolTester: s, ProtocolTester: s,
ct: ct, ct: ct,
Pss: ps, Pss: ps,
} }
} }

View file

@ -7,6 +7,7 @@ package main
import ( import (
"math/rand" "math/rand"
"reflect"
"runtime" "runtime"
"time" "time"
@ -15,7 +16,6 @@ import (
"github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters" "github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations" "github.com/ethereum/go-ethereum/p2p/simulations"
//p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
"github.com/ethereum/go-ethereum/swarm/network" "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) na := adapters.NewSimNode(id, self.Network, adapters.NewSimPipe)
addr := network.NewPeerAddrFromNodeId(id) addr := network.NewPeerAddrFromNodeId(id)
kp := network.NewKadParams() 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.NewKademlia(addr.OverlayAddr(), kp) // overlay topology driver
// to := network.NewTestOverlay(addr.OverlayAddr()) // overlay topology driver // to := network.NewTestOverlay(addr.OverlayAddr()) // overlay topology driver
hp := network.NewHiveParams() hp := network.NewHiveParams()
hp.CallInterval = 5000 hp.CallInterval = 5000
pp := network.NewHive(hp, to) // hive pp := network.NewHive(hp, to) // hive
self.hives = append(self.hives, pp) // remember hive self.hives = append(self.hives, pp) // remember hive
// bzz protocol Run function. messaging through SimPipe // bzz protocol Run function. messaging through SimPipe
services := func(p network.Peer) error { services := func(p network.Peer) error {
dp := network.NewDiscovery(p, to) dp := network.NewDiscovery(p, to)
pp.Add(dp) 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) { p.DisconnectHook(func(err error) {
pp.Remove(dp) pp.Remove(dp)
}) })
return nil 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 na.Run = network.Bzz(addr.OverlayAddr(), na, ct, services, nil, nil).Run
connect := func(s string) error { connect := func(s string) error {
return self.Connect(id, adapters.NewNodeIdFromHex(s)) 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 { func NewNetwork(network *simulations.Network) *Network {
n := &Network{ n := &Network{
// hives:
Network: network, Network: network,
} }
n.SetNaf(n.NewSimNode) n.SetNaf(n.NewSimNode)
@ -104,7 +109,7 @@ func NewNetwork(network *simulations.Network) *Network {
func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simulations.ResourceController) { func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simulations.ResourceController) {
conf.DefaultMockerConfig = simulations.DefaultMockerConfig() conf.DefaultMockerConfig = simulations.DefaultMockerConfig()
conf.DefaultMockerConfig.SwitchonRate = 100 conf.DefaultMockerConfig.SwitchonRate = 100
// conf.DefaultMockerConfig.NodesTarget = 15 conf.DefaultMockerConfig.NodesTarget = 15
conf.DefaultMockerConfig.NewConnCount = 1 conf.DefaultMockerConfig.NewConnCount = 1
conf.DefaultMockerConfig.DegreeTarget = 0 conf.DefaultMockerConfig.DegreeTarget = 0
conf.Id = "0" conf.Id = "0"
@ -112,7 +117,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
net := NewNetwork(simulations.NewNetwork(conf)) net := NewNetwork(simulations.NewNetwork(conf))
//ids := p2ptest.RandomNodeIds(10) //ids := p2ptest.RandomNodeIds(10)
ids := adapters.RandomNodeIds(10) ids := adapters.RandomNodeIds(3)
for i, id := range ids { for i, id := range ids {
net.NewNode(&simulations.NodeConfig{Id: id}) net.NewNode(&simulations.NodeConfig{Id: id})
@ -129,23 +134,45 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
} }
go func() { go func() {
for _, id := range ids { for _, id := range ids {
n := rand.Intn(1000)
time.Sleep(time.Duration(n) * time.Millisecond)
net.NewNode(&simulations.NodeConfig{Id: id}) net.NewNode(&simulations.NodeConfig{Id: id})
net.Start(id) net.Start(id)
glog.V(logger.Debug).Infof("node %v starting up", id) glog.V(logger.Debug).Infof("node %v starting up", id)
n := rand.Intn(1000) // time.Sleep(1000 * time.Millisecond)
time.Sleep(time.Duration(n) * time.Millisecond)
// net.Stop(id) // net.Stop(id)
} }
}() }()
// go func() { // for i, id := range ids {
// for _, id := range ids { // n := 3000 + i*1000
// net.Stop(id) // go func() {
// glog.V(logger.Debug).Infof("node %v shutting down", id) // for {
// n := rand.Intn(500) // // n := rand.Intn(5000)
// time.Sleep(time.Duration(n) * time.Millisecond) // // 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
// }
// }()
// } // }
// }() nodes := simulations.NewResourceContoller(
return net, nil &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 // var server