From 29bd0319dc3f08bda3840bea01597eb22943ac8d Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 28 Nov 2017 15:57:38 -0600 Subject: [PATCH 1/9] Create a SendSync Interface for the Broadcast handler --- broadcast.go | 1 + broadcast_test.go | 4 ++++ 2 files changed, 5 insertions(+) diff --git a/broadcast.go b/broadcast.go index 0785a21cc..a5dc98cc7 100644 --- a/broadcast.go +++ b/broadcast.go @@ -86,6 +86,7 @@ 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 diff --git a/broadcast_test.go b/broadcast_test.go index 8f0245c88..1331520ba 100644 --- a/broadcast_test.go +++ b/broadcast_test.go @@ -103,3 +103,7 @@ 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 +} From 9cf994bda7a38b495d08c4c332dedd52e2f8d9a6 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 28 Nov 2017 15:58:20 -0600 Subject: [PATCH 2/9] Use the Broadcast Handler's SendSync implementation rather than Gossip --- gossip/gossip.go | 28 ++-------------------------- 1 file changed, 2 insertions(+), 26 deletions(-) diff --git a/gossip/gossip.go b/gossip/gossip.go index e413410d8..39a1c019d 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -22,8 +22,6 @@ import ( "strings" "time" - "golang.org/x/sync/errgroup" - "github.com/gogo/protobuf/proto" "github.com/hashicorp/memberlist" "github.com/pilosa/pilosa" @@ -236,30 +234,8 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed // SendSync implementation of the Broadcaster interface. func (g *GossipNodeSet) SendSync(pb proto.Message) error { - msg, err := pilosa.MarshalMessage(pb) - if err != nil { - return err - } - - mlist := g.memberlist - - // Direct sends the message directly to every node. - // An error from any node raises an error on the entire operation. - // - // Gossip uses the gossip protocol to eventually deliver the message - // to every node. - var eg errgroup.Group - for _, n := range mlist.Members() { - // Don't send the message to the local node. - if n == mlist.LocalNode() { - continue - } - node := n - eg.Go(func() error { - return mlist.SendToTCP(node, msg) - }) - } - return eg.Wait() + // Use the SendSync implementation in Server. + return g.handler.SendSync(pb) } // SendAsync implementation of the Broadcaster interface. From 013d32099e84c6ddb27047cc29a3227ec8b4fd7a Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Tue, 28 Nov 2017 16:00:02 -0600 Subject: [PATCH 3/9] New /cluster/message endpoint handles all SendSync Messages --- handler.go | 119 +++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 119 insertions(+) diff --git a/handler.go b/handler.go index ad9e02e2e..563dd887a 100644 --- a/handler.go +++ b/handler.go @@ -136,6 +136,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") // TODO: Apply MethodNotAllowed statuses to all endpoints. // Ideally this would be automatic, as described in this (wontfix) ticket: @@ -1987,3 +1988,121 @@ func GetTimeStamp(data map[string]interface{}, timeField string) (int64, error) return v.Unix(), nil } + +func (h *Handler) handleClusterMessage(w http.ResponseWriter, r *http.Request) { + fmt.Println("**handleClusterMessage**") + fmt.Printf("%v\n", r.Header) + + // Verify that request is only communicating over protobufs. + if r.Header.Get("Content-Type") != "application/x-protobuf" { + fmt.Println("**unsupported media type**") + http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType) + return + } + // else if r.Header.Get("Accept") != "application/x-protobuf" { + // fmt.Println("**Not acceptable**") + // http.Error(w, "Not acceptable", http.StatusNotAcceptable) + // return + // } + + // Read entire body. + body, err := ioutil.ReadAll(r.Body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + // Marshal into request object. + pb, err := UnmarshalMessage(body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + err = h.ProcessClusterMessage(pb) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + if err := json.NewEncoder(w).Encode(defaultClusterMessageResponse{}); err != nil { + h.logger().Printf("response encoding error: %s", err) + } +} + +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: + fmt.Printf("*** Create Index %v ***\n", obj.Index) + 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 +} From e296a0d64897c2a3bc1a0ada8839d44360e2aa08 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Wed, 29 Nov 2017 15:10:11 -0600 Subject: [PATCH 4/9] Client method for the SendSync /cluster/message --- client.go | 36 ++++++++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/client.go b/client.go index e44d9c5de..a5fad55de 100644 --- a/client.go +++ b/client.go @@ -1044,6 +1044,41 @@ 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 { + msg, err := MarshalMessage(pb) + if err != nil { + return err + } + + u := uriPathToURL(ctx.Value("uri").(*URI), "/cluster/message") + req, err := http.NewRequest("POST", u.String(), bytes.NewReader(msg)) + req.Header.Set("Content-Type", "application/x-protobuf") + req.Header.Set("User-Agent", "pilosa/"+Version) + + // Execute request. + resp, err := c.HTTPClient.Do(req.WithContext(ctx)) + if err != nil { + return err + } + defer resp.Body.Close() + + // Read body. + body, err := ioutil.ReadAll(resp.Body) + if err != nil { + return err + } + + // Return error if status is not OK. + switch resp.StatusCode { + case http.StatusOK: // ok + default: + return errors.New(string(body)) + } + + return nil +} + func (c *InternalHTTPClient) clientURI(ctx context.Context) *URI { clientURI := c.defaultURI if contextURI, ok := ctx.Value("uri").(*URI); ok { @@ -1226,4 +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 } From 814bc40a16e48e8671d64a566acff2f2370f3534 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Wed, 29 Nov 2017 15:10:53 -0600 Subject: [PATCH 5/9] return the error from /cluster/message --- handler.go | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/handler.go b/handler.go index 563dd887a..0e4ebacde 100644 --- a/handler.go +++ b/handler.go @@ -484,6 +484,8 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) { }) if err != nil { h.logger().Printf("problem sending CreateIndex message: %s", err) + http.Error(w, err.Error(), http.StatusInternalServerError) + return } // Encode response. @@ -1990,20 +1992,12 @@ func GetTimeStamp(data map[string]interface{}, timeField string) (int64, error) } func (h *Handler) handleClusterMessage(w http.ResponseWriter, r *http.Request) { - fmt.Println("**handleClusterMessage**") - fmt.Printf("%v\n", r.Header) - // Verify that request is only communicating over protobufs. if r.Header.Get("Content-Type") != "application/x-protobuf" { fmt.Println("**unsupported media type**") http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType) return } - // else if r.Header.Get("Accept") != "application/x-protobuf" { - // fmt.Println("**Not acceptable**") - // http.Error(w, "Not acceptable", http.StatusNotAcceptable) - // return - // } // Read entire body. body, err := ioutil.ReadAll(r.Body) @@ -2019,6 +2013,7 @@ func (h *Handler) handleClusterMessage(w http.ResponseWriter, r *http.Request) { return } + // Forward the error message. err = h.ProcessClusterMessage(pb) if err != nil { http.Error(w, err.Error(), http.StatusBadRequest) @@ -2046,7 +2041,6 @@ func (h *Handler) ProcessClusterMessage(pb proto.Message) error { idx.SetRemoteMaxSlice(obj.Slice) } case *internal.CreateIndexMessage: - fmt.Printf("*** Create Index %v ***\n", obj.Index) opt := IndexOptions{ ColumnLabel: obj.Meta.ColumnLabel, TimeQuantum: TimeQuantum(obj.Meta.TimeQuantum), From 04790f70f2f2a71844d01a4d937017e56ea5ba1e Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Wed, 29 Nov 2017 15:12:05 -0600 Subject: [PATCH 6/9] moved the broadcast handler receive message process to the handler ProcessClusterMessage --- server.go | 83 ++++++++++++------------------------------------------- 1 file changed, 18 insertions(+), 65 deletions(-) diff --git a/server.go b/server.go index 208d0cd6f..67a9136c6 100644 --- a/server.go +++ b/server.go @@ -36,6 +36,7 @@ import ( "github.com/pilosa/pilosa/diagnostics" "github.com/pilosa/pilosa/internal" "golang.org/x/net/context" + "golang.org/x/sync/errgroup" ) // Default server settings. @@ -332,76 +333,28 @@ func (s *Server) monitorMaxSlices() { // ReceiveMessage represents an implementation of BroadcastHandler. func (s *Server) ReceiveMessage(pb proto.Message) error { - 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) + return s.Handler.ProcessClusterMessage(pb) +} + +// SendSync represents an implementation of BroadcastHandler. +func (s *Server) SendSync(pb proto.Message) error { + var eg errgroup.Group + for _, node := range s.Cluster.Nodes { + uri, err := node.URI() 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 + + // Don't forward the message to ourselves + if *s.URI != *uri { + ctx := context.WithValue(context.Background(), "uri", uri) + eg.Go(func() error { + return s.defaultClient.ClusterMessage(ctx, pb) + }) } } - return nil + + return eg.Wait() } // LocalStatus returns the state of the local node as well as the From 21072a665f6d1b17c89f7df02541d2328dde9a65 Mon Sep 17 00:00:00 2001 From: Michael Baird Date: Thu, 30 Nov 2017 22:03:24 -0600 Subject: [PATCH 7/9] refactor same node URI check in SendSync --- server.go | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/server.go b/server.go index 67a9136c6..b3e42a4c8 100644 --- a/server.go +++ b/server.go @@ -345,13 +345,15 @@ func (s *Server) SendSync(pb proto.Message) error { return err } - // Don't forward the message to ourselves - if *s.URI != *uri { - ctx := context.WithValue(context.Background(), "uri", uri) - eg.Go(func() error { - return s.defaultClient.ClusterMessage(ctx, pb) - }) + // Don't forward the message to ourselves. + if *s.URI == *uri { + continue } + + ctx := context.WithValue(context.Background(), "uri", uri) + eg.Go(func() error { + return s.defaultClient.ClusterMessage(ctx, pb) + }) } return eg.Wait() From 637111a76ca504b7f34821fcec713441cd6861ee Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 8 Dec 2017 08:20:50 -0600 Subject: [PATCH 8/9] Change Broadcaster.SyndSync to send direct via http as opposed to using memberlist's gossip broadcast. Introduce a Gossiper interface for SendAsync gossip messages. --- broadcast.go | 21 +++++++++-- broadcast_test.go | 4 -- client.go | 6 +-- gossip/gossip.go | 13 +++---- handler.go | 87 ++++--------------------------------------- server.go | 87 +++++++++++++++++++++++++++++++++++++++++-- server/server.go | 4 +- server/server_test.go | 6 ++- 8 files changed, 125 insertions(+), 103 deletions(-) diff --git a/broadcast.go b/broadcast.go index a5dc98cc7..38eb53b96 100644 --- a/broadcast.go +++ b/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 diff --git a/broadcast_test.go b/broadcast_test.go index 1331520ba..8f0245c88 100644 --- a/broadcast_test.go +++ b/broadcast_test.go @@ -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 -} diff --git a/client.go b/client.go index a5fad55de..f160d88bc 100644 --- a/client.go +++ b/client.go @@ -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 } diff --git a/gossip/gossip.go b/gossip/gossip.go index 39a1c019d..d296f7b83 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -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 { diff --git a/handler.go b/handler.go index 0e4ebacde..7a094bdc9 100644 --- a/handler.go +++ b/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 -} diff --git a/server.go b/server.go index b3e42a4c8..68e9c6a74 100644 --- a/server.go +++ b/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. diff --git a/server/server.go b/server/server.go index 93b0f184f..22bfc592a 100644 --- a/server/server.go +++ b/server/server.go @@ -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 diff --git a/server/server_test.go b/server/server_test.go index ca8ac7f09..f9f63581a 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -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 From dd390453ac207478aaaafeecd8dcdd6727e2793d Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 8 Dec 2017 16:50:17 -0600 Subject: [PATCH 9/9] improve error messages in client.SendMessage() --- client.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/client.go b/client.go index f160d88bc..f71cf57d4 100644 --- a/client.go +++ b/client.go @@ -1048,7 +1048,7 @@ func (c *InternalHTTPClient) RowAttrDiff(ctx context.Context, index, frame strin func (c *InternalHTTPClient) SendMessage(ctx context.Context, pb proto.Message) error { msg, err := MarshalMessage(pb) if err != nil { - return err + return fmt.Errorf("marshaling message: %v", err) } u := uriPathToURL(ctx.Value("uri").(*URI), "/cluster/message") @@ -1059,21 +1059,21 @@ func (c *InternalHTTPClient) SendMessage(ctx context.Context, pb proto.Message) // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { - return err + return fmt.Errorf("executing http request: %v", err) } defer resp.Body.Close() // Read body. body, err := ioutil.ReadAll(resp.Body) if err != nil { - return err + return fmt.Errorf("reading response body: %v", err) } // Return error if status is not OK. switch resp.StatusCode { case http.StatusOK: // ok default: - return errors.New(string(body)) + return fmt.Errorf("unexpected response status code: %d: %s", resp.StatusCode, body) } return nil