p2p/simulations: Refactor

Signed-off-by: Lewis Marshall <lewis@lmars.net>
This commit is contained in:
Lewis Marshall 2017-05-01 19:39:05 +01:00
parent 246eb00def
commit 12eb088e5b
31 changed files with 990 additions and 899 deletions

View file

@ -165,8 +165,6 @@ func (n *Node) Start() error {
if n.serverConfig.NodeDatabase == "" {
n.serverConfig.NodeDatabase = n.config.NodeDB()
}
running := p2p.NewServer(n.serverConfig)
log.Info("Starting peer-to-peer node", "instance", n.serverConfig.Name)
// Otherwise copy and specialize the P2P configuration
services := make(map[reflect.Type]Service)
@ -196,6 +194,8 @@ func (n *Node) Start() error {
for _, service := range services {
n.serverConfig.Protocols = append(n.serverConfig.Protocols, service.Protocols()...)
}
running := p2p.NewServer(n.serverConfig)
log.Info("Starting peer-to-peer node", "instance", n.serverConfig.Name)
if err := running.Start(); err != nil {
if errno, ok := err.(syscall.Errno); ok && datadirInUseErrnos[uint(errno)] {
return ErrDatadirUsed

View file

@ -1,300 +0,0 @@
// Copyright 2016 The go-ethereum Authors
// This file is part of the go-ethereum library.
//
// The go-ethereum library is free software: you can redistribute it and/or modify
// it under the terms of the GNU Lesser General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// The go-ethereum library is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Lesser General Public License for more details.
//
// 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/>.
package adapters
import (
"errors"
"fmt"
"sync"
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rpc"
)
func newPeer(rw MsgReadWriteCloser) *Peer {
return &Peer{
MsgReadWriteCloser: rw,
Errc: make(chan error, 1),
Connc: make(chan bool),
Readyc: make(chan bool),
}
}
type Peer struct {
MsgReadWriteCloser
Connc chan bool
Readyc chan bool
Errc chan error
}
// Network interface to retrieve protocol runner to launch upon peer
// connection
type Network interface {
GetNodeAdapter(id *NodeId) NodeAdapter
Reporter
}
// Adds close pipe to the MsgReadWriter
type MsgReadWriteCloser interface {
p2p.MsgReadWriter
Close() error
}
// SimNode is the network adapter that
type SimNode struct {
lock sync.RWMutex
Id *NodeId
network Network
service node.Service
peerMap map[discover.NodeID]int
peers []*Peer
peerFeed event.Feed
client *rpc.Client
}
func NewSimNode(id *NodeId, svc node.Service, n Network) *SimNode {
// for simplicity, only support single protocol services
if len(svc.Protocols()) != 1 {
panic("service must have a single protocol")
}
return &SimNode{
Id: id,
network: n,
service: svc,
peerMap: make(map[discover.NodeID]int),
}
}
// Addr returns the node's address
func (self *SimNode) Addr() []byte {
return []byte(self.Node().String())
}
func (self *SimNode) Node() *discover.Node {
return discover.NewNode(self.Id.NodeID, nil, 0, 0)
}
func (self *SimNode) Client() (*rpc.Client, error) {
self.lock.Lock()
defer self.lock.Unlock()
if self.client == nil {
return nil, errors.New("RPC not started")
}
return self.client, nil
}
// Start starts the RPC handler and the underlying service
func (self *SimNode) Start() error {
if err := self.startRPC(); err != nil {
return err
}
return self.service.Start(self)
}
// Stop stops the RPC handler and the underlying service
func (self *SimNode) Stop() error {
self.stopRPC()
return self.service.Stop()
}
func (self *SimNode) startRPC() error {
self.lock.Lock()
defer self.lock.Unlock()
if self.client != nil {
return errors.New("RPC already started")
}
// add SimAdminAPI and PeerAPI so that the network can call the
// AddPeer, RemovePeer and PeerEvents RPC methods
apis := append(self.service.APIs(), []rpc.API{
{
Namespace: "admin",
Version: "1.0",
Service: &SimAdminAPI{self},
},
{
Namespace: "eth",
Version: "1.0",
Service: &PeerAPI{func() p2p.Server { return self }},
},
}...)
// start the RPC handler
handler := rpc.NewServer()
for _, api := range apis {
if err := handler.RegisterName(api.Namespace, api.Service); err != nil {
return fmt.Errorf("error registering RPC: %s", err)
}
}
// create an in-process RPC client
self.client = rpc.DialInProc(handler)
return nil
}
func (self *SimNode) stopRPC() {
self.lock.Lock()
defer self.lock.Unlock()
if self.client != nil {
self.client.Close()
self.client = nil
}
}
func (self *SimNode) GetPeer(id *NodeId) *Peer {
self.lock.Lock()
defer self.lock.Unlock()
return self.getPeer(id)
}
func (self *SimNode) getPeer(id *NodeId) *Peer {
i, found := self.peerMap[id.NodeID]
if !found {
return nil
}
return self.peers[i]
}
func (self *SimNode) setPeer(id *NodeId, rw MsgReadWriteCloser) *Peer {
i, found := self.peerMap[id.NodeID]
if !found {
i = len(self.peers)
self.peerMap[id.NodeID] = i
p := newPeer(rw)
self.peers = append(self.peers, p)
return p
}
// if self.peers[i] != nil && m != nil {
// panic(fmt.Sprintf("pipe for %v already set", id))
// }
// legit reconnect reset disconnection error,
p := self.peers[i]
//p.MsgPipeRW = rw
p.MsgReadWriteCloser = rw
p.Connc = make(chan bool)
p.Readyc = make(chan bool)
return p
}
func (self *SimNode) RemovePeer(node *discover.Node) {
self.lock.Lock()
defer self.lock.Unlock()
id := &NodeId{node.ID}
peer := self.getPeer(id)
//if peer == nil || peer.MsgPipeRW == nil {
if peer == nil || peer.MsgReadWriteCloser == nil {
return
}
peer.MsgReadWriteCloser.Close()
peer.MsgReadWriteCloser = nil
// na := self.network.GetNodeAdapter(id)
// peer = na.(*SimNode).GetPeer(self.Id)
// peer.RW = nil
log.Trace(fmt.Sprintf("dropped peer %v", id))
}
func (self *SimNode) AddPeer(node *discover.Node) {
self.lock.Lock()
defer self.lock.Unlock()
id := &NodeId{node.ID}
na := self.network.GetNodeAdapter(id)
if na == nil {
panic(fmt.Sprintf("node adapter for %v is missing", id))
}
rw, rrw := p2p.MsgPipe()
// // run protocol on remote node with self as peer
peer := self.getPeer(id)
//if peer != nil && peer.MsgPipeRW != nil {
if peer != nil && peer.MsgReadWriteCloser != nil {
return
}
peer = self.setPeer(id, rrw)
close(peer.Connc)
defer close(peer.Readyc)
na.(*SimNode).RunProtocol(self, rrw, rw, peer)
// run protocol on remote node with self as peer
self.RunProtocol(na.(*SimNode), rw, rrw, peer)
}
func (self *SimNode) SubscribeEvents(ch chan *p2p.PeerEvent) event.Subscription {
return self.peerFeed.Subscribe(ch)
}
func (self *SimNode) PeerCount() int {
self.lock.Lock()
defer self.lock.Unlock()
return len(self.peers)
}
func (self *SimNode) NodeInfo() *p2p.NodeInfo {
return &p2p.NodeInfo{ID: self.Id.String()}
}
func (self *SimNode) PeersInfo() (info []*p2p.PeerInfo) {
return nil
}
func (self *SimNode) RunProtocol(node *SimNode, rw, rrw MsgReadWriteCloser, peer *Peer) {
id := node.Id
protocol := self.service.Protocols()[0]
if protocol.Run == nil {
log.Trace(fmt.Sprintf("no protocol starting on peer %v (connection with %v)", self.Id, id))
return
}
log.Trace(fmt.Sprintf("protocol starting on peer %v (connection with %v)", self.Id, id))
p := p2p.NewPeer(id.NodeID, id.Label(), []p2p.Cap{})
go func() {
self.peerFeed.Send(&p2p.PeerEvent{Type: p2p.PeerEventTypeAdd, Peer: id.NodeID})
err := protocol.Run(p, rw)
<-peer.Readyc
self.RemovePeer(node.Node())
peer.Errc <- err
log.Trace(fmt.Sprintf("protocol quit on peer %v (connection with %v broken: %v)", self.Id, id, err))
self.peerFeed.Send(&p2p.PeerEvent{Type: p2p.PeerEventTypeDrop, Peer: id.NodeID})
}()
}
// SimAdminAPI implements the AddPeer and RemovePeer RPC methods (API
// compatible with node.PrivateAdminAPI)
type SimAdminAPI struct {
*SimNode
}
func (api *SimAdminAPI) AddPeer(url string) (bool, error) {
node, err := discover.ParseNode(url)
if err != nil {
return false, fmt.Errorf("invalid enode: %v", err)
}
api.SimNode.AddPeer(node)
return true, nil
}
func (api *SimAdminAPI) RemovePeer(url string) (bool, error) {
node, err := discover.ParseNode(url)
if err != nil {
return false, fmt.Errorf("invalid enode: %v", err)
}
api.SimNode.RemovePeer(node)
return true, nil
}

View file

@ -1,44 +0,0 @@
// Copyright 2016 The go-ethereum Authors
// This file is part of the go-ethereum library.
//
// The go-ethereum library is free software: you can redistribute it and/or modify
// it under the terms of the GNU Lesser General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// The go-ethereum library is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Lesser General Public License for more details.
//
// 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/>.
package adapters
import (
// "fmt"
// "net"
// "github.com/ethereum/go-ethereum/p2p"
// "github.com/ethereum/go-ethereum/p2p/discover"
)
type RemoteReporter struct {
}
func NewRemoteReporter(url string) *RemoteReporter {
return &RemoteReporter{}
}
func (self *RemoteReporter) DidConnect(source, target *NodeId) {
self.post(true)
}
func (self *RemoteReporter) DidDisconnect(source, target *NodeId) {
self.post(true)
}
func (self *RemoteReporter) post(interface{}) {
}

View file

@ -1,80 +0,0 @@
// Copyright 2016 The go-ethereum Authors
// This file is part of the go-ethereum library.
//
// The go-ethereum library is free software: you can redistribute it and/or modify
// it under the terms of the GNU Lesser General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// The go-ethereum library is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Lesser General Public License for more details.
//
// 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/>.
package adapters
import (
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rpc"
)
const lablen = 4
type NodeId struct {
discover.NodeID
}
func NewNodeId(id []byte) *NodeId {
var n discover.NodeID
copy(n[:], id)
return &NodeId{n}
}
func NewNodeIdFromHex(s string) *NodeId {
id := discover.MustHexID(s)
return &NodeId{id}
}
type ProtoCall func(*p2p.Peer, p2p.MsgReadWriter) error
func (self *NodeId) Bytes() []byte {
return self.NodeID[:]
}
func (self *NodeId) MarshalJSON() (out []byte, err error) {
return []byte(`"` + self.String() + `"`), nil
}
func (self *NodeId) UnmarshalJSON(value []byte) error {
s := string(value)
h, err := discover.HexID(s[1 : len(s)-1])
if err != nil {
return err
}
*self = NodeId{h}
return nil
}
func (self *NodeId) Label() string {
return self.String()[:lablen]
}
type NodeAdapter interface {
Addr() []byte
Client() (*rpc.Client, error)
Start() error
Stop() error
}
type ProtocolRunner interface {
RunProtocol(id *NodeId, rw, rrw p2p.MsgReadWriter, p *Peer) error
}
type Reporter interface {
DidConnect(*NodeId, *NodeId) error
DidDisconnect(*NodeId, *NodeId) error
}

View file

@ -71,12 +71,12 @@ const (
// PeerEventTypeDrop is the type of event emitted when a peer is
// dropped from a p2p.Server
PeerEventTypeDrop PeerEventType = "drop"
// PeerEventTypeMsgSend is the type of event emitted when a
// PeerEventTypeMsgSend is the type of event emitted when a
// message is successfully sent to a peer
PeerEventTypeMsgSend PeerEventType = "msgsend"
// PeerEventTypeMsgSend is the type of event emitted when a
// PeerEventTypeMsgSend is the type of event emitted when a
// message is successfully sent to a peer
PeerEventTypeMsgRecv PeerEventType = "msgrecv"
)
@ -84,8 +84,9 @@ const (
// PeerEvent is an event emitted when peers are either added or dropped from
// a p2p.Server
type PeerEvent struct {
Type PeerEventType
Peer discover.NodeID
Type PeerEventType
Peer discover.NodeID
Error string
Label string
}

View file

@ -8,8 +8,7 @@ import (
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
)
@ -56,7 +55,7 @@ const networkId = "420"
// newProtocol sets up a protocol
// the run function here demonstrates a typical protocol using peerPool, handshake
// and messages registered to handlers
func newProtocol(pp *p2ptest.TestPeerPool) adapters.ProtoCall {
func newProtocol(pp *p2ptest.TestPeerPool) adapters.RunProtocol {
ct := NewCodeMap("test", 42, 1024, &protoHandshake{}, &hs0{}, &kill{}, &drop{})
return func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
peer := NewPeer(p, ct, rw)
@ -118,7 +117,7 @@ func newProtocol(pp *p2ptest.TestPeerPool) adapters.ProtoCall {
}
func protocolTester(t *testing.T, pp *p2ptest.TestPeerPool) *p2ptest.ProtocolTester {
conf := simulations.RandomNodeConfig()
conf := adapters.RandomNodeConfig()
return p2ptest.NewProtocolTester(t, conf.Id, 2, newProtocol(pp))
}
@ -151,12 +150,16 @@ func runProtoHandshake(t *testing.T, proto *protoHandshake, errs ...error) {
s := protocolTester(t, pp)
// TODO: make this more than one handshake
id := s.Ids[0]
s.TestExchanges(protoHandshakeExchange(id, proto)...)
if err := s.TestExchanges(protoHandshakeExchange(id, proto)...); err != nil {
t.Fatal(err)
}
var disconnects []*p2ptest.Disconnect
for i, err := range errs {
disconnects = append(disconnects, &p2ptest.Disconnect{Peer: s.Ids[i], Error: err})
}
s.TestDisconnected(disconnects...)
if err := s.TestDisconnected(disconnects...); err != nil {
t.Fatal(err)
}
}
func TestProtoHandshakeVersionMismatch(t *testing.T) {

View file

@ -779,11 +779,22 @@ func (srv *server) runPeer(p *Peer) {
srv.newPeerHook(p)
}
// broadcast peer add / drop events
srv.peerFeed.Send(&PeerEvent{Type: PeerEventTypeAdd, Peer: p.ID()})
defer srv.peerFeed.Send(&PeerEvent{Type: PeerEventTypeDrop, Peer: p.ID()})
// broadcast peer add
srv.peerFeed.Send(&PeerEvent{
Type: PeerEventTypeAdd,
Peer: p.ID(),
})
// run the protocol
remoteRequested, err := p.run()
// broadcast peer drop
srv.peerFeed.Send(&PeerEvent{
Type: PeerEventTypeDrop,
Peer: p.ID(),
Error: err.Error(),
})
// Note: run waits for existing peers to be sent on srv.delpeer
// before returning, so this send should not select on srv.quit.
srv.delpeer <- peerDrop{p, err, remoteRequested}

View file

@ -1,7 +1,7 @@
package adapters
import (
"crypto/ecdsa"
"errors"
"fmt"
"io"
"io/ioutil"
@ -9,36 +9,41 @@ import (
"os/exec"
"path/filepath"
"runtime"
"sync"
"github.com/docker/docker/pkg/reexec"
"github.com/ethereum/go-ethereum/node"
)
// DockerNode is a NodeAdapter which wraps an ExecNode but exec's the current
// binary in a docker container rather than locally
type DockerNode struct {
ExecNode
// DockerAdapter is a NodeAdapter which runs nodes inside Docker containers.
//
// A Docker image is built which contains the current binary at /bin/p2p-node
// which when executed runs the underlying service (see the description
// of the execP2PNode function for more details)
type DockerAdapter struct{}
// NewDockerAdapter builds the p2p-node Docker image containing the current
// binary and returns a DockerAdapter
func NewDockerAdapter() (*DockerAdapter, error) {
if runtime.GOOS != "linux" {
return nil, errors.New("DockerAdapter can only be used on Linux as it uses the current binary (which must be a Linux binary)")
}
if err := buildDockerImage(); err != nil {
return nil, err
}
return &DockerAdapter{}, nil
}
// NewDockerNode creates a new DockerNode, building the docker image if
// necessary
func NewDockerNode(id *NodeId, key *ecdsa.PrivateKey, service string) (*DockerNode, error) {
if runtime.GOOS != "linux" {
return nil, fmt.Errorf("NewDockerNode can only be used on Linux as it uses the current binary (which must be a Linux binary)")
}
// Name returns the name of the adapter for logging purpoeses
func (d *DockerAdapter) Name() string {
return "docker-adapter"
}
if _, exists := serviceFuncs[service]; !exists {
return nil, fmt.Errorf("unknown node service %q", service)
}
// build the docker image
var err error
dockerOnce.Do(func() {
err = buildDockerImage()
})
if err != nil {
return nil, err
// NewNode returns a new DockerNode using the given config
func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) {
if _, exists := serviceFuncs[config.Service]; !exists {
return nil, fmt.Errorf("unknown node service %q", config.Service)
}
// generate the config
@ -49,16 +54,22 @@ func NewDockerNode(id *NodeId, key *ecdsa.PrivateKey, service string) (*DockerNo
node := &DockerNode{
ExecNode: ExecNode{
ID: id,
Service: service,
ID: config.Id,
Service: config.Service,
Config: &conf,
key: key,
key: config.PrivateKey,
},
}
node.newCmd = node.dockerCommand
return node, nil
}
// DockerNode wraps an ExecNode but exec's the current binary in a docker
// container rather than locally
type DockerNode struct {
ExecNode
}
// dockerCommand returns a command which exec's the binary in a docker
// container.
//
@ -77,9 +88,6 @@ func (n *DockerNode) dockerCommand() *exec.Cmd {
// dockerImage is the name of the docker image
const dockerImage = "p2p-node"
// dockerOnce is used to build the docker image only once
var dockerOnce sync.Once
// buildDockerImage builds the docker image which is used to run devp2p nodes
// using docker.
//

View file

@ -25,24 +25,63 @@ import (
"github.com/ethereum/go-ethereum/rpc"
)
// serviceFunc returns a node.Service which can be used to boot devp2p nodes
type serviceFunc func(id *NodeId) node.Service
// ExecAdapter is a NodeAdapter which runs nodes by executing the current
// binary as a child process.
//
// An init hook is used so that the child process executes the node service
// (rather than whataver the main() function would normally do), see the
// execP2PNode function for more information.
type ExecAdapter struct {
BaseDir string
}
// serviceFuncs is a map of registered services which are used to boot devp2p
// nodes
var serviceFuncs = make(map[string]serviceFunc)
// NewExecAdapter returns an ExecAdapter which stores node data in
// subdirectories of the given base directory
func NewExecAdapter(baseDir string) *ExecAdapter {
return &ExecAdapter{BaseDir: baseDir}
}
// RegisterService registers the given serviceFunc which can then be used to
// start a devp2p node with the given name
func RegisterService(name string, f serviceFunc) {
if _, exists := serviceFuncs[name]; exists {
panic(fmt.Sprintf("node service already exists: %q", name))
// Name returns the name of the adapter for logging purpoeses
func (e *ExecAdapter) Name() string {
return "exec-adapter"
}
// NewNode returns a new ExecNode using the given config
func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) {
if _, exists := serviceFuncs[config.Service]; !exists {
return nil, fmt.Errorf("unknown node service %q", config.Service)
}
serviceFuncs[name] = f
// create the node directory using the first 12 characters of the ID
// as Unix socket paths cannot be longer than 256 characters
dir := filepath.Join(e.BaseDir, config.Id.String()[:12])
if err := os.Mkdir(dir, 0755); err != nil {
return nil, fmt.Errorf("error creating node directory: %s", err)
}
// generate the config
conf := node.DefaultConfig
conf.DataDir = filepath.Join(dir, "data")
conf.P2P.NoDiscovery = true
conf.P2P.NAT = nil
// listen on a random localhost port (we'll get the actual port after
// starting the node through the RPC admin.nodeInfo method)
conf.P2P.ListenAddr = "127.0.0.1:0"
node := &ExecNode{
ID: config.Id,
Service: config.Service,
Dir: dir,
Config: &conf,
key: config.PrivateKey,
}
node.newCmd = node.execCommand
return node, nil
}
// ExecNode is a NodeAdapter which starts the node by exec'ing the current
// binary and running a registered serviceFunc.
// binary and running a registered ServiceFunc.
//
// Communication with the node is performed using RPC over stdin / stdout
// so that we don't need access to either the node's filesystem or TCP stack
@ -61,37 +100,6 @@ type ExecNode struct {
key *ecdsa.PrivateKey
}
// NewExecNode creates a new ExecNode which will run the given service using a
// sub-directory of the given baseDir
func NewExecNode(id *NodeId, key *ecdsa.PrivateKey, service, baseDir string) (*ExecNode, error) {
if _, exists := serviceFuncs[service]; !exists {
return nil, fmt.Errorf("unknown node service %q", service)
}
// create the node directory using the first 12 characters of the ID
dir := filepath.Join(baseDir, id.String()[0:12])
if err := os.Mkdir(dir, 0755); err != nil {
return nil, fmt.Errorf("error creating node directory: %s", err)
}
// generate the config
conf := node.DefaultConfig
conf.DataDir = filepath.Join(dir, "data")
conf.P2P.ListenAddr = "127.0.0.1:0"
conf.P2P.NoDiscovery = true
conf.P2P.NAT = nil
node := &ExecNode{
ID: id,
Service: service,
Dir: dir,
Config: &conf,
key: key,
}
node.newCmd = node.execCommand
return node, nil
}
// Addr returns the node's enode URL
func (n *ExecNode) Addr() []byte {
if n.Info == nil {
@ -100,6 +108,8 @@ func (n *ExecNode) Addr() []byte {
return []byte(n.Info.Enode)
}
// Client returns an rpc.Client which can be used to communicate with the
// underlying service (it is set once the node has started)
func (n *ExecNode) Client() (*rpc.Client, error) {
return n.client, nil
}
@ -206,8 +216,9 @@ func init() {
}
// execP2PNode starts a devp2p node when the current binary is executed with
// argv[0] being "p2p-node", reading the service / ID from argv[1] / argv[2]
// and the node config from the _P2P_NODE_CONFIG environment variable
// argv[0] being "p2p-node", reading the service / ID from argv[1] / argv[2],
// the node config from the _P2P_NODE_CONFIG environment variable and the
// private key from the _P2P_NODE_KEY environment variable
func execP2PNode() {
glogger := log.NewGlogHandler(log.StreamHandler(os.Stderr, log.LogfmtFormat()))
glogger.Verbosity(log.LvlInfo)
@ -374,8 +385,8 @@ func (p *PeerAPI) PeerEvents(ctx context.Context) (*rpc.Subscription, error) {
return rpcSub, nil
}
// stdioConn wraps os.Stdin / os.Stdout with a nop Close method so we can
// use them to handle RPC messages
// stdioConn wraps os.Stdin / os.Stdout with a no-op Close method so we can
// use stdio for RPC messages
type stdioConn struct {
io.Reader
io.Writer

View file

@ -0,0 +1,318 @@
// Copyright 2016 The go-ethereum Authors
// This file is part of the go-ethereum library.
//
// The go-ethereum library is free software: you can redistribute it and/or modify
// it under the terms of the GNU Lesser General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// The go-ethereum library is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Lesser General Public License for more details.
//
// 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/>.
package adapters
import (
"errors"
"fmt"
"sync"
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rpc"
)
// SimAdapter is a NodeAdapter which creates in-memory nodes and connects them
// using an in-memory p2p.MsgReadWriter pipe
type SimAdapter struct {
mtx sync.RWMutex
nodes map[discover.NodeID]*SimNode
services map[string]ServiceFunc
}
// NewSimAdapter creates a SimAdapter which is capable of running in-memory
// nodes running any of the given services (the service to run on a particular
// node is passed to the NewNode function in the NodeConfig)
func NewSimAdapter(services map[string]ServiceFunc) *SimAdapter {
return &SimAdapter{
nodes: make(map[discover.NodeID]*SimNode),
services: services,
}
}
// Name returns the name of the adapter for logging purpoeses
func (s *SimAdapter) Name() string {
return "sim-adapter"
}
// NewNode returns a new SimNode using the given config
func (s *SimAdapter) NewNode(config *NodeConfig) (Node, error) {
s.mtx.Lock()
defer s.mtx.Unlock()
// check a node with the ID doesn't already exist
id := config.Id
if _, exists := s.nodes[id.NodeID]; exists {
return nil, fmt.Errorf("node already exists: %s", id)
}
// check the service is valid and initialize it
serviceFunc, exists := s.services[config.Service]
if !exists {
return nil, fmt.Errorf("unknown node service %q", config.Service)
}
service := serviceFunc(id)
// for simplicity, only support single protocol services (simulating
// multiple protocols on the same peer is extra effort, and we don't
// currently run any simulations which run multiple protocols)
if len(service.Protocols()) != 1 {
return nil, errors.New("service must have a single protocol")
}
node := &SimNode{
Id: id,
adapter: s,
service: service,
peers: make(map[discover.NodeID]MsgReadWriteCloser),
}
s.nodes[id.NodeID] = node
return node, nil
}
// GetNode returns the node with the given ID if it exists
func (s *SimAdapter) GetNode(id discover.NodeID) (*SimNode, bool) {
s.mtx.RLock()
defer s.mtx.RUnlock()
node, ok := s.nodes[id]
return node, ok
}
// MsgReadWriteCloser wraps a MsgReadWriter with the addition of a Close method
// so we can simulate the closing of a p2p connection (which usually happens by
/// closing the underlying TCP connection)
type MsgReadWriteCloser interface {
p2p.MsgReadWriter
Close() error
}
// SimNode is an in-memory node which connects to other SimNodes using an
// in-memory p2p.MsgReadWriter pipe, running an underlying service protocol
// directly over that pipe.
//
// It implements the p2p.Server interface so it can be used transparently
// by the underlying service.
type SimNode struct {
lock sync.RWMutex
Id *NodeId
adapter *SimAdapter
service node.Service
peers map[discover.NodeID]MsgReadWriteCloser
peerFeed event.Feed
client *rpc.Client
}
// Addr returns the node's discovery address
func (self *SimNode) Addr() []byte {
return []byte(self.Node().String())
}
// Node returns a discover.Node representing the SimNode
func (self *SimNode) Node() *discover.Node {
return discover.NewNode(self.Id.NodeID, nil, 0, 0)
}
// Client returns an rpc.Client which can be used to communicate with the
// underlying service (it is set once the node has started)
func (self *SimNode) Client() (*rpc.Client, error) {
self.lock.Lock()
defer self.lock.Unlock()
if self.client == nil {
return nil, errors.New("RPC not started")
}
return self.client, nil
}
// Start starts the RPC handler and the underlying service
func (self *SimNode) Start() error {
if err := self.startRPC(); err != nil {
return err
}
return self.service.Start(self)
}
// Stop stops the RPC handler and the underlying service
func (self *SimNode) Stop() error {
self.stopRPC()
return self.service.Stop()
}
// Service returns the underlying node.Service
func (self *SimNode) Service() node.Service {
return self.service
}
// startRPC starts an RPC server and connects to it using an in-process RPC
// client
func (self *SimNode) startRPC() error {
self.lock.Lock()
defer self.lock.Unlock()
if self.client != nil {
return errors.New("RPC already started")
}
// add SimAdminAPI and PeerAPI so that the network can call the
// AddPeer, RemovePeer and PeerEvents RPC methods
apis := append(self.service.APIs(), []rpc.API{
{
Namespace: "admin",
Version: "1.0",
Service: &SimAdminAPI{self},
},
{
Namespace: "eth",
Version: "1.0",
Service: &PeerAPI{func() p2p.Server { return self }},
},
}...)
// start the RPC handler
handler := rpc.NewServer()
for _, api := range apis {
if err := handler.RegisterName(api.Namespace, api.Service); err != nil {
return fmt.Errorf("error registering RPC: %s", err)
}
}
// create an in-process RPC client
self.client = rpc.DialInProc(handler)
return nil
}
// stopRPC closes the node's RPC client
func (self *SimNode) stopRPC() {
self.lock.Lock()
defer self.lock.Unlock()
if self.client != nil {
self.client.Close()
self.client = nil
}
}
// RemovePeer removes the given node as a peer by looking up the corresponding
// p2p.MsgReadWriter pipe and closing it (which will cause both the local
// and peer Protocol.Run functions to exit)
func (self *SimNode) RemovePeer(peer *discover.Node) {
self.lock.Lock()
defer self.lock.Unlock()
peerRW, exists := self.peers[peer.ID]
if !exists {
return
}
peerRW.Close()
delete(self.peers, peer.ID)
log.Trace(fmt.Sprintf("dropped peer %v", peer.ID))
}
// AddPeer adds the given node as a peer by creating a p2p.MsgReadWriter pipe
// and running both the local and peer's Protocol.Run function over the pipe
func (self *SimNode) AddPeer(peer *discover.Node) {
self.lock.Lock()
defer self.lock.Unlock()
if _, exists := self.peers[peer.ID]; exists {
return
}
peerNode, exists := self.adapter.GetNode(peer.ID)
if !exists {
panic(fmt.Sprintf("unknown peer: %s", peer.ID))
}
localRW, peerRW := p2p.MsgPipe()
self.peers[peer.ID] = peerRW
peerNode.RunProtocol(self, peerRW)
self.RunProtocol(peerNode, localRW)
}
// SubscribeEvents subscribes the given channel to p2p peer events
func (self *SimNode) SubscribeEvents(ch chan *p2p.PeerEvent) event.Subscription {
return self.peerFeed.Subscribe(ch)
}
// PeerCount returns the number of currently connected peers
func (self *SimNode) PeerCount() int {
self.lock.Lock()
defer self.lock.Unlock()
return len(self.peers)
}
// NodeInfo returns information about the node
func (self *SimNode) NodeInfo() *p2p.NodeInfo {
return &p2p.NodeInfo{ID: self.Id.String()}
}
// PeersInfo is a stub so that SimNode implements p2p.Server
func (self *SimNode) PeersInfo() (info []*p2p.PeerInfo) {
return nil
}
// RunProtocol runs the underlying service's protocol with the peer using the
// given p2p.MsgReadWriter, emitting peer add / drop events for peer event
// subscribers
func (self *SimNode) RunProtocol(peer *SimNode, rw p2p.MsgReadWriter) {
id := peer.Id
log.Trace(fmt.Sprintf("protocol starting on peer %v (connection with %v)", self.Id, id))
protocol := self.service.Protocols()[0]
p := p2p.NewPeer(id.NodeID, id.Label(), []p2p.Cap{})
go func() {
// emit peer add event
self.peerFeed.Send(&p2p.PeerEvent{
Type: p2p.PeerEventTypeAdd,
Peer: id.NodeID,
})
// run the protocol
err := protocol.Run(p, rw)
// remove the peer
self.RemovePeer(peer.Node())
log.Trace(fmt.Sprintf("protocol quit on peer %v (connection with %v broken: %v)", self.Id, id, err))
// emit peer drop event
self.peerFeed.Send(&p2p.PeerEvent{
Type: p2p.PeerEventTypeDrop,
Peer: id.NodeID,
Error: err.Error(),
})
}()
}
// SimAdminAPI implements the AddPeer and RemovePeer RPC methods (API
// compatible with node.PrivateAdminAPI)
type SimAdminAPI struct {
*SimNode
}
func (api *SimAdminAPI) AddPeer(url string) (bool, error) {
node, err := discover.ParseNode(url)
if err != nil {
return false, fmt.Errorf("invalid enode: %v", err)
}
api.SimNode.AddPeer(node)
return true, nil
}
func (api *SimAdminAPI) RemovePeer(url string) (bool, error) {
node, err := discover.ParseNode(url)
if err != nil {
return false, fmt.Errorf("invalid enode: %v", err)
}
api.SimNode.RemovePeer(node)
return true, nil
}

View file

@ -0,0 +1,173 @@
// Copyright 2016 The go-ethereum Authors
// This file is part of the go-ethereum library.
//
// The go-ethereum library is free software: you can redistribute it and/or modify
// it under the terms of the GNU Lesser General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// The go-ethereum library is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Lesser General Public License for more details.
//
// 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/>.
package adapters
import (
"crypto/ecdsa"
"encoding/hex"
"encoding/json"
"fmt"
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rpc"
)
// Node represents a node in a simulation network which is created by a
// NodeAdapter, for example:
//
// * SimNode - An in-memory node
// * ExecNode - A child process node
// * DockerNode - A docker container node
//
type Node interface {
// Addr returns the node's address (e.g. an Enode URL)
Addr() []byte
// Client returns the RPC client which is created once the node is
// up and running
Client() (*rpc.Client, error)
// Start starts the node
Start() error
// Stop stops the node
Stop() error
}
// NodeAdapter is an object which creates Nodes to be used in a simulation
// network
type NodeAdapter interface {
// Name returns the name of the adapter for logging purposes
Name() string
// NewNode creates a new node with the given configuration
NewNode(config *NodeConfig) (Node, error)
}
// RunProtocol is a function which runs a p2p protocol (see p2p.Protocol.Run)
type RunProtocol func(*p2p.Peer, p2p.MsgReadWriter) error
// NodeId wraps a discover.NodeID with some convenience methods
type NodeId struct {
discover.NodeID
}
func NewNodeId(id []byte) *NodeId {
var n discover.NodeID
copy(n[:], id)
return &NodeId{n}
}
func NewNodeIdFromHex(s string) *NodeId {
id := discover.MustHexID(s)
return &NodeId{id}
}
func (self *NodeId) Bytes() []byte {
return self.NodeID[:]
}
func (self *NodeId) Label() string {
return self.String()[:4]
}
// NodeConfig is the configuration used to start a node in a simulation
// network
type NodeConfig struct {
Id *NodeId
PrivateKey *ecdsa.PrivateKey
// Service is the name of the service which should be run when starting
// the node (for SimNodes it should be the name of a service contained
// in SimAdapter.services, for other nodes it should be a service
// registered by calling the RegisterService function)
Service string
}
// nodeConfigJSON is used to encode and decode NodeConfig as JSON by converting
// all fields to strings
type nodeConfigJSON struct {
Id string `json:"id"`
PrivateKey string `json:"private_key"`
Service string `json:"service"`
}
func (n *NodeConfig) MarshalJSON() ([]byte, error) {
return json.Marshal(nodeConfigJSON{
n.Id.String(),
hex.EncodeToString(crypto.FromECDSA(n.PrivateKey)),
n.Service,
})
}
func (n *NodeConfig) UnmarshalJSON(data []byte) error {
var confJSON nodeConfigJSON
if err := json.Unmarshal(data, &confJSON); err != nil {
return err
}
nodeID, err := discover.HexID(confJSON.Id)
if err != nil {
return err
}
n.Id = &NodeId{NodeID: nodeID}
key, err := hex.DecodeString(confJSON.PrivateKey)
if err != nil {
return err
}
n.PrivateKey = crypto.ToECDSA(key)
n.Service = confJSON.Service
return nil
}
// RandomNodeConfig returns node configuration with a randomly generated ID and
// PrivateKey
func RandomNodeConfig() *NodeConfig {
key, err := crypto.GenerateKey()
if err != nil {
panic("unable to generate key")
}
var id discover.NodeID
pubkey := crypto.FromECDSAPub(&key.PublicKey)
copy(id[:], pubkey[1:])
return &NodeConfig{
Id: &NodeId{NodeID: id},
PrivateKey: key,
}
}
// ServiceFunc returns a node.Service which can be used to boot devp2p nodes
type ServiceFunc func(id *NodeId) node.Service
// serviceFuncs is a map of registered services which are used to boot devp2p
// nodes
var serviceFuncs = make(map[string]ServiceFunc)
// RegisterService registers the given ServiceFunc which can then be used to
// start a devp2p node with the given name
func RegisterService(name string, f ServiceFunc) {
if _, exists := serviceFuncs[name]; exists {
panic(fmt.Sprintf("node service already exists: %q", name))
}
serviceFuncs[name] = f
}

View file

@ -1,21 +1,115 @@
package main
import (
"io/ioutil"
"os"
"runtime"
"time"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/simulations"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
)
// main() starts a simulation session which is capable of creating in-memory
// simulation networks containing nodes running a simple ping-pong protocol
func main() {
runtime.GOMAXPROCS(runtime.NumCPU())
log.Root().SetHandler(log.LvlFilterHandler(log.LvlTrace, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
c, quitc := simulations.NewSessionController(simulations.DefaultNet)
services := map[string]adapters.ServiceFunc{
"ping-pong": func(id *adapters.NodeId) node.Service {
return newPingPongService(id)
},
}
c, quitc := simulations.NewSessionController(simulations.DefaultNet(services, "ping-pong"))
simulations.StartRestApiServer("8888", c)
// wait until server shuts down
<-quitc
}
// pingPongService runs a ping-pong protocol between nodes where each node
// sends a ping to all its connected peers every 10s and receives a pong in
// return
type pingPongService struct {
id *adapters.NodeId
log log.Logger
}
func newPingPongService(id *adapters.NodeId) *pingPongService {
return &pingPongService{
id: id,
log: log.New("node.id", id),
}
}
func (p *pingPongService) Protocols() []p2p.Protocol {
return []p2p.Protocol{{
Name: "ping-pong",
Version: 1,
Length: 2,
Run: p.Run,
}}
}
func (p *pingPongService) APIs() []rpc.API {
return nil
}
func (p *pingPongService) Start(server p2p.Server) error {
p.log.Info("ping-pong service starting")
return nil
}
func (p *pingPongService) Stop() error {
p.log.Info("ping-pong service stopping")
return nil
}
const (
pingMsgCode = iota
pongMsgCode
)
// Run implements the ping-pong protocol which sends ping messages to the peer
// at 10s intervals, and responds to pings with pong messages.
func (p *pingPongService) Run(peer *p2p.Peer, rw p2p.MsgReadWriter) error {
log := p.log.New("peer.id", peer.ID())
errC := make(chan error)
go func() {
for range time.Tick(10 * time.Second) {
log.Info("sending ping")
if err := p2p.Send(rw, pingMsgCode, "PING"); err != nil {
errC <- err
return
}
}
}()
go func() {
for {
msg, err := rw.ReadMsg()
if err != nil {
errC <- err
return
}
payload, err := ioutil.ReadAll(msg.Payload)
if err != nil {
errC <- err
return
}
log.Info("received message", "msg.code", msg.Code, "msg.payload", string(payload))
if msg.Code == pingMsgCode {
log.Info("sending pong")
go p2p.Send(rw, pongMsgCode, "PONG")
}
}
}()
return <-errC
}

View file

@ -10,7 +10,7 @@ import (
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
// Journal is an instance of a guaranteed no-loss subscription to network related events

View file

@ -7,7 +7,7 @@ import (
"time"
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
func testEvents(intervals ...int) (events []*event.TypeMuxEvent) {

View file

@ -8,7 +8,7 @@ import (
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
func NewMockersController(eventer *event.TypeMux, defaultMockerConfig *MockerConfig) Controller {

View file

@ -28,20 +28,17 @@ package simulations
import (
"context"
"crypto/ecdsa"
"encoding/hex"
"encoding/json"
"fmt"
"net/http"
"reflect"
"sync"
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
)
@ -53,6 +50,7 @@ type NetworkConfig struct {
Id string
DefaultMockerConfig *MockerConfig
Backend bool
DefaultService string
}
type NetworkControl interface {
@ -241,6 +239,8 @@ func NewDebugController(journal *Journal) Controller {
// the actual logic of bringing nodes and connections up and down and
// messaging is implemented in the particular NodeAdapter interface
type Network struct {
nodeAdapter adapters.NodeAdapter
// input trigger events and other events
events *event.TypeMux // generated events a journal can subsribe to
lock sync.RWMutex
@ -250,27 +250,19 @@ type Network struct {
Conns []*Conn `json:"conns"`
quitc chan bool
conf *NetworkConfig
//
// adapters.Messenger
// node adapter function that creates the node model for
// the particular type of network from a config
naf func(*NodeConfig) adapters.NodeAdapter
}
func NewNetwork(conf *NetworkConfig) *Network {
func NewNetwork(nodeAdapter adapters.NodeAdapter, conf *NetworkConfig) *Network {
return &Network{
conf: conf,
events: &event.TypeMux{},
nodeMap: make(map[discover.NodeID]int),
connMap: make(map[string]int),
quitc: make(chan bool),
nodeAdapter: nodeAdapter,
conf: conf,
events: &event.TypeMux{},
nodeMap: make(map[discover.NodeID]int),
connMap: make(map[string]int),
quitc: make(chan bool),
}
}
func (self *Network) SetNaf(naf func(*NodeConfig) adapters.NodeAdapter) {
self.naf = naf
}
// Subscribe takes an event.TypeMux and subscibes to types
// and launches a goroutine that reads control events from an eventer Subsription channel
// and executes the events
@ -373,11 +365,11 @@ type LiveEventer interface {
}
type Node struct {
NodeConfig
adapters.Node
adapters.NodeConfig
Up bool
na adapters.NodeAdapter
controlFired bool
}
@ -510,81 +502,10 @@ func (self *NodeEvent) ToControlEvent() *NodeControlEvent {
return &NodeControlEvent{self}
}
func (self *Node) Adapter() adapters.NodeAdapter {
return self.na
}
type NodeConfig struct {
Id *adapters.NodeId
PrivateKey *ecdsa.PrivateKey
}
type nodeConfigJSON struct {
Id string `json:"id"`
PrivateKey string `json:"private_key"`
}
func (n *NodeConfig) MarshalJSON() ([]byte, error) {
return json.Marshal(nodeConfigJSON{
n.Id.String(),
hex.EncodeToString(crypto.FromECDSA(n.PrivateKey)),
})
}
func (n *NodeConfig) UnmarshalJSON(data []byte) error {
var confJSON nodeConfigJSON
if err := json.Unmarshal(data, &confJSON); err != nil {
return err
}
nodeID, err := discover.HexID(confJSON.Id)
if err != nil {
return err
}
n.Id = &adapters.NodeId{NodeID: nodeID}
key, err := hex.DecodeString(confJSON.PrivateKey)
if err != nil {
return err
}
n.PrivateKey = crypto.ToECDSA(key)
return nil
}
func RandomNodeConfig() *NodeConfig {
key, err := crypto.GenerateKey()
if err != nil {
panic("unable to generate key")
}
var id discover.NodeID
pubkey := crypto.FromECDSAPub(&key.PublicKey)
copy(id[:], pubkey[1:])
return &NodeConfig{
Id: &adapters.NodeId{NodeID: id},
PrivateKey: key,
}
}
// TODO: ignored for now
type QueryConfig struct {
Format string // "sim.update", "journal",
}
type Know struct {
Subject *adapters.NodeId `json:"subject"`
Object *adapters.NodeId `json:"object"`
// Into
// number of attempted connections
// time of attempted connections
// number of active connections during the session
// number of active connections since records began
// swap balance
}
// NewNode adds a new node to the network with a random ID
func (self *Network) NewNode() (*NodeConfig, error) {
conf := RandomNodeConfig()
func (self *Network) NewNode() (*adapters.NodeConfig, error) {
conf := adapters.RandomNodeConfig()
conf.Service = self.conf.DefaultService
if err := self.NewNodeWithConfig(conf); err != nil {
return nil, err
}
@ -593,20 +514,27 @@ func (self *Network) NewNode() (*NodeConfig, error) {
// NewNodeWithConfig adds a new node to the network with the given config
// errors if a node by the same id already exist
func (self *Network) NewNodeWithConfig(conf *NodeConfig) error {
func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) error {
self.lock.Lock()
defer self.lock.Unlock()
id := conf.Id
if conf.Service == "" {
conf.Service = self.conf.DefaultService
}
_, found := self.nodeMap[id.NodeID]
if found {
return fmt.Errorf("node %v already added", id)
}
self.nodeMap[id.NodeID] = len(self.Nodes)
na := self.naf(conf)
adapterNode, err := self.nodeAdapter.NewNode(conf)
if err != nil {
return err
}
node := &Node{
Node: adapterNode,
NodeConfig: *conf,
na: na,
}
self.Nodes = append(self.Nodes, node)
log.Trace(fmt.Sprintf("node %v created", id))
@ -655,8 +583,8 @@ func (self *Network) Start(id *adapters.NodeId) error {
if node.Up {
return fmt.Errorf("node %v already up", id)
}
log.Trace(fmt.Sprintf("starting node %v: %v adapter %v", id, node.Up, node.Adapter()))
if err := node.Adapter().Start(); err != nil {
log.Trace(fmt.Sprintf("starting node %v: %v using %v", id, node.Up, self.nodeAdapter.Name()))
if err := node.Start(); err != nil {
return err
}
node.Up = true
@ -665,7 +593,7 @@ func (self *Network) Start(id *adapters.NodeId) error {
self.events.Post(node.EmitEvent(ControlEvent))
// subscribe to peer events
client, err := node.Adapter().Client()
client, err := node.Client()
if err != nil {
return fmt.Errorf("error getting rpc client for node %v: %s", id, err)
}
@ -715,7 +643,7 @@ func (self *Network) Stop(id *adapters.NodeId) error {
if !node.Up {
return fmt.Errorf("node %v already down", id)
}
if err := node.Adapter().Stop(); err != nil {
if err := node.Stop(); err != nil {
return err
}
node.Up = false
@ -753,11 +681,11 @@ func (self *Network) Connect(oneId, otherId *adapters.NodeId) error {
var addr []byte
var client *rpc.Client
if rev {
addr = conn.one.na.Addr()
client, err = conn.other.na.Client()
addr = conn.one.Addr()
client, err = conn.other.Client()
} else {
addr = conn.other.na.Addr()
client, err = conn.one.na.Client()
addr = conn.other.Addr()
client, err = conn.one.Client()
}
if err != nil {
return err
@ -787,11 +715,11 @@ func (self *Network) Disconnect(oneId, otherId *adapters.NodeId) error {
var client *rpc.Client
var err error
if rev {
addr = conn.one.na.Addr()
client, err = conn.other.na.Client()
addr = conn.one.Addr()
client, err = conn.other.Client()
} else {
addr = conn.other.na.Addr()
client, err = conn.one.na.Client()
addr = conn.other.Addr()
client, err = conn.one.Client()
}
if err != nil {
return err
@ -840,18 +768,6 @@ func (self *Network) Send(senderid, receiverid *adapters.NodeId, msgcode uint64,
self.events.Post(msg.EmitEvent(ControlEvent))
}
// GetNodeAdapter(id) returns the NodeAdapter for node with id
// returns nil if node does not exist
func (self *Network) GetNodeAdapter(id *adapters.NodeId) adapters.NodeAdapter {
self.lock.Lock()
defer self.lock.Unlock()
node := self.getNode(id)
if node == nil {
return nil
}
return node.na
}
// GetNode retrieves the node model for the id given as arg
// returns nil if the node does not exist
func (self *Network) GetNode(id *adapters.NodeId) *Node {
@ -918,7 +834,7 @@ func (self *Network) Shutdown() {
// stop all nodes
for _, node := range self.Nodes {
log.Debug(fmt.Sprintf("stopping node %s", node.Id.Label()))
if err := node.na.Stop(); err != nil {
if err := node.Stop(); err != nil {
log.Warn(fmt.Sprintf("error stopping node %s", node.Id.Label()), "err", err)
}
}

View file

@ -11,7 +11,7 @@ import (
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
type returnHandler func(body io.Reader) (resp io.ReadSeeker, err error)
@ -83,9 +83,12 @@ func NewResourceContoller(c *ResourceHandlers) *ResourceController {
var empty = struct{}{}
func DefaultNet(conf *NetworkConfig) (NetworkControl, *ResourceController) {
net := NewNetwork(conf)
return NetworkControl(net), NewNodesController(net)
func DefaultNet(services map[string]adapters.ServiceFunc, defaultService string) func(conf *NetworkConfig) (NetworkControl, *ResourceController) {
return func(conf *NetworkConfig) (NetworkControl, *ResourceController) {
conf.DefaultService = defaultService
net := NewNetwork(adapters.NewSimAdapter(services), conf)
return NetworkControl(net), NewNodesController(net)
}
}
func NewSessionController(nethook func(*NetworkConfig) (NetworkControl, *ResourceController)) (*ResourceController, chan bool) {

View file

@ -13,7 +13,7 @@ import (
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
/***

View file

@ -4,7 +4,7 @@ import (
"context"
"time"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
// Simulation provides a framework for running actions in a simulated network

View file

@ -5,8 +5,8 @@ import (
"sync"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
type TestPeer interface {

View file

@ -1,32 +1,23 @@
package testing
import (
"errors"
"fmt"
"regexp"
"strconv"
"strings"
"sync"
"time"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
type ProtocolSession struct {
TestNodeAdapter
Ids []*adapters.NodeId
ignore []uint64
}
*adapters.SimNode
type TestMessenger interface {
ExpectMsg(uint64, interface{}) error
TriggerMsg(uint64, interface{}) error
}
type TestNodeAdapter interface {
p2p.Server
GetPeer(id *adapters.NodeId) *adapters.Peer
Ids []*adapters.NodeId
adapter *adapters.SimAdapter
events chan *p2p.PeerEvent
}
// exchanges are the basic units of protocol tests
@ -61,32 +52,22 @@ type Disconnect struct {
Error error // disconnect reason
}
func NewProtocolSession(na TestNodeAdapter, ids []*adapters.NodeId) *ProtocolSession {
ps := &ProtocolSession{
TestNodeAdapter: na,
Ids: ids,
}
return ps
}
func (self *ProtocolSession) SetIgnoreCodes(ignore ...uint64) {
self.ignore = ignore
}
// trigger sends messages from peers
func (self *ProtocolSession) trigger(trig Trigger) error {
peer := self.GetPeer(trig.Peer)
if peer == nil {
panic(fmt.Sprintf("trigger: peer %v does not exist (1- %v)", trig.Peer, len(self.Ids)))
simNode, ok := self.adapter.GetNode(trig.Peer.NodeID)
if !ok {
return fmt.Errorf("trigger: peer %v does not exist (1- %v)", trig.Peer, len(self.Ids))
}
if peer.MsgReadWriteCloser == nil {
return fmt.Errorf("trigger: peer %v unreachable", trig.Peer)
mockNode, ok := simNode.Service().(*mockNode)
if !ok {
return fmt.Errorf("trigger: peer %v is not a mock", trig.Peer)
}
errc := make(chan error)
go func() {
log.Trace(fmt.Sprintf("trigger %v (%v)....", trig.Msg, trig.Code))
errc <- p2p.Send(peer, trig.Code, trig.Msg)
errc <- mockNode.Trigger(&trig)
log.Trace(fmt.Sprintf("triggered %v (%v)", trig.Msg, trig.Code))
}()
@ -106,51 +87,21 @@ func (self *ProtocolSession) trigger(trig Trigger) error {
// expect checks an expectation
func (self *ProtocolSession) expect(exp Expect) error {
if exp.Msg == nil {
panic("no message to expect")
return errors.New("no message to expect")
}
peer := self.GetPeer(exp.Peer)
if peer == nil {
panic(fmt.Sprintf("expect: peer %v does not exist (1- %v)", exp.Peer, len(self.Ids)))
simNode, ok := self.adapter.GetNode(exp.Peer.NodeID)
if !ok {
return fmt.Errorf("trigger: peer %v does not exist (1- %v)", exp.Peer, len(self.Ids))
}
if peer.MsgReadWriteCloser== nil {
return fmt.Errorf("trigger: peer %v unreachable", exp.Peer)
mockNode, ok := simNode.Service().(*mockNode)
if !ok {
return fmt.Errorf("trigger: peer %v is not a mock", exp.Peer)
}
errc := make(chan error)
go func() {
var err error
ignored := true
log.Trace("waiting for msg", "code", exp.Code, "msg", exp.Msg)
for ignored {
ignored = false
err = p2p.ExpectMsg(peer, exp.Code, exp.Msg)
// frail, but we can't know what code expectmsg got otherwise
// can we do better error reporting in p2p.ExpectMsg()?
if err != nil {
if strings.Contains(err.Error(), "code") {
re, _ := regexp.Compile("got ([0-9]+),")
match := re.FindStringSubmatch(err.Error())
if len(match) > 1 {
for _, codetoignore := range self.ignore {
codewegot, err := strconv.ParseUint(match[1], 10, 64)
if err == nil {
if codetoignore == codewegot {
ignored = true
log.Trace("ignore msg with wrong code", "received", codewegot, "expected", exp.Code)
break
}
} else {
log.Warn("expectmsg errormsg parse error?!")
}
}
} else {
log.Warn("expectmsg errormsg parse error?!")
break
}
}
}
}
errc <- err
log.Trace(fmt.Sprintf("waiting for msg, %v", exp.Msg))
errc <- mockNode.Expect(&exp)
}()
t := exp.Timeout
@ -226,29 +177,30 @@ func (self *ProtocolSession) TestExchanges(exchanges ...Exchange) error {
}
func (self *ProtocolSession) TestDisconnected(disconnects ...*Disconnect) error {
expects := make(map[discover.NodeID]error)
for _, disconnect := range disconnects {
id := disconnect.Peer
err := disconnect.Error
peer := self.GetPeer(id)
expects[disconnect.Peer.NodeID] = disconnect.Error
}
alarm := time.NewTimer(1000 * time.Millisecond)
timeout := time.After(time.Second)
for len(expects) > 0 {
select {
case derr := <-peer.Errc:
if !((err == nil && derr == nil) || err != nil && derr != nil && err.Error() == derr.Error()) {
return fmt.Errorf("unexpected error on peer %v. expected '%v', got '%v'", id, err, derr)
case event := <-self.events:
if event.Type != p2p.PeerEventTypeDrop {
continue
}
case <-alarm.C:
return fmt.Errorf("timed out waiting for peer %v to disconnect", id)
expectErr, ok := expects[event.Peer]
if !ok {
continue
}
if !((expectErr == nil && event.Error == "") || expectErr != nil && event.Error != "" && expectErr.Error() == event.Error) {
return fmt.Errorf("unexpected error on peer %v. expected '%v', got '%v'", event.Peer, expectErr, event.Error)
}
delete(expects, event.Peer)
case <-timeout:
return fmt.Errorf("timed out waiting for peers to disconnect")
}
}
return nil
}
func (self *ProtocolSession) Stop() {
for _, id := range self.Ids {
p := self.GetPeer(id)
if p != nil && p.MsgReadWriteCloser != nil {
p.Close()
}
}
}

View file

@ -2,12 +2,14 @@ package testing
import (
"fmt"
"sync"
"testing"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
)
@ -17,32 +19,39 @@ type ProtocolTester struct {
}
func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(*p2p.Peer, p2p.MsgReadWriter) error) *ProtocolTester {
net := simulations.NewNetwork(&simulations.NetworkConfig{})
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
node := &testNode{}
if conf.Id.NodeID == id.NodeID {
log.Trace(fmt.Sprintf("adapter run function set to protocol for node %v (=%v)", conf.Id, id))
node.run = run
}
return adapters.NewSimNode(conf.Id, node, net)
services := map[string]adapters.ServiceFunc{
"test": func(id *adapters.NodeId) node.Service {
return &testNode{run}
},
"mock": func(id *adapters.NodeId) node.Service {
return newMockNode()
},
}
net.SetNaf(naf)
if err := net.NewNodeWithConfig(&simulations.NodeConfig{Id: id}); err != nil {
adapter := adapters.NewSimAdapter(services)
net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{})
if err := net.NewNodeWithConfig(&adapters.NodeConfig{Id: id, Service: "test"}); err != nil {
panic(err.Error())
}
if err := net.Start(id); err != nil {
panic(err.Error())
}
node := net.GetNodeAdapter(id).(*adapters.SimNode)
peers := make([]*simulations.NodeConfig, n)
node := net.GetNode(id).Node.(*adapters.SimNode)
peers := make([]*adapters.NodeConfig, n)
peerIDs := make([]*adapters.NodeId, n)
for i := 0; i < n; i++ {
peers[i] = simulations.RandomNodeConfig()
peers[i] = adapters.RandomNodeConfig()
peers[i].Service = "mock"
peerIDs[i] = peers[i].Id
}
ps := NewProtocolSession(node, peerIDs)
events := make(chan *p2p.PeerEvent, 1000)
node.SubscribeEvents(events)
ps := &ProtocolSession{
SimNode: node,
Ids: peerIDs,
adapter: adapter,
events: events,
}
self := &ProtocolTester{
ProtocolSession: ps,
network: net,
@ -53,7 +62,7 @@ func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(*p2p.P
return self
}
func (self *ProtocolTester) Connect(selfId *adapters.NodeId, peers ...*simulations.NodeConfig) {
func (self *ProtocolTester) Connect(selfId *adapters.NodeId, peers ...*adapters.NodeConfig) {
for _, peer := range peers {
log.Trace(fmt.Sprintf("start node %v", peer.Id))
if err := self.network.NewNodeWithConfig(peer); err != nil {
@ -91,3 +100,57 @@ func (t *testNode) Start(server p2p.Server) error {
func (t *testNode) Stop() error {
return nil
}
// mockNode is a testNode which doesn't actually run a protocol, instead
// exposing channels so that tests can manually trigger and expect certain
// messages
type mockNode struct {
testNode
trigger chan *Trigger
expect chan *Expect
err chan error
stop chan struct{}
stopOnce sync.Once
}
func newMockNode() *mockNode {
mock := &mockNode{
trigger: make(chan *Trigger),
expect: make(chan *Expect),
err: make(chan error),
stop: make(chan struct{}),
}
mock.testNode.run = mock.Run
return mock
}
// Run is a protocol run function which just loops waiting for tests to
// instruct it to either trigger or expect a message from the peer
func (m *mockNode) Run(peer *p2p.Peer, rw p2p.MsgReadWriter) error {
for {
select {
case trig := <-m.trigger:
m.err <- p2p.Send(rw, trig.Code, trig.Msg)
case exp := <-m.expect:
m.err <- p2p.ExpectMsg(rw, exp.Code, exp.Msg)
case <-m.stop:
return nil
}
}
}
func (m *mockNode) Trigger(trig *Trigger) error {
m.trigger <- trig
return <-m.err
}
func (m *mockNode) Expect(exp *Expect) error {
m.expect <- exp
return <-m.err
}
func (m *mockNode) Stop() error {
m.stopOnce.Do(func() { close(m.stop) })
return nil
}

View file

@ -23,8 +23,8 @@ import (
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
)
/*

View file

@ -5,7 +5,7 @@ import (
"testing"
"time"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
)
@ -82,7 +82,7 @@ func TestRegisterAndConnect(t *testing.T) {
},
ticker: make(chan time.Time),
}
pp.Start(s.TestNodeAdapter, tc.ping)
pp.Start(s, tc.ping)
defer pp.Stop()
tc.ticker <- time.Now()

View file

@ -23,9 +23,9 @@ import (
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/protocols"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/pot"
)

View file

@ -5,8 +5,8 @@ import (
"testing"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/protocols"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
)
@ -96,7 +96,6 @@ func (s *bzzTester) testHandshake(lhs, rhs *bzzHandshake, disconnects ...*p2ptes
} else {
peers = []*adapters.NodeId{id}
}
<-s.GetPeer(id).Connc
s.TestExchanges(bzzHandshakeExchange(lhs, rhs, id)...)
s.TestDisconnected(disconnects...)

View file

@ -12,8 +12,8 @@ import (
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/protocols"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/pot"
"github.com/ethereum/go-ethereum/rlp"
"github.com/ethereum/go-ethereum/swarm/storage"
@ -242,12 +242,12 @@ func (self *Pss) GetHandler(topic PssTopic) func([]byte, *p2p.Peer, []byte) erro
// Links a pss peer address and topic to a dedicated p2p.MsgReadWriter in the pss peerpool, and runs the specificed protocol on this p2p.MsgReadWriter and the specified peer
//
// The effect is that now we have a "virtual" protocol running on an artificial p2p.Peer, which can be looked up and piped to through Pss using swarm overlay address and topic
func (self *Pss) AddPeer(p *p2p.Peer, addr pot.Address, protocall adapters.ProtoCall, topic PssTopic, rw p2p.MsgReadWriter) error {
func (self *Pss) AddPeer(p *p2p.Peer, addr pot.Address, run adapters.RunProtocol, topic PssTopic, rw p2p.MsgReadWriter) error {
self.lock.Lock()
defer self.lock.Unlock()
self.addPeerTopic(addr, topic, rw)
go func() {
err := protocall(p, rw)
err := run(p, rw)
log.Warn(fmt.Sprintf("pss vprotocol quit on addr %v topic %v: %v", addr, topic, err))
}()
return nil

View file

@ -14,10 +14,11 @@ import (
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/protocols"
"github.com/ethereum/go-ethereum/p2p/simulations"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
)
@ -47,12 +48,11 @@ type pssTestPeer struct {
type pssTestNode struct {
*Hive
*Pss
*adapters.SimNode
id *adapters.NodeId
network *simulations.Network
trigger chan *adapters.NodeId
run adapters.ProtoCall
run adapters.RunProtocol
ct *protocols.CodeMap
expectC chan []int
ws *http.Handler
@ -70,25 +70,12 @@ func (n *pssTestNode) Remove(peer Peer) {
n.Hive.Remove(peer)
}
func (n *pssTestNode) Start() error {
return n.Hive.Start(n.SimNode, n.hiveKeepAlive)
}
func (n *pssTestNode) Stop() error {
n.Hive.Stop()
return nil
}
func (n *pssTestNode) hiveKeepAlive() <-chan time.Time {
return time.Tick(time.Second * 10)
}
func (n *pssTestNode) triggerCheck() {
go func() { n.trigger <- adapters.NewNodeId(n.Addr()) }()
}
func (n *pssTestNode) ProtoCall() adapters.ProtoCall {
return n.run
go func() { n.trigger <- n.id }()
}
func (n *pssTestNode) OverlayAddr() []byte {
@ -125,11 +112,11 @@ func newPssTestService(t *testing.T, handlefunc func(interface{}) error, testnod
}
func (self *pssTestService) Start(server p2p.Server) error {
self.node.SimNode = server.(*adapters.SimNode) // server is adapter.SimnNode now
return nil
return self.node.Hive.Start(server, self.node.hiveKeepAlive)
}
func (self *pssTestService) Stop() error {
self.node.Hive.Stop()
return nil
}
@ -322,12 +309,8 @@ func testPssFullRandom(t *testing.T, numsends int, numnodes int, numfullnodes in
topic, _ := MakeTopic(protocolName, protocolVersion)
trigger := make(chan *adapters.NodeId)
net := simulations.NewNetwork(&simulations.NetworkConfig{
Id: "0",
Backend: true,
})
testpeers := make(map[*adapters.NodeId]*pssTestPeer)
nodes := newPssSimulationTester(t, numnodes, numfullnodes, net, trigger, vct, protocolName, protocolVersion, testpeers)
net, nodes := newPssSimulationTester(t, numnodes, numfullnodes, trigger, vct, protocolName, protocolVersion, testpeers)
ids := []*adapters.NodeId{}
@ -533,12 +516,8 @@ func TestPssFullLinearEcho(t *testing.T) {
fullnodes := []*adapters.NodeId{}
trigger := make(chan *adapters.NodeId)
net := simulations.NewNetwork(&simulations.NetworkConfig{
Id: "0",
Backend: true,
})
testpeers := make(map[*adapters.NodeId]*pssTestPeer)
nodes := newPssSimulationTester(t, 3, 2, net, trigger, vct, protocolName, protocolVersion, testpeers)
net, nodes := newPssSimulationTester(t, 3, 2, trigger, vct, protocolName, protocolVersion, testpeers)
ids := []*adapters.NodeId{} // ohh risky! but the action for a specific id should come before the expect anyway
action = func(ctx context.Context) error {
@ -717,12 +696,8 @@ func TestPssFullWS(t *testing.T) {
topic, _ := MakeTopic(pingTopicName, pingTopicVersion)
trigger := make(chan *adapters.NodeId)
simnet := simulations.NewNetwork(&simulations.NetworkConfig{
Id: "0",
Backend: true,
})
testpeers := make(map[*adapters.NodeId]*pssTestPeer)
nodes := newPssSimulationTester(t, 3, 2, simnet, trigger, vct, protocolName, protocolVersion, testpeers)
simnet, nodes := newPssSimulationTester(t, 3, 2, trigger, vct, protocolName, protocolVersion, testpeers)
ids := []*adapters.NodeId{} // ohh risky! but the action for a specific id should come before the expect anyway
action = func(ctx context.Context) error {
@ -935,16 +910,16 @@ func TestPssFullWS(t *testing.T) {
// the simulation tester constructor is currently a hack to fit previous code with later stack using node.Services to start SimNodes
func newPssSimulationTester(t *testing.T, numnodes int, numfullnodes int, simnet *simulations.Network, trigger chan *adapters.NodeId, vct *protocols.CodeMap, name string, version int, testpeers map[*adapters.NodeId]*pssTestPeer) map[*adapters.NodeId]*pssTestNode {
func newPssSimulationTester(t *testing.T, numnodes int, numfullnodes int, trigger chan *adapters.NodeId, vct *protocols.CodeMap, name string, version int, testpeers map[*adapters.NodeId]*pssTestPeer) (*simulations.Network, map[*adapters.NodeId]*pssTestNode) {
topic, _ := MakeTopic(name, version)
nodes := make(map[*adapters.NodeId]*pssTestNode, numnodes)
psss := make(map[*adapters.NodeId]*Pss)
simnet.SetNaf(func(conf *simulations.NodeConfig) adapters.NodeAdapter {
var simnet *simulations.Network
serviceFunc := func(id *adapters.NodeId) node.Service {
node := &pssTestNode{
Pss: psss[conf.Id],
Pss: psss[id],
Hive: nil,
SimNode: &adapters.SimNode{},
id: conf.Id,
id: id,
network: simnet,
trigger: trigger,
ct: vct,
@ -956,33 +931,37 @@ func newPssSimulationTester(t *testing.T, numnodes int, numfullnodes int, simnet
var handlefunc func(interface{}) error
addr := NewPeerAddrFromNodeId(conf.Id)
addr := NewPeerAddrFromNodeId(id)
if testpeers[conf.Id] != nil {
handlefunc = makePssHandleProtocol(psss[conf.Id])
log.Trace(fmt.Sprintf("Making full protocol id %x addr %x (testpeers %p)", common.ByteLabel(conf.Id.Bytes()), common.ByteLabel(addr.OverlayAddr()), testpeers))
if testpeers[id] != nil {
handlefunc = makePssHandleProtocol(psss[id])
log.Trace(fmt.Sprintf("Making full protocol id %x addr %x (testpeers %p)", common.ByteLabel(id.Bytes()), common.ByteLabel(addr.OverlayAddr()), testpeers))
} else {
handlefunc = makePssHandleForward(psss[conf.Id])
handlefunc = makePssHandleForward(psss[id])
}
// protocols are now registered by invoking node services
// since adapters.SimNode implements p2p.Server, needed for the services to start, we use this as a convenience wrapper
testservice := newPssTestService(t, handlefunc, node)
svc := adapters.NewSimNode(conf.Id, testservice, simnet)
testservice.Start(svc)
// the network sim wants a adapters.NodeAdapter, so we pass back to it a SimNode
// this is the SimNode member of the testNode initialized above, but assigned through the service start
// that is so say: node == testservice.node, but we access it as a member of testservice below for clarity (to the extent that this can be clear)
nodes[conf.Id] = testservice.node
nodes[id] = testservice.node
testservice.node.apifunc = testservice.APIs
return node.SimNode
return testservice
}
adapter := adapters.NewSimAdapter(map[string]adapters.ServiceFunc{"pss": serviceFunc})
simnet = simulations.NewNetwork(adapter, &simulations.NetworkConfig{
Id: "0",
Backend: true,
})
configs := make([]*simulations.NodeConfig, numnodes)
configs := make([]*adapters.NodeConfig, numnodes)
for i := 0; i < numnodes; i++ {
configs[i] = simulations.RandomNodeConfig()
configs[i] = adapters.RandomNodeConfig()
configs[i].Service = "pss"
}
for i, conf := range configs {
addr := NewPeerAddrFromNodeId(conf.Id)
@ -1001,13 +980,15 @@ func newPssSimulationTester(t *testing.T, numnodes int, numfullnodes int, simnet
psss[conf.Id].Register(topic, pssprotocol.GetHandler())
}
simnet.NewNodeWithConfig(conf)
if err := simnet.NewNodeWithConfig(conf); err != nil {
t.Fatalf("error creating node %s: %s", conf.Id.Label(), err)
}
if err := simnet.Start(conf.Id); err != nil {
t.Fatalf("error starting node %s: %s", conf.Id.Label(), err)
}
}
return nodes
return simnet, nodes
}
func makePss(addr []byte) *Pss {

View file

@ -12,8 +12,8 @@ import (
"github.com/ethereum/go-ethereum/log"
p2pnode "github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
"github.com/ethereum/go-ethereum/swarm/network"
)
@ -22,12 +22,14 @@ import (
// service to execute
const serviceName = "discovery"
func newService(id *adapters.NodeId) p2pnode.Service {
return newNode(id)
}
func init() {
// register the discovery service which will run as a devp2p
// protocol when using the exec adapter
adapters.RegisterService(serviceName, func(id *adapters.NodeId) p2pnode.Service {
return newNode(id)
})
adapters.RegisterService(serviceName, newService)
// log.Root().SetHandler(log.LvlFilterHandler(log.LvlError, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
log.Root().SetHandler(log.LvlFilterHandler(log.LvlTrace, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
@ -43,17 +45,11 @@ func TestMain(m *testing.M) {
}
func TestDiscoverySimulationDockerAdapter(t *testing.T) {
setup := func(net *simulations.Network) {
net.SetNaf(func(conf *simulations.NodeConfig) adapters.NodeAdapter {
node, err := adapters.NewDockerNode(conf.Id, conf.PrivateKey, serviceName)
if err != nil {
panic(err)
}
return node
})
adapter, err := adapters.NewDockerAdapter()
if err != nil {
t.Fatal(err)
}
testDiscoverySimulation(t, setup)
testDiscoverySimulation(t, adapter)
}
func TestDiscoverySimulationExecAdapter(t *testing.T) {
@ -62,39 +58,23 @@ func TestDiscoverySimulationExecAdapter(t *testing.T) {
t.Fatal(err)
}
defer os.RemoveAll(baseDir)
setup := func(net *simulations.Network) {
net.SetNaf(func(conf *simulations.NodeConfig) adapters.NodeAdapter {
node, err := adapters.NewExecNode(conf.Id, conf.PrivateKey, serviceName, baseDir)
if err != nil {
panic(err)
}
return node
})
}
testDiscoverySimulation(t, setup)
testDiscoverySimulation(t, adapters.NewExecAdapter(baseDir))
}
func TestDiscoverySimulationSimAdapter(t *testing.T) {
setup := func(net *simulations.Network) {
net.SetNaf(func(conf *simulations.NodeConfig) adapters.NodeAdapter {
return adapters.NewSimNode(conf.Id, newNode(conf.Id), net)
})
}
testDiscoverySimulation(t, setup)
services := map[string]adapters.ServiceFunc{serviceName: newService}
testDiscoverySimulation(t, adapters.NewSimAdapter(services))
}
func testDiscoverySimulation(t *testing.T, setup func(net *simulations.Network)) {
func testDiscoverySimulation(t *testing.T, adapter adapters.NodeAdapter) {
// create 10 node network
nodeCount := 10
net := simulations.NewNetwork(&simulations.NetworkConfig{
Id: "0",
Backend: true,
net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{
Id: "0",
Backend: true,
DefaultService: serviceName,
})
defer net.Shutdown()
setup(net)
trigger := make(chan *adapters.NodeId)
ids := make([]*adapters.NodeId, nodeCount)
for i := 0; i < nodeCount; i++ {
@ -134,7 +114,7 @@ func testDiscoverySimulation(t *testing.T, setup func(net *simulations.Network))
default:
}
node := net.GetNodeAdapter(id)
node := net.GetNode(id)
if node == nil {
return false, fmt.Errorf("unknown node: %s", id)
}
@ -179,7 +159,7 @@ func testDiscoverySimulation(t *testing.T, setup func(net *simulations.Network))
// triggerChecks triggers a simulation step check whenever a peer is added or
// removed from the given node
func triggerChecks(trigger chan *adapters.NodeId, net *simulations.Network, id *adapters.NodeId) error {
node := net.GetNodeAdapter(id)
node := net.GetNode(id)
if node == nil {
return fmt.Errorf("unknown node: %s", id)
}

View file

@ -14,10 +14,11 @@ import (
"time"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/simulations"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
"github.com/ethereum/go-ethereum/swarm/network"
)
@ -61,8 +62,7 @@ func (self *SimNode) Stop() error {
}
// NewSimNode creates adapters for nodes in the simulation.
func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapter {
id := conf.Id
func (self *Network) NewSimNode(id *adapters.NodeId) node.Service {
addr := network.NewPeerAddrFromNodeId(id)
kp := network.NewKadParams()
@ -91,21 +91,17 @@ func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapt
ct := network.BzzCodeMap(network.DiscoveryMsgs...) // bzz protocol code map
node := &SimNode{
return &SimNode{
hive: pp,
protocol: network.Bzz(addr.OverlayAddr(), addr.UnderlayAddr(), ct, services, nil, nil),
}
return adapters.NewSimNode(id, node, self.Network)
}
func NewNetwork(net *simulations.Network) *Network {
n := &Network{
return &Network{
Network: net,
hives: make(map[discover.NodeID]*network.Hive),
}
n.SetNaf(n.NewSimNode)
return n
}
func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simulations.ResourceController) {
@ -116,7 +112,14 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
conf.DefaultMockerConfig.DegreeTarget = 0
conf.Id = "0"
conf.Backend = true
net := NewNetwork(simulations.NewNetwork(conf))
conf.DefaultService = "overlay"
net := &Network{
hives: make(map[discover.NodeID]*network.Hive),
}
services := map[string]adapters.ServiceFunc{"overlay": net.NewSimNode}
adapter := adapters.NewSimAdapter(services)
net.Network = simulations.NewNetwork(adapter, conf)
ids := make([]*adapters.NodeId, 10)
for i := 0; i < 10; i++ {

View file

@ -32,7 +32,7 @@ import (
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/adapters"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
"github.com/ethereum/go-ethereum/rpc"
"github.com/ethereum/go-ethereum/swarm/api"
httpapi "github.com/ethereum/go-ethereum/swarm/api/http"
@ -43,22 +43,21 @@ import (
// the swarm stack
type Swarm struct {
config *api.Config // swarm configuration
api *api.Api // high level api layer (fs/manifest)
dns api.Resolver // DNS registrar
storage storage.ChunkStore // internal access to storage, common interface to cloud storage backends
dpa *storage.DPA // distributed preimage archive, the local API to the storage with document level storage/retrieval support
cloud storage.CloudStore // procurement, cloud storage backend (can multi-cloud)
hive *network.Hive // the logistic manager
backend chequebook.Backend // simple blockchain Backend
config *api.Config // swarm configuration
api *api.Api // high level api layer (fs/manifest)
dns api.Resolver // DNS registrar
storage storage.ChunkStore // internal access to storage, common interface to cloud storage backends
dpa *storage.DPA // distributed preimage archive, the local API to the storage with document level storage/retrieval support
cloud storage.CloudStore // procurement, cloud storage backend (can multi-cloud)
hive *network.Hive // the logistic manager
backend chequebook.Backend // simple blockchain Backend
privateKey *ecdsa.PrivateKey
corsString string
swapEnabled bool
pssEnabled bool
pss *network.Pss
pss *network.Pss
lstore *storage.LocalStore // local store, needs to store for releasing resources after node stopped
sfs *fuse.SwarmFS // need this to cleanup all the active mounts on node exit
na *adapters.NodeAdapter
}
type SwarmAPI struct {
@ -91,7 +90,7 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api.
backend: backend,
privateKey: config.Swap.PrivateKey(),
corsString: cors,
pssEnabled: pssEnabled,
pssEnabled: pssEnabled,
}
log.Debug(fmt.Sprintf("Setting up Swarm service components"))
@ -102,20 +101,20 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api.
}
log.Debug("Set up local db access (iterator/counter)")
kp := network.NewKadParams()
to := network.NewKademlia(
common.FromHex(config.BzzKey),
kp,
)
// set up the kademlia hive
self.hive = network.NewHive(
config.HiveParams, // configuration parameters
config.HiveParams, // configuration parameters
to,
)
log.Debug(fmt.Sprintf("Set up swarm network with Kademlia hive"))
// setup cloud storage internal access layer
self.storage = storage.NewNetStore(hash, self.lstore, nil, config.StoreParams)
log.Debug("-> swarm net store shared access layer to Swarm Chunk Store")
@ -142,7 +141,7 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api.
log.Debug(fmt.Sprintf("-> Swarm Domain Name Registrar @ address %v", config.EnsRoot.Hex()))
self.api = api.NewApi(self.dpa, self.dns)
// Manifests for Smart Hosting
log.Debug(fmt.Sprintf("-> Web3 virtual server API"))
@ -180,7 +179,7 @@ func (self *Swarm) Start(net p2p.Server) error {
self.hive.Start(
net,
func () <-chan time.Time{
func() <-chan time.Time {
return time.NewTicker(time.Second).C
},
)
@ -190,14 +189,14 @@ func (self *Swarm) Start(net p2p.Server) error {
if self.pssEnabled {
pssparams := network.NewPssParams()
self.pss = network.NewPss(self.hive.Overlay, pssparams)
// for testing purposes, shold be removed in production environment!!
pingtopic, _ := network.MakeTopic("pss", 1)
self.pss.Register(pingtopic, self.pss.GetPingHandler())
log.Debug("Pss started: %v", self.pss)
}
self.dpa.Start()
log.Debug(fmt.Sprintf("Swarm DPA started"))
@ -245,7 +244,7 @@ func (self *Swarm) Protocols() []p2p.Protocol {
if self.pssEnabled {
ct.Register(&network.PssMsg{})
}
srv := func(p network.Peer) error {
if self.pssEnabled {
p.Register(&network.PssMsg{}, func(msg interface{}) error {
@ -275,7 +274,7 @@ func (self *Swarm) Protocols() []p2p.Protocol {
})
return nil
}
proto := network.Bzz(
self.hive.Overlay.GetAddr().OverlayAddr(),
self.hive.Overlay.GetAddr().UnderlayAddr(),
@ -284,7 +283,7 @@ func (self *Swarm) Protocols() []p2p.Protocol {
nil,
nil,
)
return []p2p.Protocol{*proto}
}
@ -339,12 +338,12 @@ func (self *Swarm) APIs() []rpc.API {
if self.pssEnabled {
apis = append(apis, rpc.API{
Namespace: "eth",
Version: "0.1/pss",
Service: network.NewPssApi(self.pss),
Public: true,
Version: "0.1/pss",
Service: network.NewPssApi(self.pss),
Public: true,
})
}
return apis
}
@ -389,7 +388,7 @@ func NewLocalSwarm(datadir, port string) (self *Swarm, err error) {
return
}
// serialisable info about swarm
type Info struct {
*api.Config