diff --git a/cmd/swarm/main.go b/cmd/swarm/main.go index aaf2e3ad28..e6da14c83e 100644 --- a/cmd/swarm/main.go +++ b/cmd/swarm/main.go @@ -459,7 +459,7 @@ func getPassPhrase(prompt string, i int, passwords []string) string { return password } -func injectBootnodes(srv p2p.Server, nodes []string) { +func injectBootnodes(srv *p2p.Server, nodes []string) { for _, url := range nodes { n, err := discover.ParseNode(url) if err != nil { diff --git a/cmd/wnode/main.go b/cmd/wnode/main.go index 5d5aa77ee8..f18025dff8 100644 --- a/cmd/wnode/main.go +++ b/cmd/wnode/main.go @@ -51,7 +51,7 @@ const quitCommand = "~Q" // singletons var ( - server p2p.Server + server *p2p.Server shh *whisper.Whisper done chan struct{} mailServer mailserver.WMailServer @@ -253,17 +253,19 @@ func initialize() { maxPeers = 800 } - server = p2p.NewServer(p2p.Config{ - PrivateKey: nodeid, - MaxPeers: maxPeers, - Name: common.MakeName("wnode", "5.0"), - Protocols: shh.Protocols(), - ListenAddr: *argIP, - NAT: nat.Any(), - BootstrapNodes: peers, - StaticNodes: peers, - TrustedNodes: peers, - }) + server = &p2p.Server{ + Config: p2p.Config{ + PrivateKey: nodeid, + MaxPeers: maxPeers, + Name: common.MakeName("wnode", "5.0"), + Protocols: shh.Protocols(), + ListenAddr: *argIP, + NAT: nat.Any(), + BootstrapNodes: peers, + StaticNodes: peers, + TrustedNodes: peers, + }, + } } func startServer() { diff --git a/contracts/release/release.go b/contracts/release/release.go index e0441fd41e..28a35381d4 100644 --- a/contracts/release/release.go +++ b/contracts/release/release.go @@ -94,7 +94,7 @@ func (r *ReleaseService) Protocols() []p2p.Protocol { return nil } func (r *ReleaseService) APIs() []rpc.API { return nil } // Start spawns the periodic version checker goroutine -func (r *ReleaseService) Start(server p2p.Server) error { +func (r *ReleaseService) Start(server *p2p.Server) error { go r.checker() return nil } diff --git a/eth/backend.go b/eth/backend.go index 9361cb3064..f864b1d88b 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -49,7 +49,7 @@ import ( ) type LesServer interface { - Start(srvr p2p.Server) + Start(srvr *p2p.Server) Stop() Protocols() []p2p.Protocol } @@ -362,7 +362,7 @@ func (s *Ethereum) Protocols() []p2p.Protocol { // Start implements node.Service, starting all internal goroutines needed by the // Ethereum protocol implementation. -func (s *Ethereum) Start(srvr p2p.Server) error { +func (s *Ethereum) Start(srvr *p2p.Server) error { s.netRPCService = ethapi.NewPublicNetAPI(srvr, s.NetVersion()) s.protocolManager.Start() diff --git a/ethstats/ethstats.go b/ethstats/ethstats.go index f295fb60d0..8765da8faf 100644 --- a/ethstats/ethstats.go +++ b/ethstats/ethstats.go @@ -52,7 +52,7 @@ const historyUpdateRange = 50 type Service struct { stack *node.Node // Temporary workaround, remove when API finalized - server p2p.Server // Peer-to-peer server to retrieve networking infos + server *p2p.Server // Peer-to-peer server to retrieve networking infos eth *eth.Ethereum // Full Ethereum service if monitoring a full node les *les.LightEthereum // Light Ethereum service if monitoring a light node engine consensus.Engine // Consensus engine to retrieve variadic block fields @@ -101,7 +101,7 @@ func (s *Service) Protocols() []p2p.Protocol { return nil } func (s *Service) APIs() []rpc.API { return nil } // Start implements node.Service, starting up the monitoring and reporting daemon. -func (s *Service) Start(server p2p.Server) error { +func (s *Service) Start(server *p2p.Server) error { s.server = server go s.loop() diff --git a/internal/ethapi/api.go b/internal/ethapi/api.go index 146e7f0475..f9eed87975 100644 --- a/internal/ethapi/api.go +++ b/internal/ethapi/api.go @@ -1434,12 +1434,12 @@ func (api *PrivateDebugAPI) SetHead(number hexutil.Uint64) { // PublicNetAPI offers network related RPC methods type PublicNetAPI struct { - net p2p.Server + net *p2p.Server networkVersion uint64 } // NewPublicNetAPI creates a new net API instance. -func NewPublicNetAPI(net p2p.Server, networkVersion uint64) *PublicNetAPI { +func NewPublicNetAPI(net *p2p.Server, networkVersion uint64) *PublicNetAPI { return &PublicNetAPI{net, networkVersion} } diff --git a/les/backend.go b/les/backend.go index 3dac7f15fc..646c81a7b1 100644 --- a/les/backend.go +++ b/les/backend.go @@ -185,7 +185,7 @@ func (s *LightEthereum) Protocols() []p2p.Protocol { // Start implements node.Service, starting all internal goroutines needed by the // Ethereum protocol implementation. -func (s *LightEthereum) Start(srvr p2p.Server) error { +func (s *LightEthereum) Start(srvr *p2p.Server) error { log.Warn("Light client mode is an experimental feature") s.netRPCService = ethapi.NewPublicNetAPI(srvr, s.networkId) s.protocolManager.Start(srvr) diff --git a/les/handler.go b/les/handler.go index 5af89d4fca..64023af0f5 100644 --- a/les/handler.go +++ b/les/handler.go @@ -61,10 +61,6 @@ const ( disableClientRemovePeer = false ) -type discV5Server interface { - DiscV5() *discv5.Network -} - // errIncompatibleConfig is returned if the requested protocols and configs are // not compatible (low protocol version restrictions and high requirements). var errIncompatibleConfig = errors.New("incompatible configuration") @@ -260,10 +256,10 @@ func (pm *ProtocolManager) removePeer(id string) { } } -func (pm *ProtocolManager) Start(srvr p2p.Server) { +func (pm *ProtocolManager) Start(srvr *p2p.Server) { var topicDisc *discv5.Network - if v, ok := srvr.(discV5Server); ok { - topicDisc = v.DiscV5() + if srvr != nil { + topicDisc = srvr.DiscV5 } lesTopic := discv5.Topic("LES@" + common.Bytes2Hex(pm.blockchain.Genesis().Hash().Bytes()[0:8])) if pm.lightSync { diff --git a/les/server.go b/les/server.go index 1be61ef53a..22fe59b7ac 100644 --- a/les/server.go +++ b/les/server.go @@ -68,7 +68,7 @@ func (s *LesServer) Protocols() []p2p.Protocol { } // Start starts the LES server -func (s *LesServer) Start(srvr p2p.Server) { +func (s *LesServer) Start(srvr *p2p.Server) { s.protocolManager.Start(srvr) } diff --git a/les/serverpool.go b/les/serverpool.go index 14a7354261..64fe991c63 100644 --- a/les/serverpool.go +++ b/les/serverpool.go @@ -97,7 +97,7 @@ const ( type serverPool struct { db ethdb.Database dbKey []byte - server p2p.Server + server *p2p.Server quit chan struct{} wg *sync.WaitGroup connWg sync.WaitGroup @@ -118,7 +118,7 @@ type serverPool struct { } // newServerPool creates a new serverPool instance -func newServerPool(db ethdb.Database, dbPrefix []byte, server p2p.Server, topic discv5.Topic, quit chan struct{}, wg *sync.WaitGroup) *serverPool { +func newServerPool(db ethdb.Database, dbPrefix []byte, server *p2p.Server, topic discv5.Topic, quit chan struct{}, wg *sync.WaitGroup) *serverPool { pool := &serverPool{ db: db, dbKey: append(dbPrefix, []byte(topic)...), @@ -139,11 +139,11 @@ func newServerPool(db ethdb.Database, dbPrefix []byte, server p2p.Server, topic pool.loadNodes() pool.checkDial() - if srv, ok := pool.server.(discV5Server); ok && srv.DiscV5() != nil { + if pool.server.DiscV5 != nil { pool.discSetPeriod = make(chan time.Duration, 1) pool.discNodes = make(chan *discv5.Node, 100) pool.discLookups = make(chan bool, 100) - go srv.DiscV5().SearchTopic(topic, pool.discSetPeriod, pool.discNodes, pool.discLookups) + go pool.server.DiscV5.SearchTopic(topic, pool.discSetPeriod, pool.discNodes, pool.discLookups) } go pool.eventLoop() diff --git a/node/api.go b/node/api.go index ba8cc1b8b7..1b2ef2ed90 100644 --- a/node/api.go +++ b/node/api.go @@ -375,7 +375,7 @@ func NewPublicWeb3API(stack *Node) *PublicWeb3API { // ClientVersion returns the node name func (s *PublicWeb3API) ClientVersion() string { - return s.stack.serverConfig.Name + return s.stack.Server().Name } // Sha3 applies the ethereum sha3 implementation on the input. diff --git a/node/node.go b/node/node.go index 7b4988a9cf..a89110599f 100644 --- a/node/node.go +++ b/node/node.go @@ -56,7 +56,7 @@ type Node struct { instanceDirLock storage.Storage // prevents concurrent use of instance directory serverConfig p2p.Config - server p2p.Server // Currently running P2P networking layer + server *p2p.Server // Currently running P2P networking layer serviceFuncs []ServiceConstructor // Service constructors (in dependency order) services map[reflect.Type]Service // Currently running services @@ -165,6 +165,8 @@ func (n *Node) Start() error { if n.serverConfig.NodeDatabase == "" { n.serverConfig.NodeDatabase = n.config.NodeDB() } + running := &p2p.Server{Config: 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) @@ -192,10 +194,8 @@ func (n *Node) Start() error { } // Gather the protocols and start the freshly assembled P2P server for _, service := range services { - n.serverConfig.Protocols = append(n.serverConfig.Protocols, service.Protocols()...) + running.Protocols = append(running.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 @@ -582,7 +582,7 @@ func (n *Node) RPCHandler() (*rpc.Server, error) { // Server retrieves the currently running P2P network layer. This method is meant // only to inspect fields of the currently running server, life cycle management // should be left to this Node entity. -func (n *Node) Server() p2p.Server { +func (n *Node) Server() *p2p.Server { n.lock.RLock() defer n.lock.RUnlock() diff --git a/node/node_example_test.go b/node/node_example_test.go index ddf76faea7..ee06f4065c 100644 --- a/node/node_example_test.go +++ b/node/node_example_test.go @@ -37,7 +37,7 @@ type SampleService struct{} func (s *SampleService) Protocols() []p2p.Protocol { return nil } func (s *SampleService) APIs() []rpc.API { return nil } -func (s *SampleService) Start(p2p.Server) error { fmt.Println("Service starting..."); return nil } +func (s *SampleService) Start(*p2p.Server) error { fmt.Println("Service starting..."); return nil } func (s *SampleService) Stop() error { fmt.Println("Service stopping..."); return nil } func ExampleService() { diff --git a/node/node_test.go b/node/node_test.go index a0f13ca570..2880efa619 100644 --- a/node/node_test.go +++ b/node/node_test.go @@ -154,7 +154,7 @@ func TestServiceLifeCycle(t *testing.T) { id := id // Closure for the constructor constructor := func(*ServiceContext) (Service, error) { return &InstrumentedService{ - startHook: func(p2p.Server) { started[id] = true }, + startHook: func(*p2p.Server) { started[id] = true }, stopHook: func() { stopped[id] = true }, }, nil } @@ -200,7 +200,7 @@ func TestServiceRestarts(t *testing.T) { running = false return &InstrumentedService{ - startHook: func(p2p.Server) { + startHook: func(*p2p.Server) { if running { panic("already running") } @@ -250,7 +250,7 @@ func TestServiceConstructionAbortion(t *testing.T) { id := id // Closure for the constructor constructor := func(*ServiceContext) (Service, error) { return &InstrumentedService{ - startHook: func(p2p.Server) { started[id] = true }, + startHook: func(*p2p.Server) { started[id] = true }, }, nil } if err := stack.Register(maker(constructor)); err != nil { @@ -299,7 +299,7 @@ func TestServiceStartupAbortion(t *testing.T) { id := id // Closure for the constructor constructor := func(*ServiceContext) (Service, error) { return &InstrumentedService{ - startHook: func(p2p.Server) { started[id] = true }, + startHook: func(*p2p.Server) { started[id] = true }, stopHook: func() { stopped[id] = true }, }, nil } @@ -352,7 +352,7 @@ func TestServiceTerminationGuarantee(t *testing.T) { id := id // Closure for the constructor constructor := func(*ServiceContext) (Service, error) { return &InstrumentedService{ - startHook: func(p2p.Server) { started[id] = true }, + startHook: func(*p2p.Server) { started[id] = true }, stopHook: func() { stopped[id] = true }, }, nil } diff --git a/node/service.go b/node/service.go index 05320b75ab..5e1eb0e645 100644 --- a/node/service.go +++ b/node/service.go @@ -86,7 +86,7 @@ type Service interface { // Start is called after all services have been constructed and the networking // layer was also initialized to spawn any goroutines required by the service. - Start(server p2p.Server) error + Start(server *p2p.Server) error // Stop terminates all goroutines belonging to the service, blocking until they // are all terminated. diff --git a/node/utils_test.go b/node/utils_test.go index e9c999b93c..7cdfc2b3aa 100644 --- a/node/utils_test.go +++ b/node/utils_test.go @@ -31,7 +31,7 @@ type NoopService struct{} func (s *NoopService) Protocols() []p2p.Protocol { return nil } func (s *NoopService) APIs() []rpc.API { return nil } -func (s *NoopService) Start(p2p.Server) error { return nil } +func (s *NoopService) Start(*p2p.Server) error { return nil } func (s *NoopService) Stop() error { return nil } func NewNoopService(*ServiceContext) (Service, error) { return new(NoopService), nil } @@ -57,7 +57,7 @@ type InstrumentedService struct { stop error protocolsHook func() - startHook func(p2p.Server) + startHook func(*p2p.Server) stopHook func() } @@ -74,7 +74,7 @@ func (s *InstrumentedService) APIs() []rpc.API { return s.apis } -func (s *InstrumentedService) Start(server p2p.Server) error { +func (s *InstrumentedService) Start(server *p2p.Server) error { if s.startHook != nil { s.startHook(server) } diff --git a/p2p/dial.go b/p2p/dial.go index 10f16f5166..9d58caa103 100644 --- a/p2p/dial.go +++ b/p2p/dial.go @@ -97,7 +97,7 @@ type pastDial struct { } type task interface { - Do(*server) + Do(*Server) } // A dialTask is generated for each node that is dialed. Its @@ -280,7 +280,7 @@ func (s *dialstate) taskDone(t task, now time.Time) { } } -func (t *dialTask) Do(srv *server) { +func (t *dialTask) Do(srv *Server) { if t.dest.Incomplete() { if !t.resolve(srv) { return @@ -301,7 +301,7 @@ func (t *dialTask) Do(srv *server) { // Resolve operations are throttled with backoff to avoid flooding the // discovery network with useless queries for nodes that don't exist. // The backoff delay resets when the node is found. -func (t *dialTask) resolve(srv *server) bool { +func (t *dialTask) resolve(srv *Server) bool { if srv.ntab == nil { log.Debug("Can't resolve node", "id", t.dest.ID, "err", "discovery is disabled") return false @@ -330,7 +330,7 @@ func (t *dialTask) resolve(srv *server) bool { } // dial performs the actual connection attempt. -func (t *dialTask) dial(srv *server, dest *discover.Node) bool { +func (t *dialTask) dial(srv *Server, dest *discover.Node) bool { fd, err := srv.Dialer.Dial(dest) if err != nil { log.Trace("Dial error", "task", t, "err", err) @@ -345,7 +345,7 @@ func (t *dialTask) String() string { return fmt.Sprintf("%v %x %v:%d", t.flags, t.dest.ID[:8], t.dest.IP, t.dest.TCP) } -func (t *discoverTask) Do(srv *server) { +func (t *discoverTask) Do(srv *Server) { // newTasks generates a lookup task whenever dynamic dials are // necessary. Lookups need to take some time, otherwise the // event loop spins too fast. @@ -367,7 +367,7 @@ func (t *discoverTask) String() string { return s } -func (t waitExpireTask) Do(*server) { +func (t waitExpireTask) Do(*Server) { time.Sleep(t.Duration) } func (t waitExpireTask) String() string { diff --git a/p2p/server.go b/p2p/server.go index 14a987d875..eb7c70e593 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -141,24 +141,8 @@ type Config struct { EnableMsgEvents bool } -type Server interface { - Start() error - Stop() error - SetupConn(net.Conn, connFlag, *discover.Node) - AddPeer(node *discover.Node) - RemovePeer(node *discover.Node) - SubscribeEvents(ch chan *PeerEvent) event.Subscription - PeerCount() int - NodeInfo() *NodeInfo - PeersInfo() []*PeerInfo -} - -func NewServer(conf Config) Server { - return &server{Config: conf} -} - // Server manages all peer connections. -type server struct { +type Server struct { // Config fields may not be modified while the server is running. Config @@ -267,7 +251,7 @@ func (c *conn) is(f connFlag) bool { } // Peers returns all connected peers. -func (srv *server) Peers() []*Peer { +func (srv *Server) Peers() []*Peer { var ps []*Peer select { // Note: We'd love to put this function into a variable but @@ -285,7 +269,7 @@ func (srv *server) Peers() []*Peer { } // PeerCount returns the number of connected peers. -func (srv *server) PeerCount() int { +func (srv *Server) PeerCount() int { var count int select { case srv.peerOp <- func(ps map[discover.NodeID]*Peer) { count = len(ps) }: @@ -298,7 +282,7 @@ func (srv *server) PeerCount() int { // AddPeer connects to the given node and maintains the connection until the // server is shut down. If the connection fails for any reason, the server will // attempt to reconnect the peer. -func (srv *server) AddPeer(node *discover.Node) { +func (srv *Server) AddPeer(node *discover.Node) { select { case srv.addstatic <- node: case <-srv.quit: @@ -306,7 +290,7 @@ func (srv *server) AddPeer(node *discover.Node) { } // RemovePeer disconnects from the given node -func (srv *server) RemovePeer(node *discover.Node) { +func (srv *Server) RemovePeer(node *discover.Node) { select { case srv.removestatic <- node: case <-srv.quit: @@ -314,12 +298,12 @@ func (srv *server) RemovePeer(node *discover.Node) { } // SubscribePeers subscribes the given channel to peer events -func (srv *server) SubscribeEvents(ch chan *PeerEvent) event.Subscription { +func (srv *Server) SubscribeEvents(ch chan *PeerEvent) event.Subscription { return srv.peerFeed.Subscribe(ch) } // Self returns the local node's endpoint information. -func (srv *server) Self() *discover.Node { +func (srv *Server) Self() *discover.Node { srv.lock.Lock() defer srv.lock.Unlock() @@ -329,7 +313,7 @@ func (srv *server) Self() *discover.Node { return srv.makeSelf(srv.listener, srv.ntab) } -func (srv *server) makeSelf(listener net.Listener, ntab discoverTable) *discover.Node { +func (srv *Server) makeSelf(listener net.Listener, ntab discoverTable) *discover.Node { // If the server's not running, return an empty node. // If the node is running but discovery is off, manually assemble the node infos. if ntab == nil { @@ -351,7 +335,7 @@ func (srv *server) makeSelf(listener net.Listener, ntab discoverTable) *discover // Stop terminates the server and all active peer connections. // It blocks until all active connections have been closed. -func (srv *server) Stop() error { +func (srv *Server) Stop() error { srv.lock.Lock() defer srv.lock.Unlock() if !srv.running { @@ -369,7 +353,7 @@ func (srv *server) Stop() error { // Start starts running the server. // Servers can not be re-used after stopping. -func (srv *server) Start() (err error) { +func (srv *Server) Start() (err error) { srv.lock.Lock() defer srv.lock.Unlock() if srv.running { @@ -447,7 +431,7 @@ func (srv *server) Start() (err error) { return nil } -func (srv *server) startListening() error { +func (srv *Server) startListening() error { // Launch the TCP listener. listener, err := net.Listen("tcp", srv.ListenAddr) if err != nil { @@ -476,7 +460,7 @@ type dialer interface { removeStatic(*discover.Node) } -func (srv *server) run(dialstate dialer) { +func (srv *Server) run(dialstate dialer) { defer srv.loopWG.Done() var ( peers = make(map[discover.NodeID]*Peer) @@ -617,7 +601,7 @@ running: } } -func (srv *server) protoHandshakeChecks(peers map[discover.NodeID]*Peer, c *conn) error { +func (srv *Server) protoHandshakeChecks(peers map[discover.NodeID]*Peer, c *conn) error { // Drop connections with no matching protocols. if len(srv.Protocols) > 0 && countMatchingProtocols(srv.Protocols, c.caps) == 0 { return DiscUselessPeer @@ -627,7 +611,7 @@ func (srv *server) protoHandshakeChecks(peers map[discover.NodeID]*Peer, c *conn return srv.encHandshakeChecks(peers, c) } -func (srv *server) encHandshakeChecks(peers map[discover.NodeID]*Peer, c *conn) error { +func (srv *Server) encHandshakeChecks(peers map[discover.NodeID]*Peer, c *conn) error { switch { case !c.is(trustedConn|staticDialedConn) && len(peers) >= srv.MaxPeers: return DiscTooManyPeers @@ -646,7 +630,7 @@ type tempError interface { // listenLoop runs in its own goroutine and accepts // inbound connections. -func (srv *server) listenLoop() { +func (srv *Server) listenLoop() { defer srv.loopWG.Done() log.Info("RLPx listener up", "self", srv.makeSelf(srv.listener, srv.ntab)) @@ -707,7 +691,7 @@ func (srv *server) listenLoop() { // setupConn runs the handshakes and attempts to add the connection // as a peer. It returns when the connection has been added as a peer // or the handshakes have failed. -func (srv *server) SetupConn(fd net.Conn, flags connFlag, dialDest *discover.Node) { +func (srv *Server) SetupConn(fd net.Conn, flags connFlag, dialDest *discover.Node) { // Prevent leftover pending conns from entering the handshake. srv.lock.Lock() running := srv.running @@ -767,7 +751,7 @@ func truncateName(s string) string { // checkpoint sends the conn to run, which performs the // post-handshake checks for the stage (posthandshake, addpeer). -func (srv *server) checkpoint(c *conn, stage chan<- *conn) error { +func (srv *Server) checkpoint(c *conn, stage chan<- *conn) error { select { case stage <- c: case <-srv.quit: @@ -784,7 +768,7 @@ func (srv *server) checkpoint(c *conn, stage chan<- *conn) error { // runPeer runs in its own goroutine for each peer. // it waits until the Peer logic returns and removes // the peer. -func (srv *server) runPeer(p *Peer) { +func (srv *Server) runPeer(p *Peer) { if srv.newPeerHook != nil { srv.newPeerHook(p) } @@ -825,7 +809,7 @@ type NodeInfo struct { } // NodeInfo gathers and returns a collection of metadata known about the host. -func (srv *server) NodeInfo() *NodeInfo { +func (srv *Server) NodeInfo() *NodeInfo { node := srv.Self() // Gather and assemble the generic node infos @@ -854,7 +838,7 @@ func (srv *server) NodeInfo() *NodeInfo { } // PeersInfo returns an array of metadata objects describing connected peers. -func (srv *server) PeersInfo() []*PeerInfo { +func (srv *Server) PeersInfo() []*PeerInfo { // Gather all the generic and sub-protocol specific infos infos := make([]*PeerInfo, 0, srv.PeerCount()) for _, peer := range srv.Peers() { @@ -873,6 +857,6 @@ func (srv *server) PeersInfo() []*PeerInfo { return infos } -func (srv *server) DiscV5() *discv5.Network { +func (srv *Server) DiscV5() *discv5.Network { return srv.discV5 } diff --git a/p2p/server_test.go b/p2p/server_test.go index ca8f1873ee..11dd83e5d6 100644 --- a/p2p/server_test.go +++ b/p2p/server_test.go @@ -72,7 +72,7 @@ func startTestServer(t *testing.T, id discover.NodeID, pf func(*Peer)) *Server { ListenAddr: "127.0.0.1:0", PrivateKey: newkey(), } - server := &server{ + server := &Server{ Config: config, newPeerHook: pf, newTransport: func(fd net.Conn) transport { return newTestTransport(id, fd) }, @@ -201,7 +201,7 @@ func TestServerTaskScheduling(t *testing.T) { // The Server in this test isn't actually running // because we're only interested in what run does. - srv := &server{ + srv := &Server{ Config: Config{MaxPeers: 10}, quit: make(chan struct{}), ntab: fakeTable{}, @@ -246,7 +246,7 @@ func TestServerManyTasks(t *testing.T) { } var ( - srv = &server{quit: make(chan struct{}), ntab: fakeTable{}, running: true} + srv = &Server{quit: make(chan struct{}), ntab: fakeTable{}, running: true} done = make(chan *testTask) start, end = 0, 0 ) @@ -317,7 +317,7 @@ func (t *testTask) Do(srv *Server) { // at capacity. Trusted connections should still be accepted. func TestServerAtCap(t *testing.T) { trustedID := randomID() - srv := &server{ + srv := &Server{ Config: Config{ PrivateKey: newkey(), MaxPeers: 10, @@ -420,7 +420,7 @@ func TestServerSetupConn(t *testing.T) { } for i, test := range tests { - srv := &server{ + srv := &Server{ Config: Config{ PrivateKey: srvkey, MaxPeers: 10, @@ -435,7 +435,7 @@ func TestServerSetupConn(t *testing.T) { } } p1, _ := net.Pipe() - srv.setupConn(p1, test.flags, test.dialDest) + srv.SetupConn(p1, test.flags, test.dialDest) if !reflect.DeepEqual(test.tt.closeErr, test.wantCloseErr) { t.Errorf("test %d: close error mismatch: got %q, want %q", i, test.tt.closeErr, test.wantCloseErr) } diff --git a/p2p/simulations/examples/connectivity.go b/p2p/simulations/examples/connectivity.go index 43beff951a..472054cb7a 100644 --- a/p2p/simulations/examples/connectivity.go +++ b/p2p/simulations/examples/connectivity.go @@ -99,7 +99,7 @@ func (p *pingPongService) APIs() []rpc.API { return nil } -func (p *pingPongService) Start(server p2p.Server) error { +func (p *pingPongService) Start(server *p2p.Server) error { p.log.Info("ping-pong service starting") return nil } diff --git a/p2p/simulations/http_test.go b/p2p/simulations/http_test.go index f2c31d48e7..99fb49ed8f 100644 --- a/p2p/simulations/http_test.go +++ b/p2p/simulations/http_test.go @@ -46,7 +46,7 @@ func (t *testService) APIs() []rpc.API { }} } -func (t *testService) Start(server p2p.Server) error { +func (t *testService) Start(server *p2p.Server) error { return nil } diff --git a/p2p/testing/protocolsession.go b/p2p/testing/protocolsession.go index d611cd45a5..2de849885d 100644 --- a/p2p/testing/protocolsession.go +++ b/p2p/testing/protocolsession.go @@ -13,8 +13,7 @@ import ( ) type ProtocolSession struct { - p2p.Server - + Server *p2p.Server Ids []*adapters.NodeId adapter *adapters.SimAdapter events chan *p2p.PeerEvent diff --git a/p2p/testing/protocoltester.go b/p2p/testing/protocoltester.go index 695cfd0292..3b104bd12f 100644 --- a/p2p/testing/protocoltester.go +++ b/p2p/testing/protocoltester.go @@ -62,6 +62,10 @@ func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(*p2p.P return self } +func (self *ProtocolTester) Stop() error { + return self.Server.Stop() +} + func (self *ProtocolTester) Connect(selfId *adapters.NodeId, peers ...*adapters.NodeConfig) { for _, peer := range peers { log.Trace(fmt.Sprintf("start node %v", peer.Id)) @@ -96,7 +100,7 @@ func (t *testNode) APIs() []rpc.API { return nil } -func (t *testNode) Start(server p2p.Server) error { +func (t *testNode) Start(server *p2p.Server) error { return nil } diff --git a/swarm/network/hive.go b/swarm/network/hive.go index 07c1bef6b3..8a9c3d950a 100644 --- a/swarm/network/hive.go +++ b/swarm/network/hive.go @@ -106,7 +106,7 @@ func NewHive(params *HiveParams, overlay Overlay, store Store) *Hive { // these are called on the p2p.Server which runs on the node // af() returns an arbitrary ticker channel // rw is a read writer for json configs -func (self *Hive) Start(server p2p.Server) error { +func (self *Hive) Start(server *p2p.Server) error { if self.store != nil { if err := self.loadPeers(); err != nil { return err diff --git a/swarm/network/hive_test.go b/swarm/network/hive_test.go index 473ba81c0d..6eaa48c48e 100644 --- a/swarm/network/hive_test.go +++ b/swarm/network/hive_test.go @@ -64,7 +64,7 @@ func TestRegisterAndConnect(t *testing.T) { ticker: make(chan time.Time), } pp.newTicker = func() hiveTicker { return tc } - pp.Start(s) + pp.Start(s.Server) defer pp.Stop() tc.ticker <- time.Now() diff --git a/swarm/network/protocol.go b/swarm/network/protocol.go index e4ad09f705..ac24790721 100644 --- a/swarm/network/protocol.go +++ b/swarm/network/protocol.go @@ -164,7 +164,7 @@ func (b *Bzz) APIs() []rpc.API { }} } -func (b *Bzz) Start(server p2p.Server) error { +func (b *Bzz) Start(server *p2p.Server) error { return b.Hive.Start(server) } diff --git a/swarm/network/pss_test.go b/swarm/network/pss_test.go index 502558b0a0..2fc7ddc11b 100644 --- a/swarm/network/pss_test.go +++ b/swarm/network/pss_test.go @@ -89,7 +89,7 @@ func newPssTestService(t *testing.T, handlefunc func(interface{}) error, testnod } } -func (self *pssTestService) Start(server p2p.Server) error { +func (self *pssTestService) Start(server *p2p.Server) error { return self.node.Hive.Start(server) } diff --git a/swarm/network/simulations/overlay.go b/swarm/network/simulations/overlay.go index 1fe613676b..3498be5bca 100644 --- a/swarm/network/simulations/overlay.go +++ b/swarm/network/simulations/overlay.go @@ -64,7 +64,7 @@ func af() <-chan time.Time { // Start() starts up the hive // makes SimNode implement node.Service -func (self *SimNode) Start(server p2p.Server) error { +func (self *SimNode) Start(server *p2p.Server) error { self.init() return self.hive.Start(server, af, self.rw) } diff --git a/swarm/swarm.go b/swarm/swarm.go index a8b6353dd9..3bad8b8e1e 100644 --- a/swarm/swarm.go +++ b/swarm/swarm.go @@ -161,7 +161,7 @@ Start is called when the stack is started * TODO: start subservices like sword, swear, swarmdns */ // implements the node.Service interface -func (self *Swarm) Start(net p2p.Server) error { +func (self *Swarm) Start(net *p2p.Server) error { // set chequebook if self.swapEnabled { ctx := context.Background() // The initial setup has no deadline. diff --git a/whisper/whisperv2/main.go b/whisper/whisperv2/main.go index 081e4fe438..a2e49659bc 100644 --- a/whisper/whisperv2/main.go +++ b/whisper/whisperv2/main.go @@ -48,14 +48,16 @@ func main() { shh := whisper.New() // Create an Ethereum peer to communicate through - server := p2p.NewServer(p2p.Config{ - PrivateKey: key, - MaxPeers: 10, - Name: name, - Protocols: []p2p.Protocol{shh.Protocol()}, - ListenAddr: ":30300", - NAT: nat.Any(), - }) + server := &p2p.Server{ + Config: p2p.Config{ + PrivateKey: key, + MaxPeers: 10, + Name: name, + Protocols: []p2p.Protocol{shh.Protocol()}, + ListenAddr: ":30300", + NAT: nat.Any(), + }, + } fmt.Println("Starting Ethereum peer...") if err := server.Start(); err != nil { fmt.Printf("Failed to start Ethereum peer: %v.\n", err) diff --git a/whisper/whisperv2/whisper.go b/whisper/whisperv2/whisper.go index 6d460bdee5..1d7c21bd12 100644 --- a/whisper/whisperv2/whisper.go +++ b/whisper/whisperv2/whisper.go @@ -172,7 +172,7 @@ func (self *Whisper) Send(envelope *Envelope) error { // Start implements node.Service, starting the background data propagation thread // of the Whisper protocol. -func (self *Whisper) Start(p2p.Server) error { +func (self *Whisper) Start(*p2p.Server) error { log.Info(fmt.Sprint("Whisper started")) go self.update() return nil diff --git a/whisper/whisperv5/peer_test.go b/whisper/whisperv5/peer_test.go index 9e3b66c63d..d3cd63b0b2 100644 --- a/whisper/whisperv5/peer_test.go +++ b/whisper/whisperv5/peer_test.go @@ -78,7 +78,7 @@ type TestData struct { type TestNode struct { shh *Whisper id *ecdsa.PrivateKey - server p2p.Server + server *p2p.Server filerId string } @@ -140,17 +140,19 @@ func initialize(t *testing.T) { peers = append(peers, peer) } - node.server = p2p.NewServer(p2p.Config{ - PrivateKey: node.id, - MaxPeers: NumNodes/2 + 1, - Name: name, - Protocols: node.shh.Protocols(), - ListenAddr: addr, - NAT: nat.Any(), - BootstrapNodes: peers, - StaticNodes: peers, - TrustedNodes: peers, - }) + node.server = &p2p.Server{ + Config: p2p.Config{ + PrivateKey: node.id, + MaxPeers: NumNodes/2 + 1, + Name: name, + Protocols: node.shh.Protocols(), + ListenAddr: addr, + NAT: nat.Any(), + BootstrapNodes: peers, + StaticNodes: peers, + TrustedNodes: peers, + }, + } err = node.server.Start() if err != nil { diff --git a/whisper/whisperv5/whisper.go b/whisper/whisperv5/whisper.go index 77ca20a58c..f2aad08efb 100644 --- a/whisper/whisperv5/whisper.go +++ b/whisper/whisperv5/whisper.go @@ -396,7 +396,7 @@ func (w *Whisper) Send(envelope *Envelope) error { // Start implements node.Service, starting the background data propagation thread // of the Whisper protocol. -func (w *Whisper) Start(p2p.Server) error { +func (w *Whisper) Start(*p2p.Server) error { log.Info("started whisper v." + ProtocolVersionStr) go w.update()