From 584f8b8cf4d979dd91a348c804d6d9e67505d248 Mon Sep 17 00:00:00 2001 From: Lewis Marshall Date: Sun, 7 May 2017 12:40:24 +0100 Subject: [PATCH] p2p/simulations: Refactor events Signed-off-by: Lewis Marshall --- p2p/simulations/adapters/types.go | 17 ++ p2p/simulations/events.go | 61 ++++ p2p/simulations/http.go | 30 +- p2p/simulations/journal.go | 50 ++-- p2p/simulations/journal_test.go | 34 +-- p2p/simulations/mocker.go | 33 +-- p2p/simulations/network.go | 261 ++++-------------- p2p/simulations/sim_events.go | 137 --------- p2p/simulations/simulation.go | 6 +- .../simulations/discovery/discovery_test.go | 1 - 10 files changed, 210 insertions(+), 420 deletions(-) create mode 100644 p2p/simulations/events.go delete mode 100644 p2p/simulations/sim_events.go diff --git a/p2p/simulations/adapters/types.go b/p2p/simulations/adapters/types.go index aa08301e86..5e9e34ce8d 100644 --- a/p2p/simulations/adapters/types.go +++ b/p2p/simulations/adapters/types.go @@ -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 { diff --git a/p2p/simulations/events.go b/p2p/simulations/events.go new file mode 100644 index 0000000000..8e8168c34b --- /dev/null +++ b/p2p/simulations/events.go @@ -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(" id: %s up: %t", e.Node.ID().Label(), e.Node.Up) + case EventTypeConn: + return fmt.Sprintf(" nodes: %s->%s up: %t", e.Conn.One.Label(), e.Conn.Other.Label(), e.Conn.Up) + case EventTypeMsg: + return fmt.Sprintf(" nodes: %s->%s code: %d, received: %t", e.Msg.One.Label(), e.Msg.Other.Label(), e.Msg.Code, e.Msg.Received) + default: + return "" + } +} diff --git a/p2p/simulations/http.go b/p2p/simulations/http.go index 2d3acd9b06..7809830d40 100644 --- a/p2p/simulations/http.go +++ b/p2p/simulations/http.go @@ -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 } diff --git a/p2p/simulations/journal.go b/p2p/simulations/journal.go index 4d232e3d29..a827d3957e 100644 --- a/p2p/simulations/journal.go +++ b/p2p/simulations/journal.go @@ -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 } diff --git a/p2p/simulations/journal_test.go b/p2p/simulations/journal_test.go index 8a6890a667..1c25ea7900 100644 --- a/p2p/simulations/journal_test.go +++ b/p2p/simulations/journal_test.go @@ -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) } diff --git a/p2p/simulations/mocker.go b/p2p/simulations/mocker.go index 05d5a2e017..8f03c801d6 100644 --- a/p2p/simulations/mocker.go +++ b/p2p/simulations/mocker.go @@ -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++ } diff --git a/p2p/simulations/network.go b/p2p/simulations/network.go index e24c9b5cef..b90f367b89 100644 --- a/p2p/simulations/network.go +++ b/p2p/simulations/network.go @@ -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 <-self.quitc: + for { + select { + case event, ok := <-events: + if !ok { 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)) - } - err = self.Start(ne.Node.Id) - if err != nil { - log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err)) - } - } else { - err := self.Stop(ne.Node.Id) - if err != nil { - log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err)) + if event.Control { + self.executeControlEvent(event) + } + case <-self.quitc: + return } } - 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)) +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) } - } else { - err := self.Disconnect(ce.Connection.One, ce.Connection.Other) - if err != nil { - log.Trace(fmt.Sprintf("error execute event %v: %v", ce, err)) + case EventTypeConn: + if err := self.executeConnEvent(event); err != nil { + log.Error("error executing conn event", "event", event, "err", err) } + case EventTypeMsg: + log.Warn("ignoring control msg event") } - ce.Connection.controlFired = ce.Up } -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()) - } - } 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()) - } +func (self *Network) executeNodeEvent(e *Event) error { + if !e.Node.Up { + return self.Stop(e.Node.ID()) } - if ne, ok := ev.(*NodeControlEvent); ok { - self.executeNodeEvent(ne) - } else if ce, ok := ev.(*ConnControlEvent); ok { - self.executeConnEvent(ce) + + if err := self.NewNodeWithConfig(e.Node.Config); err != nil { + return err + } + return self.Start(e.Node.ID()) +} + +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("\n", self.Up, self.Node) -} - -func (self *ConnEvent) String() string { - return fmt.Sprintf("\n", self.Up, self.Reverse, self.Connection) -} - -func (self *MsgEvent) String() string { - return fmt.Sprintf("\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() @@ -361,11 +221,12 @@ func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) error { return err } node := &Node{ - Node: adapterNode, - NodeConfig: *conf, + Node: adapterNode, + 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) } } } diff --git a/p2p/simulations/sim_events.go b/p2p/simulations/sim_events.go deleted file mode 100644 index e17091a22f..0000000000 --- a/p2p/simulations/sim_events.go +++ /dev/null @@ -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 -} diff --git a/p2p/simulations/simulation.go b/p2p/simulations/simulation.go index f38d1701f4..008905fe9e 100644 --- a/p2p/simulations/simulation.go +++ b/p2p/simulations/simulation.go @@ -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 } diff --git a/swarm/network/simulations/discovery/discovery_test.go b/swarm/network/simulations/discovery/discovery_test.go index 5a5ec53f37..5137adf91f 100644 --- a/swarm/network/simulations/discovery/discovery_test.go +++ b/swarm/network/simulations/discovery/discovery_test.go @@ -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()