mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 15:16:43 +00:00
p2p/simulations: Refactor events
Signed-off-by: Lewis Marshall <lewis@lmars.net>
This commit is contained in:
parent
aa624cdb99
commit
584f8b8cf4
10 changed files with 210 additions and 420 deletions
|
|
@ -93,6 +93,23 @@ func (self *NodeId) Label() string {
|
|||
return self.String()[:4]
|
||||
}
|
||||
|
||||
func (self *NodeId) MarshalJSON() ([]byte, error) {
|
||||
return json.Marshal(hex.EncodeToString(self.NodeID[:]))
|
||||
}
|
||||
|
||||
func (self *NodeId) UnmarshalJSON(data []byte) error {
|
||||
var s string
|
||||
if err := json.Unmarshal(data, &s); err != nil {
|
||||
return err
|
||||
}
|
||||
id, err := discover.HexID(s)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
self.NodeID = id
|
||||
return nil
|
||||
}
|
||||
|
||||
// NodeConfig is the configuration used to start a node in a simulation
|
||||
// network
|
||||
type NodeConfig struct {
|
||||
|
|
|
|||
61
p2p/simulations/events.go
Normal file
61
p2p/simulations/events.go
Normal file
|
|
@ -0,0 +1,61 @@
|
|||
package simulations
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
type EventType string
|
||||
|
||||
const (
|
||||
EventTypeNode EventType = "node"
|
||||
EventTypeConn EventType = "conn"
|
||||
EventTypeMsg EventType = "msg"
|
||||
)
|
||||
|
||||
type Event struct {
|
||||
Type EventType `json:"type"`
|
||||
Time time.Time `json:"time"`
|
||||
Control bool `json:"control"`
|
||||
|
||||
Node *Node `json:"node,omitempty"`
|
||||
Conn *Conn `json:"conn,omitempty"`
|
||||
Msg *Msg `json:"msg,omitempty"`
|
||||
}
|
||||
|
||||
func NewEvent(v interface{}) *Event {
|
||||
event := &Event{Time: time.Now()}
|
||||
switch v := v.(type) {
|
||||
case *Node:
|
||||
event.Type = EventTypeNode
|
||||
event.Node = v
|
||||
case *Conn:
|
||||
event.Type = EventTypeConn
|
||||
event.Conn = v
|
||||
case *Msg:
|
||||
event.Type = EventTypeMsg
|
||||
event.Msg = v
|
||||
default:
|
||||
panic(fmt.Sprintf("invalid event type: %T", v))
|
||||
}
|
||||
return event
|
||||
}
|
||||
|
||||
func ControlEvent(v interface{}) *Event {
|
||||
event := NewEvent(v)
|
||||
event.Control = true
|
||||
return event
|
||||
}
|
||||
|
||||
func (e *Event) String() string {
|
||||
switch e.Type {
|
||||
case EventTypeNode:
|
||||
return fmt.Sprintf("<node-event> id: %s up: %t", e.Node.ID().Label(), e.Node.Up)
|
||||
case EventTypeConn:
|
||||
return fmt.Sprintf("<conn-event> nodes: %s->%s up: %t", e.Conn.One.Label(), e.Conn.Other.Label(), e.Conn.Up)
|
||||
case EventTypeMsg:
|
||||
return fmt.Sprintf("<msg-event> nodes: %s->%s code: %d, received: %t", e.Msg.One.Label(), e.Msg.Other.Label(), e.Msg.Code, e.Msg.Received)
|
||||
default:
|
||||
return ""
|
||||
}
|
||||
}
|
||||
|
|
@ -118,7 +118,8 @@ func (s *Server) StartMocker(w http.ResponseWriter, req *http.Request) {
|
|||
}
|
||||
|
||||
func (s *Server) streamNetworkEvents(network *Network, w http.ResponseWriter) {
|
||||
sub := network.events.Subscribe(ConnectivityAllEvents...)
|
||||
events := make(chan *Event)
|
||||
sub := network.events.Subscribe(events)
|
||||
defer sub.Unsubscribe()
|
||||
|
||||
// stop the stream if the client goes away
|
||||
|
|
@ -140,23 +141,20 @@ func (s *Server) streamNetworkEvents(network *Network, w http.ResponseWriter) {
|
|||
}
|
||||
}
|
||||
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
ch := sub.Chan()
|
||||
w.Header().Set("Content-Type", "text/event-stream; charset=utf-8")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
if fw, ok := w.(http.Flusher); ok {
|
||||
fw.Flush()
|
||||
}
|
||||
for {
|
||||
select {
|
||||
case event := <-ch:
|
||||
// convert the event to a SimUpdate
|
||||
update, err := NewSimUpdate(event)
|
||||
case event := <-events:
|
||||
data, err := json.Marshal(event)
|
||||
if err != nil {
|
||||
write("error", err.Error())
|
||||
return
|
||||
}
|
||||
data, err := json.Marshal(update)
|
||||
if err != nil {
|
||||
write("error", err.Error())
|
||||
return
|
||||
}
|
||||
write("simupdate", string(data))
|
||||
write("network", string(data))
|
||||
case <-clientGone:
|
||||
return
|
||||
}
|
||||
|
|
@ -198,7 +196,7 @@ func (s *Server) StartNode(w http.ResponseWriter, req *http.Request) {
|
|||
network := req.Context().Value("network").(*Network)
|
||||
node := req.Context().Value("node").(*Node)
|
||||
|
||||
if err := network.Start(node.Id); err != nil {
|
||||
if err := network.Start(node.ID()); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
|
@ -210,7 +208,7 @@ func (s *Server) StopNode(w http.ResponseWriter, req *http.Request) {
|
|||
network := req.Context().Value("network").(*Network)
|
||||
node := req.Context().Value("node").(*Node)
|
||||
|
||||
if err := network.Stop(node.Id); err != nil {
|
||||
if err := network.Stop(node.ID()); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
|
@ -223,7 +221,7 @@ func (s *Server) ConnectNode(w http.ResponseWriter, req *http.Request) {
|
|||
node := req.Context().Value("node").(*Node)
|
||||
peer := req.Context().Value("peer").(*Node)
|
||||
|
||||
if err := network.Connect(node.Id, peer.Id); err != nil {
|
||||
if err := network.Connect(node.ID(), peer.ID()); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
|
@ -236,7 +234,7 @@ func (s *Server) DisconnectNode(w http.ResponseWriter, req *http.Request) {
|
|||
node := req.Context().Value("node").(*Node)
|
||||
peer := req.Context().Value("peer").(*Node)
|
||||
|
||||
if err := network.Disconnect(node.Id, peer.Id); err != nil {
|
||||
if err := network.Disconnect(node.ID(), peer.ID()); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,12 +16,13 @@ import (
|
|||
// (using event.TypeMux). Network components POST events to the TypeMux, which then is
|
||||
// read by the journal. Each journal belongs to a subscription.
|
||||
type Journal struct {
|
||||
Id string
|
||||
Id string `json:"id"`
|
||||
Events []*Event `json:"events"`
|
||||
|
||||
lock sync.Mutex
|
||||
counter int
|
||||
cursor int
|
||||
quitc chan bool
|
||||
Events []*event.TypeMuxEvent
|
||||
}
|
||||
|
||||
// NewJournal constructor
|
||||
|
|
@ -34,19 +35,19 @@ func NewJournal() *Journal {
|
|||
return &Journal{quitc: make(chan bool)}
|
||||
}
|
||||
|
||||
// Subscribe takes an event.TypeMux and subscibes to types
|
||||
// and launches a gorourine that appends any new event to the event log
|
||||
// used for journalling history of a network
|
||||
// Subscribe subscribes to an event.Feed and launches a gorourine that appends
|
||||
// any new event to the event log used for journalling history of a network
|
||||
// the goroutine terminates when the journal is closed
|
||||
func (self *Journal) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
||||
func (self *Journal) Subscribe(feed *event.Feed) {
|
||||
log.Info("subscribe")
|
||||
sub := eventer.Subscribe(types...)
|
||||
events := make(chan *Event)
|
||||
sub := feed.Subscribe(events)
|
||||
go func() {
|
||||
defer sub.Unsubscribe()
|
||||
for {
|
||||
select {
|
||||
case ev := <-sub.Chan():
|
||||
self.append(ev)
|
||||
case event := <-events:
|
||||
self.append(event)
|
||||
case <-self.quitc:
|
||||
return
|
||||
}
|
||||
|
|
@ -75,11 +76,12 @@ func NewJournalFromJSON(b []byte) (*Journal, error) {
|
|||
// params:
|
||||
// * acc: using acceleration factor acc
|
||||
// * journal: journal to use
|
||||
// * eventer: where to post the replayed events
|
||||
func Replay(acc float64, j *Journal, eventer *event.TypeMux) {
|
||||
f := func(d interface{}) bool {
|
||||
// * feed: where to post the replayed events
|
||||
func Replay(acc float64, j *Journal, feed *event.Feed) {
|
||||
f := func(event *Event) bool {
|
||||
// reposts the data with the eventer (the data receives a new timestamp)
|
||||
eventer.Post(d)
|
||||
event.Time = time.Now()
|
||||
feed.Send(event)
|
||||
return true
|
||||
}
|
||||
j.TimedRead(acc, f)
|
||||
|
|
@ -97,10 +99,10 @@ func (self *Journal) Close() {
|
|||
close(self.quitc)
|
||||
}
|
||||
|
||||
func (self *Journal) append(evs ...*event.TypeMuxEvent) {
|
||||
func (self *Journal) append(events ...*Event) {
|
||||
self.lock.Lock()
|
||||
defer self.lock.Unlock()
|
||||
self.Events = append(self.Events, evs...)
|
||||
self.Events = append(self.Events, events...)
|
||||
self.counter++
|
||||
}
|
||||
|
||||
|
|
@ -116,7 +118,7 @@ func (self *Journal) WaitEntries(n int) {
|
|||
}
|
||||
}
|
||||
|
||||
func (self *Journal) Read(f func(*event.TypeMuxEvent) bool) (read int) {
|
||||
func (self *Journal) Read(f func(*Event) bool) (read int) {
|
||||
self.lock.Lock()
|
||||
defer self.lock.Unlock()
|
||||
ok := true
|
||||
|
|
@ -139,20 +141,20 @@ func (self *Journal) Read(f func(*event.TypeMuxEvent) bool) (read int) {
|
|||
// NOTE: the events' timestamps are supposed to be strictly ordered otherwise
|
||||
// the call panics.
|
||||
// acc is an acceleration factor
|
||||
func (self *Journal) TimedRead(acc float64, f func(interface{}) bool) (read int) {
|
||||
func (self *Journal) TimedRead(acc float64, f func(*Event) bool) (read int) {
|
||||
var lastEvent time.Time
|
||||
timer := time.NewTimer(0)
|
||||
var data interface{}
|
||||
h := func(ev *event.TypeMuxEvent) bool {
|
||||
var event *Event
|
||||
h := func(e *Event) bool {
|
||||
// wait for the interval time passes event time
|
||||
if ev.Time.Before(lastEvent) {
|
||||
if e.Time.Before(lastEvent) {
|
||||
panic("events not ordered")
|
||||
}
|
||||
interval := ev.Time.Sub(lastEvent)
|
||||
interval := e.Time.Sub(lastEvent)
|
||||
log.Trace(fmt.Sprintf("reset timer to interval %v", interval))
|
||||
timer.Reset(time.Duration(acc) * interval)
|
||||
lastEvent = ev.Time
|
||||
data = ev.Data
|
||||
lastEvent = e.Time
|
||||
event = e
|
||||
return false
|
||||
}
|
||||
var n int
|
||||
|
|
@ -168,7 +170,7 @@ func (self *Journal) TimedRead(acc float64, f func(interface{}) bool) (read int)
|
|||
}
|
||||
}
|
||||
read += n
|
||||
if n == 0 || !f(data) {
|
||||
if n == 0 || !f(event) {
|
||||
log.Trace(fmt.Sprintf("timed read ends (read %v entries)", read))
|
||||
break
|
||||
}
|
||||
|
|
|
|||
|
|
@ -10,16 +10,14 @@ import (
|
|||
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
|
||||
)
|
||||
|
||||
func testEvents(intervals ...int) (events []*event.TypeMuxEvent) {
|
||||
func testEvents(intervals ...int) (events []*Event) {
|
||||
t := time.Now()
|
||||
for _, interval := range intervals {
|
||||
t = t.Add(time.Duration(interval) * time.Millisecond)
|
||||
events = append(events, &event.TypeMuxEvent{
|
||||
events = append(events, &Event{
|
||||
Type: EventTypeNode,
|
||||
Time: t,
|
||||
Data: interface{}(&NodeEvent{
|
||||
Type: "node",
|
||||
Action: "down",
|
||||
}),
|
||||
Node: &Node{Up: false},
|
||||
})
|
||||
}
|
||||
return events
|
||||
|
|
@ -33,8 +31,7 @@ func TestTimedRead(t *testing.T) {
|
|||
var i int
|
||||
acc := 0.5
|
||||
length := 4
|
||||
f := func(data interface{}) bool {
|
||||
_ = data.(*NodeEvent)
|
||||
f := func(event *Event) bool {
|
||||
newTimes = append(newTimes, time.Now())
|
||||
i++
|
||||
return i <= length
|
||||
|
|
@ -68,10 +65,15 @@ func testIDs() (ids []*adapters.NodeId) {
|
|||
}
|
||||
|
||||
func testJournal(ids []*adapters.NodeId) *Journal {
|
||||
eventer := &event.TypeMux{}
|
||||
eventer := &event.Feed{}
|
||||
journal := NewJournal()
|
||||
journal.Subscribe(eventer, ConnectivityEvents...)
|
||||
mockNewNodes(eventer, ids)
|
||||
journal.Subscribe(eventer)
|
||||
for _, id := range ids {
|
||||
eventer.Send(&Event{
|
||||
Type: EventTypeNode,
|
||||
Node: &Node{Config: &adapters.NodeConfig{Id: id}, Up: true},
|
||||
})
|
||||
}
|
||||
journal.WaitEntries(len(ids))
|
||||
return journal
|
||||
}
|
||||
|
|
@ -80,7 +82,7 @@ func TestSubscribe(t *testing.T) {
|
|||
ids := testIDs()
|
||||
journal := testJournal(ids)
|
||||
for i, ev := range journal.Events {
|
||||
id := ev.Data.(*NodeEvent).node.Id
|
||||
id := ev.Node.ID()
|
||||
if id != ids[i] {
|
||||
t.Fatalf("incorrect id: expected %v, got %v", id, ids[i])
|
||||
}
|
||||
|
|
@ -115,15 +117,15 @@ func TestLoadSave(t *testing.T) {
|
|||
|
||||
func TestReplay(t *testing.T) {
|
||||
_, jo := loadTestJournal(t)
|
||||
eventer := &event.TypeMux{}
|
||||
eventer := &event.Feed{}
|
||||
|
||||
journal := NewJournal()
|
||||
journal.Subscribe(eventer, ConnectivityEvents...)
|
||||
journal.Subscribe(eventer)
|
||||
|
||||
Replay(0, jo, eventer)
|
||||
for i, ev := range jo.Events {
|
||||
exp := ev.Data.(*NodeEvent).String()
|
||||
got := journal.Events[i].Data.(*NodeEvent).String()
|
||||
exp := ev.String()
|
||||
got := journal.Events[i].String()
|
||||
if exp != got {
|
||||
t.Fatalf("incorrent replayed journal entry at pos %v: expected %v, got %v", i, exp, got)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -57,7 +57,7 @@ func DefaultMockerConfig() *MockerConfig {
|
|||
// to the eventer
|
||||
// The journal using the eventer can then be read to visualise or
|
||||
// drive connections
|
||||
func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConfig) {
|
||||
func MockEvents(eventer *event.Feed, ids []*adapters.NodeId, conf *MockerConfig) {
|
||||
|
||||
var onNodes []*Node
|
||||
offNodes := ids
|
||||
|
|
@ -105,22 +105,15 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
|||
for i := 0; len(onNodes) > 0 && i < nodesDown; i++ {
|
||||
c := rand.Intn(len(onNodes))
|
||||
sn := onNodes[c]
|
||||
err := eventer.Post(sn.EmitEvent(ControlEvent))
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
eventer.Send(ControlEvent(sn))
|
||||
onNodes = append(onNodes[0:c], onNodes[c+1:]...)
|
||||
offNodes = append(offNodes, sn.Id)
|
||||
offNodes = append(offNodes, sn.ID())
|
||||
}
|
||||
var mustconnect []int
|
||||
for i := 0; len(offNodes) > 0 && i < nodesUp; i++ {
|
||||
c := rand.Intn(len(offNodes))
|
||||
sn := &Node{}
|
||||
sn.Id = offNodes[c]
|
||||
err := eventer.Post(sn.EmitEvent(ControlEvent))
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
sn := &Node{Config: &adapters.NodeConfig{Id: offNodes[c]}}
|
||||
eventer.Send(ControlEvent(sn))
|
||||
mustconnect = append(mustconnect, len(onNodes))
|
||||
onNodes = append(onNodes, sn)
|
||||
offNodes = append(offNodes[0:c], offNodes[c+1:]...)
|
||||
|
|
@ -145,7 +138,7 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
|||
m := n + rand.Intn(len(onNodes)-n)
|
||||
// m := n + 1 + rand.Intn(len(onNodes)-n-1)
|
||||
for k := m; k < len(onNodes); k++ {
|
||||
lab := ConnLabel(onNodes[n].Id, onNodes[k].Id)
|
||||
lab := ConnLabel(onNodes[n].ID(), onNodes[k].ID())
|
||||
var j int
|
||||
j, found = onConnsMap[lab]
|
||||
if found {
|
||||
|
|
@ -157,8 +150,8 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
|||
break
|
||||
}
|
||||
connected[k] = true
|
||||
caller := onNodes[n].Id
|
||||
callee := onNodes[k].Id
|
||||
caller := onNodes[n].ID()
|
||||
callee := onNodes[k].ID()
|
||||
|
||||
sc := &Conn{
|
||||
One: caller,
|
||||
|
|
@ -176,10 +169,7 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
|||
lab := ConnLabel(sc.One, sc.Other)
|
||||
onConnsMap[lab] = len(onConns)
|
||||
onConns = append(onConns, sc)
|
||||
err := eventer.Post(sc.EmitEvent(ControlEvent))
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
eventer.Send(ControlEvent(sc))
|
||||
}
|
||||
|
||||
for i := 0; len(onConns) > 0 && i < connsDown; i++ {
|
||||
|
|
@ -188,10 +178,7 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
|||
onConns = append(onConns[0:c], onConns[c+1:]...)
|
||||
lab := ConnLabel(conn.One, conn.Other)
|
||||
delete(onConnsMap, lab)
|
||||
err := eventer.Post(conn.EmitEvent(ControlEvent))
|
||||
if err != nil {
|
||||
panic(err.Error())
|
||||
}
|
||||
eventer.Send(ControlEvent(conn))
|
||||
}
|
||||
rounds++
|
||||
}
|
||||
|
|
|
|||
|
|
@ -56,12 +56,6 @@ type NetworkControl interface {
|
|||
Subscribe(*event.TypeMux, ...interface{})
|
||||
}
|
||||
|
||||
// event types related to connectivity, i.e., nodes coming on dropping off
|
||||
// and connections established and dropped
|
||||
var ConnectivityControlEvents = []interface{}{&NodeControlEvent{}, &ConnControlEvent{}, &MsgControlEvent{}}
|
||||
var ConnectivityLiveEvents = []interface{}{&NodeEvent{}, &ConnEvent{}, &MsgEvent{}}
|
||||
var ConnectivityAllEvents = append(ConnectivityControlEvents, ConnectivityLiveEvents...)
|
||||
|
||||
// Network models a p2p network
|
||||
// the actual logic of bringing nodes and connections up and down and
|
||||
// messaging is implemented in the particular NodeAdapter interface
|
||||
|
|
@ -69,7 +63,7 @@ type Network struct {
|
|||
nodeAdapter adapters.NodeAdapter
|
||||
|
||||
// input trigger events and other events
|
||||
events *event.TypeMux // generated events a journal can subsribe to
|
||||
events event.Feed // generated events a journal can subsribe to
|
||||
lock sync.RWMutex
|
||||
nodeMap map[discover.NodeID]int
|
||||
connMap map[string]int
|
||||
|
|
@ -83,125 +77,85 @@ func NewNetwork(nodeAdapter adapters.NodeAdapter, conf *NetworkConfig) *Network
|
|||
return &Network{
|
||||
nodeAdapter: nodeAdapter,
|
||||
conf: conf,
|
||||
events: &event.TypeMux{},
|
||||
nodeMap: make(map[discover.NodeID]int),
|
||||
connMap: make(map[string]int),
|
||||
quitc: make(chan bool),
|
||||
}
|
||||
}
|
||||
|
||||
// Subscribe takes an event.TypeMux and subscibes to types
|
||||
// and launches a goroutine that reads control events from an eventer Subsription channel
|
||||
// and executes the events
|
||||
func (self *Network) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
||||
// Subscribe reads control events from a channel and executes them
|
||||
func (self *Network) Subscribe(events chan *Event) {
|
||||
log.Info("subscribe")
|
||||
sub := eventer.Subscribe(types...)
|
||||
go func() {
|
||||
defer sub.Unsubscribe()
|
||||
for {
|
||||
select {
|
||||
case ev := <-sub.Chan():
|
||||
self.execute(ev)
|
||||
case event, ok := <-events:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if event.Control {
|
||||
self.executeControlEvent(event)
|
||||
}
|
||||
case <-self.quitc:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func (self *Network) executeNodeEvent(ne *NodeControlEvent) {
|
||||
if ne.Up {
|
||||
err := self.NewNodeWithConfig(&ne.Node.NodeConfig)
|
||||
if err != nil {
|
||||
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||
func (self *Network) executeControlEvent(event *Event) {
|
||||
log.Trace("execute control event", "type", event.Type, "event", event)
|
||||
switch event.Type {
|
||||
case EventTypeNode:
|
||||
if err := self.executeNodeEvent(event); err != nil {
|
||||
log.Error("error executing node event", "event", event, "err", err)
|
||||
}
|
||||
err = self.Start(ne.Node.Id)
|
||||
if err != nil {
|
||||
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||
case EventTypeConn:
|
||||
if err := self.executeConnEvent(event); err != nil {
|
||||
log.Error("error executing conn event", "event", event, "err", err)
|
||||
}
|
||||
} else {
|
||||
err := self.Stop(ne.Node.Id)
|
||||
if err != nil {
|
||||
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||
case EventTypeMsg:
|
||||
log.Warn("ignoring control msg event")
|
||||
}
|
||||
}
|
||||
ne.Node.controlFired = ne.Up
|
||||
}
|
||||
|
||||
func (self *Network) executeConnEvent(ce *ConnControlEvent) {
|
||||
if ce.Up {
|
||||
err := self.Connect(ce.Connection.One, ce.Connection.Other)
|
||||
if err != nil {
|
||||
log.Trace(fmt.Sprintf("error execute event %v: %v", ce, err))
|
||||
}
|
||||
} else {
|
||||
err := self.Disconnect(ce.Connection.One, ce.Connection.Other)
|
||||
if err != nil {
|
||||
log.Trace(fmt.Sprintf("error execute event %v: %v", ce, err))
|
||||
}
|
||||
}
|
||||
ce.Connection.controlFired = ce.Up
|
||||
func (self *Network) executeNodeEvent(e *Event) error {
|
||||
if !e.Node.Up {
|
||||
return self.Stop(e.Node.ID())
|
||||
}
|
||||
|
||||
func (self *Network) execute(in *event.TypeMuxEvent) {
|
||||
log.Trace(fmt.Sprintf("execute event %v", in))
|
||||
ev := in.Data
|
||||
if ne, ok := ev.(*NodeEvent); ok {
|
||||
if ne.Up && ne.Node.controlFired || (!ne.Up && !ne.Node.controlFired) {
|
||||
log.Trace(fmt.Sprintf("Got NodeEvent %v, but Control Event has already been applied for : %v", ne, ne.Node))
|
||||
//ignore this real event; control event already took care of this
|
||||
} else {
|
||||
self.executeNodeEvent(ne.ToControlEvent())
|
||||
if err := self.NewNodeWithConfig(e.Node.Config); err != nil {
|
||||
return err
|
||||
}
|
||||
} else if ce, ok := ev.(*ConnEvent); ok {
|
||||
if ce.Up && ce.Connection.controlFired || (!ce.Up && !ce.Connection.controlFired) {
|
||||
log.Trace(fmt.Sprintf("Got ConnEvent %v, but Control Event has already been applied for : %v", ce, ce.Connection))
|
||||
//ignore this real event; control event already took care of this
|
||||
} else {
|
||||
self.executeConnEvent(ce.ToControlEvent())
|
||||
return self.Start(e.Node.ID())
|
||||
}
|
||||
}
|
||||
if ne, ok := ev.(*NodeControlEvent); ok {
|
||||
self.executeNodeEvent(ne)
|
||||
} else if ce, ok := ev.(*ConnControlEvent); ok {
|
||||
self.executeConnEvent(ce)
|
||||
|
||||
func (self *Network) executeConnEvent(e *Event) error {
|
||||
if e.Conn.Up {
|
||||
return self.Connect(e.Conn.One, e.Conn.Other)
|
||||
} else {
|
||||
log.Trace(fmt.Sprintf("event: %#v", ev))
|
||||
panic("unhandled event")
|
||||
return self.Disconnect(e.Conn.One, e.Conn.Other)
|
||||
}
|
||||
}
|
||||
|
||||
// Events returns the output eventer of the Network.
|
||||
func (self *Network) Events() *event.TypeMux {
|
||||
return self.events
|
||||
}
|
||||
|
||||
type EventType int
|
||||
|
||||
const (
|
||||
ControlEvent EventType = iota
|
||||
LiveEvent
|
||||
)
|
||||
|
||||
type EventEmitter interface {
|
||||
EmitEvent()
|
||||
}
|
||||
|
||||
type LiveEventer interface {
|
||||
ToControlEvent()
|
||||
func (self *Network) Events() *event.Feed {
|
||||
return &self.events
|
||||
}
|
||||
|
||||
type Node struct {
|
||||
adapters.Node
|
||||
adapters.NodeConfig
|
||||
adapters.Node `json:"-"`
|
||||
|
||||
Up bool
|
||||
Config *adapters.NodeConfig `json:"config"`
|
||||
Up bool `json:"up"`
|
||||
|
||||
controlFired bool
|
||||
}
|
||||
|
||||
func (self *Node) ID() *adapters.NodeId {
|
||||
return self.Config.Id
|
||||
}
|
||||
|
||||
func (self *Node) String() string {
|
||||
return fmt.Sprintf("Node %v", self.Id.Label())
|
||||
return fmt.Sprintf("Node %v", self.ID().Label())
|
||||
}
|
||||
|
||||
// active connections are represented by the Node entry object so that
|
||||
|
|
@ -236,100 +190,6 @@ func (self *Msg) String() string {
|
|||
return fmt.Sprintf("Msg(%d) %v->%v", self.Code, self.One.Label(), self.Other.Label())
|
||||
}
|
||||
|
||||
type NodeEvent struct {
|
||||
Node *Node
|
||||
Up bool
|
||||
}
|
||||
|
||||
type ConnEvent struct {
|
||||
Connection *Conn
|
||||
Up bool
|
||||
Reverse bool
|
||||
}
|
||||
|
||||
type MsgEvent struct {
|
||||
Message *Msg
|
||||
}
|
||||
|
||||
type NodeControlEvent struct {
|
||||
*NodeEvent
|
||||
}
|
||||
|
||||
type ConnControlEvent struct {
|
||||
*ConnEvent
|
||||
}
|
||||
|
||||
type MsgControlEvent struct {
|
||||
*MsgEvent
|
||||
}
|
||||
|
||||
func (self *NodeEvent) String() string {
|
||||
return fmt.Sprintf("<Up: %v, Data: %v>\n", self.Up, self.Node)
|
||||
}
|
||||
|
||||
func (self *ConnEvent) String() string {
|
||||
return fmt.Sprintf("<Up: %v, Reverse: %v, Data: %v>\n", self.Up, self.Reverse, self.Connection)
|
||||
}
|
||||
|
||||
func (self *MsgEvent) String() string {
|
||||
return fmt.Sprintf("<Msg: %v>\n", self.Message)
|
||||
}
|
||||
|
||||
func (self *Node) EmitEvent(eventType EventType) interface{} {
|
||||
evt := &NodeEvent{
|
||||
Node: self,
|
||||
Up: self.Up,
|
||||
}
|
||||
if eventType == ControlEvent {
|
||||
return &NodeControlEvent{
|
||||
evt,
|
||||
}
|
||||
} else {
|
||||
return evt
|
||||
}
|
||||
}
|
||||
|
||||
func (self *Conn) EmitEvent(eventType EventType) interface{} {
|
||||
evt := &ConnEvent{
|
||||
Connection: self,
|
||||
Up: self.Up,
|
||||
Reverse: self.Reverse,
|
||||
}
|
||||
|
||||
if eventType == ControlEvent {
|
||||
return &ConnControlEvent{
|
||||
evt,
|
||||
}
|
||||
} else {
|
||||
return evt
|
||||
}
|
||||
}
|
||||
|
||||
func (self *Msg) EmitEvent(eventType EventType) interface{} {
|
||||
evt := &MsgEvent{
|
||||
Message: self,
|
||||
}
|
||||
if eventType == ControlEvent {
|
||||
return &MsgControlEvent{
|
||||
evt,
|
||||
}
|
||||
} else {
|
||||
return evt
|
||||
}
|
||||
}
|
||||
|
||||
func (self *MsgEvent) ToControlEvent() *MsgControlEvent {
|
||||
return &MsgControlEvent{self}
|
||||
}
|
||||
|
||||
func (self *ConnEvent) ToControlEvent() *ConnControlEvent {
|
||||
return &ConnControlEvent{self}
|
||||
}
|
||||
|
||||
func (self *NodeEvent) ToControlEvent() *NodeControlEvent {
|
||||
return &NodeControlEvent{self}
|
||||
}
|
||||
|
||||
// NewNode adds a new node to the network with a random ID
|
||||
func (self *Network) NewNode() (*adapters.NodeConfig, error) {
|
||||
conf := adapters.RandomNodeConfig()
|
||||
|
|
@ -362,10 +222,11 @@ func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) error {
|
|||
}
|
||||
node := &Node{
|
||||
Node: adapterNode,
|
||||
NodeConfig: *conf,
|
||||
Config: conf,
|
||||
}
|
||||
self.Nodes = append(self.Nodes, node)
|
||||
log.Trace(fmt.Sprintf("node %v created", id))
|
||||
self.events.Send(ControlEvent(node))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -418,7 +279,7 @@ func (self *Network) Start(id *adapters.NodeId) error {
|
|||
node.Up = true
|
||||
log.Info(fmt.Sprintf("started node %v: %v", id, node.Up))
|
||||
|
||||
self.events.Post(node.EmitEvent(ControlEvent))
|
||||
self.events.Send(NewEvent(node))
|
||||
|
||||
// subscribe to peer events
|
||||
client, err := node.Client()
|
||||
|
|
@ -485,7 +346,7 @@ func (self *Network) Stop(id *adapters.NodeId) error {
|
|||
node.Up = false
|
||||
log.Info(fmt.Sprintf("stop node %v: %v", id, node.Up))
|
||||
|
||||
self.events.Post(node.EmitEvent(ControlEvent))
|
||||
self.events.Send(ControlEvent(node))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -526,7 +387,7 @@ func (self *Network) Connect(oneId, otherId *adapters.NodeId) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
self.events.Post(conn.EmitEvent(ControlEvent))
|
||||
self.events.Send(ControlEvent(conn))
|
||||
return client.Call(nil, "admin_addPeer", string(addr))
|
||||
}
|
||||
|
||||
|
|
@ -560,7 +421,7 @@ func (self *Network) Disconnect(oneId, otherId *adapters.NodeId) error {
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
self.events.Post(conn.EmitEvent(ControlEvent))
|
||||
self.events.Send(ControlEvent(conn))
|
||||
return client.Call(nil, "admin_removePeer", string(addr))
|
||||
}
|
||||
|
||||
|
|
@ -575,7 +436,7 @@ func (self *Network) DidConnect(one, other *adapters.NodeId) error {
|
|||
conn.Reverse = conn.One.NodeID != one.NodeID
|
||||
conn.Up = true
|
||||
// connection event posted
|
||||
self.events.Post(conn.EmitEvent(LiveEvent))
|
||||
self.events.Send(NewEvent(conn))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -589,7 +450,7 @@ func (self *Network) DidDisconnect(one, other *adapters.NodeId) error {
|
|||
}
|
||||
conn.Reverse = conn.One.NodeID != one.NodeID
|
||||
conn.Up = false
|
||||
self.events.Post(conn.EmitEvent(LiveEvent))
|
||||
self.events.Send(NewEvent(conn))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -601,7 +462,7 @@ func (self *Network) Send(senderid, receiverid *adapters.NodeId, msgcode uint64,
|
|||
Code: msgcode,
|
||||
}
|
||||
//self.GetNode(senderid).na.(*adapters.SimNode).GetPeer(receiverid).SendMsg(msgcode, protomsg) // phew!
|
||||
self.events.Post(msg.EmitEvent(ControlEvent))
|
||||
self.events.Send(ControlEvent(msg))
|
||||
}
|
||||
|
||||
func (self *Network) DidSend(sender, receiver *adapters.NodeId, msgcode uint64) error {
|
||||
|
|
@ -611,7 +472,7 @@ func (self *Network) DidSend(sender, receiver *adapters.NodeId, msgcode uint64)
|
|||
Code: msgcode,
|
||||
Received: false,
|
||||
}
|
||||
self.events.Post(msg.EmitEvent(LiveEvent))
|
||||
self.events.Send(NewEvent(msg))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -622,7 +483,7 @@ func (self *Network) DidReceive(sender, receiver *adapters.NodeId, msgcode uint6
|
|||
Code: msgcode,
|
||||
Received: true,
|
||||
}
|
||||
self.events.Post(msg.EmitEvent(LiveEvent))
|
||||
self.events.Send(NewEvent(msg))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -697,9 +558,9 @@ func (self *Network) Shutdown() {
|
|||
|
||||
// stop all nodes
|
||||
for _, node := range self.Nodes {
|
||||
log.Debug(fmt.Sprintf("stopping node %s", node.Id.Label()))
|
||||
log.Debug(fmt.Sprintf("stopping node %s", node.ID().Label()))
|
||||
if err := node.Stop(); err != nil {
|
||||
log.Warn(fmt.Sprintf("error stopping node %s", node.Id.Label()), "err", err)
|
||||
log.Warn(fmt.Sprintf("error stopping node %s", node.ID().Label()), "err", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,137 +0,0 @@
|
|||
package simulations
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/ethereum/go-ethereum/event"
|
||||
)
|
||||
|
||||
// TODO: to implement simulation global behav
|
||||
type SimConfig struct {
|
||||
}
|
||||
|
||||
type SimData struct {
|
||||
Id string `json:"id"`
|
||||
Source string `json:"source,omitempty"`
|
||||
Target string `json:"target,omitempty"`
|
||||
Up bool `json:"up"`
|
||||
}
|
||||
|
||||
type SimElement struct {
|
||||
Data *SimData `json:"data"`
|
||||
Classes string `json:"classes,omitempty"`
|
||||
Group string `json:"group"`
|
||||
Control bool `json:"control"`
|
||||
// selected: false, // whether the element is selected (default false)
|
||||
// selectable: true, // whether the selection state is mutable (default true)
|
||||
// locked: false, // when locked a node's position is immutable (default false)
|
||||
// grabbable: true, // whether the node can be grabbed and moved by the user
|
||||
}
|
||||
|
||||
type SimUpdate struct {
|
||||
Add []*SimElement `json:"add"`
|
||||
Remove []*SimElement `json:"remove"`
|
||||
Message []*SimElement `json:"message"`
|
||||
}
|
||||
|
||||
func NewSimUpdate(e *event.TypeMuxEvent) (*SimUpdate, error) {
|
||||
var update SimUpdate
|
||||
var el *SimElement
|
||||
entry := e.Data
|
||||
|
||||
switch entry.(type) {
|
||||
case *NodeControlEvent, *NodeEvent:
|
||||
var data *SimData
|
||||
var control bool
|
||||
nce, ok := entry.(*NodeControlEvent)
|
||||
if ok {
|
||||
data = &SimData{Id: nce.Node.Id.String()}
|
||||
data.Up = nce.Up
|
||||
control = true
|
||||
} else {
|
||||
ne := entry.(*NodeEvent)
|
||||
data = &SimData{Id: ne.Node.Id.String()}
|
||||
data.Up = ne.Up
|
||||
control = false
|
||||
}
|
||||
el = &SimElement{Group: "nodes", Data: data}
|
||||
el.Control = control
|
||||
if el.Data.Up {
|
||||
update.Add = append(update.Add, el)
|
||||
} else {
|
||||
update.Remove = append(update.Remove, el)
|
||||
}
|
||||
case *MsgControlEvent, *MsgEvent:
|
||||
var control bool
|
||||
var msg *Msg
|
||||
mce, ok := entry.(*MsgControlEvent)
|
||||
if ok {
|
||||
msg = mce.Message
|
||||
control = true
|
||||
} else {
|
||||
me := entry.(*MsgEvent)
|
||||
msg = me.Message
|
||||
control = false
|
||||
}
|
||||
id := ConnLabel(msg.One, msg.Other)
|
||||
var source, target string
|
||||
source = msg.One.String()
|
||||
target = msg.Other.String()
|
||||
el = &SimElement{Group: "msgs", Data: &SimData{Id: id, Source: source, Target: target}}
|
||||
el.Data.Up = true
|
||||
el.Control = control
|
||||
update.Message = append(update.Message, el)
|
||||
case *ConnControlEvent, *ConnEvent:
|
||||
var control bool
|
||||
var conn *Conn
|
||||
var up bool
|
||||
cce, ok := entry.(*ConnControlEvent)
|
||||
if ok {
|
||||
conn = cce.Connection
|
||||
up = cce.Up
|
||||
control = true
|
||||
} else {
|
||||
ce := entry.(*ConnEvent)
|
||||
conn = ce.Connection
|
||||
up = ce.Up
|
||||
control = false
|
||||
}
|
||||
// mutually exclusive directed edge (caller -> callee)
|
||||
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 = &SimElement{Group: "edges", Data: &SimData{Id: id, Source: source, Target: target}}
|
||||
el.Control = control
|
||||
el.Data.Up = up
|
||||
if up {
|
||||
update.Add = append(update.Add, el)
|
||||
} else {
|
||||
update.Remove = append(update.Remove, el)
|
||||
}
|
||||
default:
|
||||
return nil, fmt.Errorf("unknown event type: %T", entry)
|
||||
}
|
||||
|
||||
return &update, nil
|
||||
}
|
||||
|
||||
func UpdateSim(conf *SimConfig, j *Journal) (*SimUpdate, error) {
|
||||
var update SimUpdate
|
||||
j.Read(func(e *event.TypeMuxEvent) bool {
|
||||
u, err := NewSimUpdate(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
|
||||
})
|
||||
return &update, nil
|
||||
}
|
||||
|
|
@ -68,11 +68,11 @@ func (s *Simulation) Run(ctx context.Context, step *Step) (result *StepResult) {
|
|||
func (s *Simulation) watchNetwork(result *StepResult) func() {
|
||||
stop := make(chan struct{})
|
||||
done := make(chan struct{})
|
||||
sub := s.network.Events().Subscribe(ConnectivityAllEvents...)
|
||||
events := make(chan *Event)
|
||||
sub := s.network.Events().Subscribe(events)
|
||||
go func() {
|
||||
defer close(done)
|
||||
defer sub.Unsubscribe()
|
||||
events := sub.Chan()
|
||||
for {
|
||||
select {
|
||||
case event := <-events:
|
||||
|
|
@ -128,5 +128,5 @@ type StepResult struct {
|
|||
Passes map[*adapters.NodeId]time.Time
|
||||
|
||||
// NetworkEvents are the network events which occurred during the step
|
||||
NetworkEvents []interface{}
|
||||
NetworkEvents []*Event
|
||||
}
|
||||
|
|
|
|||
|
|
@ -62,7 +62,6 @@ func testDiscoverySimulation(t *testing.T, adapter adapters.NodeAdapter) {
|
|||
nodeCount := 10
|
||||
net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{
|
||||
Id: "0",
|
||||
Backend: true,
|
||||
DefaultService: serviceName,
|
||||
})
|
||||
defer net.Shutdown()
|
||||
|
|
|
|||
Loading…
Reference in a new issue