This commit is contained in:
nolash 2017-05-18 00:21:11 +02:00 committed by Lewis Marshall
parent aaeb8497a0
commit d219b0edb0
14 changed files with 225 additions and 234 deletions

View file

@ -9,6 +9,7 @@ import (
"os/exec" "os/exec"
"path/filepath" "path/filepath"
"runtime" "runtime"
"strings"
"github.com/docker/docker/pkg/reexec" "github.com/docker/docker/pkg/reexec"
"github.com/ethereum/go-ethereum/node" "github.com/ethereum/go-ethereum/node"
@ -42,10 +43,8 @@ func (d *DockerAdapter) Name() string {
// NewNode returns a new DockerNode using the given config // NewNode returns a new DockerNode using the given config
func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) { func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) {
for _, name := range config.Services { if _, exists := serviceFuncs[config.Service]; !exists {
if _, exists := serviceFuncs[name]; !exists { return nil, fmt.Errorf("unknown node service %q", config.Service)
return nil, fmt.Errorf("unknown node service %q", name)
}
} }
// generate the config // generate the config
@ -62,7 +61,6 @@ func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) {
ExecNode: ExecNode{ ExecNode: ExecNode{
ID: config.Id, ID: config.Id,
Config: conf, Config: conf,
Services: config.Services,
}, },
} }
node.newCmd = node.dockerCommand node.newCmd = node.dockerCommand
@ -85,7 +83,7 @@ func (n *DockerNode) dockerCommand() *exec.Cmd {
"sh", "-c", "sh", "-c",
fmt.Sprintf( fmt.Sprintf(
`exec docker run --interactive --env _P2P_NODE_CONFIG="${_P2P_NODE_CONFIG}" --env _P2P_NODE_KEY="${_P2P_NODE_KEY}" %s p2p-node %s %s`, `exec docker run --interactive --env _P2P_NODE_CONFIG="${_P2P_NODE_CONFIG}" --env _P2P_NODE_KEY="${_P2P_NODE_KEY}" %s p2p-node %s %s`,
dockerImage, n.Services[0], n.ID.String(), dockerImage, strings.Join(n.Services, " "), n.ID.String(),
), ),
) )
} }

View file

@ -46,10 +46,8 @@ func (e *ExecAdapter) Name() string {
// NewNode returns a new ExecNode using the given config // NewNode returns a new ExecNode using the given config
func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) { func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) {
for _, name := range config.Services { if _, exists := serviceFuncs[config.Service]; !exists {
if _, exists := serviceFuncs[name]; !exists { return nil, fmt.Errorf("unknown node service %q", config.Service)
return nil, fmt.Errorf("unknown node service %q", name)
}
} }
// create the node directory using the first 12 characters of the ID // create the node directory using the first 12 characters of the ID
@ -77,7 +75,6 @@ func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) {
ID: config.Id, ID: config.Id,
Dir: dir, Dir: dir,
Config: conf, Config: conf,
Services: config.Services,
} }
node.newCmd = node.execCommand node.newCmd = node.execCommand
return node, nil return node, nil
@ -91,11 +88,11 @@ func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) {
// (so for example we can run the node in a remote Docker container and // (so for example we can run the node in a remote Docker container and
// still communicate with it). // still communicate with it).
type ExecNode struct { type ExecNode struct {
ID *NodeId ID *NodeId
Dir string Dir string
Config *execNodeConfig Config *execNodeConfig
Cmd *exec.Cmd Cmd *exec.Cmd
Info *p2p.NodeInfo Info *p2p.NodeInfo
Services []string Services []string
client *rpc.Client client *rpc.Client
@ -168,7 +165,6 @@ func (n *ExecNode) Start(snapshot []byte) (err error) {
return nil return nil
} }
func (n *ExecNode) GetService(name string) node.Service { func (n *ExecNode) GetService(name string) node.Service {
return nil return nil
} }
@ -287,7 +283,7 @@ func execP2PNode() {
if !exists { if !exists {
log.Crit(fmt.Sprintf("unknown node service %q", serviceName)) log.Crit(fmt.Sprintf("unknown node service %q", serviceName))
} }
service := serviceFunc(id, conf.Snapshot) services := serviceFunc(id, conf.Snapshot)
// use explicit IP address in ListenAddr so that Enode URL is usable // use explicit IP address in ListenAddr so that Enode URL is usable
if strings.HasPrefix(conf.Stack.P2P.ListenAddr, ":") { if strings.HasPrefix(conf.Stack.P2P.ListenAddr, ":") {
@ -304,7 +300,7 @@ func execP2PNode() {
} }
// start the devp2p stack // start the devp2p stack
stack, err := startP2PNode(&conf.Stack, service) stack, err := startP2PNode(&conf.Stack, services)
if err != nil { if err != nil {
log.Crit("error starting p2p node", "err", err) log.Crit("error starting p2p node", "err", err)
} }
@ -330,17 +326,20 @@ func execP2PNode() {
stack.Wait() stack.Wait()
} }
func startP2PNode(conf *node.Config, service node.Service) (*node.Node, error) { func startP2PNode(conf *node.Config, services []node.Service) (*node.Node, error) {
stack, err := node.New(conf) stack, err := node.New(conf)
if err != nil { if err != nil {
return nil, err return nil, err
} }
constructor := func(ctx *node.ServiceContext) (node.Service, error) { for _, svc := range services {
return &snapshotService{service}, nil constructor := func(ctx *node.ServiceContext) (node.Service, error) {
} return &snapshotService{svc}, nil
if err := stack.Register(constructor); err != nil { }
return nil, err if err := stack.Register(constructor); err != nil {
return nil, err
}
} }
if err := stack.Start(); err != nil { if err := stack.Start(); err != nil {
return nil, err return nil, err
} }

View file

@ -21,6 +21,7 @@ import (
"fmt" "fmt"
"math" "math"
"net" "net"
"reflect"
"sync" "sync"
"github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/event"
@ -36,7 +37,6 @@ import (
type SimAdapter struct { type SimAdapter struct {
mtx sync.RWMutex mtx sync.RWMutex
nodes map[discover.NodeID]*SimNode nodes map[discover.NodeID]*SimNode
services map[string]ServiceFunc
} }
// NewSimAdapter creates a SimAdapter which is capable of running in-memory // NewSimAdapter creates a SimAdapter which is capable of running in-memory
@ -45,7 +45,6 @@ type SimAdapter struct {
func NewSimAdapter(services map[string]ServiceFunc) *SimAdapter { func NewSimAdapter(services map[string]ServiceFunc) *SimAdapter {
return &SimAdapter{ return &SimAdapter{
nodes: make(map[discover.NodeID]*SimNode), nodes: make(map[discover.NodeID]*SimNode),
services: services,
} }
} }
@ -101,28 +100,18 @@ func (s *SimAdapter) NewNode(config *NodeConfig) (Node, error) {
return nil, err return nil, err
} }
servicefuncs := make(map[string]ServiceFunc) for _, service := range serviceFuncs[config.Service](id, nil) {
for name, servicefunc := range s.services {
service := servicefunc(id, nil)
/*if err := n.Register(func(ctx *node.ServiceContext) (node.Service, error) {
return service, err
}); err != nil {
return nil, err
}*/
for _, proto := range service.Protocols() { for _, proto := range service.Protocols() {
nodeprotos = append(nodeprotos, proto) nodeprotos = append(nodeprotos, proto)
} }
servicefuncs[name] = servicefunc
} }
simnode := &SimNode{ simnode := &SimNode{
//node: n,
Id: id, Id: id,
serviceFuncs: servicefuncs, serviceFunc: serviceFuncs[config.Service],
adapter: s, adapter: s,
config: config, config: config,
running: make(map[string]node.Service), running: []node.Service{},
} }
s.nodes[id.NodeID] = simnode s.nodes[id.NodeID] = simnode
return simnode, nil return simnode, nil
@ -161,11 +150,11 @@ type SimNode struct {
Id *NodeId Id *NodeId
config *NodeConfig config *NodeConfig
adapter *SimAdapter adapter *SimAdapter
serviceFuncs map[string]ServiceFunc serviceFunc ServiceFunc
node *node.Node node *node.Node
running map[string]node.Service
client *rpc.Client client *rpc.Client
rpcMux *rpcMux rpcMux *rpcMux
running []node.Service
} }
// Addr returns the node's discovery address // Addr returns the node's discovery address
@ -224,14 +213,16 @@ func (self *SimNode) Start(snapshot []byte) error {
services := []node.ServiceConstructor{} services := []node.ServiceConstructor{}
sf := self.serviceFunc(self.Id, snapshot)
// so we can control the order of the services if we need for i, _ := range sf {
for _, name := range self.config.Services { service := sf[i]
service := self.serviceFuncs[name](self.Id, snapshot) sc := func(ctx *node.ServiceContext) (node.Service, error) {
services = append(services, func(ctx *node.ServiceContext) (node.Service, error) {
self.running[name] = service
return service, nil return service, nil
}) }
log.Debug(fmt.Sprintf("servicefunc yield: %v %p %p", reflect.TypeOf(sf[i]), sf[i], sc))
services = append(services, sc)
self.running = append(self.running, sf[i])
} }
node, err := node.New(&node.Config{ node, err := node.New(&node.Config{
@ -249,7 +240,7 @@ func (self *SimNode) Start(snapshot []byte) error {
} }
for _, service := range services { for _, service := range services {
log.Debug("registering service", "service", service) log.Debug(fmt.Sprintf("service %v", service))
if err := node.Register(service); err != nil { if err := node.Register(service); err != nil {
return err return err
} }
@ -290,6 +281,19 @@ func (self *SimNode) Stop() error {
return nil return nil
} }
// Service returns the underlying running node.Service matching the supplied servuce type
func (self *SimNode) Service(servicetype interface{}) node.Service {
self.lock.Lock()
defer self.lock.Unlock()
typ := reflect.TypeOf(servicetype)
for _, service := range self.running {
if reflect.TypeOf(service) == typ {
return service
}
}
return nil
}
func (self *SimNode) Server() *p2p.Server { func (self *SimNode) Server() *p2p.Server {
self.lock.Lock() self.lock.Lock()
defer self.lock.Unlock() defer self.lock.Unlock()
@ -299,12 +303,6 @@ func (self *SimNode) Server() *p2p.Server {
return self.node.Server() return self.node.Server()
} }
// Service returns a underlying node.Service of the speficied type
func (self *SimNode) GetService(servicename string) node.Service {
log.Warn("retrieving service", "name", servicename)
return self.running[servicename]
}
func (self *SimNode) SubscribeEvents(ch chan *p2p.PeerEvent) event.Subscription { func (self *SimNode) SubscribeEvents(ch chan *p2p.PeerEvent) event.Subscription {
srv := self.Server() srv := self.Server()
if srv == nil { if srv == nil {

View file

@ -62,9 +62,6 @@ type Node interface {
// Snapshot creates a snapshot of the running service // Snapshot creates a snapshot of the running service
Snapshot() ([]byte, error) Snapshot() ([]byte, error)
// Gets a service by name
GetService(string) node.Service
} }
// NodeAdapter is an object which creates Nodes to be used in a simulation // NodeAdapter is an object which creates Nodes to be used in a simulation
@ -130,11 +127,11 @@ type NodeConfig struct {
// Name is a human friendly name for the node like "node01" // Name is a human friendly name for the node like "node01"
Name string Name string
// Services is the name of the services which should be run when starting // Service is the name of the services which should be run when starting
// the node (for SimNodes it should be the names of services contained // the node (for SimNodes it should be the names of services contained
// in SimAdapter.services, for other nodes it should be services // in SimAdapter.services, for other nodes it should be services
// registered by calling the RegisterService function) // registered by calling the RegisterService function)
Services []string Service string
} }
// nodeConfigJSON is used to encode and decode NodeConfig as JSON by converting // nodeConfigJSON is used to encode and decode NodeConfig as JSON by converting
@ -143,13 +140,13 @@ type nodeConfigJSON struct {
Id string `json:"id"` Id string `json:"id"`
PrivateKey string `json:"private_key"` PrivateKey string `json:"private_key"`
Name string `json:"name"` Name string `json:"name"`
Services []string `json:"services"` Service string `json:"service"`
} }
func (n *NodeConfig) MarshalJSON() ([]byte, error) { func (n *NodeConfig) MarshalJSON() ([]byte, error) {
confJSON := nodeConfigJSON{ confJSON := nodeConfigJSON{
Name: n.Name, Name: n.Name,
Services: n.Services, Service: n.Service,
} }
if n.Id != nil { if n.Id != nil {
confJSON.Id = n.Id.String() confJSON.Id = n.Id.String()
@ -183,7 +180,7 @@ func (n *NodeConfig) UnmarshalJSON(data []byte) error {
} }
n.Name = confJSON.Name n.Name = confJSON.Name
n.Services = confJSON.Services n.Service = confJSON.Service
return nil return nil
} }

View file

@ -190,7 +190,7 @@ func (self *Msg) String() string {
// NewNode adds a new node to the network with a random ID // NewNode adds a new node to the network with a random ID
func (self *Network) NewNode() (*Node, error) { func (self *Network) NewNode() (*Node, error) {
conf := adapters.RandomNodeConfig() conf := adapters.RandomNodeConfig()
conf.Services = append(conf.Services, self.DefaultService) conf.Service = self.DefaultService
return self.NewNodeWithConfig(conf) return self.NewNodeWithConfig(conf)
} }
@ -203,8 +203,8 @@ func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) (*Node, error)
if conf.Name == "" { if conf.Name == "" {
conf.Name = fmt.Sprintf("node%02d", len(self.Nodes)+1) conf.Name = fmt.Sprintf("node%02d", len(self.Nodes)+1)
} }
if len(conf.Services) == 0 { if conf.Service == "" {
conf.Services = append(conf.Services, self.DefaultService) conf.Service = self.DefaultService
} }
_, found := self.nodeMap[id.NodeID] _, found := self.nodeMap[id.NodeID]

View file

@ -57,7 +57,7 @@ func (self *ProtocolSession) trigger(trig Trigger) error {
if !ok { if !ok {
return fmt.Errorf("trigger: peer %v does not exist (1- %v)", trig.Peer, len(self.Ids)) return fmt.Errorf("trigger: peer %v does not exist (1- %v)", trig.Peer, len(self.Ids))
} }
mockNode, ok := simNode.GetService("mock").(*mockNode) mockNode, ok := simNode.Service(&mockNode{}).(*mockNode)
if !ok { if !ok {
return fmt.Errorf("trigger: peer %v is not a mock", trig.Peer) return fmt.Errorf("trigger: peer %v is not a mock", trig.Peer)
} }
@ -92,7 +92,7 @@ func (self *ProtocolSession) expect(exp Expect) error {
if !ok { if !ok {
return fmt.Errorf("trigger: peer %v does not exist (1- %v)", exp.Peer, len(self.Ids)) return fmt.Errorf("trigger: peer %v does not exist (1- %v)", exp.Peer, len(self.Ids))
} }
mockNode, ok := simNode.GetService("mock").(*mockNode) mockNode, ok := simNode.Service(&mockNode{}).(*mockNode)
if !ok { if !ok {
return fmt.Errorf("trigger: peer %v is not a mock", exp.Peer) return fmt.Errorf("trigger: peer %v is not a mock", exp.Peer)
} }

View file

@ -18,24 +18,20 @@ type ProtocolTester struct {
network *simulations.Network network *simulations.Network
} }
func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, moreservices adapters.Services, run func(*p2p.Peer, p2p.MsgReadWriter) error) *ProtocolTester { func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(*p2p.Peer, p2p.MsgReadWriter) error) *ProtocolTester {
//func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(*p2p.Peer, p2p.MsgReadWriter) error) *ProtocolTester { //func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(*p2p.Peer, p2p.MsgReadWriter) error) *ProtocolTester {
moreservicesstring := []string{} services := adapters.Services {
services := map[string]adapters.ServiceFunc{ "test": func(id *adapters.NodeId, _ []byte) []node.Service {
"test": func(id *adapters.NodeId, _ []byte) node.Service { return []node.Service{&testNode{run}}
return &testNode{run}
}, },
"mock": func(id *adapters.NodeId, _ []byte) node.Service { "mock": func(id *adapters.NodeId, _ []byte) []node.Service {
return newMockNode() return []node.Service{newMockNode()}
}, },
} }
for name, service := range moreservices { adapters.RegisterServices(services)
services[name] = service
moreservicesstring = append(moreservicesstring, name)
}
adapter := adapters.NewSimAdapter(services) adapter := adapters.NewSimAdapter(services)
net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{}) net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{})
if _, err := net.NewNodeWithConfig(&adapters.NodeConfig{Id: id, Services: append(moreservicesstring, "test")}); err != nil { if _, err := net.NewNodeWithConfig(&adapters.NodeConfig{Id: id, Service: "test"}); err != nil {
panic(err.Error()) panic(err.Error())
} }
if err := net.Start(id); err != nil { if err := net.Start(id); err != nil {
@ -47,8 +43,7 @@ func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, moreservices ad
peerIDs := make([]*adapters.NodeId, n) peerIDs := make([]*adapters.NodeId, n)
for i := 0; i < n; i++ { for i := 0; i < n; i++ {
peers[i] = adapters.RandomNodeConfig() peers[i] = adapters.RandomNodeConfig()
peers[i].Services = moreservicesstring peers[i].Service = "mock"
peers[i].Services = append(peers[i].Services, "mock")
peerIDs[i] = peers[i].Id peerIDs[i] = peers[i].Id
} }
events := make(chan *p2p.PeerEvent, 1000) events := make(chan *p2p.PeerEvent, 1000)

View file

@ -80,7 +80,7 @@ func NewHiveParams() *HiveParams {
type Hive struct { type Hive struct {
*HiveParams // settings *HiveParams // settings
Overlay // the overlay topology driver Overlay // the overlay topology driver
store Store store StateStore
// bookkeeping // bookkeeping
lock sync.Mutex lock sync.Mutex
@ -93,7 +93,7 @@ type Hive struct {
// Hive constructor embeds both arguments // Hive constructor embeds both arguments
// HiveParams: config parameters // HiveParams: config parameters
// Overlay: Topology Driver Interface // Overlay: Topology Driver Interface
func NewHive(params *HiveParams, overlay Overlay, store Store) *Hive { func NewHive(params *HiveParams, overlay Overlay, store StateStore) *Hive {
return &Hive{ return &Hive{
HiveParams: params, HiveParams: params,
Overlay: overlay, Overlay: overlay,

View file

@ -85,7 +85,7 @@ type Conn interface {
} }
// TODO: implement store for exec nodes // TODO: implement store for exec nodes
type Store interface { type StateStore interface {
Load(string) ([]byte, error) Load(string) ([]byte, error)
Save(string, []byte) error Save(string, []byte) error
} }
@ -106,7 +106,7 @@ type Bzz struct {
} }
// NewBzz is the swarm protocol constructor // NewBzz is the swarm protocol constructor
func NewBzz(config *BzzConfig, kad Overlay, store Store) *Bzz { func NewBzz(config *BzzConfig, kad Overlay, store StateStore) *Bzz {
return &Bzz{ return &Bzz{
Hive: NewHive(config.HiveParams, kad, store), Hive: NewHive(config.HiveParams, kad, store),
localAddr: &bzzAddr{config.OverlayAddr, config.UnderlayAddr}, localAddr: &bzzAddr{config.OverlayAddr, config.UnderlayAddr},

View file

@ -22,33 +22,14 @@ import (
"github.com/ethereum/go-ethereum/swarm/network" "github.com/ethereum/go-ethereum/swarm/network"
) )
type simStore struct {
m map[string][]byte
}
func (self *simStore) Load(s string) ([]byte, error) {
return self.m[s], nil
}
func (self *simStore) Save(s string, data []byte) error {
self.m[s] = data
return nil
}
func NewSimStore() *simStore {
return &simStore{
make(map[string][]byte),
}
}
type Simulation struct { type Simulation struct {
mtx sync.Mutex mtx sync.Mutex
stores map[discover.NodeID]*simStore stores map[discover.NodeID]*adapters.stateStore
} }
func NewSimulation() *Simulation { func NewSimulation() *Simulation {
return &Simulation{ return &Simulation{
stores: make(map[discover.NodeID]*simStore), stores: make(map[discover.NodeID]*adapters.stateStore),
} }
} }

View file

@ -1,11 +1,13 @@
package pss package pss
import ( import (
"fmt"
"io/ioutil" "io/ioutil"
"os" "os"
"time" "time"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/protocols" "github.com/ethereum/go-ethereum/p2p/protocols"
"github.com/ethereum/go-ethereum/swarm/network" "github.com/ethereum/go-ethereum/swarm/network"
"github.com/ethereum/go-ethereum/swarm/storage" "github.com/ethereum/go-ethereum/swarm/storage"
@ -87,3 +89,17 @@ func newPssPingMsg(ps *Pss, spec *protocols.Spec, topic PssTopic, senderaddr []b
return pssmsg return pssmsg
} }
func newPssPingProtocol(handler func (interface{}) error) *p2p.Protocol {
return &p2p.Protocol{
Name: pssPingProtocol.Name,
Version: pssPingProtocol.Version,
Length: uint64(pssPingProtocol.MaxMsgSize),
Run: func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
pp := protocols.NewPeer(p, rw, pssPingProtocol)
log.Trace(fmt.Sprintf("running pss vprotocol on peer %v", p))
err := pp.Run(handler)
return err
},
}
}

View file

@ -191,16 +191,16 @@ func (self *Pss) Protocols() []p2p.Protocol {
Name: pssTransportProtocol.Name, Name: pssTransportProtocol.Name,
Version: pssTransportProtocol.Version, Version: pssTransportProtocol.Version,
Length: pssTransportProtocol.Length(), Length: pssTransportProtocol.Length(),
Run: func(p *p2p.Peer, rw p2p.MsgReadWriter) error { Run: self.Run,
pp := protocols.NewPeer(p, rw, pssTransportProtocol)
err := pp.Run(self.handlePssMsg)
log.Warn("pss protocol peer returned", "peer", p, "err", err)
return nil
},
}, },
} }
} }
func (self *Pss) Run(p *p2p.Peer, rw p2p.MsgReadWriter) error {
pp := protocols.NewPeer(p, rw, pssTransportProtocol)
return pp.Run(self.handlePssMsg)
}
func (self *Pss) APIs() []rpc.API { func (self *Pss) APIs() []rpc.API {
return []rpc.API{ return []rpc.API{
rpc.API { rpc.API {
@ -218,24 +218,24 @@ func (self *Pss) APIs() []rpc.API {
// a topic allows for multiple handlers // a topic allows for multiple handlers
// returns a deregister function which needs to be called to deregister the handler // returns a deregister function which needs to be called to deregister the handler
// (similar to event.Subscription.Unsubscribe()) // (similar to event.Subscription.Unsubscribe())
func (self *Pss) Register(topic PssTopic, handler pssHandler) func() { func (self *Pss) Register(topic *PssTopic, handler pssHandler) func() {
self.lock.Lock() self.lock.Lock()
defer self.lock.Unlock() defer self.lock.Unlock()
handlers := self.handlers[topic] handlers := self.handlers[*topic]
if handlers == nil { if handlers == nil {
handlers = make(map[*pssHandler]bool) handlers = make(map[*pssHandler]bool)
self.handlers[topic] = handlers self.handlers[*topic] = handlers
} }
handlers[&handler] = true handlers[&handler] = true
return func() { self.deregister(topic, &handler) } return func() { self.deregister(topic, &handler) }
} }
func (self *Pss) deregister(topic PssTopic, h *pssHandler) { func (self *Pss) deregister(topic *PssTopic, h *pssHandler) {
self.lock.Lock() self.lock.Lock()
defer self.lock.Unlock() defer self.lock.Unlock()
handlers := self.handlers[topic] handlers := self.handlers[*topic]
if len(handlers) == 1 { if len(handlers) == 1 {
delete(self.handlers, topic) delete(self.handlers, *topic)
return return
} }
delete(handlers, h) delete(handlers, h)
@ -497,14 +497,17 @@ type PssProtocol struct {
} }
// Constructor // Constructor
func NewPssProtocol(pss *Pss, topic *PssTopic, spec *protocols.Spec, targetprotocol *p2p.Protocol) *PssProtocol { //func RegisterPssProtocol(pss *Pss, topic *PssTopic, spec *protocols.Spec, targetprotocol *p2p.Protocol) *PssProtocol {
func RegisterPssProtocol(pss *Pss, topic *PssTopic, spec *protocols.Spec, targetprotocol *p2p.Protocol) error {
pp := &PssProtocol{ pp := &PssProtocol{
Pss: pss, Pss: pss,
proto: targetprotocol, proto: targetprotocol,
topic: topic, topic: topic,
spec: spec, spec: spec,
} }
return pp pss.Register(topic, pp.handle)
//return pp
return nil
} }
func (self *PssProtocol) handle(msg []byte, p *p2p.Peer, senderAddr []byte) error { func (self *PssProtocol) handle(msg []byte, p *p2p.Peer, senderAddr []byte) error {

View file

@ -21,12 +21,10 @@ import (
) )
const ( const (
pssServiceName = "pss" pssServiceName = "pss"
bzzServiceName = "bzz" bzzServiceName = "bzz"
) )
var topic PssTopic = NewTopic(pssPingProtocol.Name, int(pssPingProtocol.Version))
var services = newServices() var services = newServices()
func init() { func init() {
@ -49,9 +47,9 @@ func TestPssCache(t *testing.T) {
fwdaddr := network.RandomAddr() fwdaddr := network.RandomAddr()
msg := &PssMsg{ msg := &PssMsg{
Payload: &PssEnvelope{ Payload: &PssEnvelope{
TTL: 0, TTL: 0,
From: oaddr, From: oaddr,
Topic: topic, Topic: pssPingTopic,
Payload: data, Payload: data,
}, },
To: to, To: to,
@ -59,9 +57,9 @@ func TestPssCache(t *testing.T) {
msgtwo := &PssMsg{ msgtwo := &PssMsg{
Payload: &PssEnvelope{ Payload: &PssEnvelope{
TTL: 0, TTL: 0,
From: oaddr, From: oaddr,
Topic: topic, Topic: pssPingTopic,
Payload: datatwo, Payload: datatwo,
}, },
To: to, To: to,
@ -144,7 +142,7 @@ func TestPssRegisterHandler(t *testing.T) {
} }
return nil return nil
} }
deregister := ps.Register(topic, checkMsg) deregister := ps.Register(&topic, checkMsg)
pssmsg := &PssMsg{Payload: NewPssEnvelope(from.OAddr, topic, payload)} pssmsg := &PssMsg{Payload: NewPssEnvelope(from.OAddr, topic, payload)}
err = ps.Process(pssmsg) err = ps.Process(pssmsg)
if err != nil { if err != nil {
@ -156,7 +154,7 @@ func TestPssRegisterHandler(t *testing.T) {
if err == nil || err.Error() == expErr { if err == nil || err.Error() == expErr {
t.Fatalf("unhandled topic expected '%v', got '%v'", expErr, err) t.Fatalf("unhandled topic expected '%v', got '%v'", expErr, err)
} }
deregister2 := ps.Register(topic, func(msg []byte, p *p2p.Peer, sender []byte) error { i++; return nil }) deregister2 := ps.Register(&topic, func(msg []byte, p *p2p.Peer, sender []byte) error { i++; return nil })
err = ps.Process(pssmsg) err = ps.Process(pssmsg)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
@ -176,33 +174,42 @@ func TestPssRegisterHandler(t *testing.T) {
func TestPssSimpleLinear(t *testing.T) { func TestPssSimpleLinear(t *testing.T) {
nodeconfig := adapters.RandomNodeConfig() nodeconfig := adapters.RandomNodeConfig()
addr := network.NewAddrFromNodeId(nodeconfig.Id) addr := network.NewAddrFromNodeId(nodeconfig.Id)
_ = p2ptest.NewTestPeerPool()
ps := newTestPss(addr.OAddr) ps := newTestPss(addr.OAddr)
ps.Register(pssPingTopic, pssPingHandler)
pt := p2ptest.NewProtocolTester(t, nodeconfig.Id, 2, newServices(), ps.Protocols()[0].Run)
msg := newPssPingMsg(ps, pssPingProtocol, pssPingTopic, []byte{1,2,3}) ping := &pssPing{
quitC: make(chan struct{}),
}
err := RegisterPssProtocol(ps, &pssPingTopic, pssPingProtocol, newPssPingProtocol(ping.pssPingHandler))
if err != nil {
t.Fatalf("Failed to register virtual protocol in pss: %v", err)
}
pt := p2ptest.NewProtocolTester(t, nodeconfig.Id, 2, ps.Run)
msg := newPssPingMsg(ps, pssPingProtocol, pssPingTopic, []byte{1, 2, 3})
exchange := p2ptest.Exchange{ exchange := p2ptest.Exchange{
Expects: []p2ptest.Expect{ Expects: []p2ptest.Expect{
p2ptest.Expect{ p2ptest.Expect{
Code: 0, Code: 0,
Msg: msg, Msg: msg,
Peer: pt.Ids[1], Peer: pt.Ids[0],
},
}, },
Triggers: []p2ptest.Trigger{ },
p2ptest.Trigger{ Triggers: []p2ptest.Trigger{
Code: 0, p2ptest.Trigger{
Msg: msg, Code: 0,
Peer: pt.Ids[1], Msg: msg,
}, Peer: pt.Ids[1],
}, },
} },
}
pt.TestExchanges(exchange) pt.TestExchanges(exchange)
} }
func TestPssFullRandom10_5_5(t *testing.T) { func TestPssFullRandom10_5_5(t *testing.T) {
adapter := adapters.NewSimAdapter(services) adapter := adapters.NewSimAdapter(services)
testPssFullRandom(t, adapter, 10, 5, 5) testPssFullRandom(t, adapter, 10, 5, 5)
@ -210,10 +217,11 @@ func TestPssFullRandom10_5_5(t *testing.T) {
func testPssFullRandom(t *testing.T, adapter adapters.NodeAdapter, nodecount int, fullnodecount int, msgcount int) { func testPssFullRandom(t *testing.T, adapter adapters.NodeAdapter, nodecount int, fullnodecount int, msgcount int) {
var lastid *adapters.NodeId = nil var lastid *adapters.NodeId = nil
nodeCount := 5 nodeCount := 5
net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{ net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{
Id: "0", Id: "0",
DefaultService: bzzServiceName, DefaultService: "psstest",
}) })
defer net.Shutdown() defer net.Shutdown()
@ -222,7 +230,7 @@ func testPssFullRandom(t *testing.T, adapter adapters.NodeAdapter, nodecount int
for i := 0; i < nodeCount; i++ { for i := 0; i < nodeCount; i++ {
nodeconfig := adapters.RandomNodeConfig() nodeconfig := adapters.RandomNodeConfig()
nodeconfig.Services = []string{"bzz", "pss"} nodeconfig.Service = "psstest"
node, err := net.NewNodeWithConfig(nodeconfig) node, err := net.NewNodeWithConfig(nodeconfig)
if err != nil { if err != nil {
t.Fatalf("error starting node: %s", err) t.Fatalf("error starting node: %s", err)
@ -274,9 +282,9 @@ func testPssFullRandom(t *testing.T, adapter adapters.NodeAdapter, nodecount int
if lastid != nil { if lastid != nil {
//msg := pssPingMsg{Created: time.Now(),} //msg := pssPingMsg{Created: time.Now(),}
client.CallContext(context.Background(), nil, "pss_sendRaw", topic, PssAPIMsg{ client.CallContext(context.Background(), nil, "pss_sendRaw", pssPingTopic, PssAPIMsg{
Addr: lastid.Bytes(), Addr: lastid.Bytes(),
Msg: []byte{1,2,3}, Msg: []byte{1, 2, 3},
}) })
} }
lastid = id lastid = id
@ -313,81 +321,77 @@ func triggerChecks(trigger chan *adapters.NodeId, net *simulations.Network, id *
if node == nil { if node == nil {
return fmt.Errorf("unknown node: %s", id) return fmt.Errorf("unknown node: %s", id)
} }
go func(){ go func() {
time.Sleep(time.Second) time.Sleep(time.Second)
trigger <- id trigger <- id
}() }()
/* /*
client, err := node.Client() client, err := node.Client()
if err != nil { if err != nil {
return err return err
}
events := make(chan PssAPIMsg)
sub, err := client.Subscribe(context.Background(), "pss", events, "newMsg", topic)
if err != nil {
return fmt.Errorf("error getting peer events for node %v: %s", id, err)
}
go func() {
defer sub.Unsubscribe()
for {
select {
case msg := <-events:
log.Warn("pss rpc got msg", "msg", msg)
trigger <- id
case err := <-sub.Err():
if err != nil {
log.Error(fmt.Sprintf("error getting peer events for node %v", id), "err", err)
}
return
}
} }
}() events := make(chan PssAPIMsg)
sub, err := client.Subscribe(context.Background(), "pss", events, "newMsg", topic)
if err != nil {
return fmt.Errorf("error getting peer events for node %v: %s", id, err)
}
go func() {
defer sub.Unsubscribe()
for {
select {
case msg := <-events:
log.Warn("pss rpc got msg", "msg", msg)
trigger <- id
case err := <-sub.Err():
if err != nil {
log.Error(fmt.Sprintf("error getting peer events for node %v", id), "err", err)
}
return
}
}
}()
*/ */
return nil return nil
} }
func newServices() adapters.Services { func newServices() adapters.Services {
return func(id *adapters.NodeId, snapshot []byte) []node.Service { return adapters.Services{
// setup hive "psstest": func(id *adapters.NodeId, snapshot []byte) []node.Service {
addr := network.NewAddrFromNodeId(id) addr := network.NewAddrFromNodeId(id)
config := &network.BzzConfig{ kadparams := network.NewKadParams()
OverlayAddr: addr.Over(), kadparams.MinProxBinSize = 2
UnderlayAddr: addr.Under(), kadparams.MaxBinSize = 3
KadParams: network.NewKadParams(), kadparams.MinBinSize = 1
HiveParams: network.NewHiveParams(), kadparams.MaxRetries = 1000
} kadparams.RetryExponent = 2
kadparams.RetryInterval = 1000000
kademlia := network.NewKademlia(addr.OAddr, kadparams)
config.KadParams.MinProxBinSize = 2 config := &network.BzzConfig{
config.KadParams.MaxBinSize = 3 OverlayAddr: addr.Over(),
config.KadParams.MinBinSize = 1 UnderlayAddr: addr.Under(),
config.KadParams.MaxRetries = 1000 HiveParams: network.NewHiveParams(),
config.KadParams.RetryExponent = 2 }
config.KadParams.RetryInterval = 1000000
config.HiveParams.KeepAliveInterval = time.Second config.HiveParams.KeepAliveInterval = time.Second
network.NewBzz(config) cachedir, err := ioutil.TempDir("", "pss-cache")
if err != nil {
log.Error("create pss cache tmpdir failed", "error", err)
return nil
}
dpa, err := storage.NewLocalDPA(cachedir)
if err != nil {
log.Error("local dpa creation failed", "error", err)
return nil
}
pssp := NewPssParams()
// pss setup return []node.Service{network.NewBzz(config, kademlia, adapters.NewSimStateStore()), NewPss(kademlia, dpa, pssp)}
cachedir, err := ioutil.TempDir("", "pss-cache") },
if err != nil {
log.Error("create pss cache tmpdir failed", "error", err)
return nil
}
dpa, err := storage.NewLocalDPA(cachedir)
if err != nil {
log.Error("local dpa creation failed", "error", err)
return nil
}
pssp := NewPssParams()
for bzzs[id] == nil {
time.Sleep(time.Microsecond * 100)
}
return NewPss(bzzs[id].Kademlia, dpa, pssp)
} }
} }
/* /*

View file

@ -44,7 +44,7 @@ func (pssapi *PssAPI) NewMsg(ctx context.Context, topic PssTopic) (*rpc.Subscrip
} }
return nil return nil
} }
deregf := pssapi.Pss.Register(topic, handler) deregf := pssapi.Pss.Register(&topic, handler)
go func() { go func() {
defer deregf() defer deregf()