mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-12 07:41:02 +00:00
get gossip using serializer stuff, remove proto and internal
This commit is contained in:
parent
2cec75e399
commit
6309d3b7f7
5 changed files with 57 additions and 127 deletions
4
api.go
4
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 {
|
||||
|
|
|
|||
48
broadcast.go
48
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
|
|
@ -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),
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue