mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-12 23:51:03 +00:00
add lock protection around gossip.memberlist and holder.indexes
This commit is contained in:
parent
0342da92bb
commit
2afbdba1c5
2 changed files with 15 additions and 0 deletions
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue