featurebase/encoding/proto/proto.go
Seebs 3b696da34a plugins and precomputed data
So in some cases, when we do a query, the results of one
part of the query are innately shared-across-nodes; for
instance, a hypothetical Distinct query. More generally,
we allow cross-index queries; calls can have "index=foo"
in them.

This patch lets us handle that without duplicating that
query all over. Before we actually start doing the
separate calls, we run the query once from the coordinating
node, then patch the results in, and send relevant subsets
over to each client, etcetera. Also provides slightly
friendlier (and I hope faster) support for converting
bitmaps to/from sets of rows.

We also add an extension interface, and some fancy stuff
to let us define new calls, which use this. They're sort
of tied together because the first extension I wanted to
implement needed precomputed calls. The extension API
lets us create extensions using `pkg/plugin` (with all its
associated limitations, unfortunately), then query them
at load time for functionality.

This also implies some revamping of the argument
validation for PQL, like verifying that functions exist
and knowing things about their argument types.

So basically this is an overly intrusive patch, and would
be better as separate patches, but they're hard to detangle.

add trivial execution-time profiling

What if you could ?profile=true on a query and get some
numbers back? That'd be really cool.

We already have tracing/spans, but right now, those only generate
any data if you have something set up for them to trace to. Add a
fancy wrapper that lets us generate our own tracing data, and dump
it into the request response, if ?profile=true.

add a sample extension, add missing features to extension interface

Implement a naive probabilistic filter extension as an example of
what an extension looks like. In the process, discover multiple
omissions in the bitmap API. Well, I did *say* it was experimental.
2019-11-12 12:14:29 -06:00

1369 lines
35 KiB
Go

// 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 proto
import (
"fmt"
"sort"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/internal"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pkg/errors"
)
// Serializer implements pilosa.Serializer for protobufs.
type Serializer struct{}
// Marshal turns pilosa messages into protobuf serialized bytes.
func (Serializer) Marshal(m pilosa.Message) ([]byte, error) {
pm := encodeToProto(m)
if pm == nil {
return nil, errors.New("passed invalid pilosa.Message")
}
buf, err := proto.Marshal(pm)
return buf, errors.Wrap(err, "marshalling")
}
// Unmarshal takes byte slices and protobuf deserializes them into a pilosa Message.
func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error {
switch mt := m.(type) {
case *pilosa.CreateShardMessage:
msg := &internal.CreateShardMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling CreateShardMessage")
}
decodeCreateShardMessage(msg, mt)
return nil
case *pilosa.CreateIndexMessage:
msg := &internal.CreateIndexMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling CreateIndexMessage")
}
decodeCreateIndexMessage(msg, mt)
return nil
case *pilosa.DeleteIndexMessage:
msg := &internal.DeleteIndexMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling DeleteIndexMessage")
}
decodeDeleteIndexMessage(msg, mt)
return nil
case *pilosa.CreateFieldMessage:
msg := &internal.CreateFieldMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling CreateFieldMessage")
}
decodeCreateFieldMessage(msg, mt)
return nil
case *pilosa.DeleteFieldMessage:
msg := &internal.DeleteFieldMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling DeleteFieldMessage")
}
decodeDeleteFieldMessage(msg, mt)
return nil
case *pilosa.DeleteAvailableShardMessage:
msg := &internal.DeleteAvailableShardMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling DeleteAvailableShardMessage")
}
decodeDeleteAvailableShardMessage(msg, mt)
return nil
case *pilosa.CreateViewMessage:
msg := &internal.CreateViewMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling CreateViewMessage")
}
decodeCreateViewMessage(msg, mt)
return nil
case *pilosa.DeleteViewMessage:
msg := &internal.DeleteViewMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling DeleteViewMessage")
}
decodeDeleteViewMessage(msg, mt)
return nil
case *pilosa.ClusterStatus:
msg := &internal.ClusterStatus{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ClusterStatus")
}
decodeClusterStatus(msg, mt)
return nil
case *pilosa.ResizeInstruction:
msg := &internal.ResizeInstruction{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ResizeInstruction")
}
decodeResizeInstruction(msg, mt)
return nil
case *pilosa.ResizeInstructionComplete:
msg := &internal.ResizeInstructionComplete{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ResizeInstructionComplete")
}
decodeResizeInstructionComplete(msg, mt)
return nil
case *pilosa.SetCoordinatorMessage:
msg := &internal.SetCoordinatorMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling SetCoordinatorMessage")
}
decodeSetCoordinatorMessage(msg, mt)
return nil
case *pilosa.UpdateCoordinatorMessage:
msg := &internal.UpdateCoordinatorMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling UpdateCoordinatorMessage")
}
decodeUpdateCoordinatorMessage(msg, mt)
return nil
case *pilosa.NodeStateMessage:
msg := &internal.NodeStateMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling NodeStateMessage")
}
decodeNodeStateMessage(msg, mt)
return nil
case *pilosa.RecalculateCaches:
msg := &internal.RecalculateCaches{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling RecalculateCaches")
}
decodeRecalculateCaches(msg, mt)
return nil
case *pilosa.NodeEvent:
msg := &internal.NodeEventMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling NodeEvent")
}
decodeNodeEventMessage(msg, mt)
return nil
case *pilosa.NodeStatus:
msg := &internal.NodeStatus{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling NodeStatus")
}
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
case *pilosa.QueryRequest:
msg := &internal.QueryRequest{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling QueryRequest")
}
decodeQueryRequest(msg, mt)
return nil
case *pilosa.QueryResponse:
msg := &internal.QueryResponse{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling QueryResponse")
}
decodeQueryResponse(msg, mt)
return nil
case *pilosa.ImportRequest:
msg := &internal.ImportRequest{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ImportRequest")
}
decodeImportRequest(msg, mt)
return nil
case *pilosa.ImportValueRequest:
msg := &internal.ImportValueRequest{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ImportValueRequest")
}
decodeImportValueRequest(msg, mt)
return nil
case *pilosa.ImportRoaringRequest:
msg := &internal.ImportRoaringRequest{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ImportRoaringRequest")
}
decodeImportRoaringRequest(msg, mt)
return nil
case *pilosa.ImportResponse:
msg := &internal.ImportResponse{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling ImportResponse")
}
decodeImportResponse(msg, mt)
return nil
case *pilosa.BlockDataRequest:
msg := &internal.BlockDataRequest{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling BlockDataRequest")
}
decodeBlockDataRequest(msg, mt)
return nil
case *pilosa.BlockDataResponse:
msg := &internal.BlockDataResponse{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling BlockDataResponse")
}
decodeBlockDataResponse(msg, mt)
return nil
case *pilosa.TranslateKeysRequest:
msg := &internal.TranslateKeysRequest{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling TranslateKeysRequest")
}
decodeTranslateKeysRequest(msg, mt)
return nil
case *pilosa.TranslateKeysResponse:
msg := &internal.TranslateKeysResponse{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling TranslateKeysResponse")
}
decodeTranslateKeysResponse(msg, mt)
return nil
default:
panic(fmt.Sprintf("unhandled pilosa.Message of type %T: %#v", mt, m))
}
}
func encodeToProto(m pilosa.Message) proto.Message {
switch mt := m.(type) {
case *pilosa.CreateShardMessage:
return encodeCreateShardMessage(mt)
case *pilosa.CreateIndexMessage:
return encodeCreateIndexMessage(mt)
case *pilosa.DeleteIndexMessage:
return encodeDeleteIndexMessage(mt)
case *pilosa.CreateFieldMessage:
return encodeCreateFieldMessage(mt)
case *pilosa.DeleteFieldMessage:
return encodeDeleteFieldMessage(mt)
case *pilosa.DeleteAvailableShardMessage:
return encodeDeleteAvailableShardMessage(mt)
case *pilosa.CreateViewMessage:
return encodeCreateViewMessage(mt)
case *pilosa.DeleteViewMessage:
return encodeDeleteViewMessage(mt)
case *pilosa.ClusterStatus:
return encodeClusterStatus(mt)
case *pilosa.ResizeInstruction:
return encodeResizeInstruction(mt)
case *pilosa.ResizeInstructionComplete:
return encodeResizeInstructionComplete(mt)
case *pilosa.SetCoordinatorMessage:
return encodeSetCoordinatorMessage(mt)
case *pilosa.UpdateCoordinatorMessage:
return encodeUpdateCoordinatorMessage(mt)
case *pilosa.NodeStateMessage:
return encodeNodeStateMessage(mt)
case *pilosa.RecalculateCaches:
return encodeRecalculateCaches(mt)
case *pilosa.NodeEvent:
return encodeNodeEventMessage(mt)
case *pilosa.NodeStatus:
return encodeNodeStatus(mt)
case *pilosa.Node:
return encodeNode(mt)
case *pilosa.QueryRequest:
return encodeQueryRequest(mt)
case *pilosa.QueryResponse:
return encodeQueryResponse(mt)
case *pilosa.ImportRequest:
return encodeImportRequest(mt)
case *pilosa.ImportValueRequest:
return encodeImportValueRequest(mt)
case *pilosa.ImportRoaringRequest:
return encodeImportRoaringRequest(mt)
case *pilosa.ImportResponse:
return encodeImportResponse(mt)
case *pilosa.BlockDataRequest:
return encodeBlockDataRequest(mt)
case *pilosa.BlockDataResponse:
return encodeBlockDataResponse(mt)
case *pilosa.TranslateKeysRequest:
return encodeTranslateKeysRequest(mt)
case *pilosa.TranslateKeysResponse:
return encodeTranslateKeysResponse(mt)
}
return nil
}
func encodeBlockDataRequest(m *pilosa.BlockDataRequest) *internal.BlockDataRequest {
return &internal.BlockDataRequest{
Index: m.Index,
Field: m.Field,
View: m.View,
Shard: m.Shard,
Block: m.Block,
}
}
func encodeBlockDataResponse(m *pilosa.BlockDataResponse) *internal.BlockDataResponse {
return &internal.BlockDataResponse{
RowIDs: m.RowIDs,
ColumnIDs: m.ColumnIDs,
}
}
func encodeImportResponse(m *pilosa.ImportResponse) *internal.ImportResponse {
return &internal.ImportResponse{
Err: m.Err,
}
}
func encodeImportRequest(m *pilosa.ImportRequest) *internal.ImportRequest {
return &internal.ImportRequest{
Index: m.Index,
Field: m.Field,
Shard: m.Shard,
RowIDs: m.RowIDs,
ColumnIDs: m.ColumnIDs,
RowKeys: m.RowKeys,
ColumnKeys: m.ColumnKeys,
Timestamps: m.Timestamps,
}
}
func encodeImportValueRequest(m *pilosa.ImportValueRequest) *internal.ImportValueRequest {
return &internal.ImportValueRequest{
Index: m.Index,
Field: m.Field,
Shard: m.Shard,
ColumnIDs: m.ColumnIDs,
ColumnKeys: m.ColumnKeys,
Values: m.Values,
FloatValues: m.FloatValues,
}
}
func encodeImportRoaringRequest(m *pilosa.ImportRoaringRequest) *internal.ImportRoaringRequest {
views := make([]*internal.ImportRoaringRequestView, len(m.Views))
i := 0
for viewName, viewData := range m.Views {
views[i] = &internal.ImportRoaringRequestView{
Name: viewName,
Data: viewData,
}
i += 1
}
return &internal.ImportRoaringRequest{
Clear: m.Clear,
Views: views,
}
}
func encodeQueryRequest(m *pilosa.QueryRequest) *internal.QueryRequest {
r := &internal.QueryRequest{
Query: m.Query,
Shards: m.Shards,
ColumnAttrs: m.ColumnAttrs,
Remote: m.Remote,
ExcludeRowAttrs: m.ExcludeRowAttrs,
ExcludeColumns: m.ExcludeColumns,
EmbeddedData: make([]*internal.Row, len(m.EmbeddedData)),
}
for i := range m.EmbeddedData {
r.EmbeddedData[i] = encodeRow(m.EmbeddedData[i])
}
return r
}
func encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse {
pb := &internal.QueryResponse{
Results: make([]*internal.QueryResult, len(m.Results)),
ColumnAttrSets: encodeColumnAttrSets(m.ColumnAttrSets),
}
for i := range m.Results {
pb.Results[i] = &internal.QueryResult{}
switch result := m.Results[i].(type) {
case pilosa.SignedRow:
pb.Results[i].Type = queryResultTypeSignedRow
pb.Results[i].SignedRow = encodeSignedRow(result)
case *pilosa.Row:
pb.Results[i].Type = queryResultTypeRow
pb.Results[i].Row = encodeRow(result)
case []pilosa.Pair:
pb.Results[i].Type = queryResultTypePairs
pb.Results[i].Pairs = encodePairs(result)
case pilosa.ValCount:
pb.Results[i].Type = queryResultTypeValCount
pb.Results[i].ValCount = encodeValCount(result)
case uint64:
pb.Results[i].Type = queryResultTypeUint64
pb.Results[i].N = result
case bool:
pb.Results[i].Type = queryResultTypeBool
pb.Results[i].Changed = result
case pilosa.RowIDs:
pb.Results[i].Type = queryResultTypeRowIDs
pb.Results[i].RowIDs = result
case []pilosa.GroupCount:
pb.Results[i].Type = queryResultTypeGroupCounts
pb.Results[i].GroupCounts = encodeGroupCounts(result)
case pilosa.RowIdentifiers:
pb.Results[i].Type = queryResultTypeRowIdentifiers
pb.Results[i].RowIdentifiers = encodeRowIdentifiers(result)
case pilosa.Pair:
pb.Results[i].Type = queryResultTypePair
pb.Results[i].Pairs = []*internal.Pair{encodePair(result)}
case nil:
pb.Results[i].Type = queryResultTypeNil
default:
panic(fmt.Errorf("unknown type: %T", m.Results[i]))
}
}
if m.Err != nil {
pb.Err = m.Err.Error()
}
return pb
}
func encodeResizeInstruction(m *pilosa.ResizeInstruction) *internal.ResizeInstruction {
return &internal.ResizeInstruction{
JobID: m.JobID,
Node: encodeNode(m.Node),
Coordinator: encodeNode(m.Coordinator),
Sources: encodeResizeSources(m.Sources),
NodeStatus: encodeNodeStatus(m.NodeStatus),
ClusterStatus: encodeClusterStatus(m.ClusterStatus),
}
}
func encodeResizeSources(srcs []*pilosa.ResizeSource) []*internal.ResizeSource {
new := make([]*internal.ResizeSource, 0, len(srcs))
for _, src := range srcs {
new = append(new, encodeResizeSource(src))
}
return new
}
func encodeResizeSource(m *pilosa.ResizeSource) *internal.ResizeSource {
return &internal.ResizeSource{
Node: encodeNode(m.Node),
Index: m.Index,
Field: m.Field,
View: m.View,
Shard: m.Shard,
}
}
func encodeSchema(m *pilosa.Schema) *internal.Schema {
return &internal.Schema{
Indexes: encodeIndexInfos(m.Indexes),
}
}
func encodeIndexInfos(idxs []*pilosa.IndexInfo) []*internal.Index {
new := make([]*internal.Index, 0, len(idxs))
for _, idx := range idxs {
new = append(new, encodeIndexInfo(idx))
}
return new
}
func encodeIndexInfo(idx *pilosa.IndexInfo) *internal.Index {
return &internal.Index{
Name: idx.Name,
Fields: encodeFieldInfos(idx.Fields),
}
}
func encodeFieldInfos(fs []*pilosa.FieldInfo) []*internal.Field {
new := make([]*internal.Field, 0, len(fs))
for _, f := range fs {
new = append(new, encodeFieldInfo(f))
}
return new
}
func encodeFieldInfo(f *pilosa.FieldInfo) *internal.Field {
ifield := &internal.Field{
Name: f.Name,
Meta: encodeFieldOptions(&f.Options),
Views: make([]string, 0, len(f.Views)),
}
for _, viewinfo := range f.Views {
ifield.Views = append(ifield.Views, viewinfo.Name)
}
return ifield
}
func encodeFieldOptions(o *pilosa.FieldOptions) *internal.FieldOptions {
if o == nil {
return nil
}
return &internal.FieldOptions{
Type: o.Type,
CacheType: o.CacheType,
CacheSize: o.CacheSize,
Min: o.Min,
Max: o.Max,
Base: o.Base,
Scale: o.Scale,
BitDepth: uint64(o.BitDepth),
TimeQuantum: string(o.TimeQuantum),
Keys: o.Keys,
}
}
// encodeNodes converts a slice of Nodes into its internal representation.
func encodeNodes(a []*pilosa.Node) []*internal.Node {
other := make([]*internal.Node, len(a))
for i := range a {
other[i] = encodeNode(a[i])
}
return other
}
// encodeNode converts a Node into its internal representation.
func encodeNode(n *pilosa.Node) *internal.Node {
return &internal.Node{
ID: n.ID,
URI: encodeURI(n.URI),
IsCoordinator: n.IsCoordinator,
State: n.State,
}
}
func encodeURI(u pilosa.URI) *internal.URI {
return &internal.URI{
Scheme: u.Scheme,
Host: u.Host,
Port: uint32(u.Port),
}
}
func encodeClusterStatus(m *pilosa.ClusterStatus) *internal.ClusterStatus {
return &internal.ClusterStatus{
State: m.State,
ClusterID: m.ClusterID,
Nodes: encodeNodes(m.Nodes),
}
}
func encodeCreateShardMessage(m *pilosa.CreateShardMessage) *internal.CreateShardMessage {
return &internal.CreateShardMessage{
Index: m.Index,
Field: m.Field,
Shard: m.Shard,
}
}
func encodeCreateIndexMessage(m *pilosa.CreateIndexMessage) *internal.CreateIndexMessage {
return &internal.CreateIndexMessage{
Index: m.Index,
Meta: encodeIndexMeta(m.Meta),
}
}
func encodeIndexMeta(m *pilosa.IndexOptions) *internal.IndexMeta {
return &internal.IndexMeta{
Keys: m.Keys,
TrackExistence: m.TrackExistence,
}
}
func encodeDeleteIndexMessage(m *pilosa.DeleteIndexMessage) *internal.DeleteIndexMessage {
return &internal.DeleteIndexMessage{
Index: m.Index,
}
}
func encodeCreateFieldMessage(m *pilosa.CreateFieldMessage) *internal.CreateFieldMessage {
return &internal.CreateFieldMessage{
Index: m.Index,
Field: m.Field,
Meta: encodeFieldOptions(m.Meta),
}
}
func encodeDeleteFieldMessage(m *pilosa.DeleteFieldMessage) *internal.DeleteFieldMessage {
return &internal.DeleteFieldMessage{
Index: m.Index,
Field: m.Field,
}
}
func encodeDeleteAvailableShardMessage(m *pilosa.DeleteAvailableShardMessage) *internal.DeleteAvailableShardMessage {
return &internal.DeleteAvailableShardMessage{
Index: m.Index,
Field: m.Field,
ShardID: m.ShardID,
}
}
func encodeCreateViewMessage(m *pilosa.CreateViewMessage) *internal.CreateViewMessage {
return &internal.CreateViewMessage{
Index: m.Index,
Field: m.Field,
View: m.View,
}
}
func encodeDeleteViewMessage(m *pilosa.DeleteViewMessage) *internal.DeleteViewMessage {
return &internal.DeleteViewMessage{
Index: m.Index,
Field: m.Field,
View: m.View,
}
}
func encodeResizeInstructionComplete(m *pilosa.ResizeInstructionComplete) *internal.ResizeInstructionComplete {
return &internal.ResizeInstructionComplete{
JobID: m.JobID,
Node: encodeNode(m.Node),
Error: m.Error,
}
}
func encodeSetCoordinatorMessage(m *pilosa.SetCoordinatorMessage) *internal.SetCoordinatorMessage {
return &internal.SetCoordinatorMessage{
New: encodeNode(m.New),
}
}
func encodeUpdateCoordinatorMessage(m *pilosa.UpdateCoordinatorMessage) *internal.UpdateCoordinatorMessage {
return &internal.UpdateCoordinatorMessage{
New: encodeNode(m.New),
}
}
func encodeNodeStateMessage(m *pilosa.NodeStateMessage) *internal.NodeStateMessage {
return &internal.NodeStateMessage{
NodeID: m.NodeID,
State: m.State,
}
}
func encodeNodeEventMessage(m *pilosa.NodeEvent) *internal.NodeEventMessage {
return &internal.NodeEventMessage{
Event: uint32(m.Event),
Node: encodeNode(m.Node),
}
}
func encodeNodeStatus(m *pilosa.NodeStatus) *internal.NodeStatus {
return &internal.NodeStatus{
Node: encodeNode(m.Node),
Indexes: encodeIndexStatuses(m.Indexes),
Schema: encodeSchema(m.Schema),
}
}
func encodeIndexStatus(m *pilosa.IndexStatus) *internal.IndexStatus {
return &internal.IndexStatus{
Name: m.Name,
Fields: encodeFieldStatuses(m.Fields),
}
}
func encodeIndexStatuses(a []*pilosa.IndexStatus) []*internal.IndexStatus {
other := make([]*internal.IndexStatus, len(a))
for i := range a {
other[i] = encodeIndexStatus(a[i])
}
return other
}
func encodeFieldStatus(m *pilosa.FieldStatus) *internal.FieldStatus {
return &internal.FieldStatus{
Name: m.Name,
AvailableShards: m.AvailableShards.Slice(),
}
}
func encodeFieldStatuses(a []*pilosa.FieldStatus) []*internal.FieldStatus {
other := make([]*internal.FieldStatus, len(a))
for i := range a {
other[i] = encodeFieldStatus(a[i])
}
return other
}
func encodeRecalculateCaches(*pilosa.RecalculateCaches) *internal.RecalculateCaches {
return &internal.RecalculateCaches{}
}
func encodeTranslateKeysResponse(response *pilosa.TranslateKeysResponse) *internal.TranslateKeysResponse {
return &internal.TranslateKeysResponse{
IDs: response.IDs,
}
}
func encodeTranslateKeysRequest(request *pilosa.TranslateKeysRequest) *internal.TranslateKeysRequest {
return &internal.TranslateKeysRequest{
Index: request.Index,
Field: request.Field,
Keys: request.Keys,
}
}
func decodeResizeInstruction(ri *internal.ResizeInstruction, m *pilosa.ResizeInstruction) {
m.JobID = ri.JobID
m.Node = &pilosa.Node{}
decodeNode(ri.Node, m.Node)
m.Coordinator = &pilosa.Node{}
decodeNode(ri.Coordinator, m.Coordinator)
m.Sources = make([]*pilosa.ResizeSource, len(ri.Sources))
decodeResizeSources(ri.Sources, m.Sources)
m.NodeStatus = &pilosa.NodeStatus{}
decodeNodeStatus(ri.NodeStatus, m.NodeStatus)
m.ClusterStatus = &pilosa.ClusterStatus{}
decodeClusterStatus(ri.ClusterStatus, m.ClusterStatus)
}
func decodeResizeSources(srcs []*internal.ResizeSource, m []*pilosa.ResizeSource) {
for i := range srcs {
m[i] = &pilosa.ResizeSource{}
decodeResizeSource(srcs[i], m[i])
}
}
func decodeResizeSource(rs *internal.ResizeSource, m *pilosa.ResizeSource) {
m.Node = &pilosa.Node{}
decodeNode(rs.Node, m.Node)
m.Index = rs.Index
m.Field = rs.Field
m.View = rs.View
m.Shard = rs.Shard
}
func decodeSchema(s *internal.Schema, m *pilosa.Schema) {
m.Indexes = make([]*pilosa.IndexInfo, len(s.Indexes))
decodeIndexes(s.Indexes, m.Indexes)
}
func decodeIndexes(idxs []*internal.Index, m []*pilosa.IndexInfo) {
for i := range idxs {
m[i] = &pilosa.IndexInfo{}
decodeIndex(idxs[i], m[i])
}
}
func decodeIndex(idx *internal.Index, m *pilosa.IndexInfo) {
m.Name = idx.Name
m.Fields = make([]*pilosa.FieldInfo, len(idx.Fields))
decodeFields(idx.Fields, m.Fields)
}
func decodeFields(fs []*internal.Field, m []*pilosa.FieldInfo) {
for i := range fs {
m[i] = &pilosa.FieldInfo{}
decodeField(fs[i], m[i])
}
}
func decodeField(f *internal.Field, m *pilosa.FieldInfo) {
m.Name = f.Name
m.Options = pilosa.FieldOptions{}
decodeFieldOptions(f.Meta, &m.Options)
m.Views = make([]*pilosa.ViewInfo, 0, len(f.Views))
for _, viewname := range f.Views {
m.Views = append(m.Views, &pilosa.ViewInfo{Name: viewname})
}
}
func decodeFieldOptions(options *internal.FieldOptions, m *pilosa.FieldOptions) {
m.Type = options.Type
m.CacheType = options.CacheType
m.CacheSize = options.CacheSize
m.Min = options.Min
m.Max = options.Max
m.Base = options.Base
m.Scale = options.Scale
m.BitDepth = uint(options.BitDepth)
m.TimeQuantum = pilosa.TimeQuantum(options.TimeQuantum)
m.Keys = options.Keys
}
func decodeNodes(a []*internal.Node, m []*pilosa.Node) {
for i := range a {
m[i] = &pilosa.Node{}
decodeNode(a[i], m[i])
}
}
func decodeClusterStatus(cs *internal.ClusterStatus, m *pilosa.ClusterStatus) {
m.State = cs.State
m.ClusterID = cs.ClusterID
m.Nodes = make([]*pilosa.Node, len(cs.Nodes))
decodeNodes(cs.Nodes, m.Nodes)
}
func decodeNode(node *internal.Node, m *pilosa.Node) {
m.ID = node.ID
decodeURI(node.URI, &m.URI)
m.IsCoordinator = node.IsCoordinator
m.State = node.State
}
func decodeURI(i *internal.URI, m *pilosa.URI) {
m.Scheme = i.Scheme
m.Host = i.Host
m.Port = uint16(i.Port)
}
func decodeCreateShardMessage(pb *internal.CreateShardMessage, m *pilosa.CreateShardMessage) {
m.Index = pb.Index
m.Field = pb.Field
m.Shard = pb.Shard
}
func decodeCreateIndexMessage(pb *internal.CreateIndexMessage, m *pilosa.CreateIndexMessage) {
m.Index = pb.Index
m.Meta = &pilosa.IndexOptions{}
decodeIndexMeta(pb.Meta, m.Meta)
}
func decodeIndexMeta(pb *internal.IndexMeta, m *pilosa.IndexOptions) {
m.Keys = pb.Keys
m.TrackExistence = pb.TrackExistence
}
func decodeDeleteIndexMessage(pb *internal.DeleteIndexMessage, m *pilosa.DeleteIndexMessage) {
m.Index = pb.Index
}
func decodeCreateFieldMessage(pb *internal.CreateFieldMessage, m *pilosa.CreateFieldMessage) {
m.Index = pb.Index
m.Field = pb.Field
m.Meta = &pilosa.FieldOptions{}
decodeFieldOptions(pb.Meta, m.Meta)
}
func decodeDeleteFieldMessage(pb *internal.DeleteFieldMessage, m *pilosa.DeleteFieldMessage) {
m.Index = pb.Index
m.Field = pb.Field
}
func decodeDeleteAvailableShardMessage(pb *internal.DeleteAvailableShardMessage, m *pilosa.DeleteAvailableShardMessage) {
m.Index = pb.Index
m.Field = pb.Field
m.ShardID = pb.ShardID
}
func decodeCreateViewMessage(pb *internal.CreateViewMessage, m *pilosa.CreateViewMessage) {
m.Index = pb.Index
m.Field = pb.Field
m.View = pb.View
}
func decodeDeleteViewMessage(pb *internal.DeleteViewMessage, m *pilosa.DeleteViewMessage) {
m.Index = pb.Index
m.Field = pb.Field
m.View = pb.View
}
func decodeResizeInstructionComplete(pb *internal.ResizeInstructionComplete, m *pilosa.ResizeInstructionComplete) {
m.JobID = pb.JobID
m.Node = &pilosa.Node{}
decodeNode(pb.Node, m.Node)
m.Error = pb.Error
}
func decodeSetCoordinatorMessage(pb *internal.SetCoordinatorMessage, m *pilosa.SetCoordinatorMessage) {
m.New = &pilosa.Node{}
decodeNode(pb.New, m.New)
}
func decodeUpdateCoordinatorMessage(pb *internal.UpdateCoordinatorMessage, m *pilosa.UpdateCoordinatorMessage) {
m.New = &pilosa.Node{}
decodeNode(pb.New, m.New)
}
func decodeNodeStateMessage(pb *internal.NodeStateMessage, m *pilosa.NodeStateMessage) {
m.NodeID = pb.NodeID
m.State = pb.State
}
func decodeNodeEventMessage(pb *internal.NodeEventMessage, m *pilosa.NodeEvent) {
m.Event = pilosa.NodeEventType(pb.Event)
m.Node = &pilosa.Node{}
decodeNode(pb.Node, m.Node)
}
func decodeNodeStatus(pb *internal.NodeStatus, m *pilosa.NodeStatus) {
m.Node = &pilosa.Node{}
m.Indexes = decodeIndexStatuses(pb.Indexes)
m.Schema = &pilosa.Schema{}
decodeSchema(pb.Schema, m.Schema)
}
func decodeIndexStatuses(a []*internal.IndexStatus) []*pilosa.IndexStatus {
m := make([]*pilosa.IndexStatus, 0)
for i := range a {
m = append(m, &pilosa.IndexStatus{})
decodeIndexStatus(a[i], m[i])
}
return m
}
func decodeIndexStatus(pb *internal.IndexStatus, m *pilosa.IndexStatus) {
m.Name = pb.Name
m.Fields = decodeFieldStatuses(pb.Fields)
}
func decodeFieldStatuses(a []*internal.FieldStatus) []*pilosa.FieldStatus {
m := make([]*pilosa.FieldStatus, 0)
for i := range a {
m = append(m, &pilosa.FieldStatus{})
decodeFieldStatus(a[i], m[i])
}
return m
}
func decodeFieldStatus(pb *internal.FieldStatus, m *pilosa.FieldStatus) {
m.Name = pb.Name
m.AvailableShards = roaring.NewBitmap(pb.AvailableShards...)
}
func decodeRecalculateCaches(pb *internal.RecalculateCaches, m *pilosa.RecalculateCaches) {}
func decodeQueryRequest(pb *internal.QueryRequest, m *pilosa.QueryRequest) {
m.Query = pb.Query
m.Shards = pb.Shards
m.ColumnAttrs = pb.ColumnAttrs
m.Remote = pb.Remote
m.ExcludeRowAttrs = pb.ExcludeRowAttrs
m.ExcludeColumns = pb.ExcludeColumns
m.EmbeddedData = make([]*pilosa.Row, len(pb.EmbeddedData))
for i := range pb.EmbeddedData {
m.EmbeddedData[i] = decodeRow(pb.EmbeddedData[i])
}
}
func decodeImportRequest(pb *internal.ImportRequest, m *pilosa.ImportRequest) {
m.Index = pb.Index
m.Field = pb.Field
m.Shard = pb.Shard
m.RowIDs = pb.RowIDs
m.ColumnIDs = pb.ColumnIDs
m.RowKeys = pb.RowKeys
m.ColumnKeys = pb.ColumnKeys
m.Timestamps = pb.Timestamps
}
func decodeImportValueRequest(pb *internal.ImportValueRequest, m *pilosa.ImportValueRequest) {
m.Index = pb.Index
m.Field = pb.Field
m.Shard = pb.Shard
m.ColumnIDs = pb.ColumnIDs
m.ColumnKeys = pb.ColumnKeys
m.Values = pb.Values
m.FloatValues = pb.FloatValues
}
func decodeImportRoaringRequest(pb *internal.ImportRoaringRequest, m *pilosa.ImportRoaringRequest) {
views := map[string][]byte{}
for _, view := range pb.Views {
views[view.Name] = view.Data
}
m.Clear = pb.Clear
m.Views = views
}
func decodeImportResponse(pb *internal.ImportResponse, m *pilosa.ImportResponse) {
m.Err = pb.Err
}
func decodeBlockDataRequest(pb *internal.BlockDataRequest, m *pilosa.BlockDataRequest) {
m.Index = pb.Index
m.Field = pb.Field
m.View = pb.View
m.Shard = pb.Shard
m.Block = pb.Block
}
func decodeBlockDataResponse(pb *internal.BlockDataResponse, m *pilosa.BlockDataResponse) {
m.RowIDs = pb.RowIDs
m.ColumnIDs = pb.ColumnIDs
}
func decodeQueryResponse(pb *internal.QueryResponse, m *pilosa.QueryResponse) {
m.ColumnAttrSets = make([]*pilosa.ColumnAttrSet, len(pb.ColumnAttrSets))
decodeColumnAttrSets(pb.ColumnAttrSets, m.ColumnAttrSets)
if pb.Err == "" {
m.Err = nil
} else {
m.Err = errors.New(pb.Err)
}
m.Results = make([]interface{}, len(pb.Results))
decodeQueryResults(pb.Results, m.Results)
}
func decodeColumnAttrSets(pb []*internal.ColumnAttrSet, m []*pilosa.ColumnAttrSet) {
for i := range pb {
m[i] = &pilosa.ColumnAttrSet{}
decodeColumnAttrSet(pb[i], m[i])
}
}
func decodeColumnAttrSet(pb *internal.ColumnAttrSet, m *pilosa.ColumnAttrSet) {
m.ID = pb.ID
m.Key = pb.Key
m.Attrs = decodeAttrs(pb.Attrs)
}
func decodeQueryResults(pb []*internal.QueryResult, m []interface{}) {
for i := range pb {
m[i] = decodeQueryResult(pb[i])
}
}
func decodeTranslateKeysRequest(pb *internal.TranslateKeysRequest, m *pilosa.TranslateKeysRequest) {
m.Index = pb.Index
m.Field = pb.Field
m.Keys = pb.Keys
}
func decodeTranslateKeysResponse(pb *internal.TranslateKeysResponse, m *pilosa.TranslateKeysResponse) {
m.IDs = pb.IDs
}
// QueryResult types.
const (
queryResultTypeNil uint32 = iota
queryResultTypeRow
queryResultTypePairs
queryResultTypeValCount
queryResultTypeUint64
queryResultTypeBool
queryResultTypeRowIDs
queryResultTypeGroupCounts
queryResultTypeRowIdentifiers
queryResultTypePair
queryResultTypeSignedRow
)
func decodeQueryResult(pb *internal.QueryResult) interface{} {
switch pb.Type {
case queryResultTypeSignedRow:
return decodeSignedRow(pb.SignedRow)
case queryResultTypeRow:
return decodeRow(pb.Row)
case queryResultTypePairs:
return decodePairs(pb.Pairs)
case queryResultTypeValCount:
return decodeValCount(pb.ValCount)
case queryResultTypeUint64:
return pb.N
case queryResultTypeBool:
return pb.Changed
case queryResultTypeNil:
return nil
case queryResultTypeRowIDs:
return pilosa.RowIDs(pb.RowIDs)
case queryResultTypeRowIdentifiers:
return decodeRowIdentifiers(pb.RowIdentifiers)
case queryResultTypeGroupCounts:
return decodeGroupCounts(pb.GroupCounts)
case queryResultTypePair:
return decodePair(pb.Pairs[0])
}
panic(fmt.Sprintf("unknown type: %d", pb.Type))
}
// DecodeRow converts r from its internal representation.
func decodeRow(pr *internal.Row) *pilosa.Row {
if pr == nil {
return pilosa.NewRow()
}
var r *pilosa.Row
if len(pr.Roaring) > 0 {
r = pilosa.NewRowFromRoaring(pr.Roaring)
} else {
r = pilosa.NewRow()
for _, v := range pr.Columns {
r.SetBit(v)
}
}
r.Attrs = decodeAttrs(pr.Attrs)
r.Keys = pr.Keys
return r
}
func decodeSignedRow(pr *internal.SignedRow) pilosa.SignedRow {
if pr == nil {
return pilosa.SignedRow{}
}
r := pilosa.SignedRow{
Pos: decodeRow(pr.Pos),
Neg: decodeRow(pr.Neg),
}
return r
}
func decodeAttrs(pb []*internal.Attr) map[string]interface{} {
m := make(map[string]interface{}, len(pb))
for i := range pb {
key, value := decodeAttr(pb[i])
m[key] = value
}
return m
}
const (
attrTypeString = 1
attrTypeInt = 2
attrTypeBool = 3
attrTypeFloat = 4
)
func decodeAttr(attr *internal.Attr) (key string, value interface{}) {
switch attr.Type {
case attrTypeString:
return attr.Key, attr.StringValue
case attrTypeInt:
return attr.Key, attr.IntValue
case attrTypeBool:
return attr.Key, attr.BoolValue
case attrTypeFloat:
return attr.Key, attr.FloatValue
default:
return attr.Key, nil
}
}
func decodeRowIdentifiers(a *internal.RowIdentifiers) *pilosa.RowIdentifiers {
return &pilosa.RowIdentifiers{
Rows: a.Rows,
Keys: a.Keys,
}
}
func decodeGroupCounts(a []*internal.GroupCount) []pilosa.GroupCount {
other := make([]pilosa.GroupCount, len(a))
for i := range a {
other[i] = pilosa.GroupCount{
Group: decodeFieldRows(a[i].Group),
Count: a[i].Count,
}
}
return other
}
func decodeFieldRows(a []*internal.FieldRow) []pilosa.FieldRow {
other := make([]pilosa.FieldRow, len(a))
for i := range a {
fr := a[i]
other[i].Field = fr.Field
if fr.RowKey == "" {
other[i].RowID = fr.RowID
} else {
other[i].RowKey = fr.RowKey
}
}
return other
}
func decodePairs(a []*internal.Pair) []pilosa.Pair {
other := make([]pilosa.Pair, len(a))
for i := range a {
other[i] = decodePair(a[i])
}
return other
}
func decodePair(pb *internal.Pair) pilosa.Pair {
return pilosa.Pair{
ID: pb.ID,
Key: pb.Key,
Count: pb.Count,
}
}
func decodeValCount(pb *internal.ValCount) pilosa.ValCount {
return pilosa.ValCount{
Val: pb.Val,
Count: pb.Count,
}
}
func encodeColumnAttrSets(a []*pilosa.ColumnAttrSet) []*internal.ColumnAttrSet {
other := make([]*internal.ColumnAttrSet, len(a))
for i := range a {
other[i] = encodeColumnAttrSet(a[i])
}
return other
}
func encodeColumnAttrSet(set *pilosa.ColumnAttrSet) *internal.ColumnAttrSet {
return &internal.ColumnAttrSet{
ID: set.ID,
Key: set.Key,
Attrs: encodeAttrs(set.Attrs),
}
}
func encodeSignedRow(r pilosa.SignedRow) *internal.SignedRow {
ir := &internal.SignedRow{
Pos: encodeRow(r.Pos),
Neg: encodeRow(r.Neg),
}
return ir
}
func encodeRow(r *pilosa.Row) *internal.Row {
if r == nil {
return nil
}
ir := &internal.Row{
Keys: r.Keys,
Attrs: encodeAttrs(r.Attrs),
}
if false {
ir.Columns = r.Columns()
} else {
ir.Roaring = r.Roaring()
}
return ir
}
func encodeRowIdentifiers(r pilosa.RowIdentifiers) *internal.RowIdentifiers {
return &internal.RowIdentifiers{
Rows: r.Rows,
Keys: r.Keys,
//Attrs: encodeAttrs(r.Attrs),
}
}
func encodeGroupCounts(counts []pilosa.GroupCount) []*internal.GroupCount {
result := make([]*internal.GroupCount, len(counts))
for i := range counts {
result[i] = &internal.GroupCount{
Group: encodeFieldRows(counts[i].Group),
Count: counts[i].Count,
}
}
return result
}
func encodeFieldRows(a []pilosa.FieldRow) []*internal.FieldRow {
other := make([]*internal.FieldRow, len(a))
for i := range a {
fr := a[i]
if fr.RowKey == "" {
other[i] = &internal.FieldRow{
Field: fr.Field,
RowID: fr.RowID,
}
} else {
other[i] = &internal.FieldRow{
Field: fr.Field,
RowKey: fr.RowKey,
}
}
}
return other
}
func encodePairs(a pilosa.Pairs) []*internal.Pair {
other := make([]*internal.Pair, len(a))
for i := range a {
other[i] = encodePair(a[i])
}
return other
}
func encodePair(p pilosa.Pair) *internal.Pair {
return &internal.Pair{
ID: p.ID,
Key: p.Key,
Count: p.Count,
}
}
func encodeValCount(vc pilosa.ValCount) *internal.ValCount {
return &internal.ValCount{
Val: vc.Val,
Count: vc.Count,
}
}
func encodeAttrs(m map[string]interface{}) []*internal.Attr {
keys := make([]string, 0, len(m))
for k := range m {
keys = append(keys, k)
}
sort.Strings(keys)
a := make([]*internal.Attr, len(keys))
for i := range keys {
a[i] = encodeAttr(keys[i], m[keys[i]])
}
return a
}
// encodeAttr converts a key/value pair into an Attr internal representation.
func encodeAttr(key string, value interface{}) *internal.Attr {
pb := &internal.Attr{Key: key}
switch value := value.(type) {
case string:
pb.Type = attrTypeString
pb.StringValue = value
case float64:
pb.Type = attrTypeFloat
pb.FloatValue = value
case uint64:
pb.Type = attrTypeInt
pb.IntValue = int64(value)
case int64:
pb.Type = attrTypeInt
pb.IntValue = value
case bool:
pb.Type = attrTypeBool
pb.BoolValue = value
}
return pb
}