Merge pull request #109 from ethersphere/network-testing-framework-events

p2p/simulations: Send current nodes & conns in event stream
This commit is contained in:
Viktor Trón 2017-06-29 10:14:26 +02:00 committed by GitHub
commit 66f96dad09
3 changed files with 64 additions and 10 deletions

View file

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

View file

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

View file

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