From 4b397c49c5d39ac0b34506deb01d7a027f123975 Mon Sep 17 00:00:00 2001 From: Lewis Marshall Date: Wed, 28 Jun 2017 17:54:37 +0100 Subject: [PATCH] p2p/simulations: Send current nodes & conns in event stream Signed-off-by: Lewis Marshall --- p2p/simulations/cmd/p2psim/main.go | 8 ++++- p2p/simulations/http.go | 48 +++++++++++++++++++++++++----- p2p/simulations/http_test.go | 18 +++++++++-- 3 files changed, 64 insertions(+), 10 deletions(-) diff --git a/p2p/simulations/cmd/p2psim/main.go b/p2p/simulations/cmd/p2psim/main.go index b511fb5953..923b2985ec 100644 --- a/p2p/simulations/cmd/p2psim/main.go +++ b/p2p/simulations/cmd/p2psim/main.go @@ -65,6 +65,12 @@ func main() { Name: "events", Usage: "stream network events", Action: streamNetwork, + Flags: []cli.Flag{ + cli.BoolFlag{ + Name: "current", + Usage: "get existing nodes and conns first", + }, + }, }, { Name: "snapshot", @@ -176,7 +182,7 @@ func streamNetwork(ctx *cli.Context) error { return cli.ShowCommandHelp(ctx, ctx.Command.Name) } events := make(chan *simulations.Event) - sub, err := client.SubscribeNetwork(events) + sub, err := client.SubscribeNetwork(events, ctx.Bool("current")) if err != nil { return err } diff --git a/p2p/simulations/http.go b/p2p/simulations/http.go index c9140412e9..9df33d6d9a 100644 --- a/p2p/simulations/http.go +++ b/p2p/simulations/http.go @@ -68,9 +68,10 @@ func (c *Client) LoadSnapshot(snap *Snapshot) error { } // SubscribeNetwork subscribes to network events which are sent from the server -// as a server-sent-events stream -func (c *Client) SubscribeNetwork(events chan *Event) (event.Subscription, error) { - req, err := http.NewRequest("GET", fmt.Sprintf("%s/events", c.URL), nil) +// as a server-sent-events stream, optionally receiving events for existing +// nodes and connections +func (c *Client) SubscribeNetwork(events chan *Event, current bool) (event.Subscription, error) { + req, err := http.NewRequest("GET", fmt.Sprintf("%s/events?current=%t", c.URL, current), nil) if err != nil { return nil, err } @@ -366,6 +367,17 @@ func (s *Server) StreamNetworkEvents(w http.ResponseWriter, req *http.Request) { fw.Flush() } } + writeEvent := func(event *Event) error { + data, err := json.Marshal(event) + if err != nil { + return err + } + write("network", string(data)) + return nil + } + writeErr := func(err error) { + write("error", err.Error()) + } w.Header().Set("Content-Type", "text/event-stream; charset=utf-8") w.WriteHeader(http.StatusOK) @@ -373,15 +385,37 @@ func (s *Server) StreamNetworkEvents(w http.ResponseWriter, req *http.Request) { if fw, ok := w.(http.Flusher); ok { fw.Flush() } + + // optionally send the existing nodes and connections + if req.URL.Query().Get("current") == "true" { + snap, err := s.network.Snapshot() + if err != nil { + writeErr(err) + return + } + for _, node := range snap.Nodes { + event := NewEvent(&node.Node) + if err := writeEvent(event); err != nil { + writeErr(err) + return + } + } + for _, conn := range snap.Conns { + event := NewEvent(&conn) + if err := writeEvent(event); err != nil { + writeErr(err) + return + } + } + } + for { select { case event := <-events: - data, err := json.Marshal(event) - if err != nil { - write("error", err.Error()) + if err := writeEvent(event); err != nil { + writeErr(err) return } - write("network", string(data)) case <-clientGone: return } diff --git a/p2p/simulations/http_test.go b/p2p/simulations/http_test.go index f9955db5c0..bb30b5c8b3 100644 --- a/p2p/simulations/http_test.go +++ b/p2p/simulations/http_test.go @@ -145,7 +145,7 @@ func TestHTTPNetwork(t *testing.T) { // subscribe to events so we can check them later client := NewClient(s.URL) events := make(chan *Event, 100) - sub, err := client.SubscribeNetwork(events) + sub, err := client.SubscribeNetwork(events, false) if err != nil { t.Fatalf("error subscribing to network events: %s", err) } @@ -213,6 +213,20 @@ func TestHTTPNetwork(t *testing.T) { x.connEvent(nodeIDs[0], nodeIDs[1], false), x.connEvent(nodeIDs[0], nodeIDs[1], true), ) + + // reconnect the stream and check we get the current nodes and conns + events = make(chan *Event, 100) + sub, err = client.SubscribeNetwork(events, true) + if err != nil { + t.Fatalf("error subscribing to network events: %s", err) + } + defer sub.Unsubscribe() + x = &expectEvents{t, events, sub} + x.expect( + x.nodeEvent(nodeIDs[0], true), + x.nodeEvent(nodeIDs[1], true), + x.connEvent(nodeIDs[0], nodeIDs[1], true), + ) } type expectEvents struct { @@ -411,7 +425,7 @@ func TestHTTPSnapshot(t *testing.T) { // subscribe to events so we can check them later events := make(chan *Event, 100) - sub, err := client.SubscribeNetwork(events) + sub, err := client.SubscribeNetwork(events, false) if err != nil { t.Fatalf("error subscribing to network events: %s", err) }