p2p/simulations: Move PeerEvents to admin RPC namespace

Signed-off-by: Lewis Marshall <lewis@lmars.net>
This commit is contained in:
Lewis Marshall 2017-05-06 15:38:49 +01:00
parent e4ac128657
commit cdab665475
5 changed files with 79 additions and 82 deletions

View file

@ -17,6 +17,7 @@
package node package node
import ( import (
"context"
"fmt" "fmt"
"strings" "strings"
"time" "time"
@ -25,6 +26,7 @@ import (
"github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rpc"
"github.com/rcrowley/go-metrics" "github.com/rcrowley/go-metrics"
) )
@ -73,6 +75,44 @@ func (api *PrivateAdminAPI) RemovePeer(url string) (bool, error) {
return true, nil return true, nil
} }
// PeerEvents creates an RPC subscription which receives peer events from the
// node's p2p.Server
func (api *PrivateAdminAPI) PeerEvents(ctx context.Context) (*rpc.Subscription, error) {
// Make sure the server is running, fail otherwise
server := api.node.Server()
if server == nil {
return nil, ErrNodeStopped
}
// Create the subscription
notifier, supported := rpc.NotifierFromContext(ctx)
if !supported {
return nil, rpc.ErrNotificationsUnsupported
}
rpcSub := notifier.CreateSubscription()
go func() {
events := make(chan *p2p.PeerEvent)
sub := server.SubscribeEvents(events)
defer sub.Unsubscribe()
for {
select {
case event := <-events:
notifier.Notify(rpcSub.ID, event)
case <-sub.Err():
return
case <-rpcSub.Err():
return
case <-notifier.Closed():
return
}
}
}()
return rpcSub, nil
}
// StartRPC starts the HTTP RPC API server. // StartRPC starts the HTTP RPC API server.
func (api *PrivateAdminAPI) StartRPC(host *string, port *int, cors *string, apis *string) (bool, error) { func (api *PrivateAdminAPI) StartRPC(host *string, port *int, cors *string, apis *string) (bool, error) {
api.node.lock.Lock() api.node.lock.Lock()

View file

@ -315,22 +315,10 @@ func startP2PNode(conf *node.Config, service node.Service) (*node.Node, error) {
if err != nil { if err != nil {
return nil, err return nil, err
} }
constructor := func(ctx *node.ServiceContext) (node.Service, error) {
constructor := func(s node.Service) node.ServiceConstructor { return service, nil
return func(ctx *node.ServiceContext) (node.Service, error) {
return s, nil
}
} }
if err := stack.Register(constructor); err != nil {
// register the peer events API
//
// TODO: move this to node.PrivateAdminAPI once the following is merged:
// https://github.com/ethereum/go-ethereum/pull/13885
if err := stack.Register(constructor(&PeerAPI{stack.Server})); err != nil {
return nil, err
}
if err := stack.Register(constructor(service)); err != nil {
return nil, err return nil, err
} }
if err := stack.Start(); err != nil { if err := stack.Start(); err != nil {
@ -339,65 +327,6 @@ func startP2PNode(conf *node.Config, service node.Service) (*node.Node, error) {
return stack, nil return stack, nil
} }
// PeerAPI is used to expose peer events under the "eth" RPC namespace.
//
// TODO: move this to node.PrivateAdminAPI and expose under the "admin"
// namespace once the following is merged:
// https://github.com/ethereum/go-ethereum/pull/13885
type PeerAPI struct {
server func() p2p.Server
}
func (p *PeerAPI) Protocols() []p2p.Protocol {
return nil
}
func (p *PeerAPI) APIs() []rpc.API {
return []rpc.API{{
Namespace: "eth",
Version: "1.0",
Service: p,
}}
}
func (p *PeerAPI) Start(p2p.Server) error {
return nil
}
func (p *PeerAPI) Stop() error {
return nil
}
// PeerEvents creates an RPC sunscription which receives peer events from the
// underlying p2p.Server
func (p *PeerAPI) PeerEvents(ctx context.Context) (*rpc.Subscription, error) {
notifier, supported := rpc.NotifierFromContext(ctx)
if !supported {
return &rpc.Subscription{}, rpc.ErrNotificationsUnsupported
}
rpcSub := notifier.CreateSubscription()
go func() {
events := make(chan *p2p.PeerEvent)
sub := p.server().SubscribeEvents(events)
defer sub.Unsubscribe()
for {
select {
case event := <-events:
notifier.Notify(rpcSub.ID, event)
case <-rpcSub.Err():
return
case <-notifier.Closed():
return
}
}
}()
return rpcSub, nil
}
// stdioConn wraps os.Stdin / os.Stdout with a no-op Close method so we can // stdioConn wraps os.Stdin / os.Stdout with a no-op Close method so we can
// use stdio for RPC messages // use stdio for RPC messages
type stdioConn struct { type stdioConn struct {

View file

@ -17,6 +17,7 @@
package adapters package adapters
import ( import (
"context"
"errors" "errors"
"fmt" "fmt"
"sync" "sync"
@ -184,7 +185,7 @@ func (self *SimNode) startRPC() error {
return errors.New("RPC already started") return errors.New("RPC already started")
} }
// add SimAdminAPI and PeerAPI so that the network can call the // add SimAdminAPI so that the network can call the
// AddPeer, RemovePeer and PeerEvents RPC methods // AddPeer, RemovePeer and PeerEvents RPC methods
apis := append(self.service.APIs(), []rpc.API{ apis := append(self.service.APIs(), []rpc.API{
{ {
@ -192,11 +193,6 @@ func (self *SimNode) startRPC() error {
Version: "1.0", Version: "1.0",
Service: &SimAdminAPI{self}, Service: &SimAdminAPI{self},
}, },
{
Namespace: "eth",
Version: "1.0",
Service: &PeerAPI{func() p2p.Server { return self }},
},
}...) }...)
// start the RPC handler // start the RPC handler
@ -356,3 +352,35 @@ func (api *SimAdminAPI) RemovePeer(url string) (bool, error) {
api.SimNode.RemovePeer(node) api.SimNode.RemovePeer(node)
return true, nil return true, nil
} }
// PeerEvents creates an RPC subscription which receives peer events from the
// underlying p2p.Server
func (api *SimAdminAPI) PeerEvents(ctx context.Context) (*rpc.Subscription, error) {
notifier, supported := rpc.NotifierFromContext(ctx)
if !supported {
return &rpc.Subscription{}, rpc.ErrNotificationsUnsupported
}
rpcSub := notifier.CreateSubscription()
go func() {
events := make(chan *p2p.PeerEvent)
sub := api.SubscribeEvents(events)
defer sub.Unsubscribe()
for {
select {
case event := <-events:
notifier.Notify(rpcSub.ID, event)
case <-sub.Err():
return
case <-rpcSub.Err():
return
case <-notifier.Closed():
return
}
}
}()
return rpcSub, nil
}

View file

@ -426,7 +426,7 @@ func (self *Network) Start(id *adapters.NodeId) error {
return fmt.Errorf("error getting rpc client for node %v: %s", id, err) return fmt.Errorf("error getting rpc client for node %v: %s", id, err)
} }
events := make(chan *p2p.PeerEvent) events := make(chan *p2p.PeerEvent)
sub, err := client.EthSubscribe(context.Background(), events, "peerEvents") sub, err := client.Subscribe(context.Background(), "admin", events, "peerEvents")
if err != nil { if err != nil {
return fmt.Errorf("error getting peer events for node %v: %s", id, err) return fmt.Errorf("error getting peer events for node %v: %s", id, err)
} }

View file

@ -159,7 +159,7 @@ func triggerChecks(trigger chan *adapters.NodeId, net *simulations.Network, id *
return err return err
} }
events := make(chan *p2p.PeerEvent) events := make(chan *p2p.PeerEvent)
sub, err := client.EthSubscribe(context.Background(), events, "peerEvents") sub, err := client.Subscribe(context.Background(), "admin", events, "peerEvents")
if err != nil { if err != nil {
return fmt.Errorf("error getting peer events for node %v: %s", id, err) return fmt.Errorf("error getting peer events for node %v: %s", id, err)
} }