Merge pull request #103 from ethersphere/network-testing-framework-p2psim

Simulation framework improvements
This commit is contained in:
Lewis Marshall 2017-06-21 09:34:45 +02:00 committed by GitHub
commit 3899c3e086
6 changed files with 156 additions and 98 deletions

View file

@ -13,6 +13,7 @@ import (
"github.com/docker/docker/pkg/reexec"
"github.com/ethereum/go-ethereum/node"
"github.com/ethereum/go-ethereum/p2p/discover"
)
// DockerAdapter is a NodeAdapter which runs nodes inside Docker containers.
@ -20,7 +21,9 @@ import (
// 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{}
type DockerAdapter struct {
ExecAdapter
}
// NewDockerAdapter builds the p2p-node Docker image containing the current
// binary and returns a DockerAdapter
@ -33,7 +36,11 @@ func NewDockerAdapter() (*DockerAdapter, error) {
return nil, err
}
return &DockerAdapter{}, nil
return &DockerAdapter{
ExecAdapter{
nodes: make(map[discover.NodeID]*ExecNode),
},
}, nil
}
// Name returns the name of the adapter for logging purpoeses
@ -58,17 +65,22 @@ func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) {
Node: config,
}
conf.Stack.DataDir = "/data"
conf.Stack.WSHost = "0.0.0.0"
conf.Stack.WSOrigins = []string{"*"}
conf.Stack.WSExposeAll = true
conf.Stack.P2P.EnableMsgEvents = true
conf.Stack.P2P.NoDiscovery = true
conf.Stack.P2P.NAT = nil
node := &DockerNode{
ExecNode: ExecNode{
ID: config.ID,
Config: conf,
ID: config.ID,
Config: conf,
adapter: &d.ExecAdapter,
},
}
node.newCmd = node.dockerCommand
d.ExecAdapter.nodes[node.ID] = &node.ExecNode
return node, nil
}

View file

@ -183,25 +183,18 @@ func (n *ExecNode) Start(snapshots map[string][]byte) (err error) {
// read the WebSocket address from the stderr logs
var wsAddr string
errC := make(chan error)
wsAddrC := make(chan string)
go func() {
s := bufio.NewScanner(stderrR)
for s.Scan() {
if strings.Contains(s.Text(), "WebSocket endpoint opened:") {
wsAddr = wsAddrPattern.FindString(s.Text())
break
wsAddrC <- wsAddrPattern.FindString(s.Text())
}
}
select {
case errC <- s.Err():
default:
}
}()
select {
case err := <-errC:
if err != nil {
return fmt.Errorf("error reading WebSocket address from stderr: %s", err)
} else if wsAddr == "" {
case wsAddr = <-wsAddrC:
if wsAddr == "" {
return errors.New("failed to read WebSocket address from stderr")
}
case <-time.After(10 * time.Second):
@ -354,17 +347,24 @@ func execP2PNode() {
conf.Stack.P2P.PrivateKey = conf.Node.PrivateKey
// use explicit IP address in ListenAddr so that Enode URL is usable
if strings.HasPrefix(conf.Stack.P2P.ListenAddr, ":") {
externalIP := func() string {
addrs, err := net.InterfaceAddrs()
if err != nil {
log.Crit("error getting IP address", "err", err)
}
for _, addr := range addrs {
if ip, ok := addr.(*net.IPNet); ok && !ip.IP.IsLoopback() {
conf.Stack.P2P.ListenAddr = ip.IP.String() + conf.Stack.P2P.ListenAddr
break
return ip.IP.String()
}
}
log.Crit("unable to determine explicit IP address")
return ""
}
if strings.HasPrefix(conf.Stack.P2P.ListenAddr, ":") {
conf.Stack.P2P.ListenAddr = externalIP() + conf.Stack.P2P.ListenAddr
}
if conf.Stack.WSHost == "0.0.0.0" {
conf.Stack.WSHost = externalIP()
}
// initialize the devp2p stack

View file

@ -26,20 +26,26 @@ package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"strings"
"text/tabwriter"
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/p2p"
"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"
"gopkg.in/urfave/cli.v1"
)
var client *simulations.Client
var (
client *simulations.Client
networkID string
)
func main() {
app := cli.NewApp()
@ -73,9 +79,14 @@ func main() {
Action: createNetwork,
Flags: []cli.Flag{
cli.StringFlag{
Name: "config",
Value: "{}",
Usage: "JSON encoded network config",
Name: "id",
Value: "",
Usage: "network ID",
},
cli.StringFlag{
Name: "default-service",
Value: "",
Usage: "default service",
},
},
},
@ -109,59 +120,83 @@ func main() {
Name: "node",
Usage: "manage simulation nodes",
Action: listNodes,
Flags: []cli.Flag{
cli.StringFlag{
Name: "network",
Usage: "simulation network",
},
},
Before: func(ctx *cli.Context) error {
networkID = ctx.GlobalString("network")
if networkID == "" {
networkID = os.Getenv("P2PSIM_NETWORK")
}
if networkID == "" {
return errors.New("missing network, set with --network or P2PSIM_NETWORK")
}
return nil
},
Subcommands: []cli.Command{
{
Name: "list",
ArgsUsage: "<network>",
Usage: "list nodes",
Action: listNodes,
Name: "list",
Usage: "list nodes",
Action: listNodes,
},
{
Name: "create",
ArgsUsage: "<network>",
Usage: "create a node",
Action: createNode,
Name: "create",
Usage: "create a node",
Action: createNode,
Flags: []cli.Flag{
cli.StringFlag{
Name: "config",
Value: "{}",
Usage: "JSON encoded node config",
Name: "name",
Value: "",
Usage: "node name",
},
cli.StringFlag{
Name: "services",
Value: "",
Usage: "node services (comma separated)",
},
cli.StringFlag{
Name: "key",
Value: "",
Usage: "node private key (hex encoded)",
},
},
},
{
Name: "show",
ArgsUsage: "<network> <node>",
ArgsUsage: "<node>",
Usage: "show node information",
Action: showNode,
},
{
Name: "start",
ArgsUsage: "<network> <node>",
ArgsUsage: "<node>",
Usage: "start a node",
Action: startNode,
},
{
Name: "stop",
ArgsUsage: "<network> <node>",
ArgsUsage: "<node>",
Usage: "stop a node",
Action: stopNode,
},
{
Name: "connect",
ArgsUsage: "<network> <node> <peer>",
ArgsUsage: "<node> <peer>",
Usage: "connect a node to a peer node",
Action: connectNode,
},
{
Name: "disconnect",
ArgsUsage: "<network> <node> <peer>",
ArgsUsage: "<node> <peer>",
Usage: "disconnect a node from a peer node",
Action: disconnectNode,
},
{
Name: "rpc",
ArgsUsage: "<network> <node> <method> [<args>]",
ArgsUsage: "<node> <method> [<args>]",
Usage: "call a node RPC method",
Action: rpcNode,
Flags: []cli.Flag{
@ -198,9 +233,9 @@ func createNetwork(ctx *cli.Context) error {
if len(ctx.Args()) != 0 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
config := &simulations.NetworkConfig{}
if err := json.Unmarshal([]byte(ctx.String("config")), config); err != nil {
return err
config := &simulations.NetworkConfig{
ID: ctx.String("id"),
DefaultService: ctx.String("default-service"),
}
network, err := client.CreateNetwork(config)
if err != nil {
@ -280,11 +315,9 @@ func loadSnapshot(ctx *cli.Context) error {
}
func listNodes(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 1 {
if len(ctx.Args()) != 0 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodes, err := client.GetNodes(networkID)
if err != nil {
return err
@ -307,14 +340,22 @@ func protocolList(node *p2p.NodeInfo) []string {
}
func createNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 1 {
if len(ctx.Args()) != 0 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
config := &adapters.NodeConfig{}
if err := json.Unmarshal([]byte(ctx.String("config")), config); err != nil {
return err
config := &adapters.NodeConfig{
Name: ctx.String("name"),
}
if key := ctx.String("key"); key != "" {
privKey, err := crypto.HexToECDSA(key)
if err != nil {
return err
}
config.ID = discover.PubkeyID(&privKey.PublicKey)
config.PrivateKey = privKey
}
if services := ctx.String("services"); services != "" {
config.Services = strings.Split(services, ",")
}
node, err := client.CreateNode(networkID, config)
if err != nil {
@ -326,11 +367,10 @@ func createNode(ctx *cli.Context) error {
func showNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 2 {
if len(args) != 1 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodeName := args[1]
nodeName := args[0]
node, err := client.GetNode(networkID, nodeName)
if err != nil {
return err
@ -352,11 +392,10 @@ func showNode(ctx *cli.Context) error {
func startNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 2 {
if len(args) != 1 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodeName := args[1]
nodeName := args[0]
if err := client.StartNode(networkID, nodeName); err != nil {
return err
}
@ -366,11 +405,10 @@ func startNode(ctx *cli.Context) error {
func stopNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 2 {
if len(args) != 1 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodeName := args[1]
nodeName := args[0]
if err := client.StopNode(networkID, nodeName); err != nil {
return err
}
@ -380,12 +418,11 @@ func stopNode(ctx *cli.Context) error {
func connectNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 3 {
if len(args) != 2 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodeName := args[1]
peerName := args[2]
nodeName := args[0]
peerName := args[1]
if err := client.ConnectNode(networkID, nodeName, peerName); err != nil {
return err
}
@ -395,12 +432,11 @@ func connectNode(ctx *cli.Context) error {
func disconnectNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) != 3 {
if len(args) != 2 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodeName := args[1]
peerName := args[2]
nodeName := args[0]
peerName := args[1]
if err := client.DisconnectNode(networkID, nodeName, peerName); err != nil {
return err
}
@ -410,12 +446,11 @@ func disconnectNode(ctx *cli.Context) error {
func rpcNode(ctx *cli.Context) error {
args := ctx.Args()
if len(args) < 3 {
if len(args) < 2 {
return cli.ShowCommandHelp(ctx, ctx.Command.Name)
}
networkID := args[0]
nodeName := args[1]
method := args[2]
nodeName := args[0]
method := args[1]
rpcClient, err := client.RPCClient(context.Background(), networkID, nodeName)
if err != nil {
return err

View file

@ -67,7 +67,9 @@ func main() {
}
log.Info("starting simulation server on 0.0.0.0:8888...")
http.ListenAndServe(":8888", simulations.NewServer(config))
if err := http.ListenAndServe(":8888", simulations.NewServer(config)); err != nil {
log.Crit("error starting simulation server", "err", err)
}
}
// pingPongService runs a ping-pong protocol between nodes where each node

View file

@ -10,17 +10,18 @@ main() {
fi
info "creating the example network"
p2psim network create --config '{"id": "example", "default_service": "ping-pong"}'
export P2PSIM_NETWORK="example"
p2psim network create --id "${P2PSIM_NETWORK}"
info "creating 10 nodes"
for i in $(seq 1 10); do
p2psim node create "example"
p2psim node start "example" "$(node_name $i)"
p2psim node create --name "$(node_name $i)" --services "ping-pong"
p2psim node start "$(node_name $i)"
done
info "connecting node01 to all other nodes"
for i in $(seq 2 10); do
p2psim node connect "example" "node01" "$(node_name $i)"
p2psim node connect "node01" "$(node_name $i)"
done
info "done"

View file

@ -201,10 +201,21 @@ func (self *Network) NewNode() (*Node, error) {
func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) (*Node, error) {
self.lock.Lock()
defer self.lock.Unlock()
if conf.ID == (discover.NodeID{}) {
c := adapters.RandomNodeConfig()
conf.ID = c.ID
conf.PrivateKey = c.PrivateKey
}
id := conf.ID
if node := self.getNode(id); node != nil {
return nil, fmt.Errorf("node already exists: %q", id)
}
if conf.Name == "" {
conf.Name = fmt.Sprintf("node%02d", len(self.Nodes)+1)
}
if node := self.getNodeByName(conf.Name); node != nil {
return nil, fmt.Errorf("node already exists: %q", conf.Name)
}
if len(conf.Services) == 0 {
conf.Services = []string{self.DefaultService}
}
@ -326,7 +337,16 @@ func (self *Network) startWithSnapshots(id discover.NodeID, snapshots map[string
}
func (self *Network) watchPeerEvents(id discover.NodeID, events chan *p2p.PeerEvent, sub event.Subscription) {
defer sub.Unsubscribe()
defer func() {
sub.Unsubscribe()
// assume the node is now down
self.lock.Lock()
node := self.getNode(id)
node.Up = false
self.lock.Unlock()
self.events.Send(NewEvent(node))
}()
for {
select {
case event, ok := <-events:
@ -336,21 +356,13 @@ func (self *Network) watchPeerEvents(id discover.NodeID, events chan *p2p.PeerEv
peer := event.Peer
switch event.Type {
case p2p.PeerEventTypeAdd:
if err := self.DidConnect(id, peer); err != nil {
log.Error(fmt.Sprintf("error generating connection up event %s => %s", id.TerminalString(), peer.TerminalString()), "err", err)
}
self.DidConnect(id, peer)
case p2p.PeerEventTypeDrop:
if err := self.DidDisconnect(id, peer); err != nil {
log.Error(fmt.Sprintf("error generating connection down event %s => %s", id.TerminalString(), peer.TerminalString()), "err", err)
}
self.DidDisconnect(id, peer)
case p2p.PeerEventTypeMsgSend:
if err := self.DidSend(id, peer, *event.MsgCode); err != nil {
log.Error(fmt.Sprintf("error generating msg send event %s => %s", id.TerminalString(), peer.TerminalString()), "err", err)
}
self.DidSend(id, peer, *event.MsgCode)
case p2p.PeerEventTypeMsgRecv:
if err := self.DidReceive(peer, id, *event.MsgCode); err != nil {
log.Error(fmt.Sprintf("error generating msg receive event %s => %s", peer.TerminalString(), id.TerminalString()), "err", err)
}
self.DidReceive(peer, id, *event.MsgCode)
}
case err := <-sub.Err():
if err != nil {
@ -528,6 +540,10 @@ func (self *Network) GetNode(id discover.NodeID) *Node {
func (self *Network) GetNodeByName(name string) *Node {
self.lock.Lock()
defer self.lock.Unlock()
return self.getNodeByName(name)
}
func (self *Network) getNodeByName(name string) *Node {
for _, node := range self.Nodes {
if node.Config.Name == name {
return node
@ -589,14 +605,6 @@ func (self *Network) getConn(oneID, otherID discover.NodeID) *Conn {
}
func (self *Network) Shutdown() {
// disconnect all nodes
for _, conn := range self.Conns {
log.Debug(fmt.Sprintf("disconnecting %s from %s", conn.One.TerminalString(), conn.Other.TerminalString()))
if err := self.Disconnect(conn.One, conn.Other); err != nil {
log.Warn(fmt.Sprintf("error disconnecting %s from %s", conn.One.TerminalString(), conn.Other.TerminalString()), "err", err)
}
}
// stop all nodes
for _, node := range self.Nodes {
log.Debug(fmt.Sprintf("stopping node %s", node.ID().TerminalString()))