From dee06a92d9ddf120b38e79999acb8ed53a325721 Mon Sep 17 00:00:00 2001 From: Anton Kaliaev Date: Tue, 21 Mar 2017 17:26:54 +0400 Subject: [PATCH 1/5] [tm-monitor] rewrite eventmeter to use go-kit/log --- tm-monitor/eventmeter/eventmeter.go | 20 +++++++++++++------- tm-monitor/glide.yaml | 1 - tm-monitor/mock/mock.go | 2 ++ tm-monitor/monitor/node.go | 2 ++ 4 files changed, 17 insertions(+), 8 deletions(-) diff --git a/tm-monitor/eventmeter/eventmeter.go b/tm-monitor/eventmeter/eventmeter.go index 5626ba134..dea4c1b1f 100644 --- a/tm-monitor/eventmeter/eventmeter.go +++ b/tm-monitor/eventmeter/eventmeter.go @@ -6,16 +6,14 @@ import ( "sync" "time" + "github.com/go-kit/kit/log" "github.com/gorilla/websocket" + "github.com/pkg/errors" metrics "github.com/rcrowley/go-metrics" events "github.com/tendermint/go-events" client "github.com/tendermint/go-rpc/client" - log15 "github.com/tendermint/log15" ) -// Log allows you to set your own logger. -var Log log15.Logger - //------------------------------------------------------ // Generic system to subscribe to events and record their frequency //------------------------------------------------------ @@ -96,6 +94,8 @@ type EventMeter struct { unmarshalEvent EventUnmarshalFunc quit chan struct{} + + logger log.Logger } func NewEventMeter(addr string, unmarshalEvent EventUnmarshalFunc) *EventMeter { @@ -106,10 +106,16 @@ func NewEventMeter(addr string, unmarshalEvent EventUnmarshalFunc) *EventMeter { receivedPong: true, unmarshalEvent: unmarshalEvent, quit: make(chan struct{}), + logger: log.NewNopLogger(), } return em } +// SetLogger lets you set your own logger +func (em *EventMeter) SetLogger(l log.Logger) { + em.logger = l +} + func (em *EventMeter) String() string { return em.wsc.Address } @@ -225,11 +231,11 @@ func (em *EventMeter) receiveRoutine() { select { case <-pingTicker.C: if pingAttempts, err = em.pingForLatency(pingAttempts); err != nil { - Log.Error("Failed to write ping message on websocket", err) + em.logger.Log("err", errors.Wrap(err, "Failed to write ping message on websocket")) em.StopAndReconnect() return } else if pingAttempts >= maxPingsPerPong { - Log.Error(fmt.Sprintf("Have not received a pong in %v", time.Duration(pingAttempts)*pingTime)) + em.logger.Log("err", errors.Errorf("Have not received a pong in %v", time.Duration(pingAttempts)*pingTime)) em.StopAndReconnect() return } @@ -240,7 +246,7 @@ func (em *EventMeter) receiveRoutine() { } eventID, data, err := em.unmarshalEvent(r) if err != nil { - Log.Error(err.Error()) + em.logger.Log("err", err) continue } if eventID != "" { diff --git a/tm-monitor/glide.yaml b/tm-monitor/glide.yaml index 14bc82fbc..f9bd303d1 100644 --- a/tm-monitor/glide.yaml +++ b/tm-monitor/glide.yaml @@ -15,7 +15,6 @@ import: - package: github.com/tendermint/go-rpc subpackages: - client -- package: github.com/tendermint/log15 - package: github.com/go-kit/kit subpackages: - log diff --git a/tm-monitor/mock/mock.go b/tm-monitor/mock/mock.go index e765dbe19..b614e5953 100644 --- a/tm-monitor/mock/mock.go +++ b/tm-monitor/mock/mock.go @@ -4,6 +4,7 @@ import ( "log" "reflect" + gokitlog "github.com/go-kit/kit/log" ctypes "github.com/tendermint/tendermint/rpc/core/types" em "github.com/tendermint/tools/tm-monitor/eventmeter" ) @@ -16,6 +17,7 @@ type EventMeter struct { func (e *EventMeter) Start() error { return nil } func (e *EventMeter) Stop() {} +func (e *EventMeter) SetLogger(l gokitlog.Logger) {} func (e *EventMeter) RegisterLatencyCallback(cb em.LatencyCallbackFunc) { e.latencyCallback = cb } func (e *EventMeter) RegisterDisconnectCallback(cb em.DisconnectCallbackFunc) { e.disconnectCallback = cb diff --git a/tm-monitor/monitor/node.go b/tm-monitor/monitor/node.go index 8f3071d12..72d404a4f 100644 --- a/tm-monitor/monitor/node.go +++ b/tm-monitor/monitor/node.go @@ -98,6 +98,7 @@ func (n *Node) NotifyAboutDisconnects(ch chan<- bool) { // SetLogger lets you set your own logger func (n *Node) SetLogger(l log.Logger) { n.logger = l + n.em.SetLogger(l) } func (n *Node) Start() error { @@ -265,6 +266,7 @@ type eventMeter interface { RegisterDisconnectCallback(em.DisconnectCallbackFunc) Subscribe(string, em.EventCallbackFunc) error Unsubscribe(string) error + SetLogger(l log.Logger) } // UnmarshalEvent unmarshals a json event From 6e00ce9bbd3e3438b4fb30a6f7782693ba00a39a Mon Sep 17 00:00:00 2001 From: Anton Kaliaev Date: Tue, 21 Mar 2017 20:35:40 +0400 Subject: [PATCH 2/5] [tm-monitor] fix blocking issue as you can see the mistake is that we listen for quit instead of closing it. --- tm-monitor/eventmeter/eventmeter.go | 3 ++- tm-monitor/monitor/node.go | 7 +------ 2 files changed, 3 insertions(+), 7 deletions(-) diff --git a/tm-monitor/eventmeter/eventmeter.go b/tm-monitor/eventmeter/eventmeter.go index dea4c1b1f..d32882fb4 100644 --- a/tm-monitor/eventmeter/eventmeter.go +++ b/tm-monitor/eventmeter/eventmeter.go @@ -140,8 +140,9 @@ func (em *EventMeter) Start() error { return nil } +// Stop stops the EventMeter. func (em *EventMeter) Stop() { - <-em.quit + close(em.quit) em.RegisterDisconnectCallback(nil) // so we don't try and reconnect em.wsc.Stop() // close(wsc.Quit) diff --git a/tm-monitor/monitor/node.go b/tm-monitor/monitor/node.go index 72d404a4f..308428a2b 100644 --- a/tm-monitor/monitor/node.go +++ b/tm-monitor/monitor/node.go @@ -121,12 +121,7 @@ func (n *Node) Start() error { func (n *Node) Stop() { n.Online = false - n.em.RegisterLatencyCallback(nil) - n.em.Unsubscribe(tmtypes.EventStringNewBlockHeader()) - n.em.RegisterDisconnectCallback(nil) - - // FIXME stop blocks at event_meter.go:140 - // n.em.Stop() + n.em.Stop() close(n.quit) } From c053c1523125df7c0c213a6ad1d1aebdcc142c05 Mon Sep 17 00:00:00 2001 From: Anton Kaliaev Date: Tue, 21 Mar 2017 20:37:52 +0400 Subject: [PATCH 3/5] [tm-monitor] only restart EventMeter --- tm-monitor/monitor/node.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tm-monitor/monitor/node.go b/tm-monitor/monitor/node.go index 308428a2b..85ffc1492 100644 --- a/tm-monitor/monitor/node.go +++ b/tm-monitor/monitor/node.go @@ -182,7 +182,7 @@ func (n *Node) RestartBackOff() error { d := time.Duration(math.Exp2(float64(attempt))) time.Sleep(d * time.Second) - if err := n.Start(); err != nil { + if err := n.em.Start(); err != nil { n.logger.Log("err", errors.Wrap(err, "restart failed")) } else { // TODO: authenticate pubkey From 3044f66ba90694927fb22ea5267de2a90bb3281b Mon Sep 17 00:00:00 2001 From: Anton Kaliaev Date: Tue, 21 Mar 2017 20:39:48 +0400 Subject: [PATCH 4/5] [tm-monitor] now EventMeter can be restarted multiple times (Refs #6) with one caveat: go-common and go-rpc need to be updated as well --- tm-monitor/eventmeter/eventmeter.go | 42 +++++++++++++++++++++-------- tm-monitor/monitor/node.go | 4 +-- 2 files changed, 33 insertions(+), 13 deletions(-) diff --git a/tm-monitor/eventmeter/eventmeter.go b/tm-monitor/eventmeter/eventmeter.go index d32882fb4..82cd7a186 100644 --- a/tm-monitor/eventmeter/eventmeter.go +++ b/tm-monitor/eventmeter/eventmeter.go @@ -136,20 +136,28 @@ func (em *EventMeter) Start() error { } return nil }) + + em.quit = make(chan struct{}) go em.receiveRoutine() - return nil + + return em.resubscribe() } // Stop stops the EventMeter. func (em *EventMeter) Stop() { close(em.quit) - em.RegisterDisconnectCallback(nil) // so we don't try and reconnect - em.wsc.Stop() // close(wsc.Quit) + if em.wsc.IsRunning() { + em.wsc.Stop() + } } -func (em *EventMeter) StopAndReconnect() { - em.wsc.Stop() +// StopAndCallDisconnectCallback stops the EventMeter and calls +// disconnectCallback if present. +func (em *EventMeter) StopAndCallDisconnectCallback() { + if em.wsc.IsRunning() { + em.wsc.Stop() + } em.mtx.Lock() defer em.mtx.Unlock() @@ -223,38 +231,50 @@ func (em *EventMeter) RegisterDisconnectCallback(f DisconnectCallbackFunc) { //------------------------------------------------------ +func (em *EventMeter) resubscribe() error { + for eventID, _ := range em.events { + if err := em.wsc.Subscribe(eventID); err != nil { + return err + } + } + return nil +} + func (em *EventMeter) receiveRoutine() { pingTime := time.Second * 1 pingTicker := time.NewTicker(pingTime) pingAttempts := 0 // if this hits maxPingsPerPong we kill the conn + var err error for { select { case <-pingTicker.C: if pingAttempts, err = em.pingForLatency(pingAttempts); err != nil { - em.logger.Log("err", errors.Wrap(err, "Failed to write ping message on websocket")) - em.StopAndReconnect() + em.logger.Log("err", errors.Wrap(err, "failed to write ping message on websocket")) + em.StopAndCallDisconnectCallback() return } else if pingAttempts >= maxPingsPerPong { em.logger.Log("err", errors.Errorf("Have not received a pong in %v", time.Duration(pingAttempts)*pingTime)) - em.StopAndReconnect() + em.StopAndCallDisconnectCallback() return } case r := <-em.wsc.ResultsCh: if r == nil { - em.StopAndReconnect() + em.logger.Log("err", errors.New("Expected some event, received nil")) + em.StopAndCallDisconnectCallback() return } eventID, data, err := em.unmarshalEvent(r) if err != nil { - em.logger.Log("err", err) + em.logger.Log("err", errors.Wrap(err, "failed to unmarshal event")) continue } if eventID != "" { em.updateMetric(eventID, data) } case <-em.wsc.Quit: - em.StopAndReconnect() + em.logger.Log("err", errors.New("WSClient closed unexpectedly")) + em.StopAndCallDisconnectCallback() return case <-em.quit: return diff --git a/tm-monitor/monitor/node.go b/tm-monitor/monitor/node.go index 85ffc1492..b8f873dcc 100644 --- a/tm-monitor/monitor/node.go +++ b/tm-monitor/monitor/node.go @@ -162,7 +162,7 @@ func disconnectCallback(n *Node) em.DisconnectCallbackFunc { n.disconnectCh <- true } - if err := n.RestartBackOff(); err != nil { + if err := n.RestartEventMeterBackoff(); err != nil { n.logger.Log("err", errors.Wrap(err, "restart failed")) } else { n.Online = true @@ -175,7 +175,7 @@ func disconnectCallback(n *Node) em.DisconnectCallbackFunc { } } -func (n *Node) RestartBackOff() error { +func (n *Node) RestartEventMeterBackoff() error { attempt := 0 for { From 9442a069a30cac0e8566518d5e5498528003013f Mon Sep 17 00:00:00 2001 From: Anton Kaliaev Date: Tue, 28 Mar 2017 13:51:14 +0400 Subject: [PATCH 5/5] [tm-monitor] use BaseService.Reset method --- tm-monitor/eventmeter/eventmeter.go | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tm-monitor/eventmeter/eventmeter.go b/tm-monitor/eventmeter/eventmeter.go index 82cd7a186..314e039d3 100644 --- a/tm-monitor/eventmeter/eventmeter.go +++ b/tm-monitor/eventmeter/eventmeter.go @@ -105,7 +105,6 @@ func NewEventMeter(addr string, unmarshalEvent EventUnmarshalFunc) *EventMeter { timer: metrics.NewTimer(), receivedPong: true, unmarshalEvent: unmarshalEvent, - quit: make(chan struct{}), logger: log.NewNopLogger(), } return em @@ -121,6 +120,10 @@ func (em *EventMeter) String() string { } func (em *EventMeter) Start() error { + if _, err := em.wsc.Reset(); err != nil { + return err + } + if _, err := em.wsc.Start(); err != nil { return err }