mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-11 07:11:02 +00:00
Change Broadcaster.SyndSync to send direct via http
as opposed to using memberlist's gossip broadcast. Introduce a Gossiper interface for SendAsync gossip messages.
This commit is contained in:
parent
21072a665f
commit
637111a76c
8 changed files with 125 additions and 103 deletions
21
broadcast.go
21
broadcast.go
|
|
@ -65,6 +65,7 @@ type Broadcaster interface {
|
|||
|
||||
func init() {
|
||||
NopBroadcaster = &nopBroadcaster{}
|
||||
NopGossiper = &nopGossiper{}
|
||||
}
|
||||
|
||||
// NopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
|
|
@ -73,12 +74,12 @@ var NopBroadcaster Broadcaster
|
|||
type nopBroadcaster struct{}
|
||||
|
||||
// SendSync A no-op implemenetation of Broadcaster SendSync method.
|
||||
func (c *nopBroadcaster) SendSync(pb proto.Message) error {
|
||||
func (n *nopBroadcaster) SendSync(pb proto.Message) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// SendAsync A no-op implemenetation of Broadcaster SendAsync method.
|
||||
func (c *nopBroadcaster) SendAsync(pb proto.Message) error {
|
||||
func (n *nopBroadcaster) SendAsync(pb proto.Message) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -86,7 +87,6 @@ func (c *nopBroadcaster) SendAsync(pb proto.Message) error {
|
|||
// handle broadcast messages. (Hint: this is implemented by pilosa.Server)
|
||||
type BroadcastHandler interface {
|
||||
ReceiveMessage(pb proto.Message) error
|
||||
SendSync(pb proto.Message) error
|
||||
}
|
||||
|
||||
// BroadcastReceiver is the interface for the object which will listen for and
|
||||
|
|
@ -107,6 +107,21 @@ func (n *nopBroadcastReceiver) Start(b BroadcastHandler) error { return nil }
|
|||
// NopBroadcastReceiver is a no-op implementation of the BroadcastReceiver.
|
||||
var NopBroadcastReceiver = &nopBroadcastReceiver{}
|
||||
|
||||
// Gossiper is an interface for sharing messages via gossip.
|
||||
type Gossiper interface {
|
||||
SendAsync(pb proto.Message) error
|
||||
}
|
||||
|
||||
// NopBroadcaster represents a Broadcaster that doesn't do anything.
|
||||
var NopGossiper Gossiper
|
||||
|
||||
type nopGossiper struct{}
|
||||
|
||||
// SendAsync A no-op implemenetation of Gossiper SendAsync method.
|
||||
func (n *nopGossiper) SendAsync(pb proto.Message) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Broadcast message types.
|
||||
const (
|
||||
MessageTypeCreateSlice = 1
|
||||
|
|
|
|||
|
|
@ -103,7 +103,3 @@ func (h *SimpleBroadcastHandler) ReceiveMessage(pb proto.Message) error {
|
|||
h.receivedMessage = pb.(proto.Message)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *SimpleBroadcastHandler) SendSync(pb proto.Message) error {
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1044,8 +1044,8 @@ func (c *InternalHTTPClient) RowAttrDiff(ctx context.Context, index, frame strin
|
|||
return rsp.Attrs, nil
|
||||
}
|
||||
|
||||
// ClusterMessage posts a Gossip message synchronously.
|
||||
func (c *InternalHTTPClient) ClusterMessage(ctx context.Context, pb proto.Message) error {
|
||||
// SendMessage posts a message synchronously.
|
||||
func (c *InternalHTTPClient) SendMessage(ctx context.Context, pb proto.Message) error {
|
||||
msg, err := MarshalMessage(pb)
|
||||
if err != nil {
|
||||
return err
|
||||
|
|
@ -1261,5 +1261,5 @@ type InternalClient interface {
|
|||
BlockData(ctx context.Context, index, frame, view string, slice uint64, block int) ([]uint64, []uint64, error)
|
||||
ColumnAttrDiff(ctx context.Context, index string, blks []AttrBlock) (map[uint64]map[string]interface{}, error)
|
||||
RowAttrDiff(ctx context.Context, index, frame string, blks []AttrBlock) (map[uint64]map[string]interface{}, error)
|
||||
ClusterMessage(ctx context.Context, pb proto.Message) error
|
||||
SendMessage(ctx context.Context, pb proto.Message) error
|
||||
}
|
||||
|
|
|
|||
|
|
@ -28,6 +28,11 @@ import (
|
|||
"github.com/pilosa/pilosa/internal"
|
||||
)
|
||||
|
||||
// Ensure GossipNodeSet implements interfaces.
|
||||
var _ pilosa.BroadcastReceiver = &GossipNodeSet{}
|
||||
var _ pilosa.Gossiper = &GossipNodeSet{}
|
||||
var _ memberlist.Delegate = &GossipNodeSet{}
|
||||
|
||||
// GossipNodeSet represents a gossip implementation of NodeSet using memberlist
|
||||
// GossipNodeSet also represents a gossip implementation of pilosa.Broadcaster
|
||||
// GossipNodeSet also represents an implementation of memberlist.Delegate
|
||||
|
|
@ -232,13 +237,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed
|
|||
return g, nil
|
||||
}
|
||||
|
||||
// SendSync implementation of the Broadcaster interface.
|
||||
func (g *GossipNodeSet) SendSync(pb proto.Message) error {
|
||||
// Use the SendSync implementation in Server.
|
||||
return g.handler.SendSync(pb)
|
||||
}
|
||||
|
||||
// SendAsync implementation of the Broadcaster interface.
|
||||
// SendAsync implementation of the Gossiper interface.
|
||||
func (g *GossipNodeSet) SendAsync(pb proto.Message) error {
|
||||
msg, err := pilosa.MarshalMessage(pb)
|
||||
if err != nil {
|
||||
|
|
|
|||
87
handler.go
87
handler.go
|
|
@ -51,9 +51,10 @@ import (
|
|||
|
||||
// Handler represents an HTTP handler.
|
||||
type Handler struct {
|
||||
Holder *Holder
|
||||
Broadcaster Broadcaster
|
||||
StatusHandler StatusHandler
|
||||
Holder *Holder
|
||||
Broadcaster Broadcaster
|
||||
BroadcastHandler BroadcastHandler
|
||||
StatusHandler StatusHandler
|
||||
|
||||
// Local hostname & cluster configuration.
|
||||
URI *URI
|
||||
|
|
@ -136,7 +137,7 @@ func NewRouter(handler *Handler) *mux.Router {
|
|||
router.HandleFunc("/status", handler.handleGetStatus).Methods("GET")
|
||||
router.HandleFunc("/version", handler.handleGetVersion).Methods("GET")
|
||||
router.HandleFunc("/recalculate-caches", handler.handleRecalculateCaches).Methods("POST")
|
||||
router.HandleFunc("/cluster/message", handler.handleClusterMessage).Methods("POST")
|
||||
router.HandleFunc("/cluster/message", handler.handlePostClusterMessage).Methods("POST")
|
||||
|
||||
// TODO: Apply MethodNotAllowed statuses to all endpoints.
|
||||
// Ideally this would be automatic, as described in this (wontfix) ticket:
|
||||
|
|
@ -1991,7 +1992,7 @@ func GetTimeStamp(data map[string]interface{}, timeField string) (int64, error)
|
|||
return v.Unix(), nil
|
||||
}
|
||||
|
||||
func (h *Handler) handleClusterMessage(w http.ResponseWriter, r *http.Request) {
|
||||
func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Request) {
|
||||
// Verify that request is only communicating over protobufs.
|
||||
if r.Header.Get("Content-Type") != "application/x-protobuf" {
|
||||
fmt.Println("**unsupported media type**")
|
||||
|
|
@ -2014,7 +2015,7 @@ func (h *Handler) handleClusterMessage(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
|
||||
// Forward the error message.
|
||||
err = h.ProcessClusterMessage(pb)
|
||||
err = h.BroadcastHandler.ReceiveMessage(pb)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
|
|
@ -2026,77 +2027,3 @@ func (h *Handler) handleClusterMessage(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
|
||||
type defaultClusterMessageResponse struct{}
|
||||
|
||||
// ProcessClusterMessage Process Cluster messages from API handler and Gossip BroadcastHandler.
|
||||
func (h *Handler) ProcessClusterMessage(pb proto.Message) error {
|
||||
switch obj := pb.(type) {
|
||||
case *internal.CreateSliceMessage:
|
||||
idx := h.Holder.Index(obj.Index)
|
||||
if idx == nil {
|
||||
return fmt.Errorf("Local Index not found: %s", obj.Index)
|
||||
}
|
||||
if obj.IsInverse {
|
||||
idx.SetRemoteMaxInverseSlice(obj.Slice)
|
||||
} else {
|
||||
idx.SetRemoteMaxSlice(obj.Slice)
|
||||
}
|
||||
case *internal.CreateIndexMessage:
|
||||
opt := IndexOptions{
|
||||
ColumnLabel: obj.Meta.ColumnLabel,
|
||||
TimeQuantum: TimeQuantum(obj.Meta.TimeQuantum),
|
||||
}
|
||||
_, err := h.Holder.CreateIndex(obj.Index, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.DeleteIndexMessage:
|
||||
if err := h.Holder.DeleteIndex(obj.Index); err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.CreateFrameMessage:
|
||||
idx := h.Holder.Index(obj.Index)
|
||||
if idx == nil {
|
||||
return fmt.Errorf("Local Index not found: %s", obj.Index)
|
||||
}
|
||||
opt := FrameOptions{
|
||||
RowLabel: obj.Meta.RowLabel,
|
||||
InverseEnabled: obj.Meta.InverseEnabled,
|
||||
RangeEnabled: obj.Meta.RangeEnabled,
|
||||
CacheType: obj.Meta.CacheType,
|
||||
CacheSize: obj.Meta.CacheSize,
|
||||
TimeQuantum: TimeQuantum(obj.Meta.TimeQuantum),
|
||||
Fields: decodeFields(obj.Meta.Fields),
|
||||
}
|
||||
_, err := idx.CreateFrame(obj.Frame, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.DeleteFrameMessage:
|
||||
idx := h.Holder.Index(obj.Index)
|
||||
if err := idx.DeleteFrame(obj.Frame); err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.CreateInputDefinitionMessage:
|
||||
idx := h.Holder.Index(obj.Index)
|
||||
if idx == nil {
|
||||
return fmt.Errorf("Local Index not found: %s", obj.Index)
|
||||
}
|
||||
idx.CreateInputDefinition(obj.Definition)
|
||||
case *internal.DeleteInputDefinitionMessage:
|
||||
idx := h.Holder.Index(obj.Index)
|
||||
err := idx.DeleteInputDefinition(obj.Name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.DeleteViewMessage:
|
||||
f := h.Holder.Frame(obj.Index, obj.Frame)
|
||||
if f == nil {
|
||||
return fmt.Errorf("Local Frame not found: %s", obj.Frame)
|
||||
}
|
||||
err := f.DeleteView(obj.View)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
87
server.go
87
server.go
|
|
@ -46,6 +46,11 @@ const (
|
|||
DefaultDiagnosticServer = "https://diagnostics.pilosa.com/v0/diagnostics"
|
||||
)
|
||||
|
||||
// Ensure Server implements interfaces.
|
||||
var _ Broadcaster = &Server{}
|
||||
var _ BroadcastHandler = &Server{}
|
||||
var _ StatusHandler = &Server{}
|
||||
|
||||
// Server represents a holder wrapped by a running HTTP server.
|
||||
type Server struct {
|
||||
ln net.Listener
|
||||
|
|
@ -59,6 +64,7 @@ type Server struct {
|
|||
Handler *Handler
|
||||
Broadcaster Broadcaster
|
||||
BroadcastReceiver BroadcastReceiver
|
||||
Gossiper Gossiper
|
||||
RemoteClient *http.Client
|
||||
|
||||
// Cluster configuration.
|
||||
|
|
@ -182,6 +188,7 @@ func (s *Server) Open() error {
|
|||
|
||||
// Initialize HTTP handler.
|
||||
s.Handler.Broadcaster = s.Broadcaster
|
||||
s.Handler.BroadcastHandler = s
|
||||
s.Handler.StatusHandler = s
|
||||
s.Handler.URI = s.URI
|
||||
s.Handler.Cluster = s.Cluster
|
||||
|
|
@ -333,10 +340,79 @@ func (s *Server) monitorMaxSlices() {
|
|||
|
||||
// ReceiveMessage represents an implementation of BroadcastHandler.
|
||||
func (s *Server) ReceiveMessage(pb proto.Message) error {
|
||||
return s.Handler.ProcessClusterMessage(pb)
|
||||
switch obj := pb.(type) {
|
||||
case *internal.CreateSliceMessage:
|
||||
idx := s.Holder.Index(obj.Index)
|
||||
if idx == nil {
|
||||
return fmt.Errorf("Local Index not found: %s", obj.Index)
|
||||
}
|
||||
if obj.IsInverse {
|
||||
idx.SetRemoteMaxInverseSlice(obj.Slice)
|
||||
} else {
|
||||
idx.SetRemoteMaxSlice(obj.Slice)
|
||||
}
|
||||
case *internal.CreateIndexMessage:
|
||||
opt := IndexOptions{
|
||||
ColumnLabel: obj.Meta.ColumnLabel,
|
||||
TimeQuantum: TimeQuantum(obj.Meta.TimeQuantum),
|
||||
}
|
||||
_, err := s.Holder.CreateIndex(obj.Index, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.DeleteIndexMessage:
|
||||
if err := s.Holder.DeleteIndex(obj.Index); err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.CreateFrameMessage:
|
||||
idx := s.Holder.Index(obj.Index)
|
||||
if idx == nil {
|
||||
return fmt.Errorf("Local Index not found: %s", obj.Index)
|
||||
}
|
||||
opt := FrameOptions{
|
||||
RowLabel: obj.Meta.RowLabel,
|
||||
InverseEnabled: obj.Meta.InverseEnabled,
|
||||
RangeEnabled: obj.Meta.RangeEnabled,
|
||||
CacheType: obj.Meta.CacheType,
|
||||
CacheSize: obj.Meta.CacheSize,
|
||||
TimeQuantum: TimeQuantum(obj.Meta.TimeQuantum),
|
||||
Fields: decodeFields(obj.Meta.Fields),
|
||||
}
|
||||
_, err := idx.CreateFrame(obj.Frame, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.DeleteFrameMessage:
|
||||
idx := s.Holder.Index(obj.Index)
|
||||
if err := idx.DeleteFrame(obj.Frame); err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.CreateInputDefinitionMessage:
|
||||
idx := s.Holder.Index(obj.Index)
|
||||
if idx == nil {
|
||||
return fmt.Errorf("Local Index not found: %s", obj.Index)
|
||||
}
|
||||
idx.CreateInputDefinition(obj.Definition)
|
||||
case *internal.DeleteInputDefinitionMessage:
|
||||
idx := s.Holder.Index(obj.Index)
|
||||
err := idx.DeleteInputDefinition(obj.Name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.DeleteViewMessage:
|
||||
f := s.Holder.Frame(obj.Index, obj.Frame)
|
||||
if f == nil {
|
||||
return fmt.Errorf("Local Frame not found: %s", obj.Frame)
|
||||
}
|
||||
err := f.DeleteView(obj.View)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// SendSync represents an implementation of BroadcastHandler.
|
||||
// SendSync represents an implementation of Broadcaster.
|
||||
func (s *Server) SendSync(pb proto.Message) error {
|
||||
var eg errgroup.Group
|
||||
for _, node := range s.Cluster.Nodes {
|
||||
|
|
@ -352,13 +428,18 @@ func (s *Server) SendSync(pb proto.Message) error {
|
|||
|
||||
ctx := context.WithValue(context.Background(), "uri", uri)
|
||||
eg.Go(func() error {
|
||||
return s.defaultClient.ClusterMessage(ctx, pb)
|
||||
return s.defaultClient.SendMessage(ctx, pb)
|
||||
})
|
||||
}
|
||||
|
||||
return eg.Wait()
|
||||
}
|
||||
|
||||
// SendAsync represents an implementation of Broadcaster.
|
||||
func (s *Server) SendAsync(pb proto.Message) error {
|
||||
return s.Gossiper.SendAsync(pb)
|
||||
}
|
||||
|
||||
// LocalStatus returns the state of the local node as well as the
|
||||
// holder (indexes/frames) according to the local node.
|
||||
// In a gossip implementation, memberlist.Delegate.LocalState() uses this.
|
||||
|
|
|
|||
|
|
@ -226,12 +226,14 @@ func (m *Command) SetupServer() error {
|
|||
return err
|
||||
}
|
||||
m.Server.Cluster.NodeSet = gossipNodeSet
|
||||
m.Server.Broadcaster = gossipNodeSet
|
||||
m.Server.Broadcaster = m.Server
|
||||
m.Server.BroadcastReceiver = gossipNodeSet
|
||||
m.Server.Gossiper = gossipNodeSet
|
||||
case pilosa.ClusterStatic, pilosa.ClusterNone:
|
||||
m.Server.Broadcaster = pilosa.NopBroadcaster
|
||||
m.Server.Cluster.NodeSet = pilosa.NewStaticNodeSet()
|
||||
m.Server.BroadcastReceiver = pilosa.NopBroadcastReceiver
|
||||
m.Server.Gossiper = pilosa.NopGossiper
|
||||
err := m.Server.Cluster.NodeSet.(*pilosa.StaticNodeSet).Join(m.Server.Cluster.Nodes)
|
||||
if err != nil {
|
||||
return err
|
||||
|
|
|
|||
|
|
@ -412,7 +412,8 @@ func TestMain_SendReceiveMessage(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
m0.Server.Cluster.NodeSet = gossipNodeSet0
|
||||
m0.Server.Broadcaster = gossipNodeSet0
|
||||
m0.Server.Broadcaster = m0.Server
|
||||
m0.Server.Gossiper = gossipNodeSet0
|
||||
m0.Server.Handler.Broadcaster = m0.Server.Broadcaster
|
||||
m0.Server.Holder.Broadcaster = m0.Server.Broadcaster
|
||||
m0.Server.BroadcastReceiver = gossipNodeSet0
|
||||
|
|
@ -437,7 +438,8 @@ func TestMain_SendReceiveMessage(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
m1.Server.Cluster.NodeSet = gossipNodeSet1
|
||||
m1.Server.Broadcaster = gossipNodeSet1
|
||||
m1.Server.Broadcaster = m1.Server
|
||||
m1.Server.Gossiper = gossipNodeSet1
|
||||
m1.Server.Handler.Broadcaster = m1.Server.Broadcaster
|
||||
m1.Server.Holder.Broadcaster = m1.Server.Broadcaster
|
||||
m1.Server.BroadcastReceiver = gossipNodeSet1
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue