Merge pull request #86 from ethersphere/network-testing-framework-rpc-subscribe

p2p/simulations: Move PeerEvents to admin RPC namespace
This commit is contained in:
Viktor Trón 2017-05-09 11:55:52 +02:00 committed by GitHub
commit aa624cdb99
7 changed files with 121 additions and 87 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

@ -357,19 +357,24 @@ func (c *Client) BatchCallContext(ctx context.Context, b []BatchElem) error {
return err return err
} }
// EthSubscribe calls the "eth_subscribe" method with the given arguments, // EthSubscribe registers a subscripion under the "eth" namespace.
func (c *Client) EthSubscribe(ctx context.Context, channel interface{}, args ...interface{}) (*ClientSubscription, error) {
return c.Subscribe(ctx, "eth", channel, args)
}
// Subscribe calls the "<namespace>_subscribe" method with the given arguments,
// registering a subscription. Server notifications for the subscription are // registering a subscription. Server notifications for the subscription are
// sent to the given channel. The element type of the channel must match the // sent to the given channel. The element type of the channel must match the
// expected type of content returned by the subscription. // expected type of content returned by the subscription.
// //
// The context argument cancels the RPC request that sets up the subscription but has no // The context argument cancels the RPC request that sets up the subscription but has no
// effect on the subscription after EthSubscribe has returned. // effect on the subscription after Subscribe has returned.
// //
// Slow subscribers will be dropped eventually. Client buffers up to 8000 notifications // Slow subscribers will be dropped eventually. Client buffers up to 8000 notifications
// before considering the subscriber dead. The subscription Err channel will receive // before considering the subscriber dead. The subscription Err channel will receive
// ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure // ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure
// that the channel usually has at least one reader to prevent this issue. // that the channel usually has at least one reader to prevent this issue.
func (c *Client) EthSubscribe(ctx context.Context, channel interface{}, args ...interface{}) (*ClientSubscription, error) { func (c *Client) Subscribe(ctx context.Context, namespace string, channel interface{}, args ...interface{}) (*ClientSubscription, error) {
// Check type of channel first. // Check type of channel first.
chanVal := reflect.ValueOf(channel) chanVal := reflect.ValueOf(channel)
if chanVal.Kind() != reflect.Chan || chanVal.Type().ChanDir()&reflect.SendDir == 0 { if chanVal.Kind() != reflect.Chan || chanVal.Type().ChanDir()&reflect.SendDir == 0 {
@ -382,14 +387,14 @@ func (c *Client) EthSubscribe(ctx context.Context, channel interface{}, args ...
return nil, ErrNotificationsUnsupported return nil, ErrNotificationsUnsupported
} }
msg, err := c.newMessage("eth"+subscribeMethodSuffix, args...) msg, err := c.newMessage(namespace+subscribeMethodSuffix, args...)
if err != nil { if err != nil {
return nil, err return nil, err
} }
op := &requestOp{ op := &requestOp{
ids: []json.RawMessage{msg.ID}, ids: []json.RawMessage{msg.ID},
resp: make(chan *jsonrpcMessage), resp: make(chan *jsonrpcMessage),
sub: newClientSubscription(c, "eth", chanVal), sub: newClientSubscription(c, namespace, chanVal),
} }
// Send the subscription request. // Send the subscription request.

View file

@ -251,6 +251,38 @@ func TestClientSubscribe(t *testing.T) {
} }
} }
func TestClientSubscribeCustomNamespace(t *testing.T) {
namespace := "custom"
server := newTestServer(namespace, new(NotificationTestService))
defer server.Stop()
client := DialInProc(server)
defer client.Close()
nc := make(chan int)
count := 10
sub, err := client.Subscribe(context.Background(), namespace, nc, "someSubscription", count, 0)
if err != nil {
t.Fatal("can't subscribe:", err)
}
for i := 0; i < count; i++ {
if val := <-nc; val != i {
t.Fatalf("value mismatch: got %d, want %d", val, i)
}
}
sub.Unsubscribe()
select {
case v := <-nc:
t.Fatal("received value after unsubscribe:", v)
case err := <-sub.Err():
if err != nil {
t.Fatalf("Err returned a non-nil error after explicit unsubscribe: %q", err)
}
case <-time.After(1 * time.Second):
t.Fatalf("subscription not closed within 1s after unsubscribe")
}
}
// In this test, the connection drops while EthSubscribe is // In this test, the connection drops while EthSubscribe is
// waiting for a response. // waiting for a response.
func TestClientSubscribeClose(t *testing.T) { func TestClientSubscribeClose(t *testing.T) {

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)
} }