mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 07:06:42 +00:00
swarm/network: WIP kademlia network sim tests
This commit is contained in:
parent
e26d672438
commit
11aeaec420
4 changed files with 32 additions and 11 deletions
|
|
@ -271,7 +271,7 @@ func (self *Peer) Send(msg interface{}) error {
|
||||||
if !found {
|
if !found {
|
||||||
return errorf(ErrInvalidMsgType, "%v", code)
|
return errorf(ErrInvalidMsgType, "%v", code)
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("=> %v (%d)", msg, code)
|
glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg)
|
||||||
err := self.m.SendMsg(uint64(code), msg)
|
err := self.m.SendMsg(uint64(code), msg)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
err = errorf(ErrWrite, "(msg code: %v): %v", code, err)
|
err = errorf(ErrWrite, "(msg code: %v): %v", code, err)
|
||||||
|
|
@ -321,7 +321,7 @@ func (self *Peer) handleIncoming() (interface{}, error) {
|
||||||
if err := msg.Decode(val.Interface()); err != nil {
|
if err := msg.Decode(val.Interface()); err != nil {
|
||||||
return nil, errorf(ErrDecode, "<= %v: %v", msg, err)
|
return nil, errorf(ErrDecode, "<= %v: %v", msg, err)
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("<= %v %v (%d)", req, typ, msg.Code)
|
glog.V(logger.Detail).Infof("<= %v FROM %v %v %v", msg, self.ID(), req, typ)
|
||||||
|
|
||||||
// call the registered handler callbacks
|
// call the registered handler callbacks
|
||||||
// a registered callback take the decoded message as argument as an interface
|
// a registered callback take the decoded message as argument as an interface
|
||||||
|
|
|
||||||
|
|
@ -26,7 +26,8 @@ func (self *discPeer) NotifyPeer(p Peer, po uint8) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
resp := &peersMsg{
|
resp := &peersMsg{
|
||||||
Peers: []*peerAddr{p.(*discPeer).Peer.(*bzzPeer).peerAddr},
|
//Peers: []*peerAddr{p.(*discPeer).Peer.(*bzzPeer).peerAddr},
|
||||||
|
Peers: []*peerAddr{&peerAddr{OAddr: p.OverlayAddr(), UAddr: p.UnderlayAddr()}}, // perhaps the PeerAddr interface is unnecessary generalization
|
||||||
}
|
}
|
||||||
return p.Send(resp)
|
return p.Send(resp)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -73,7 +73,8 @@ func NewKadParams() *KadParams {
|
||||||
MinProxBinSize: 2,
|
MinProxBinSize: 2,
|
||||||
MinBinSize: 2,
|
MinBinSize: 2,
|
||||||
MaxBinSize: 4,
|
MaxBinSize: 4,
|
||||||
RetryInterval: 42000000000,
|
//RetryInterval: 42000000000,
|
||||||
|
RetryInterval: 420000000,
|
||||||
MaxRetries: 42,
|
MaxRetries: 42,
|
||||||
RetryExponent: 2,
|
RetryExponent: 2,
|
||||||
}
|
}
|
||||||
|
|
@ -136,6 +137,7 @@ func (self *KadPeer) String() string {
|
||||||
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 not callable", kp.PeerAddr)
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -149,6 +151,7 @@ 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
|
||||||
|
|
@ -313,6 +316,7 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) {
|
||||||
})
|
})
|
||||||
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)
|
||||||
|
|
|
||||||
|
|
@ -50,8 +50,8 @@ func (self *SimNode) Stop() error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *SimNode) RunProtocol(id *adapters.NodeId, rw, rrw p2p.MsgReadWriter, runc chan bool) error {
|
func (self *SimNode) RunProtocol(id *adapters.NodeId, rw, rrw p2p.MsgReadWriter, peer *adapters.Peer) error {
|
||||||
return self.NodeAdapter.(adapters.ProtocolRunner).RunProtocol(id, rw, rrw, runc)
|
return self.NodeAdapter.(adapters.ProtocolRunner).RunProtocol(id, rw, rrw, peer)
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewSimNode creates adapters for nodes in the simulation.
|
// NewSimNode creates adapters for nodes in the simulation.
|
||||||
|
|
@ -59,13 +59,29 @@ func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapt
|
||||||
id := conf.Id
|
id := conf.Id
|
||||||
na := adapters.NewSimNode(id, self.Network, adapters.NewSimPipe)
|
na := adapters.NewSimNode(id, self.Network, adapters.NewSimPipe)
|
||||||
addr := network.NewPeerAddrFromNodeId(id)
|
addr := network.NewPeerAddrFromNodeId(id)
|
||||||
to := network.NewKademlia(addr.OverlayAddr(), nil) // overlay topology driver
|
kp := network.NewKadParams()
|
||||||
|
kp.MinProxBinSize = 4
|
||||||
|
kp.MinBinSize = 4
|
||||||
|
to := network.NewKademlia(addr.OverlayAddr(), kp) // overlay topology driver
|
||||||
// to := network.NewTestOverlay(addr.OverlayAddr()) // overlay topology driver
|
// to := network.NewTestOverlay(addr.OverlayAddr()) // overlay topology driver
|
||||||
pp := network.NewHive(network.NewHiveParams(), to) // hive
|
hp := network.NewHiveParams()
|
||||||
|
hp.CallInterval = 5000
|
||||||
|
pp := network.NewHive(hp, to) // hive
|
||||||
self.hives = append(self.hives, pp) // remember hive
|
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 {
|
||||||
|
dp := network.NewDiscovery(p, to)
|
||||||
|
pp.Add(dp)
|
||||||
|
glog.V(logger.Detail).Infof("kademlia on %v", p)
|
||||||
|
p.DisconnectHook(func(err error) {
|
||||||
|
pp.Remove(dp)
|
||||||
|
})
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
ct := network.BzzCodeMap(network.HiveMsgs...) // bzz protocol code map
|
ct := network.BzzCodeMap(network.HiveMsgs...) // bzz protocol code map
|
||||||
na.Run = network.Bzz(addr.OverlayAddr(), pp, na, ct, nil).Run
|
na.Run = network.Bzz(addr.OverlayAddr(), na, ct, services, nil, nil).Run
|
||||||
connect := func(s string) error {
|
connect := func(s string) error {
|
||||||
return self.Connect(id, adapters.NewNodeIdFromHex(s))
|
return self.Connect(id, adapters.NewNodeIdFromHex(s))
|
||||||
}
|
}
|
||||||
|
|
@ -95,7 +111,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
|
||||||
conf.Backend = true
|
conf.Backend = true
|
||||||
net := NewNetwork(simulations.NewNetwork(conf))
|
net := NewNetwork(simulations.NewNetwork(conf))
|
||||||
|
|
||||||
ids := p2ptest.RandomNodeIds(5)
|
ids := p2ptest.RandomNodeIds(10)
|
||||||
|
|
||||||
for i, id := range ids {
|
for i, id := range ids {
|
||||||
net.NewNode(&simulations.NodeConfig{Id: id})
|
net.NewNode(&simulations.NodeConfig{Id: id})
|
||||||
|
|
@ -134,7 +150,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
|
||||||
// var server
|
// var server
|
||||||
func main() {
|
func main() {
|
||||||
runtime.GOMAXPROCS(runtime.NumCPU())
|
runtime.GOMAXPROCS(runtime.NumCPU())
|
||||||
glog.SetV(logger.Info)
|
glog.SetV(logger.Detail)
|
||||||
glog.SetToStderr(true)
|
glog.SetToStderr(true)
|
||||||
|
|
||||||
c, quitc := simulations.NewSessionController(nethook)
|
c, quitc := simulations.NewSessionController(nethook)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue