diff --git a/api.go b/api.go index 09f0531d5..cae91b20f 100644 --- a/api.go +++ b/api.go @@ -149,6 +149,10 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er return resp, nil } +func (api *API) Holder() *Holder { + return api.server.Holder() +} + // readColumnAttrSets returns a list of column attribute objects by id. func (api *API) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet, error) { if index == nil { diff --git a/broadcast.go b/broadcast.go index 1da6ca57f..00835db64 100644 --- a/broadcast.go +++ b/broadcast.go @@ -16,7 +16,6 @@ package pilosa import ( "fmt" - "reflect" "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" @@ -78,48 +77,13 @@ const ( messageTypeNodeStatus ) -// MarshalMessage encodes the protobuf message into a byte slice. -func MarshalMessage(m proto.Message) ([]byte, error) { - var typ uint8 - switch obj := m.(type) { - case *internal.CreateShardMessage: - typ = messageTypeCreateShard - case *internal.CreateIndexMessage: - typ = messageTypeCreateIndex - case *internal.DeleteIndexMessage: - typ = messageTypeDeleteIndex - case *internal.CreateFieldMessage: - typ = messageTypeCreateField - case *internal.DeleteFieldMessage: - typ = messageTypeDeleteField - case *internal.CreateViewMessage: - typ = messageTypeCreateView - case *internal.DeleteViewMessage: - typ = messageTypeDeleteView - case *internal.ClusterStatus: - typ = messageTypeClusterStatus - case *internal.ResizeInstruction: - typ = messageTypeResizeInstruction - case *internal.ResizeInstructionComplete: - typ = messageTypeResizeInstructionComplete - case *internal.SetCoordinatorMessage: - typ = messageTypeSetCoordinator - case *internal.UpdateCoordinatorMessage: - typ = messageTypeUpdateCoordinator - case *internal.NodeStateMessage: - typ = messageTypeNodeState - case *internal.RecalculateCaches: - typ = messageTypeRecalculateCaches - case *internal.NodeEventMessage: - typ = messageTypeNodeEvent - case *internal.NodeStatus: - typ = messageTypeNodeStatus - default: - return nil, fmt.Errorf("message type not implemented for marshalling: %s", reflect.TypeOf(obj)) - } - buf, err := proto.Marshal(m) +// MarshalInternalMessage serializes the pilosa message and adds pilosa internal +// type info which is used by the internal messaging stuff. +func MarshalInternalMessage(m Message, s Serializer) ([]byte, error) { + typ := getMessageType(m) + buf, err := s.Marshal(m) if err != nil { - return nil, errors.Wrap(err, "marshalling") + return nil, errors.Wrap(err, "marshaling") } return append([]byte{typ}, buf...), nil } diff --git a/broadcast_test.go b/broadcast_test.go deleted file mode 100644 index 415228718..000000000 --- a/broadcast_test.go +++ /dev/null @@ -1,51 +0,0 @@ -// Copyright 2017 Pilosa Corp. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package pilosa_test - -import ( - "reflect" - "testing" - - "github.com/gogo/protobuf/proto" - "github.com/pilosa/pilosa" - "github.com/pilosa/pilosa/internal" -) - -// Ensure a message can be marshaled and unmarshaled. -func TestMessage_Marshal(t *testing.T) { - - testMessageMarshal(t, &internal.CreateShardMessage{ - Index: "i", - Shard: 8, - }) - - testMessageMarshal(t, &internal.DeleteIndexMessage{ - Index: "i", - }) -} - -func testMessageMarshal(t *testing.T, m proto.Message) { - marshalled, err := pilosa.MarshalMessage(m) - if err != nil { - t.Fatal(err) - } - unmarshalled, err := pilosa.UnmarshalMessage(marshalled) - if err != nil { - t.Fatal(err) - } - if !reflect.DeepEqual(unmarshalled, m) { - t.Fatalf("unexpected message marshalling: %s", unmarshalled) - } -} diff --git a/encoding/proto/proto.go b/encoding/proto/proto.go index 1c83ba4ea..10cc13dc4 100644 --- a/encoding/proto/proto.go +++ b/encoding/proto/proto.go @@ -153,6 +153,14 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { } decodeNodeStatus(msg, mt) return nil + case *pilosa.Node: + msg := &internal.Node{} + err := proto.Unmarshal(buf, msg) + if err != nil { + return errors.Wrap(err, "unmarshaling Node") + } + decodeNode(msg, mt) + return nil default: panic(fmt.Sprintf("unhandled pilosa.Message of type %T: %#v", mt, m)) } @@ -192,6 +200,8 @@ func encodeToProto(m pilosa.Message) proto.Message { return encodeNodeEventMessage(mt) case *pilosa.NodeStatus: return encodeNodeStatus(mt) + case *pilosa.Node: + return encodeNode(mt) } return nil } @@ -199,8 +209,8 @@ func encodeToProto(m pilosa.Message) proto.Message { func encodeResizeInstruction(m *pilosa.ResizeInstruction) *internal.ResizeInstruction { return &internal.ResizeInstruction{ JobID: m.JobID, - Node: EncodeNode(m.Node), - Coordinator: EncodeNode(m.Coordinator), + Node: encodeNode(m.Node), + Coordinator: encodeNode(m.Coordinator), Sources: encodeResizeSources(m.Sources), Schema: encodeSchema(m.Schema), ClusterStatus: encodeClusterStatus(m.ClusterStatus), @@ -217,7 +227,7 @@ func encodeResizeSources(srcs []*pilosa.ResizeSource) []*internal.ResizeSource { func encodeResizeSource(m *pilosa.ResizeSource) *internal.ResizeSource { return &internal.ResizeSource{ - Node: EncodeNode(m.Node), + Node: encodeNode(m.Node), Index: m.Index, Field: m.Field, View: m.View, @@ -286,13 +296,13 @@ func encodeFieldOptions(o *pilosa.FieldOptions) *internal.FieldOptions { func EncodeNodes(a []*pilosa.Node) []*internal.Node { other := make([]*internal.Node, len(a)) for i := range a { - other[i] = EncodeNode(a[i]) + other[i] = encodeNode(a[i]) } return other } -// EncodeNode converts a Node into its internal representation. -func EncodeNode(n *pilosa.Node) *internal.Node { +// encodeNode converts a Node into its internal representation. +func encodeNode(n *pilosa.Node) *internal.Node { return &internal.Node{ ID: n.ID, URI: n.URI.Encode(), @@ -368,20 +378,20 @@ func encodeDeleteViewMessage(m *pilosa.DeleteViewMessage) *internal.DeleteViewMe func encodeResizeInstructionComplete(m *pilosa.ResizeInstructionComplete) *internal.ResizeInstructionComplete { return &internal.ResizeInstructionComplete{ JobID: m.JobID, - Node: EncodeNode(m.Node), + Node: encodeNode(m.Node), Error: m.Error, } } func encodeSetCoordinatorMessage(m *pilosa.SetCoordinatorMessage) *internal.SetCoordinatorMessage { return &internal.SetCoordinatorMessage{ - New: EncodeNode(m.New), + New: encodeNode(m.New), } } func encodeUpdateCoordinatorMessage(m *pilosa.UpdateCoordinatorMessage) *internal.UpdateCoordinatorMessage { return &internal.UpdateCoordinatorMessage{ - New: EncodeNode(m.New), + New: encodeNode(m.New), } } @@ -395,13 +405,13 @@ func encodeNodeStateMessage(m *pilosa.NodeStateMessage) *internal.NodeStateMessa func encodeNodeEventMessage(m *pilosa.NodeEvent) *internal.NodeEventMessage { return &internal.NodeEventMessage{ Event: uint32(m.Event), - Node: EncodeNode(m.Node), + Node: encodeNode(m.Node), } } func encodeNodeStatus(m *pilosa.NodeStatus) *internal.NodeStatus { return &internal.NodeStatus{ - Node: EncodeNode(m.Node), + Node: encodeNode(m.Node), MaxShards: &internal.MaxShards{Standard: m.MaxShards}, Schema: encodeSchema(m.Schema), } diff --git a/gossip/gossip.go b/gossip/gossip.go index cd1d9ad0b..251851077 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -26,10 +26,9 @@ import ( "sync" "time" - "github.com/gogo/protobuf/proto" "github.com/hashicorp/memberlist" "github.com/pilosa/pilosa" - "github.com/pilosa/pilosa/internal" + "github.com/pilosa/pilosa/encoding/proto" "github.com/pilosa/pilosa/toml" "github.com/pkg/errors" ) @@ -44,8 +43,9 @@ type GossipMemberSet struct { broadcasts *memberlist.TransmitLimitedQueue - papi *pilosa.API - config *gossipConfig + papi *pilosa.API + serializer pilosa.Serializer + config *gossipConfig Logger pilosa.Logger @@ -150,8 +150,9 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption { func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...GossipMemberSetOption) (*GossipMemberSet, error) { host := api.Node().URI.GetHost() g := &GossipMemberSet{ - papi: api, - Logger: pilosa.NopLogger, + papi: api, + serializer: proto.Serializer{}, + Logger: pilosa.NopLogger, } // options @@ -222,7 +223,7 @@ func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...GossipMemberSetO // NodeMeta implementation of the memberlist.Delegate interface. func (g *GossipMemberSet) NodeMeta(limit int) []byte { - buf, err := proto.Marshal(pilosa.EncodeNode(g.papi.Node())) + buf, err := g.serializer.Marshal(g.papi.Node()) if err != nil { g.Logger.Printf("marshal message error: %s", err) return []byte{} @@ -248,14 +249,14 @@ 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 := &internal.NodeStatus{ - Node: pilosa.EncodeNode(g.papi.Node()), - MaxShards: &internal.MaxShards{Standard: g.papi.MaxShards(context.Background())}, - Schema: &internal.Schema{Indexes: pilosa.EncodeIndexes(g.papi.Schema(context.Background()))}, + m := &pilosa.NodeStatus{ + Node: g.papi.Node(), + MaxShards: g.papi.MaxShards(context.Background()), + Schema: &pilosa.Schema{Indexes: g.papi.Holder().Schema()}, } // Marshal nodestate data to bytes. - buf, err := pilosa.MarshalMessage(pb) + buf, err := pilosa.MarshalInternalMessage(m, g.serializer) if err != nil { g.Logger.Printf("error marshalling nodestate data, err=%s", err) return []byte{} @@ -278,8 +279,9 @@ func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) { // Care must be taken that events are processed in a timely manner from // the channel, since this delegate will block until an event can be sent. type gossipEventReceiver struct { - ch chan memberlist.NodeEvent - papi *pilosa.API + ch chan memberlist.NodeEvent + papi *pilosa.API + serializer pilosa.Serializer logger *log.Logger } @@ -287,9 +289,10 @@ type gossipEventReceiver struct { // newGossipEventReceiver returns a new instance of GossipEventReceiver. func newGossipEventReceiver(logger *log.Logger, papi *pilosa.API) *gossipEventReceiver { ger := &gossipEventReceiver{ - ch: make(chan memberlist.NodeEvent, 1), - logger: logger, - papi: papi, + ch: make(chan memberlist.NodeEvent, 1), + logger: logger, + papi: papi, + serializer: proto.Serializer{}, } go ger.listen() return ger @@ -323,16 +326,16 @@ func (g *gossipEventReceiver) listen() { } // Get the node from the event.Node meta data. - var n internal.Node - if err := proto.Unmarshal(e.Node.Meta, &n); err != nil { - panic("failed to unmarshal event node meta data") + var n pilosa.Node + if err := g.serializer.Unmarshal(e.Node.Meta, &n); err != nil { + panic("failed to unmarshal event node meta into node") } - ne := &internal.NodeEventMessage{ - Event: uint32(nodeEventType), + ne := &pilosa.NodeEvent{ + Event: nodeEventType, Node: &n, } - buf, err := pilosa.MarshalMessage(ne) + buf, err := pilosa.MarshalInternalMessage(ne, g.serializer) if err != nil { panic(err) }