From 2afbdba1c537a7c0872271a119bb35d8724e7c12 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 13 Dec 2017 08:22:58 -0600 Subject: [PATCH] add lock protection around gossip.memberlist and holder.indexes --- gossip/gossip.go | 11 +++++++++++ holder.go | 4 ++++ 2 files changed, 15 insertions(+) diff --git a/gossip/gossip.go b/gossip/gossip.go index 9d5845315..35f0056c8 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -20,6 +20,7 @@ import ( "log" "os" "strings" + "sync" "time" "golang.org/x/sync/errgroup" @@ -37,6 +38,7 @@ var _ memberlist.Delegate = &GossipMemberSet{} // GossipMemberSet represents a gossip implementation of MemberSet using memberlist. type GossipMemberSet struct { + mu sync.RWMutex memberlist *memberlist.Memberlist handler pilosa.BroadcastHandler @@ -51,6 +53,9 @@ type GossipMemberSet struct { // Nodes implements the MemberSet interface and returns a list of nodes in the cluster. func (g *GossipMemberSet) Nodes() []*pilosa.Node { + g.mu.RLock() + defer g.mu.RUnlock() + a := make([]*pilosa.Node, 0, g.memberlist.NumMembers()) for _, n := range g.memberlist.Members() { uri, _ := pilosa.NewURIFromAddress(n.Name) @@ -78,13 +83,17 @@ func (g *GossipMemberSet) Open() error { } err := error(nil) + g.mu.Lock() g.memberlist, err = memberlist.Create(g.config.memberlistConfig) + g.mu.Unlock() if err != nil { return err } g.broadcasts = &memberlist.TransmitLimitedQueue{ NumNodes: func() int { + g.mu.RLock() + defer g.mu.RUnlock() return g.memberlist.NumMembers() }, RetransmitMult: 3, @@ -98,7 +107,9 @@ func (g *GossipMemberSet) Open() error { // attach to gossip seed node nodes := []*pilosa.Node{&pilosa.Node{URI: *uri}} //TODO: support a list of seeds + g.mu.RLock() err = g.joinWithRetry(pilosa.NodeSet(pilosa.Nodes(nodes).URIs()).ToHostPortStrings()) + g.mu.RUnlock() if err != nil { return err } diff --git a/holder.go b/holder.go index 4c3d476d2..9a1d793d5 100644 --- a/holder.go +++ b/holder.go @@ -156,7 +156,9 @@ func (h *Holder) Open() error { } return fmt.Errorf("open index: name=%s, err=%s", index.Name(), err) } + h.mu.Lock() h.indexes[index.Name()] = index + h.mu.Unlock() } h.logger().Printf("open holder: complete") @@ -190,6 +192,8 @@ func (h *Holder) Close() error { // This is used to determine if the rebalancing of data is necessary // when a node joins the cluster. func (h *Holder) HasData() bool { + h.mu.RLock() + defer h.mu.RUnlock() return h.hasData || len(h.indexes) > 0 }