diff --git a/api.go b/api.go index 33fff083d..3fca21cab 100644 --- a/api.go +++ b/api.go @@ -331,8 +331,8 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin } // Validate that this handler owns the shard. - if !api.cluster.ownsShard(api.LocalID(), indexName, shard) { - api.server.logger.Printf("node %s does not own shard %d of index %s", api.LocalID(), shard, indexName) + if !api.cluster.ownsShard(api.Node().ID, indexName, shard) { + api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName) return ErrClusterDoesNotOwnShard } @@ -477,6 +477,10 @@ func (api *API) Hosts(ctx context.Context) []*Node { return api.cluster.Nodes } +func (api *API) Node() *Node { + return api.server.Node() +} + // RecalculateCaches forces all TopN caches to be updated. Used mainly for integration tests. func (api *API) RecalculateCaches(ctx context.Context) error { if err := api.validate(apiRecalculateCaches); err != nil { @@ -517,15 +521,10 @@ func (api *API) ClusterMessage(ctx context.Context, reqBody io.Reader) error { return nil } -// LocalID returns the current node's ID. -func (api *API) LocalID() string { - return api.cluster.Node.ID -} - // Schema returns information about each index in Pilosa including which fields // and views they contain. -func (api *API) Schema(ctx context.Context) []*IndexInfo { - return api.holder.Schema() +func (api *API) Schema(ctx context.Context) []*Index { + return api.holder.Indexes() } // Views returns the views in the given field. @@ -720,8 +719,8 @@ func (api *API) LongQueryTime() time.Duration { func (api *API) indexField(indexName string, fieldName string, shard uint64) (*Index, *Field, error) { // Validate that this handler owns the shard. - if !api.cluster.ownsShard(api.LocalID(), indexName, shard) { - api.server.logger.Printf("node %s does not own shard %d of index %s", api.LocalID(), shard, indexName) + if !api.cluster.ownsShard(api.Node().ID, indexName, shard) { + api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName) return nil, nil, ErrClusterDoesNotOwnShard } @@ -755,7 +754,7 @@ func (api *API) SetCoordinator(ctx context.Context, id string) (oldNode, newNode } // If the new coordinator is this node, do the SetCoordinator directly. - if newNode.ID == api.LocalID() { + if newNode.ID == api.Node().ID { return oldNode, newNode, api.cluster.setCoordinator(newNode) } diff --git a/broadcast.go b/broadcast.go index 37452292a..19e927c18 100644 --- a/broadcast.go +++ b/broadcast.go @@ -65,6 +65,7 @@ const ( messageTypeNodeState messageTypeRecalculateCaches messageTypeNodeEvent + messageTypeNodeStatus ) // MarshalMessage encodes the protobuf message into a byte slice. @@ -101,6 +102,8 @@ func MarshalMessage(m proto.Message) ([]byte, error) { 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)) } @@ -114,7 +117,6 @@ func MarshalMessage(m proto.Message) ([]byte, error) { // UnmarshalMessage decodes the byte slice into a protobuf message. func UnmarshalMessage(buf []byte) (proto.Message, error) { typ, buf := buf[0], buf[1:] - var m proto.Message switch typ { case messageTypeCreateShard: @@ -147,6 +149,8 @@ func UnmarshalMessage(buf []byte) (proto.Message, error) { m = &internal.RecalculateCaches{} case messageTypeNodeEvent: m = &internal.NodeEventMessage{} + case messageTypeNodeStatus: + m = &internal.NodeStatus{} default: return nil, fmt.Errorf("invalid message type: %d", typ) } diff --git a/field.go b/field.go index 2e4897f91..76be9cf3f 100644 --- a/field.go +++ b/field.go @@ -1073,6 +1073,21 @@ func (f *Field) ImportValue(columnIDs []uint64, values []int64) error { return nil } +func (f *Field) MarshalJSON() ([]byte, error) { + thing := struct { + Name string + Options FieldOptions + Views []*viewInfo + }{ + Name: f.Name(), + Options: f.Options(), + } + for _, viewname := range f.viewNames() { + thing.Views = append(thing.Views, &viewInfo{Name: viewname}) + } + return json.Marshal(thing) +} + // encodeFields converts a into its internal representation. func encodeFields(a []*Field) []*internal.Field { other := make([]*internal.Field, len(a)) diff --git a/gossip/gossip.go b/gossip/gossip.go index 4f562af5e..3dc7c202b 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -15,6 +15,8 @@ package gossip import ( + "bytes" + "context" "fmt" "io/ioutil" "log" @@ -42,8 +44,8 @@ type GossipMemberSet struct { broadcasts *memberlist.TransmitLimitedQueue - pserver pilosa.MemberServer - config *gossipConfig + papi *pilosa.API + config *gossipConfig Logger pilosa.Logger @@ -145,9 +147,10 @@ func WithLogger(logger *log.Logger) GossipMemberSetOption { } // NewGossipMemberSet returns a new instance of GossipMemberSet based on options. -func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) { - host := s.Node().URI.Host() +func NewGossipMemberSet(cfg Config, api *pilosa.API, options ...GossipMemberSetOption) (*GossipMemberSet, error) { + host := api.Node().URI.Host() g := &GossipMemberSet{ + papi: api, Logger: pilosa.NopLogger, } @@ -157,7 +160,7 @@ func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSet return nil, errors.Wrap(err, "executing option") } } - ger := newGossipEventReceiver(g.logger, s) + ger := newGossipEventReceiver(g.logger, api) g.gossipEventReceiver = ger if g.transport == nil { @@ -189,11 +192,11 @@ func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSet // memberlist config conf := memberlist.DefaultWANConfig() conf.Transport = g.transport.Net - conf.Name = s.Node().ID - conf.BindAddr = s.Node().URI.Host() + conf.Name = api.Node().ID + conf.BindAddr = api.Node().URI.Host() conf.BindPort = port conf.AdvertisePort = port - conf.AdvertiseAddr = hostToIP(s.Node().URI.Host()) + conf.AdvertiseAddr = hostToIP(api.Node().URI.Host()) // conf.TCPTimeout = time.Duration(cfg.StreamTimeout) conf.SuspicionMult = cfg.SuspicionMult @@ -214,14 +217,12 @@ func NewGossipMemberSet(cfg Config, s *pilosa.Server, options ...GossipMemberSet gossipSeeds: cfg.Seeds, } - g.pserver = s - return g, nil } // NodeMeta implementation of the memberlist.Delegate interface. func (g *GossipMemberSet) NodeMeta(limit int) []byte { - buf, err := proto.Marshal(pilosa.EncodeNode(g.pserver.Node())) + buf, err := proto.Marshal(pilosa.EncodeNode(g.papi.Node())) if err != nil { g.Logger.Printf("marshal message error: %s", err) return []byte{} @@ -232,14 +233,9 @@ func (g *GossipMemberSet) NodeMeta(limit int) []byte { // NotifyMsg implementation of the memberlist.Delegate interface // called when a user-data message is received. func (g *GossipMemberSet) NotifyMsg(b []byte) { - m, err := pilosa.UnmarshalMessage(b) + err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(b)) if err != nil { - g.Logger.Printf("unmarshal message error: %s", err) - return - } - if err := g.pserver.ReceiveMessage(m); err != nil { - g.Logger.Printf("receive message error: %s", err) - return + g.Logger.Printf("cluster message error: %s", err) } } @@ -252,14 +248,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, err := g.pserver.LocalStatus() - if err != nil { - g.Logger.Printf("error getting local state, err=%s", err) - return []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()))}, } // Marshal nodestate data to bytes. - buf, err := proto.Marshal(pb) + buf, err := pilosa.MarshalMessage(pb) if err != nil { g.Logger.Printf("error marshalling nodestate data, err=%s", err) return []byte{} @@ -270,13 +266,7 @@ func (g *GossipMemberSet) LocalState(join bool) []byte { // MergeRemoteState implementation of the memberlist.Delegate interface // receive and process the remote side's LocalState. func (g *GossipMemberSet) MergeRemoteState(buf []byte, join bool) { - // Unmarshal nodestate data. - var pb internal.NodeStatus - if err := proto.Unmarshal(buf, &pb); err != nil { - g.Logger.Printf("error unmarshalling nodestate data, err=%s", err) - return - } - err := g.pserver.HandleRemoteStatus(&pb) + err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(buf)) if err != nil { g.Logger.Printf("merge state error: %s", err) } @@ -288,18 +278,18 @@ 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 - eventHandler *pilosa.Server + ch chan memberlist.NodeEvent + papi *pilosa.API logger *log.Logger } // newGossipEventReceiver returns a new instance of GossipEventReceiver. -func newGossipEventReceiver(logger *log.Logger, pserver *pilosa.Server) *gossipEventReceiver { +func newGossipEventReceiver(logger *log.Logger, papi *pilosa.API) *gossipEventReceiver { ger := &gossipEventReceiver{ - ch: make(chan memberlist.NodeEvent, 1), - logger: logger, - eventHandler: pserver, + ch: make(chan memberlist.NodeEvent, 1), + logger: logger, + papi: papi, } go ger.listen() return ger @@ -342,7 +332,11 @@ func (g *gossipEventReceiver) listen() { Event: uint32(nodeEventType), Node: &n, } - if err := g.eventHandler.ReceiveMessage(ne); err != nil { + buf, err := pilosa.MarshalMessage(ne) + if err != nil { + panic(err) + } + if err := g.papi.ClusterMessage(context.Background(), bytes.NewBuffer(buf)); err != nil { g.logger.Printf("receive event error: %s", err) } } diff --git a/holder.go b/holder.go index bc255abba..e599d21f2 100644 --- a/holder.go +++ b/holder.go @@ -267,7 +267,7 @@ func (h *Holder) encodeMaxShards() *internal.MaxShards { // encodeSchema creates an internal representation of schema. func (h *Holder) encodeSchema() *internal.Schema { return &internal.Schema{ - Indexes: encodeIndexes(h.Indexes()), + Indexes: EncodeIndexes(h.Indexes()), } } diff --git a/http/handler.go b/http/handler.go index 9d6b9e080..9195f05e9 100644 --- a/http/handler.go +++ b/http/handler.go @@ -358,9 +358,7 @@ func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { } schema := h.API.Schema(r.Context()) - if err := json.NewEncoder(w).Encode(getSchemaResponse{ - Indexes: schema, - }); err != nil { + if err := json.NewEncoder(w).Encode(map[string]interface{}{"indexes": schema}); err != nil { h.Logger.Printf("write schema response error: %s", err) } } @@ -374,7 +372,7 @@ func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) { status := getStatusResponse{ State: h.API.State(), Nodes: h.API.Hosts(r.Context()), - LocalID: h.API.LocalID(), + LocalID: h.API.Node().ID, } if err := json.NewEncoder(w).Encode(status); err != nil { h.Logger.Printf("write status response error: %s", err) @@ -467,7 +465,7 @@ func (h *Handler) handleGetIndex(w http.ResponseWriter, r *http.Request) { } indexName := mux.Vars(r)["index"] for _, idx := range h.API.Schema(r.Context()) { - if idx.Name == indexName { + if idx.Name() == indexName { if err := json.NewEncoder(w).Encode(idx); err != nil { h.Logger.Printf("write response error: %s", err) } diff --git a/index.go b/index.go index d1a6a8ee5..98f50eced 100644 --- a/index.go +++ b/index.go @@ -15,6 +15,7 @@ package pilosa import ( + "encoding/json" "fmt" "io/ioutil" "os" @@ -75,6 +76,22 @@ func NewIndex(path, name string) (*Index, error) { }, nil } +func (i *Index) MarshalJSON() ([]byte, error) { + fields := make([]*Field, 0, len(i.fields)) + for _, f := range i.fields { + + fields = append(fields, f) + } + thing := struct { + Name string + Fields []*Field + }{ + Name: i.name, + Fields: fields, + } + return json.Marshal(thing) +} + // Name returns name of the index. func (i *Index) Name() string { return i.name } @@ -386,8 +403,8 @@ func (p indexInfoSlice) Swap(i, j int) { p[i], p[j] = p[j], p[i] } func (p indexInfoSlice) Len() int { return len(p) } func (p indexInfoSlice) Less(i, j int) bool { return p[i].Name < p[j].Name } -// encodeIndexes converts a into its internal representation. -func encodeIndexes(a []*Index) []*internal.Index { +// EncodeIndexes converts a into its internal representation. +func EncodeIndexes(a []*Index) []*internal.Index { other := make([]*internal.Index, len(a)) for i := range a { other[i] = encodeIndex(a[i]) diff --git a/server.go b/server.go index 7fbd94066..14f2be179 100644 --- a/server.go +++ b/server.go @@ -511,6 +511,8 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { s.holder.RecalculateCaches() case *internal.NodeEventMessage: s.cluster.ReceiveEvent(DecodeNodeEvent(obj)) + case *internal.NodeStatus: + s.HandleRemoteStatus(pb) } return nil diff --git a/server/server.go b/server/server.go index 164dc94e3..7781ed4bc 100644 --- a/server/server.go +++ b/server/server.go @@ -317,7 +317,7 @@ func (m *Command) SetupNetworking() error { gossipMemberSet, err := gossip.NewGossipMemberSet( m.Config.Gossip, - m.Server, + m.API, gossip.WithLogger(m.logger.Logger()), gossip.WithTransport(m.gossipTransport), )