diff --git a/node/api.go b/node/api.go index 70aa59d844..ba8cc1b8b7 100644 --- a/node/api.go +++ b/node/api.go @@ -17,6 +17,7 @@ package node import ( + "context" "fmt" "strings" "time" @@ -25,6 +26,7 @@ import ( "github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p/discover" + "github.com/ethereum/go-ethereum/rpc" "github.com/rcrowley/go-metrics" ) @@ -73,6 +75,44 @@ func (api *PrivateAdminAPI) RemovePeer(url string) (bool, error) { 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. func (api *PrivateAdminAPI) StartRPC(host *string, port *int, cors *string, apis *string) (bool, error) { api.node.lock.Lock() diff --git a/p2p/simulations/adapters/exec.go b/p2p/simulations/adapters/exec.go index 13715cc100..fd31e5ef70 100644 --- a/p2p/simulations/adapters/exec.go +++ b/p2p/simulations/adapters/exec.go @@ -315,22 +315,10 @@ func startP2PNode(conf *node.Config, service node.Service) (*node.Node, error) { if err != nil { return nil, err } - - constructor := func(s node.Service) node.ServiceConstructor { - return func(ctx *node.ServiceContext) (node.Service, error) { - return s, nil - } + constructor := func(ctx *node.ServiceContext) (node.Service, error) { + return service, 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 { + if err := stack.Register(constructor); err != nil { return nil, err } if err := stack.Start(); err != nil { @@ -339,65 +327,6 @@ func startP2PNode(conf *node.Config, service node.Service) (*node.Node, error) { 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 // use stdio for RPC messages type stdioConn struct { diff --git a/p2p/simulations/adapters/inproc.go b/p2p/simulations/adapters/inproc.go index 559e3ea234..fbabe9c5f7 100644 --- a/p2p/simulations/adapters/inproc.go +++ b/p2p/simulations/adapters/inproc.go @@ -17,6 +17,7 @@ package adapters import ( + "context" "errors" "fmt" "sync" @@ -184,7 +185,7 @@ func (self *SimNode) startRPC() error { 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 apis := append(self.service.APIs(), []rpc.API{ { @@ -192,11 +193,6 @@ func (self *SimNode) startRPC() error { Version: "1.0", Service: &SimAdminAPI{self}, }, - { - Namespace: "eth", - Version: "1.0", - Service: &PeerAPI{func() p2p.Server { return self }}, - }, }...) // start the RPC handler @@ -356,3 +352,35 @@ func (api *SimAdminAPI) RemovePeer(url string) (bool, error) { api.SimNode.RemovePeer(node) 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 +} diff --git a/p2p/simulations/network.go b/p2p/simulations/network.go index 95a0198ee5..e24c9b5cef 100644 --- a/p2p/simulations/network.go +++ b/p2p/simulations/network.go @@ -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) } 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 { return fmt.Errorf("error getting peer events for node %v: %s", id, err) } diff --git a/rpc/client.go b/rpc/client.go index 55ea1783a8..8b42416670 100644 --- a/rpc/client.go +++ b/rpc/client.go @@ -357,19 +357,24 @@ func (c *Client) BatchCallContext(ctx context.Context, b []BatchElem) error { 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 "_subscribe" method with the given arguments, // registering a subscription. Server notifications for the subscription are // sent to the given channel. The element type of the channel must match the // expected type of content returned by the subscription. // // 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 // before considering the subscriber dead. The subscription Err channel will receive // ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure // 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. chanVal := reflect.ValueOf(channel) 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 } - msg, err := c.newMessage("eth"+subscribeMethodSuffix, args...) + msg, err := c.newMessage(namespace+subscribeMethodSuffix, args...) if err != nil { return nil, err } op := &requestOp{ ids: []json.RawMessage{msg.ID}, resp: make(chan *jsonrpcMessage), - sub: newClientSubscription(c, "eth", chanVal), + sub: newClientSubscription(c, namespace, chanVal), } // Send the subscription request. diff --git a/rpc/client_test.go b/rpc/client_test.go index 10d74670b7..4f354d389e 100644 --- a/rpc/client_test.go +++ b/rpc/client_test.go @@ -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 // waiting for a response. func TestClientSubscribeClose(t *testing.T) { diff --git a/swarm/network/simulations/discovery/discovery_test.go b/swarm/network/simulations/discovery/discovery_test.go index f905cb8ac5..5a5ec53f37 100644 --- a/swarm/network/simulations/discovery/discovery_test.go +++ b/swarm/network/simulations/discovery/discovery_test.go @@ -159,7 +159,7 @@ func triggerChecks(trigger chan *adapters.NodeId, net *simulations.Network, id * return err } 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 { return fmt.Errorf("error getting peer events for node %v: %s", id, err) }