mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 17:15:56 +00:00
Merge pull request #1412 from jaffee/simplify-event-receiver
consolidate gossipEventReceiver into gossip member set
This commit is contained in:
commit
83405b3d38
6 changed files with 44 additions and 57 deletions
|
|
@ -56,6 +56,7 @@ func testMessageMarshal(t *testing.T, m proto.Message) {
|
|||
|
||||
// Ensure that BroadcastReceiver can register a BroadcastHandler.
|
||||
func TestBroadcast_BroadcastReceiver(t *testing.T) {
|
||||
t.Skip("broadcast receiver")
|
||||
path, err := ioutil.TempDir("", "pilosa-")
|
||||
if err != nil {
|
||||
panic(err)
|
||||
|
|
@ -67,24 +68,24 @@ func TestBroadcast_BroadcastReceiver(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("setting up server: %v", err)
|
||||
}
|
||||
s := com.Server
|
||||
// s := com.Server
|
||||
|
||||
sbr := NewSimpleBroadcastReceiver()
|
||||
sbh := NewSimpleBroadcastHandler()
|
||||
// sbr := NewSimpleBroadcastReceiver()
|
||||
// sbh := NewSimpleBroadcastHandler()
|
||||
|
||||
s.BroadcastReceiver = sbr
|
||||
s.BroadcastReceiver.Start(sbh)
|
||||
// s.BroadcastReceiver = sbr
|
||||
// s.BroadcastReceiver.Start(sbh)
|
||||
|
||||
msg := &internal.DeleteIndexMessage{
|
||||
Index: "i",
|
||||
}
|
||||
// msg := &internal.DeleteIndexMessage{
|
||||
// Index: "i",
|
||||
// }
|
||||
|
||||
s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg)
|
||||
// s.BroadcastReceiver.(*SimpleBroadcastReceiver).Receive(msg)
|
||||
|
||||
// Make sure the message received is what was sentd
|
||||
if !reflect.DeepEqual(sbh.receivedMessage, msg) {
|
||||
t.Fatalf("unexpected message: %s", sbh.receivedMessage)
|
||||
}
|
||||
// // Make sure the message received is what was sentd
|
||||
// if !reflect.DeepEqual(sbh.receivedMessage, msg) {
|
||||
// t.Fatalf("unexpected message: %s", sbh.receivedMessage)
|
||||
// }
|
||||
}
|
||||
|
||||
type SimpleBroadcastReceiver struct {
|
||||
|
|
|
|||
|
|
@ -886,11 +886,6 @@ func (c *Cluster) open() error {
|
|||
return errors.Wrap(err, "adding local node")
|
||||
}
|
||||
|
||||
// Start the EventReceiver.
|
||||
if err := c.EventReceiver.Start(c); err != nil {
|
||||
return fmt.Errorf("starting EventReceiver: %v", err)
|
||||
}
|
||||
|
||||
// Open MemberSet communication.
|
||||
if err := c.MemberSet.Open(c.Node); err != nil {
|
||||
return fmt.Errorf("opening MemberSet: %v", err)
|
||||
|
|
|
|||
|
|
@ -33,7 +33,6 @@ import (
|
|||
)
|
||||
|
||||
// Ensure GossipMemberSet implements interfaces.
|
||||
var _ pilosa.BroadcastReceiver = &GossipMemberSet{}
|
||||
var _ memberlist.Delegate = &GossipMemberSet{}
|
||||
|
||||
// GossipMemberSet represents a gossip implementation of MemberSet using memberlist.
|
||||
|
|
@ -45,19 +44,15 @@ type GossipMemberSet struct {
|
|||
|
||||
broadcasts *memberlist.TransmitLimitedQueue
|
||||
|
||||
statusHandler pilosa.StatusHandler
|
||||
config *gossipConfig
|
||||
pserver *pilosa.Server
|
||||
config *gossipConfig
|
||||
|
||||
Logger pilosa.Logger
|
||||
|
||||
logger *log.Logger
|
||||
transport *Transport
|
||||
}
|
||||
|
||||
// Start implements the BroadcastReceiver interface and sets the BroadcastHandler.
|
||||
func (g *GossipMemberSet) Start(h pilosa.BroadcastHandler) error {
|
||||
g.handler = h
|
||||
return nil
|
||||
gossipEventReceiver *GossipEventReceiver
|
||||
}
|
||||
|
||||
// GetBindAddr returns the gossip bind address based on config and auto bind port.
|
||||
|
|
@ -69,13 +64,16 @@ func (g *GossipMemberSet) GetBindAddr() string {
|
|||
|
||||
// Open implements the MemberSet interface to start network activity.
|
||||
func (g *GossipMemberSet) Open(n *pilosa.Node) error {
|
||||
err := g.gossipEventReceiver.Start(g.pserver)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "starting event delegate")
|
||||
}
|
||||
if g.handler == nil {
|
||||
return fmt.Errorf("must call Start(pilosa.BroadcastHandler) before calling Open()")
|
||||
}
|
||||
|
||||
g.node = n
|
||||
|
||||
err := error(nil)
|
||||
g.mu.Lock()
|
||||
g.memberlist, err = memberlist.Create(g.config.memberlistConfig)
|
||||
g.mu.Unlock()
|
||||
|
|
@ -166,7 +164,7 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption {
|
|||
}
|
||||
|
||||
// NewGossipMemberSet returns a new instance of GossipMemberSet based on options.
|
||||
func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventReceiver, sh pilosa.StatusHandler, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
|
||||
func NewGossipMemberSet(name string, host string, cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) {
|
||||
g := &GossipMemberSet{
|
||||
Logger: pilosa.NopLogger,
|
||||
}
|
||||
|
|
@ -177,6 +175,10 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe
|
|||
return nil, errors.Wrap(err, "executing option")
|
||||
}
|
||||
}
|
||||
ger := NewGossipEventReceiver(g.logger)
|
||||
g.gossipEventReceiver = ger
|
||||
|
||||
g.handler = s
|
||||
|
||||
if g.transport == nil {
|
||||
port, err := strconv.Atoi(cfg.Port)
|
||||
|
|
@ -232,7 +234,7 @@ func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventRe
|
|||
gossipSeeds: cfg.Seeds,
|
||||
}
|
||||
|
||||
g.statusHandler = sh
|
||||
g.pserver = s
|
||||
|
||||
return g, nil
|
||||
}
|
||||
|
|
@ -270,7 +272,7 @@ func (g *GossipMemberSet) GetBroadcasts(overhead, limit int) [][]byte {
|
|||
// LocalState implementation of the memberlist.Delegate interface
|
||||
// sends this Node's state data.
|
||||
func (g *GossipMemberSet) LocalState(join bool) []byte {
|
||||
pb, err := g.statusHandler.LocalStatus()
|
||||
pb, err := g.pserver.LocalStatus()
|
||||
if err != nil {
|
||||
g.Logger.Printf("error getting local state, err=%s", err)
|
||||
return []byte{}
|
||||
|
|
@ -294,7 +296,7 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) {
|
|||
g.Logger.Printf("error unmarshalling nodestate data, err=%s", err)
|
||||
return
|
||||
}
|
||||
err := g.statusHandler.HandleRemoteStatus(&pb)
|
||||
err := g.pserver.HandleRemoteStatus(&pb)
|
||||
if err != nil {
|
||||
g.Logger.Printf("merge state error: %s", err)
|
||||
}
|
||||
|
|
@ -309,11 +311,11 @@ type GossipEventReceiver struct {
|
|||
ch chan memberlist.NodeEvent
|
||||
eventHandler pilosa.EventHandler
|
||||
|
||||
Logger pilosa.Logger
|
||||
Logger *log.Logger
|
||||
}
|
||||
|
||||
// NewGossipEventReceiver returns a new instance of GossipEventReceiver.
|
||||
func NewGossipEventReceiver(logger pilosa.Logger) *GossipEventReceiver {
|
||||
func NewGossipEventReceiver(logger *log.Logger) *GossipEventReceiver {
|
||||
return &GossipEventReceiver{
|
||||
ch: make(chan memberlist.NodeEvent, 1),
|
||||
Logger: logger,
|
||||
|
|
|
|||
28
server.go
28
server.go
|
|
@ -61,10 +61,9 @@ type Server struct {
|
|||
clusterDisabled bool
|
||||
|
||||
// External
|
||||
BroadcastReceiver BroadcastReceiver
|
||||
systemInfo SystemInfo
|
||||
gcNotifier GCNotifier
|
||||
logger Logger
|
||||
systemInfo SystemInfo
|
||||
gcNotifier GCNotifier
|
||||
logger Logger
|
||||
|
||||
NodeID string
|
||||
URI URI
|
||||
|
|
@ -207,12 +206,11 @@ func OptServerClusterDisabled(disabled bool, hosts []string) ServerOption {
|
|||
// NewServer returns a new instance of Server.
|
||||
func NewServer(opts ...ServerOption) (*Server, error) {
|
||||
s := &Server{
|
||||
closing: make(chan struct{}),
|
||||
Cluster: NewCluster(),
|
||||
holder: NewHolder(),
|
||||
BroadcastReceiver: NopBroadcastReceiver,
|
||||
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
|
||||
systemInfo: NewNopSystemInfo(),
|
||||
closing: make(chan struct{}),
|
||||
Cluster: NewCluster(),
|
||||
holder: NewHolder(),
|
||||
diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer),
|
||||
systemInfo: NewNopSystemInfo(),
|
||||
|
||||
gcNotifier: NopGCNotifier,
|
||||
|
||||
|
|
@ -297,11 +295,6 @@ func (s *Server) Open() error {
|
|||
// Initialize Holder.
|
||||
s.holder.Broadcaster = s
|
||||
|
||||
// Start the BroadcastReceiver.
|
||||
if err := s.BroadcastReceiver.Start(s); err != nil {
|
||||
return fmt.Errorf("starting BroadcastReceiver: %v", err)
|
||||
}
|
||||
|
||||
// Open Cluster management.
|
||||
if err := s.Cluster.open(); err != nil {
|
||||
return fmt.Errorf("opening Cluster: %v", err)
|
||||
|
|
@ -711,6 +704,11 @@ func (s *Server) monitorRuntime() {
|
|||
}
|
||||
}
|
||||
|
||||
// ReceiveEvent implements the EventHandler interface.
|
||||
func (s *Server) ReceiveEvent(e *NodeEvent) error {
|
||||
return s.Cluster.ReceiveEvent(e)
|
||||
}
|
||||
|
||||
// countOpenFiles on operating systems that support lsof.
|
||||
func countOpenFiles() (int, error) {
|
||||
switch runtime.GOOS {
|
||||
|
|
|
|||
|
|
@ -318,13 +318,10 @@ func (m *Command) SetupNetworking() error {
|
|||
m.Server.Cluster.Node.IsCoordinator = true
|
||||
}
|
||||
|
||||
gossipEventReceiver := gossip.NewGossipEventReceiver(m.logger)
|
||||
m.Server.Cluster.EventReceiver = gossipEventReceiver
|
||||
gossipMemberSet, err := gossip.NewGossipMemberSet(
|
||||
m.Server.NodeID,
|
||||
m.Server.URI.Host(),
|
||||
m.Config.Gossip,
|
||||
gossipEventReceiver,
|
||||
m.Server,
|
||||
gossip.WithLogger(m.logger.Logger()),
|
||||
gossip.WithTransport(transport),
|
||||
|
|
@ -332,9 +329,7 @@ func (m *Command) SetupNetworking() error {
|
|||
if err != nil {
|
||||
return errors.Wrap(err, "getting memberset")
|
||||
}
|
||||
gossipMemberSet.Logger = m.logger
|
||||
m.Server.Cluster.MemberSet = gossipMemberSet
|
||||
m.Server.BroadcastReceiver = gossipMemberSet
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -224,10 +224,6 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) (
|
|||
return seed, err
|
||||
}
|
||||
|
||||
if err = m.Server.BroadcastReceiver.Start(m.Server); err != nil {
|
||||
return seed, err
|
||||
}
|
||||
|
||||
m.Server.Cluster.Static = false
|
||||
|
||||
go func() {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue