diff --git a/cluster.go b/cluster.go index c7f5c4193..ed5318917 100644 --- a/cluster.go +++ b/cluster.go @@ -405,6 +405,8 @@ func (c *Cluster) setState(state string) { if c.state == ClusterStateResizing { doCleanup = true } + default: + panic(fmt.Sprintf("invalid cluster state: %s", state)) } c.state = state diff --git a/gossip/gossip.go b/gossip/gossip.go index 646bd63cb..0949feb16 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -346,12 +346,16 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) { type GossipEventReceiver struct { ch chan memberlist.NodeEvent eventHandler pilosa.EventHandler + + // The writer for any logging. + LogOutput io.Writer } // NewGossipEventReceiver returns a new instance of GossipEventReceiver. -func NewGossipEventReceiver() *GossipEventReceiver { +func NewGossipEventReceiver(logOutput io.Writer) *GossipEventReceiver { return &GossipEventReceiver{ - ch: make(chan memberlist.NodeEvent, 1), + ch: make(chan memberlist.NodeEvent, 1), + LogOutput: logOutput, } } @@ -374,6 +378,11 @@ func (g *GossipEventReceiver) Start(h pilosa.EventHandler) error { return nil } +// logger returns a logger for the GossipEventReceiver. +func (g *GossipEventReceiver) logger() *log.Logger { + return log.New(g.LogOutput, "", log.LstdFlags) +} + func (g *GossipEventReceiver) listen() { var nodeEventType pilosa.NodeEventType for { @@ -400,7 +409,9 @@ func (g *GossipEventReceiver) listen() { Event: nodeEventType, Node: node, } - _ = g.eventHandler.ReceiveEvent(ne) + if err := g.eventHandler.ReceiveEvent(ne); err != nil { + g.logger().Printf("receive event error: %s", err) + } } } diff --git a/server/cluster_test.go b/server/cluster_test.go index 65d8844d3..5ec805f81 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -54,7 +54,7 @@ func TestMain_SendReceiveMessage(t *testing.T) { m0.Server.Cluster.Coordinator = m0.Server.URI m0.Server.Cluster.Topology = &pilosa.Topology{NodeIDs: []string{m0.Server.NodeID, m1.Server.NodeID}} - m0.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver() + m0.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m0.Server.LogOutput) gossipMemberSet0, err := gossip.NewGossipMemberSet(m0.Server.URI.HostPort(), m0.Config, m0.Server) if err != nil { t.Fatal(err) @@ -81,7 +81,7 @@ func TestMain_SendReceiveMessage(t *testing.T) { m1.Config.Gossip.Seeds = gossipMemberSet0.Seeds() m1.Server.Cluster.Coordinator = m0.Server.URI - m1.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver() + m1.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m1.Server.LogOutput) gossipMemberSet1, err := gossip.NewGossipMemberSet(m1.Server.URI.HostPort(), m1.Config, m1.Server) if err != nil { t.Fatal(err) diff --git a/server/server.go b/server/server.go index 7b623461a..9803f7be2 100644 --- a/server/server.go +++ b/server/server.go @@ -260,7 +260,7 @@ func (m *Command) SetupNetworking() error { m.Server.NodeID = m.Server.LoadNodeID() - m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver() + m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m.Server.LogOutput) gossipMemberSet, err := gossip.NewGossipMemberSetWithTransport(m.Server.NodeID, m.Config, transport, m.Server) if err != nil { return err