mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 22:56:43 +00:00
p2p/simulations: event streams
This commit is contained in:
parent
5b4702304a
commit
20e91384bf
3 changed files with 139 additions and 61 deletions
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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: <event>
|
||||
// data: <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 {
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue