From 20e91384bf9e6c7fb76f6198e969f6de4c8a1967 Mon Sep 17 00:00:00 2001 From: Lewis Marshall Date: Thu, 6 Apr 2017 00:26:50 +0100 Subject: [PATCH] p2p/simulations: event streams --- p2p/simulations/cytoscape.go | 123 +++++++++++++++-------------- p2p/simulations/network.go | 59 +++++++++++++- p2p/simulations/rest_api_server.go | 18 ++++- 3 files changed, 139 insertions(+), 61 deletions(-) diff --git a/p2p/simulations/cytoscape.go b/p2p/simulations/cytoscape.go index 06571a0c72..37f3fe7956 100644 --- a/p2p/simulations/cytoscape.go +++ b/p2p/simulations/cytoscape.go @@ -1,7 +1,7 @@ package simulations import ( - // "fmt" + "fmt" "github.com/ethereum/go-ethereum/event" ) @@ -29,67 +29,72 @@ type CyElement struct { type CyUpdate struct { Add []*CyElement `json:"add"` - Remove []string `json:"remove"` - Message []string `json:"message"` + Remove []*CyElement `json:"remove"` + Message []*CyElement `json:"message"` +} + +func NewCyUpdate(e *event.TypeMuxEvent) (*CyUpdate, error) { + var update CyUpdate + var el *CyElement + entry := e.Data + var action string + if ev, ok := entry.(*NodeEvent); ok { + el = &CyElement{Group: "nodes", Data: &CyData{Id: ev.node.Id.String()}} + action = ev.Action + } else if ev, ok := entry.(*MsgEvent); ok { + msg := ev.msg + id := ConnLabel(msg.One, msg.Other) + var source, target string + source = msg.One.String() + target = msg.Other.String() + el = &CyElement{Group: "msgs", Data: &CyData{Id: id, Source: source, Target: target}} + action = ev.Action + } else if ev, ok := entry.(*ConnEvent); ok { + // mutually exclusive directed edge (caller -> callee) + conn := ev.conn + id := ConnLabel(conn.One, conn.Other) + var source, target string + if conn.Reverse { + source = conn.Other.String() + target = conn.One.String() + } else { + source = conn.One.String() + target = conn.Other.String() + } + el = &CyElement{Group: "edges", Data: &CyData{Id: id, Source: source, Target: target}} + action = ev.Action + } else { + return nil, fmt.Errorf("unknown event type: %T", entry) + } + + switch action { + case "up": + el.Data.Up = true + update.Add = append(update.Add, el) + case "down": + el.Data.Up = false + update.Remove = append(update.Remove, el) + case "msg": + el.Data.Up = true + update.Message = append(update.Message, el) + default: + return nil, fmt.Errorf("unknown action: %q", action) + } + + return &update, nil } func UpdateCy(conf *CyConfig, j *Journal) (*CyUpdate, error) { - added := []*CyElement{} - removed := []string{} - messaged := []string{} - var el *CyElement - update := func(e *event.TypeMuxEvent) bool { - entry := e.Data - var action string - if ev, ok := entry.(*NodeEvent); ok { - el = &CyElement{Group: "nodes", Data: &CyData{Id: ev.node.Id.Label()}} - action = ev.Action - } else if ev, ok := entry.(*MsgEvent); ok { - msg := ev.msg - id := ConnLabel(msg.One, msg.Other) - var source, target string - source = msg.One.Label() - target = msg.Other.Label() - el = &CyElement{Group: "msgs", Data: &CyData{Id: id, Source: source, Target: target}} - action = ev.Action - } else if ev, ok := entry.(*ConnEvent); ok { - // mutually exclusive directed edge (caller -> callee) - conn := ev.conn - id := ConnLabel(conn.One, conn.Other) - var source, target string - if conn.Reverse { - source = conn.Other.Label() - target = conn.One.Label() - } else { - source = conn.One.Label() - target = conn.Other.Label() - } - el = &CyElement{Group: "edges", Data: &CyData{Id: id, Source: source, Target: target}} - action = ev.Action - } else { - panic("unknown event type") - } - - switch action { - case "up": - el.Data.Up = true - added = append(added, el) - case "down": - el.Data.Up = false - removed = append(removed, el.Data.Id) - case "msg": - el.Data.Up = true - messaged = append(messaged, el.Data.Id) - default: - panic("unknown action") + var update CyUpdate + j.Read(func(e *event.TypeMuxEvent) bool { + u, err := NewCyUpdate(e) + if err != nil { + panic(err.Error()) } + update.Add = append(update.Add, u.Add...) + update.Remove = append(update.Remove, u.Remove...) + update.Message = append(update.Message, u.Message...) return true - } - j.Read(update) - - return &CyUpdate{ - Add: added, - Remove: removed, - Message: messaged, - }, nil + }) + return &update, nil } diff --git a/p2p/simulations/network.go b/p2p/simulations/network.go index 52002d0842..1cc481c244 100644 --- a/p2p/simulations/network.go +++ b/p2p/simulations/network.go @@ -27,7 +27,9 @@ package simulations import ( + "encoding/json" "fmt" + "net/http" "reflect" "sync" @@ -57,6 +59,61 @@ type NetworkControl interface { // and connections established and dropped var ConnectivityEvents = []interface{}{&NodeEvent{}, &ConnEvent{}, &MsgEvent{}} +type NetworkController struct { + *ResourceController + + events *event.TypeMux +} + +// ServeStream subscribes to network events and sends them to the client as a +// stream of Server-Sent-Events, with each event being a JSON encoded CyUpdate +// object +func (n *NetworkController) ServeStream(w http.ResponseWriter, req *http.Request) { + sub := n.events.Subscribe(ConnectivityEvents...) + defer sub.Unsubscribe() + + // stop the stream if the client goes away + var clientGone <-chan bool + if cn, ok := w.(http.CloseNotifier); ok { + clientGone = cn.CloseNotify() + } + + // write writes the given event and data to the stream like: + // + // event: + // data: + // + write := func(event, data string) { + fmt.Fprintf(w, "event: %s\n", event) + fmt.Fprintf(w, "data: %s\n\n", data) + if fw, ok := w.(http.Flusher); ok { + fw.Flush() + } + } + + w.Header().Set("Content-Type", "text/event-stream") + ch := sub.Chan() + for { + select { + case event := <-ch: + // convert the event to a CyUpdate + update, err := NewCyUpdate(event) + if err != nil { + write("error", err.Error()) + return + } + data, err := json.Marshal(update) + if err != nil { + write("error", err.Error()) + return + } + write("cyupdate", string(data)) + case <-clientGone: + return + } + } +} + // NewNetworkController creates a ResourceController responding to GET and DELETE methods // it embeds a mockers controller, a journal player, node and connection contollers. // @@ -108,7 +165,7 @@ func NewNetworkController(net NetworkControl, nodesController *ResourceControlle self.SetResource("mockevents", NewMockersController(eventer, conf.DefaultMockerConfig)) self.SetResource("journals", NewJournalPlayersController(eventer)) - return Controller(self) + return &NetworkController{ResourceController: self, events: net.Events()} } func NewNodesController(net *Network) *ResourceController { diff --git a/p2p/simulations/rest_api_server.go b/p2p/simulations/rest_api_server.go index b8c2cd5bf5..fe4237a4b0 100644 --- a/p2p/simulations/rest_api_server.go +++ b/p2p/simulations/rest_api_server.go @@ -16,6 +16,10 @@ type Controller interface { SetResource(id string, c Controller) } +type StreamController interface { + ServeStream(http.ResponseWriter, *http.Request) +} + // starts up http server func StartRestApiServer(port string, c Controller) { serveMux := http.NewServeMux() @@ -35,7 +39,6 @@ func handle(w http.ResponseWriter, r *http.Request, c Controller) { requestURL := r.URL log.Debug(fmt.Sprintf("HTTP %s request URL: '%s', Host: '%s', Path: '%s', Referer: '%s', Accept: '%s'", r.Method, r.RequestURI, requestURL.Host, requestURL.Path, r.Referer(), r.Header.Get("Accept"))) uri := requestURL.Path - w.Header().Set("Content-Type", "text/json") w.Header().Set("Access-Control-Allow-Origin", "*") w.Header().Set("Access-Control-Allow-Methods", "GET, POST, PUT, DELETE, OPTIONS") defer r.Body.Close() @@ -51,6 +54,18 @@ func handle(w http.ResponseWriter, r *http.Request, c Controller) { return } } + + // if the request is for a stream, call c.ServeStream + if r.Header.Get("Accept") == "text/event-stream" { + streamer, ok := c.(StreamController) + if !ok { + http.Error(w, "stream not supported", http.StatusBadRequest) + return + } + streamer.ServeStream(w, r) + return + } + handler, err := c.Handle(r.Method) if err != nil { http.Error(w, fmt.Sprintf("method %v not allowed (%v)", r.Method, err), http.StatusMethodNotAllowed) @@ -62,5 +77,6 @@ func handle(w http.ResponseWriter, r *http.Request, c Controller) { http.Error(w, fmt.Sprintf("handler error: %v", err), http.StatusBadRequest) return } + w.Header().Set("Content-Type", "text/json") http.ServeContent(w, r, "", time.Now(), response) }