From 79940cf0777b204005c3defd1e5afbc61035ef52 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 3 Apr 2020 13:08:05 -0500 Subject: [PATCH] encoding/proto: allow distinct serializers We want to be able to control whether or not we use roaring to serialize Rows, which means serializers have to be able to be distinct. We also make corresponding changes to http/handler.go to have it use the exported serializers directly rather than the API's serializer (which is always the base protobuf serializer right now, and if it weren't, that would be bad because we were assuming it was). When we're accepting protobuf from a pilosa server, flag that we'll accept roaring bitmaps as opposed to the naive column representation. --- encoding/proto/proto.go | 650 ++++++++++++++++++++-------------------- http/client.go | 7 + http/handler.go | 38 ++- 3 files changed, 362 insertions(+), 333 deletions(-) diff --git a/encoding/proto/proto.go b/encoding/proto/proto.go index 0fef659e0..1d46327b4 100644 --- a/encoding/proto/proto.go +++ b/encoding/proto/proto.go @@ -28,11 +28,16 @@ import ( ) // Serializer implements pilosa.Serializer for protobufs. -type Serializer struct{} +type Serializer struct { + RoaringRows bool +} + +var DefaultSerializer = Serializer{} +var RoaringSerializer = Serializer{RoaringRows: true} // Marshal turns pilosa messages into protobuf serialized bytes. -func (Serializer) Marshal(m pilosa.Message) ([]byte, error) { - pm := encodeToProto(m) +func (s Serializer) Marshal(m pilosa.Message) ([]byte, error) { + pm := s.encodeToProto(m) if pm == nil { return nil, errors.New("passed invalid pilosa.Message") } @@ -41,7 +46,7 @@ func (Serializer) Marshal(m pilosa.Message) ([]byte, error) { } // Unmarshal takes byte slices and protobuf deserializes them into a pilosa Message. -func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { +func (s Serializer) Unmarshal(buf []byte, m pilosa.Message) error { switch mt := m.(type) { case *pilosa.CreateShardMessage: msg := &internal.CreateShardMessage{} @@ -49,7 +54,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling CreateShardMessage") } - decodeCreateShardMessage(msg, mt) + s.decodeCreateShardMessage(msg, mt) return nil case *pilosa.CreateIndexMessage: msg := &internal.CreateIndexMessage{} @@ -57,7 +62,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling CreateIndexMessage") } - decodeCreateIndexMessage(msg, mt) + s.decodeCreateIndexMessage(msg, mt) return nil case *pilosa.DeleteIndexMessage: msg := &internal.DeleteIndexMessage{} @@ -65,7 +70,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling DeleteIndexMessage") } - decodeDeleteIndexMessage(msg, mt) + s.decodeDeleteIndexMessage(msg, mt) return nil case *pilosa.CreateFieldMessage: msg := &internal.CreateFieldMessage{} @@ -73,7 +78,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling CreateFieldMessage") } - decodeCreateFieldMessage(msg, mt) + s.decodeCreateFieldMessage(msg, mt) return nil case *pilosa.DeleteFieldMessage: msg := &internal.DeleteFieldMessage{} @@ -81,7 +86,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling DeleteFieldMessage") } - decodeDeleteFieldMessage(msg, mt) + s.decodeDeleteFieldMessage(msg, mt) return nil case *pilosa.DeleteAvailableShardMessage: msg := &internal.DeleteAvailableShardMessage{} @@ -89,7 +94,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling DeleteAvailableShardMessage") } - decodeDeleteAvailableShardMessage(msg, mt) + s.decodeDeleteAvailableShardMessage(msg, mt) return nil case *pilosa.CreateViewMessage: msg := &internal.CreateViewMessage{} @@ -97,7 +102,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling CreateViewMessage") } - decodeCreateViewMessage(msg, mt) + s.decodeCreateViewMessage(msg, mt) return nil case *pilosa.DeleteViewMessage: msg := &internal.DeleteViewMessage{} @@ -105,7 +110,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling DeleteViewMessage") } - decodeDeleteViewMessage(msg, mt) + s.decodeDeleteViewMessage(msg, mt) return nil case *pilosa.ClusterStatus: msg := &internal.ClusterStatus{} @@ -113,7 +118,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ClusterStatus") } - decodeClusterStatus(msg, mt) + s.decodeClusterStatus(msg, mt) return nil case *pilosa.ResizeInstruction: msg := &internal.ResizeInstruction{} @@ -121,7 +126,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ResizeInstruction") } - decodeResizeInstruction(msg, mt) + s.decodeResizeInstruction(msg, mt) return nil case *pilosa.ResizeInstructionComplete: msg := &internal.ResizeInstructionComplete{} @@ -129,7 +134,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ResizeInstructionComplete") } - decodeResizeInstructionComplete(msg, mt) + s.decodeResizeInstructionComplete(msg, mt) return nil case *pilosa.SetCoordinatorMessage: msg := &internal.SetCoordinatorMessage{} @@ -137,7 +142,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling SetCoordinatorMessage") } - decodeSetCoordinatorMessage(msg, mt) + s.decodeSetCoordinatorMessage(msg, mt) return nil case *pilosa.UpdateCoordinatorMessage: msg := &internal.UpdateCoordinatorMessage{} @@ -145,7 +150,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling UpdateCoordinatorMessage") } - decodeUpdateCoordinatorMessage(msg, mt) + s.decodeUpdateCoordinatorMessage(msg, mt) return nil case *pilosa.NodeStateMessage: msg := &internal.NodeStateMessage{} @@ -153,7 +158,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling NodeStateMessage") } - decodeNodeStateMessage(msg, mt) + s.decodeNodeStateMessage(msg, mt) return nil case *pilosa.RecalculateCaches: msg := &internal.RecalculateCaches{} @@ -161,7 +166,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling RecalculateCaches") } - decodeRecalculateCaches(msg, mt) + s.decodeRecalculateCaches(msg, mt) return nil case *pilosa.NodeEvent: msg := &internal.NodeEventMessage{} @@ -169,7 +174,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling NodeEvent") } - decodeNodeEventMessage(msg, mt) + s.decodeNodeEventMessage(msg, mt) return nil case *pilosa.NodeStatus: msg := &internal.NodeStatus{} @@ -177,7 +182,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling NodeStatus") } - decodeNodeStatus(msg, mt) + s.decodeNodeStatus(msg, mt) return nil case *pilosa.Node: msg := &internal.Node{} @@ -185,7 +190,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling Node") } - decodeNode(msg, mt) + s.decodeNode(msg, mt) return nil case *pilosa.QueryRequest: msg := &internal.QueryRequest{} @@ -193,7 +198,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling QueryRequest") } - decodeQueryRequest(msg, mt) + s.decodeQueryRequest(msg, mt) return nil case *pilosa.QueryResponse: msg := &internal.QueryResponse{} @@ -201,7 +206,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling QueryResponse") } - decodeQueryResponse(msg, mt) + s.decodeQueryResponse(msg, mt) return nil case *pilosa.ImportRequest: msg := &internal.ImportRequest{} @@ -209,7 +214,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ImportRequest") } - decodeImportRequest(msg, mt) + s.decodeImportRequest(msg, mt) return nil case *pilosa.ImportValueRequest: msg := &internal.ImportValueRequest{} @@ -217,7 +222,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ImportValueRequest") } - decodeImportValueRequest(msg, mt) + s.decodeImportValueRequest(msg, mt) return nil case *pilosa.ImportRoaringRequest: msg := &internal.ImportRoaringRequest{} @@ -225,7 +230,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ImportRoaringRequest") } - decodeImportRoaringRequest(msg, mt) + s.decodeImportRoaringRequest(msg, mt) return nil case *pilosa.ImportColumnAttrsRequest: msg := &internal.ImportColumnAttrsRequest{} @@ -233,7 +238,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ImportColumnAttrsRequest") } - decodeImportColumnAttrsRequest(msg, mt) + s.decodeImportColumnAttrsRequest(msg, mt) return nil case *pilosa.ImportResponse: msg := &internal.ImportResponse{} @@ -241,7 +246,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling ImportResponse") } - decodeImportResponse(msg, mt) + s.decodeImportResponse(msg, mt) return nil case *pilosa.BlockDataRequest: msg := &internal.BlockDataRequest{} @@ -249,7 +254,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling BlockDataRequest") } - decodeBlockDataRequest(msg, mt) + s.decodeBlockDataRequest(msg, mt) return nil case *pilosa.BlockDataResponse: msg := &internal.BlockDataResponse{} @@ -257,7 +262,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling BlockDataResponse") } - decodeBlockDataResponse(msg, mt) + s.decodeBlockDataResponse(msg, mt) return nil case *pilosa.TranslateKeysRequest: msg := &internal.TranslateKeysRequest{} @@ -265,7 +270,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling TranslateKeysRequest") } - decodeTranslateKeysRequest(msg, mt) + s.decodeTranslateKeysRequest(msg, mt) return nil case *pilosa.TranslateKeysResponse: msg := &internal.TranslateKeysResponse{} @@ -273,7 +278,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling TranslateKeysResponse") } - decodeTranslateKeysResponse(msg, mt) + s.decodeTranslateKeysResponse(msg, mt) return nil case *pilosa.TranslateIDsRequest: msg := &internal.TranslateIDsRequest{} @@ -281,7 +286,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling TranslateIDsRequest") } - decodeTranslateIDsRequest(msg, mt) + s.decodeTranslateIDsRequest(msg, mt) return nil case *pilosa.TranslateIDsResponse: msg := &internal.TranslateIDsResponse{} @@ -289,7 +294,7 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { if err != nil { return errors.Wrap(err, "unmarshaling TranslateIDsResponse") } - decodeTranslateIDsResponse(msg, mt) + s.decodeTranslateIDsResponse(msg, mt) return nil case *pilosa.TransactionMessage: msg := &internal.TransactionMessage{} @@ -304,77 +309,77 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { } } -func encodeToProto(m pilosa.Message) proto.Message { +func (s Serializer) encodeToProto(m pilosa.Message) proto.Message { switch mt := m.(type) { case *pilosa.CreateShardMessage: - return encodeCreateShardMessage(mt) + return s.encodeCreateShardMessage(mt) case *pilosa.CreateIndexMessage: - return encodeCreateIndexMessage(mt) + return s.encodeCreateIndexMessage(mt) case *pilosa.DeleteIndexMessage: - return encodeDeleteIndexMessage(mt) + return s.encodeDeleteIndexMessage(mt) case *pilosa.CreateFieldMessage: - return encodeCreateFieldMessage(mt) + return s.encodeCreateFieldMessage(mt) case *pilosa.DeleteFieldMessage: - return encodeDeleteFieldMessage(mt) + return s.encodeDeleteFieldMessage(mt) case *pilosa.DeleteAvailableShardMessage: - return encodeDeleteAvailableShardMessage(mt) + return s.encodeDeleteAvailableShardMessage(mt) case *pilosa.CreateViewMessage: - return encodeCreateViewMessage(mt) + return s.encodeCreateViewMessage(mt) case *pilosa.DeleteViewMessage: - return encodeDeleteViewMessage(mt) + return s.encodeDeleteViewMessage(mt) case *pilosa.ClusterStatus: - return encodeClusterStatus(mt) + return s.encodeClusterStatus(mt) case *pilosa.ResizeInstruction: - return encodeResizeInstruction(mt) + return s.encodeResizeInstruction(mt) case *pilosa.ResizeInstructionComplete: - return encodeResizeInstructionComplete(mt) + return s.encodeResizeInstructionComplete(mt) case *pilosa.SetCoordinatorMessage: - return encodeSetCoordinatorMessage(mt) + return s.encodeSetCoordinatorMessage(mt) case *pilosa.UpdateCoordinatorMessage: - return encodeUpdateCoordinatorMessage(mt) + return s.encodeUpdateCoordinatorMessage(mt) case *pilosa.NodeStateMessage: - return encodeNodeStateMessage(mt) + return s.encodeNodeStateMessage(mt) case *pilosa.RecalculateCaches: - return encodeRecalculateCaches(mt) + return s.encodeRecalculateCaches(mt) case *pilosa.NodeEvent: - return encodeNodeEventMessage(mt) + return s.encodeNodeEventMessage(mt) case *pilosa.NodeStatus: - return encodeNodeStatus(mt) + return s.encodeNodeStatus(mt) case *pilosa.Node: - return encodeNode(mt) + return s.encodeNode(mt) case *pilosa.QueryRequest: - return encodeQueryRequest(mt) + return s.encodeQueryRequest(mt) case *pilosa.QueryResponse: - return encodeQueryResponse(mt) + return s.encodeQueryResponse(mt) case *pilosa.ImportRequest: - return encodeImportRequest(mt) + return s.encodeImportRequest(mt) case *pilosa.ImportValueRequest: - return encodeImportValueRequest(mt) + return s.encodeImportValueRequest(mt) case *pilosa.ImportRoaringRequest: - return encodeImportRoaringRequest(mt) + return s.encodeImportRoaringRequest(mt) case *pilosa.ImportColumnAttrsRequest: - return encodeImportColumnAttrsRequest(mt) + return s.encodeImportColumnAttrsRequest(mt) case *pilosa.ImportResponse: - return encodeImportResponse(mt) + return s.encodeImportResponse(mt) case *pilosa.BlockDataRequest: - return encodeBlockDataRequest(mt) + return s.encodeBlockDataRequest(mt) case *pilosa.BlockDataResponse: - return encodeBlockDataResponse(mt) + return s.encodeBlockDataResponse(mt) case *pilosa.TranslateKeysRequest: - return encodeTranslateKeysRequest(mt) + return s.encodeTranslateKeysRequest(mt) case *pilosa.TranslateKeysResponse: - return encodeTranslateKeysResponse(mt) + return s.encodeTranslateKeysResponse(mt) case *pilosa.TranslateIDsRequest: - return encodeTranslateIDsRequest(mt) + return s.encodeTranslateIDsRequest(mt) case *pilosa.TranslateIDsResponse: - return encodeTranslateIDsResponse(mt) + return s.encodeTranslateIDsResponse(mt) case *pilosa.TransactionMessage: - return encodeTransactionMessage(mt) + return s.encodeTransactionMessage(mt) } return nil } -func encodeBlockDataRequest(m *pilosa.BlockDataRequest) *internal.BlockDataRequest { +func (s Serializer) encodeBlockDataRequest(m *pilosa.BlockDataRequest) *internal.BlockDataRequest { return &internal.BlockDataRequest{ Index: m.Index, Field: m.Field, @@ -383,20 +388,20 @@ func encodeBlockDataRequest(m *pilosa.BlockDataRequest) *internal.BlockDataReque Block: m.Block, } } -func encodeBlockDataResponse(m *pilosa.BlockDataResponse) *internal.BlockDataResponse { +func (s Serializer) encodeBlockDataResponse(m *pilosa.BlockDataResponse) *internal.BlockDataResponse { return &internal.BlockDataResponse{ RowIDs: m.RowIDs, ColumnIDs: m.ColumnIDs, } } -func encodeImportResponse(m *pilosa.ImportResponse) *internal.ImportResponse { +func (s Serializer) encodeImportResponse(m *pilosa.ImportResponse) *internal.ImportResponse { return &internal.ImportResponse{ Err: m.Err, } } -func encodeImportRequest(m *pilosa.ImportRequest) *internal.ImportRequest { +func (s Serializer) encodeImportRequest(m *pilosa.ImportRequest) *internal.ImportRequest { return &internal.ImportRequest{ Index: m.Index, Field: m.Field, @@ -409,7 +414,7 @@ func encodeImportRequest(m *pilosa.ImportRequest) *internal.ImportRequest { } } -func encodeImportValueRequest(m *pilosa.ImportValueRequest) *internal.ImportValueRequest { +func (s Serializer) encodeImportValueRequest(m *pilosa.ImportValueRequest) *internal.ImportValueRequest { return &internal.ImportValueRequest{ Index: m.Index, Field: m.Field, @@ -422,7 +427,7 @@ func encodeImportValueRequest(m *pilosa.ImportValueRequest) *internal.ImportValu } } -func encodeImportRoaringRequest(m *pilosa.ImportRoaringRequest) *internal.ImportRoaringRequest { +func (s Serializer) encodeImportRoaringRequest(m *pilosa.ImportRoaringRequest) *internal.ImportRoaringRequest { views := make([]*internal.ImportRoaringRequestView, len(m.Views)) i := 0 for viewName, viewData := range m.Views { @@ -440,7 +445,7 @@ func encodeImportRoaringRequest(m *pilosa.ImportRoaringRequest) *internal.Import } } -func encodeImportColumnAttrsRequest(m *pilosa.ImportColumnAttrsRequest) *internal.ImportColumnAttrsRequest { +func (s Serializer) encodeImportColumnAttrsRequest(m *pilosa.ImportColumnAttrsRequest) *internal.ImportColumnAttrsRequest { return &internal.ImportColumnAttrsRequest{ Index: m.Index, Shard: m.Shard, @@ -450,7 +455,7 @@ func encodeImportColumnAttrsRequest(m *pilosa.ImportColumnAttrsRequest) *interna } } -func encodeQueryRequest(m *pilosa.QueryRequest) *internal.QueryRequest { +func (s Serializer) encodeQueryRequest(m *pilosa.QueryRequest) *internal.QueryRequest { r := &internal.QueryRequest{ Query: m.Query, Shards: m.Shards, @@ -461,15 +466,15 @@ func encodeQueryRequest(m *pilosa.QueryRequest) *internal.QueryRequest { EmbeddedData: make([]*internal.Row, len(m.EmbeddedData)), } for i := range m.EmbeddedData { - r.EmbeddedData[i] = encodeRow(m.EmbeddedData[i]) + r.EmbeddedData[i] = s.encodeRow(m.EmbeddedData[i]) } return r } -func encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse { +func (s Serializer) encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse { pb := &internal.QueryResponse{ Results: make([]*internal.QueryResult, len(m.Results)), - ColumnAttrSets: encodeColumnAttrSets(m.ColumnAttrSets), + ColumnAttrSets: s.encodeColumnAttrSets(m.ColumnAttrSets), } for i := range m.Results { @@ -478,19 +483,19 @@ func encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse { switch result := m.Results[i].(type) { case pilosa.SignedRow: pb.Results[i].Type = queryResultTypeSignedRow - pb.Results[i].SignedRow = encodeSignedRow(result) + pb.Results[i].SignedRow = s.encodeSignedRow(result) case *pilosa.Row: pb.Results[i].Type = queryResultTypeRow - pb.Results[i].Row = encodeRow(result) + pb.Results[i].Row = s.encodeRow(result) case []pilosa.Pair: pb.Results[i].Type = queryResultTypePairs - pb.Results[i].Pairs = encodePairs(result) + pb.Results[i].Pairs = s.encodePairs(result) case *pilosa.PairsField: pb.Results[i].Type = queryResultTypePairsField - pb.Results[i].PairsField = encodePairsField(result) + pb.Results[i].PairsField = s.encodePairsField(result) case pilosa.ValCount: pb.Results[i].Type = queryResultTypeValCount - pb.Results[i].ValCount = encodeValCount(result) + pb.Results[i].ValCount = s.encodeValCount(result) case uint64: pb.Results[i].Type = queryResultTypeUint64 pb.Results[i].N = result @@ -502,16 +507,16 @@ func encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse { pb.Results[i].RowIDs = result case []pilosa.GroupCount: pb.Results[i].Type = queryResultTypeGroupCounts - pb.Results[i].GroupCounts = encodeGroupCounts(result) + pb.Results[i].GroupCounts = s.encodeGroupCounts(result) case pilosa.RowIdentifiers: pb.Results[i].Type = queryResultTypeRowIdentifiers - pb.Results[i].RowIdentifiers = encodeRowIdentifiers(result) + pb.Results[i].RowIdentifiers = s.encodeRowIdentifiers(result) case pilosa.Pair: pb.Results[i].Type = queryResultTypePair - pb.Results[i].Pairs = []*internal.Pair{encodePair(result)} + pb.Results[i].Pairs = []*internal.Pair{s.encodePair(result)} case pilosa.PairField: pb.Results[i].Type = queryResultTypePairField - pb.Results[i].Pairs = []*internal.Pair{encodePairField(result)} + pb.Results[i].Pairs = []*internal.Pair{s.encodePairField(result)} case nil: pb.Results[i].Type = queryResultTypeNil default: @@ -526,29 +531,29 @@ func encodeQueryResponse(m *pilosa.QueryResponse) *internal.QueryResponse { return pb } -func encodeResizeInstruction(m *pilosa.ResizeInstruction) *internal.ResizeInstruction { +func (s Serializer) encodeResizeInstruction(m *pilosa.ResizeInstruction) *internal.ResizeInstruction { return &internal.ResizeInstruction{ JobID: m.JobID, - Node: encodeNode(m.Node), - Coordinator: encodeNode(m.Coordinator), - Sources: encodeResizeSources(m.Sources), - TranslationSources: encodeTranslationResizeSources(m.TranslationSources), - NodeStatus: encodeNodeStatus(m.NodeStatus), - ClusterStatus: encodeClusterStatus(m.ClusterStatus), + Node: s.encodeNode(m.Node), + Coordinator: s.encodeNode(m.Coordinator), + Sources: s.encodeResizeSources(m.Sources), + TranslationSources: s.encodeTranslationResizeSources(m.TranslationSources), + NodeStatus: s.encodeNodeStatus(m.NodeStatus), + ClusterStatus: s.encodeClusterStatus(m.ClusterStatus), } } -func encodeResizeSources(srcs []*pilosa.ResizeSource) []*internal.ResizeSource { +func (s Serializer) encodeResizeSources(srcs []*pilosa.ResizeSource) []*internal.ResizeSource { new := make([]*internal.ResizeSource, 0, len(srcs)) for _, src := range srcs { - new = append(new, encodeResizeSource(src)) + new = append(new, s.encodeResizeSource(src)) } return new } -func encodeResizeSource(m *pilosa.ResizeSource) *internal.ResizeSource { +func (s Serializer) encodeResizeSource(m *pilosa.ResizeSource) *internal.ResizeSource { return &internal.ResizeSource{ - Node: encodeNode(m.Node), + Node: s.encodeNode(m.Node), Index: m.Index, Field: m.Field, View: m.View, @@ -556,56 +561,56 @@ func encodeResizeSource(m *pilosa.ResizeSource) *internal.ResizeSource { } } -func encodeTranslationResizeSources(srcs []*pilosa.TranslationResizeSource) []*internal.TranslationResizeSource { +func (s Serializer) encodeTranslationResizeSources(srcs []*pilosa.TranslationResizeSource) []*internal.TranslationResizeSource { new := make([]*internal.TranslationResizeSource, 0, len(srcs)) for _, src := range srcs { - new = append(new, encodeTranslationResizeSource(src)) + new = append(new, s.encodeTranslationResizeSource(src)) } return new } -func encodeTranslationResizeSource(m *pilosa.TranslationResizeSource) *internal.TranslationResizeSource { +func (s Serializer) encodeTranslationResizeSource(m *pilosa.TranslationResizeSource) *internal.TranslationResizeSource { return &internal.TranslationResizeSource{ - Node: encodeNode(m.Node), + Node: s.encodeNode(m.Node), Index: m.Index, PartitionID: int32(m.PartitionID), } } -func encodeSchema(m *pilosa.Schema) *internal.Schema { +func (s Serializer) encodeSchema(m *pilosa.Schema) *internal.Schema { return &internal.Schema{ - Indexes: encodeIndexInfos(m.Indexes), + Indexes: s.encodeIndexInfos(m.Indexes), } } -func encodeIndexInfos(idxs []*pilosa.IndexInfo) []*internal.Index { +func (s Serializer) encodeIndexInfos(idxs []*pilosa.IndexInfo) []*internal.Index { new := make([]*internal.Index, 0, len(idxs)) for _, idx := range idxs { - new = append(new, encodeIndexInfo(idx)) + new = append(new, s.encodeIndexInfo(idx)) } return new } -func encodeIndexInfo(idx *pilosa.IndexInfo) *internal.Index { +func (s Serializer) encodeIndexInfo(idx *pilosa.IndexInfo) *internal.Index { return &internal.Index{ Name: idx.Name, - Options: encodeIndexMeta(&idx.Options), - Fields: encodeFieldInfos(idx.Fields), + Options: s.encodeIndexMeta(&idx.Options), + Fields: s.encodeFieldInfos(idx.Fields), } } -func encodeFieldInfos(fs []*pilosa.FieldInfo) []*internal.Field { +func (s Serializer) encodeFieldInfos(fs []*pilosa.FieldInfo) []*internal.Field { new := make([]*internal.Field, 0, len(fs)) for _, f := range fs { - new = append(new, encodeFieldInfo(f)) + new = append(new, s.encodeFieldInfo(f)) } return new } -func encodeFieldInfo(f *pilosa.FieldInfo) *internal.Field { +func (s Serializer) encodeFieldInfo(f *pilosa.FieldInfo) *internal.Field { ifield := &internal.Field{ Name: f.Name, - Meta: encodeFieldOptions(&f.Options), + Meta: s.encodeFieldOptions(&f.Options), Views: make([]string, 0, len(f.Views)), } @@ -615,7 +620,7 @@ func encodeFieldInfo(f *pilosa.FieldInfo) *internal.Field { return ifield } -func encodeFieldOptions(o *pilosa.FieldOptions) *internal.FieldOptions { +func (s Serializer) encodeFieldOptions(o *pilosa.FieldOptions) *internal.FieldOptions { if o == nil { return nil } @@ -634,26 +639,26 @@ func encodeFieldOptions(o *pilosa.FieldOptions) *internal.FieldOptions { } } -// encodeNodes converts a slice of Nodes into its internal representation. -func encodeNodes(a []*pilosa.Node) []*internal.Node { +// s.encodeNodes converts a slice of Nodes into its internal representation. +func (s Serializer) encodeNodes(a []*pilosa.Node) []*internal.Node { other := make([]*internal.Node, len(a)) for i := range a { - other[i] = encodeNode(a[i]) + other[i] = s.encodeNode(a[i]) } return other } -// encodeNode converts a Node into its internal representation. -func encodeNode(n *pilosa.Node) *internal.Node { +// s.encodeNode converts a Node into its internal representation. +func (s Serializer) encodeNode(n *pilosa.Node) *internal.Node { return &internal.Node{ ID: n.ID, - URI: encodeURI(n.URI), + URI: s.encodeURI(n.URI), IsCoordinator: n.IsCoordinator, State: n.State, } } -func encodeURI(u pilosa.URI) *internal.URI { +func (s Serializer) encodeURI(u pilosa.URI) *internal.URI { return &internal.URI{ Scheme: u.Scheme, Host: u.Host, @@ -661,15 +666,15 @@ func encodeURI(u pilosa.URI) *internal.URI { } } -func encodeClusterStatus(m *pilosa.ClusterStatus) *internal.ClusterStatus { +func (s Serializer) encodeClusterStatus(m *pilosa.ClusterStatus) *internal.ClusterStatus { return &internal.ClusterStatus{ State: m.State, ClusterID: m.ClusterID, - Nodes: encodeNodes(m.Nodes), + Nodes: s.encodeNodes(m.Nodes), } } -func encodeCreateShardMessage(m *pilosa.CreateShardMessage) *internal.CreateShardMessage { +func (s Serializer) encodeCreateShardMessage(m *pilosa.CreateShardMessage) *internal.CreateShardMessage { return &internal.CreateShardMessage{ Index: m.Index, Field: m.Field, @@ -677,42 +682,42 @@ func encodeCreateShardMessage(m *pilosa.CreateShardMessage) *internal.CreateShar } } -func encodeCreateIndexMessage(m *pilosa.CreateIndexMessage) *internal.CreateIndexMessage { +func (s Serializer) encodeCreateIndexMessage(m *pilosa.CreateIndexMessage) *internal.CreateIndexMessage { return &internal.CreateIndexMessage{ Index: m.Index, - Meta: encodeIndexMeta(m.Meta), + Meta: s.encodeIndexMeta(m.Meta), } } -func encodeIndexMeta(m *pilosa.IndexOptions) *internal.IndexMeta { +func (s Serializer) encodeIndexMeta(m *pilosa.IndexOptions) *internal.IndexMeta { return &internal.IndexMeta{ Keys: m.Keys, TrackExistence: m.TrackExistence, } } -func encodeDeleteIndexMessage(m *pilosa.DeleteIndexMessage) *internal.DeleteIndexMessage { +func (s Serializer) encodeDeleteIndexMessage(m *pilosa.DeleteIndexMessage) *internal.DeleteIndexMessage { return &internal.DeleteIndexMessage{ Index: m.Index, } } -func encodeCreateFieldMessage(m *pilosa.CreateFieldMessage) *internal.CreateFieldMessage { +func (s Serializer) encodeCreateFieldMessage(m *pilosa.CreateFieldMessage) *internal.CreateFieldMessage { return &internal.CreateFieldMessage{ Index: m.Index, Field: m.Field, - Meta: encodeFieldOptions(m.Meta), + Meta: s.encodeFieldOptions(m.Meta), } } -func encodeDeleteFieldMessage(m *pilosa.DeleteFieldMessage) *internal.DeleteFieldMessage { +func (s Serializer) encodeDeleteFieldMessage(m *pilosa.DeleteFieldMessage) *internal.DeleteFieldMessage { return &internal.DeleteFieldMessage{ Index: m.Index, Field: m.Field, } } -func encodeDeleteAvailableShardMessage(m *pilosa.DeleteAvailableShardMessage) *internal.DeleteAvailableShardMessage { +func (s Serializer) encodeDeleteAvailableShardMessage(m *pilosa.DeleteAvailableShardMessage) *internal.DeleteAvailableShardMessage { return &internal.DeleteAvailableShardMessage{ Index: m.Index, Field: m.Field, @@ -720,7 +725,7 @@ func encodeDeleteAvailableShardMessage(m *pilosa.DeleteAvailableShardMessage) *i } } -func encodeCreateViewMessage(m *pilosa.CreateViewMessage) *internal.CreateViewMessage { +func (s Serializer) encodeCreateViewMessage(m *pilosa.CreateViewMessage) *internal.CreateViewMessage { return &internal.CreateViewMessage{ Index: m.Index, Field: m.Field, @@ -728,7 +733,7 @@ func encodeCreateViewMessage(m *pilosa.CreateViewMessage) *internal.CreateViewMe } } -func encodeDeleteViewMessage(m *pilosa.DeleteViewMessage) *internal.DeleteViewMessage { +func (s Serializer) encodeDeleteViewMessage(m *pilosa.DeleteViewMessage) *internal.DeleteViewMessage { return &internal.DeleteViewMessage{ Index: m.Index, Field: m.Field, @@ -736,83 +741,83 @@ func encodeDeleteViewMessage(m *pilosa.DeleteViewMessage) *internal.DeleteViewMe } } -func encodeResizeInstructionComplete(m *pilosa.ResizeInstructionComplete) *internal.ResizeInstructionComplete { +func (s Serializer) encodeResizeInstructionComplete(m *pilosa.ResizeInstructionComplete) *internal.ResizeInstructionComplete { return &internal.ResizeInstructionComplete{ JobID: m.JobID, - Node: encodeNode(m.Node), + Node: s.encodeNode(m.Node), Error: m.Error, } } -func encodeSetCoordinatorMessage(m *pilosa.SetCoordinatorMessage) *internal.SetCoordinatorMessage { +func (s Serializer) encodeSetCoordinatorMessage(m *pilosa.SetCoordinatorMessage) *internal.SetCoordinatorMessage { return &internal.SetCoordinatorMessage{ - New: encodeNode(m.New), + New: s.encodeNode(m.New), } } -func encodeUpdateCoordinatorMessage(m *pilosa.UpdateCoordinatorMessage) *internal.UpdateCoordinatorMessage { +func (s Serializer) encodeUpdateCoordinatorMessage(m *pilosa.UpdateCoordinatorMessage) *internal.UpdateCoordinatorMessage { return &internal.UpdateCoordinatorMessage{ - New: encodeNode(m.New), + New: s.encodeNode(m.New), } } -func encodeNodeStateMessage(m *pilosa.NodeStateMessage) *internal.NodeStateMessage { +func (s Serializer) encodeNodeStateMessage(m *pilosa.NodeStateMessage) *internal.NodeStateMessage { return &internal.NodeStateMessage{ NodeID: m.NodeID, State: m.State, } } -func encodeNodeEventMessage(m *pilosa.NodeEvent) *internal.NodeEventMessage { +func (s Serializer) encodeNodeEventMessage(m *pilosa.NodeEvent) *internal.NodeEventMessage { return &internal.NodeEventMessage{ Event: uint32(m.Event), - Node: encodeNode(m.Node), + Node: s.encodeNode(m.Node), } } -func encodeNodeStatus(m *pilosa.NodeStatus) *internal.NodeStatus { +func (s Serializer) encodeNodeStatus(m *pilosa.NodeStatus) *internal.NodeStatus { return &internal.NodeStatus{ - Node: encodeNode(m.Node), - Indexes: encodeIndexStatuses(m.Indexes), - Schema: encodeSchema(m.Schema), + Node: s.encodeNode(m.Node), + Indexes: s.encodeIndexStatuses(m.Indexes), + Schema: s.encodeSchema(m.Schema), } } -func encodeIndexStatus(m *pilosa.IndexStatus) *internal.IndexStatus { +func (s Serializer) encodeIndexStatus(m *pilosa.IndexStatus) *internal.IndexStatus { return &internal.IndexStatus{ Name: m.Name, - Fields: encodeFieldStatuses(m.Fields), + Fields: s.encodeFieldStatuses(m.Fields), } } -func encodeIndexStatuses(a []*pilosa.IndexStatus) []*internal.IndexStatus { +func (s Serializer) encodeIndexStatuses(a []*pilosa.IndexStatus) []*internal.IndexStatus { other := make([]*internal.IndexStatus, len(a)) for i := range a { - other[i] = encodeIndexStatus(a[i]) + other[i] = s.encodeIndexStatus(a[i]) } return other } -func encodeFieldStatus(m *pilosa.FieldStatus) *internal.FieldStatus { +func (s Serializer) encodeFieldStatus(m *pilosa.FieldStatus) *internal.FieldStatus { return &internal.FieldStatus{ Name: m.Name, AvailableShards: m.AvailableShards.Slice(), } } -func encodeFieldStatuses(a []*pilosa.FieldStatus) []*internal.FieldStatus { +func (s Serializer) encodeFieldStatuses(a []*pilosa.FieldStatus) []*internal.FieldStatus { other := make([]*internal.FieldStatus, len(a)) for i := range a { - other[i] = encodeFieldStatus(a[i]) + other[i] = s.encodeFieldStatus(a[i]) } return other } -func encodeRecalculateCaches(*pilosa.RecalculateCaches) *internal.RecalculateCaches { +func (s Serializer) encodeRecalculateCaches(*pilosa.RecalculateCaches) *internal.RecalculateCaches { return &internal.RecalculateCaches{} } -func encodeTranslateKeysRequest(request *pilosa.TranslateKeysRequest) *internal.TranslateKeysRequest { +func (s Serializer) encodeTranslateKeysRequest(request *pilosa.TranslateKeysRequest) *internal.TranslateKeysRequest { return &internal.TranslateKeysRequest{ Index: request.Index, Field: request.Field, @@ -820,13 +825,13 @@ func encodeTranslateKeysRequest(request *pilosa.TranslateKeysRequest) *internal. } } -func encodeTranslateKeysResponse(response *pilosa.TranslateKeysResponse) *internal.TranslateKeysResponse { +func (s Serializer) encodeTranslateKeysResponse(response *pilosa.TranslateKeysResponse) *internal.TranslateKeysResponse { return &internal.TranslateKeysResponse{ IDs: response.IDs, } } -func encodeTranslateIDsRequest(request *pilosa.TranslateIDsRequest) *internal.TranslateIDsRequest { +func (s Serializer) encodeTranslateIDsRequest(request *pilosa.TranslateIDsRequest) *internal.TranslateIDsRequest { return &internal.TranslateIDsRequest{ Index: request.Index, Field: request.Field, @@ -834,20 +839,20 @@ func encodeTranslateIDsRequest(request *pilosa.TranslateIDsRequest) *internal.Tr } } -func encodeTranslateIDsResponse(response *pilosa.TranslateIDsResponse) *internal.TranslateIDsResponse { +func (s Serializer) encodeTranslateIDsResponse(response *pilosa.TranslateIDsResponse) *internal.TranslateIDsResponse { return &internal.TranslateIDsResponse{ Keys: response.Keys, } } -func encodeTransactionMessage(msg *pilosa.TransactionMessage) *internal.TransactionMessage { +func (s Serializer) encodeTransactionMessage(msg *pilosa.TransactionMessage) *internal.TransactionMessage { return &internal.TransactionMessage{ Action: msg.Action, - Transaction: encodeTransaction(msg.Transaction), + Transaction: s.encodeTransaction(msg.Transaction), } } -func encodeTransaction(trns *pilosa.Transaction) *internal.Transaction { +func (s Serializer) encodeTransaction(trns *pilosa.Transaction) *internal.Transaction { if trns == nil { return nil } @@ -856,111 +861,111 @@ func encodeTransaction(trns *pilosa.Transaction) *internal.Transaction { Active: trns.Active, Exclusive: trns.Exclusive, Timeout: int64(trns.Timeout), - Deadline: encodeTransactionDeadline(trns.Deadline), - Stats: encodeTransactionStats(trns.Stats), + Deadline: s.encodeTransactionDeadline(trns.Deadline), + Stats: s.encodeTransactionStats(trns.Stats), } } -func encodeTransactionDeadline(deadline time.Time) int64 { +func (s Serializer) encodeTransactionDeadline(deadline time.Time) int64 { if deadline.Year() > 2262 || deadline.Year() < 1678 { return 0 } return deadline.UnixNano() } -func encodeTransactionStats(stats pilosa.TransactionStats) *internal.TransactionStats { +func (s Serializer) encodeTransactionStats(stats pilosa.TransactionStats) *internal.TransactionStats { return &internal.TransactionStats{} } -func decodeResizeInstruction(ri *internal.ResizeInstruction, m *pilosa.ResizeInstruction) { +func (s Serializer) decodeResizeInstruction(ri *internal.ResizeInstruction, m *pilosa.ResizeInstruction) { m.JobID = ri.JobID m.Node = &pilosa.Node{} - decodeNode(ri.Node, m.Node) + s.decodeNode(ri.Node, m.Node) m.Coordinator = &pilosa.Node{} - decodeNode(ri.Coordinator, m.Coordinator) + s.decodeNode(ri.Coordinator, m.Coordinator) m.Sources = make([]*pilosa.ResizeSource, len(ri.Sources)) - decodeResizeSources(ri.Sources, m.Sources) + s.decodeResizeSources(ri.Sources, m.Sources) m.TranslationSources = make([]*pilosa.TranslationResizeSource, len(ri.TranslationSources)) - decodeTranslationResizeSources(ri.TranslationSources, m.TranslationSources) + s.decodeTranslationResizeSources(ri.TranslationSources, m.TranslationSources) m.NodeStatus = &pilosa.NodeStatus{} - decodeNodeStatus(ri.NodeStatus, m.NodeStatus) + s.decodeNodeStatus(ri.NodeStatus, m.NodeStatus) m.ClusterStatus = &pilosa.ClusterStatus{} - decodeClusterStatus(ri.ClusterStatus, m.ClusterStatus) + s.decodeClusterStatus(ri.ClusterStatus, m.ClusterStatus) } -func decodeResizeSources(srcs []*internal.ResizeSource, m []*pilosa.ResizeSource) { +func (s Serializer) decodeResizeSources(srcs []*internal.ResizeSource, m []*pilosa.ResizeSource) { for i := range srcs { m[i] = &pilosa.ResizeSource{} - decodeResizeSource(srcs[i], m[i]) + s.decodeResizeSource(srcs[i], m[i]) } } -func decodeResizeSource(rs *internal.ResizeSource, m *pilosa.ResizeSource) { +func (s Serializer) decodeResizeSource(rs *internal.ResizeSource, m *pilosa.ResizeSource) { m.Node = &pilosa.Node{} - decodeNode(rs.Node, m.Node) + s.decodeNode(rs.Node, m.Node) m.Index = rs.Index m.Field = rs.Field m.View = rs.View m.Shard = rs.Shard } -func decodeTranslationResizeSources(srcs []*internal.TranslationResizeSource, m []*pilosa.TranslationResizeSource) { +func (s Serializer) decodeTranslationResizeSources(srcs []*internal.TranslationResizeSource, m []*pilosa.TranslationResizeSource) { for i := range srcs { m[i] = &pilosa.TranslationResizeSource{} - decodeTranslationResizeSource(srcs[i], m[i]) + s.decodeTranslationResizeSource(srcs[i], m[i]) } } -func decodeTranslationResizeSource(rs *internal.TranslationResizeSource, m *pilosa.TranslationResizeSource) { +func (s Serializer) decodeTranslationResizeSource(rs *internal.TranslationResizeSource, m *pilosa.TranslationResizeSource) { m.Node = &pilosa.Node{} - decodeNode(rs.Node, m.Node) + s.decodeNode(rs.Node, m.Node) m.Index = rs.Index m.PartitionID = int(rs.PartitionID) } -func decodeSchema(s *internal.Schema, m *pilosa.Schema) { - m.Indexes = make([]*pilosa.IndexInfo, len(s.Indexes)) - decodeIndexes(s.Indexes, m.Indexes) +func (s Serializer) decodeSchema(sc *internal.Schema, m *pilosa.Schema) { + m.Indexes = make([]*pilosa.IndexInfo, len(sc.Indexes)) + s.decodeIndexes(sc.Indexes, m.Indexes) } -func decodeIndexes(idxs []*internal.Index, m []*pilosa.IndexInfo) { +func (s Serializer) decodeIndexes(idxs []*internal.Index, m []*pilosa.IndexInfo) { for i := range idxs { m[i] = &pilosa.IndexInfo{} - decodeIndex(idxs[i], m[i]) + s.decodeIndex(idxs[i], m[i]) } } -func decodeIndex(idx *internal.Index, m *pilosa.IndexInfo) { +func (s Serializer) decodeIndex(idx *internal.Index, m *pilosa.IndexInfo) { m.Name = idx.Name m.Options = pilosa.IndexOptions{} - decodeIndexMeta(idx.Options, &m.Options) + s.decodeIndexMeta(idx.Options, &m.Options) m.Fields = make([]*pilosa.FieldInfo, len(idx.Fields)) - decodeFields(idx.Fields, m.Fields) + s.decodeFields(idx.Fields, m.Fields) } -func decodeFields(fs []*internal.Field, m []*pilosa.FieldInfo) { +func (s Serializer) decodeFields(fs []*internal.Field, m []*pilosa.FieldInfo) { for i := range fs { m[i] = &pilosa.FieldInfo{} - decodeField(fs[i], m[i]) + s.decodeField(fs[i], m[i]) } } -func decodeField(f *internal.Field, m *pilosa.FieldInfo) { +func (s Serializer) decodeField(f *internal.Field, m *pilosa.FieldInfo) { m.Name = f.Name m.Options = pilosa.FieldOptions{} - decodeFieldOptions(f.Meta, &m.Options) + s.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) { +func (s Serializer) decodeFieldOptions(options *internal.FieldOptions, m *pilosa.FieldOptions) { m.Type = options.Type m.CacheType = options.CacheType m.CacheSize = options.CacheSize - decodeDecimal(options.Min, &m.Min) - decodeDecimal(options.Max, &m.Max) + s.decodeDecimal(options.Min, &m.Min) + s.decodeDecimal(options.Max, &m.Max) m.Base = options.Base m.Scale = options.Scale m.BitDepth = uint(options.BitDepth) @@ -969,157 +974,158 @@ func decodeFieldOptions(options *internal.FieldOptions, m *pilosa.FieldOptions) m.ForeignIndex = options.ForeignIndex } -func decodeDecimal(d *internal.Decimal, m *pql.Decimal) { +func (s Serializer) decodeDecimal(d *internal.Decimal, m *pql.Decimal) { m.Value = d.Value m.Scale = d.Scale } -func decodeNodes(a []*internal.Node, m []*pilosa.Node) { +func (s Serializer) decodeNodes(a []*internal.Node, m []*pilosa.Node) { for i := range a { m[i] = &pilosa.Node{} - decodeNode(a[i], m[i]) + s.decodeNode(a[i], m[i]) } } -func decodeClusterStatus(cs *internal.ClusterStatus, m *pilosa.ClusterStatus) { +func (s Serializer) 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) + s.decodeNodes(cs.Nodes, m.Nodes) } -func decodeNode(node *internal.Node, m *pilosa.Node) { +func (s Serializer) decodeNode(node *internal.Node, m *pilosa.Node) { m.ID = node.ID - decodeURI(node.URI, &m.URI) + s.decodeURI(node.URI, &m.URI) m.IsCoordinator = node.IsCoordinator m.State = node.State } -func decodeURI(i *internal.URI, m *pilosa.URI) { +func (s Serializer) 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) { +func (s Serializer) 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) { +func (s Serializer) decodeCreateIndexMessage(pb *internal.CreateIndexMessage, m *pilosa.CreateIndexMessage) { m.Index = pb.Index m.Meta = &pilosa.IndexOptions{} - decodeIndexMeta(pb.Meta, m.Meta) + s.decodeIndexMeta(pb.Meta, m.Meta) } -func decodeIndexMeta(pb *internal.IndexMeta, m *pilosa.IndexOptions) { +func (s Serializer) decodeIndexMeta(pb *internal.IndexMeta, m *pilosa.IndexOptions) { if pb != nil { m.Keys = pb.Keys m.TrackExistence = pb.TrackExistence } } -func decodeDeleteIndexMessage(pb *internal.DeleteIndexMessage, m *pilosa.DeleteIndexMessage) { +func (s Serializer) decodeDeleteIndexMessage(pb *internal.DeleteIndexMessage, m *pilosa.DeleteIndexMessage) { m.Index = pb.Index } -func decodeCreateFieldMessage(pb *internal.CreateFieldMessage, m *pilosa.CreateFieldMessage) { +func (s Serializer) decodeCreateFieldMessage(pb *internal.CreateFieldMessage, m *pilosa.CreateFieldMessage) { m.Index = pb.Index m.Field = pb.Field m.Meta = &pilosa.FieldOptions{} - decodeFieldOptions(pb.Meta, m.Meta) + s.decodeFieldOptions(pb.Meta, m.Meta) } -func decodeDeleteFieldMessage(pb *internal.DeleteFieldMessage, m *pilosa.DeleteFieldMessage) { +func (s Serializer) decodeDeleteFieldMessage(pb *internal.DeleteFieldMessage, m *pilosa.DeleteFieldMessage) { m.Index = pb.Index m.Field = pb.Field } -func decodeDeleteAvailableShardMessage(pb *internal.DeleteAvailableShardMessage, m *pilosa.DeleteAvailableShardMessage) { +func (s Serializer) 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) { +func (s Serializer) 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) { +func (s Serializer) 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) { +func (s Serializer) decodeResizeInstructionComplete(pb *internal.ResizeInstructionComplete, m *pilosa.ResizeInstructionComplete) { m.JobID = pb.JobID m.Node = &pilosa.Node{} - decodeNode(pb.Node, m.Node) + s.decodeNode(pb.Node, m.Node) m.Error = pb.Error } -func decodeSetCoordinatorMessage(pb *internal.SetCoordinatorMessage, m *pilosa.SetCoordinatorMessage) { +func (s Serializer) decodeSetCoordinatorMessage(pb *internal.SetCoordinatorMessage, m *pilosa.SetCoordinatorMessage) { m.New = &pilosa.Node{} - decodeNode(pb.New, m.New) + s.decodeNode(pb.New, m.New) } -func decodeUpdateCoordinatorMessage(pb *internal.UpdateCoordinatorMessage, m *pilosa.UpdateCoordinatorMessage) { +func (s Serializer) decodeUpdateCoordinatorMessage(pb *internal.UpdateCoordinatorMessage, m *pilosa.UpdateCoordinatorMessage) { m.New = &pilosa.Node{} - decodeNode(pb.New, m.New) + s.decodeNode(pb.New, m.New) } -func decodeNodeStateMessage(pb *internal.NodeStateMessage, m *pilosa.NodeStateMessage) { +func (s Serializer) decodeNodeStateMessage(pb *internal.NodeStateMessage, m *pilosa.NodeStateMessage) { m.NodeID = pb.NodeID m.State = pb.State } -func decodeNodeEventMessage(pb *internal.NodeEventMessage, m *pilosa.NodeEvent) { +func (s Serializer) decodeNodeEventMessage(pb *internal.NodeEventMessage, m *pilosa.NodeEvent) { m.Event = pilosa.NodeEventType(pb.Event) m.Node = &pilosa.Node{} - decodeNode(pb.Node, m.Node) + s.decodeNode(pb.Node, m.Node) } -func decodeNodeStatus(pb *internal.NodeStatus, m *pilosa.NodeStatus) { +func (s Serializer) decodeNodeStatus(pb *internal.NodeStatus, m *pilosa.NodeStatus) { m.Node = &pilosa.Node{} - m.Indexes = decodeIndexStatuses(pb.Indexes) + m.Indexes = s.decodeIndexStatuses(pb.Indexes) m.Schema = &pilosa.Schema{} - decodeSchema(pb.Schema, m.Schema) + s.decodeSchema(pb.Schema, m.Schema) } -func decodeIndexStatuses(a []*internal.IndexStatus) []*pilosa.IndexStatus { +func (s Serializer) 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]) + s.decodeIndexStatus(a[i], m[i]) } return m } -func decodeIndexStatus(pb *internal.IndexStatus, m *pilosa.IndexStatus) { +func (s Serializer) decodeIndexStatus(pb *internal.IndexStatus, m *pilosa.IndexStatus) { m.Name = pb.Name - m.Fields = decodeFieldStatuses(pb.Fields) + m.Fields = s.decodeFieldStatuses(pb.Fields) } -func decodeFieldStatuses(a []*internal.FieldStatus) []*pilosa.FieldStatus { +func (s Serializer) 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]) + s.decodeFieldStatus(a[i], m[i]) } return m } -func decodeFieldStatus(pb *internal.FieldStatus, m *pilosa.FieldStatus) { +func (s Serializer) 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 (s Serializer) decodeRecalculateCaches(pb *internal.RecalculateCaches, m *pilosa.RecalculateCaches) { +} -func decodeQueryRequest(pb *internal.QueryRequest, m *pilosa.QueryRequest) { +func (s Serializer) decodeQueryRequest(pb *internal.QueryRequest, m *pilosa.QueryRequest) { m.Query = pb.Query m.Shards = pb.Shards m.ColumnAttrs = pb.ColumnAttrs @@ -1128,11 +1134,11 @@ func decodeQueryRequest(pb *internal.QueryRequest, m *pilosa.QueryRequest) { m.ExcludeColumns = pb.ExcludeColumns m.EmbeddedData = make([]*pilosa.Row, len(pb.EmbeddedData)) for i := range pb.EmbeddedData { - m.EmbeddedData[i] = decodeRow(pb.EmbeddedData[i]) + m.EmbeddedData[i] = s.decodeRow(pb.EmbeddedData[i]) } } -func decodeImportRequest(pb *internal.ImportRequest, m *pilosa.ImportRequest) { +func (s Serializer) decodeImportRequest(pb *internal.ImportRequest, m *pilosa.ImportRequest) { m.Index = pb.Index m.Field = pb.Field m.Shard = pb.Shard @@ -1143,7 +1149,7 @@ func decodeImportRequest(pb *internal.ImportRequest, m *pilosa.ImportRequest) { m.Timestamps = pb.Timestamps } -func decodeImportValueRequest(pb *internal.ImportValueRequest, m *pilosa.ImportValueRequest) { +func (s Serializer) decodeImportValueRequest(pb *internal.ImportValueRequest, m *pilosa.ImportValueRequest) { m.Index = pb.Index m.Field = pb.Field m.Shard = pb.Shard @@ -1154,7 +1160,7 @@ func decodeImportValueRequest(pb *internal.ImportValueRequest, m *pilosa.ImportV m.StringValues = pb.StringValues } -func decodeImportRoaringRequest(pb *internal.ImportRoaringRequest, m *pilosa.ImportRoaringRequest) { +func (s Serializer) decodeImportRoaringRequest(pb *internal.ImportRoaringRequest, m *pilosa.ImportRoaringRequest) { views := map[string][]byte{} for _, view := range pb.Views { views[view.Name] = view.Data @@ -1165,7 +1171,7 @@ func decodeImportRoaringRequest(pb *internal.ImportRoaringRequest, m *pilosa.Imp m.Views = views } -func decodeImportColumnAttrsRequest(pb *internal.ImportColumnAttrsRequest, m *pilosa.ImportColumnAttrsRequest) { +func (s Serializer) decodeImportColumnAttrsRequest(pb *internal.ImportColumnAttrsRequest, m *pilosa.ImportColumnAttrsRequest) { m.Index = pb.Index m.Shard = pb.Shard m.AttrKey = pb.AttrKey @@ -1173,11 +1179,11 @@ func decodeImportColumnAttrsRequest(pb *internal.ImportColumnAttrsRequest, m *pi m.ColumnIDs = pb.ColumnIDs } -func decodeImportResponse(pb *internal.ImportResponse, m *pilosa.ImportResponse) { +func (s Serializer) decodeImportResponse(pb *internal.ImportResponse, m *pilosa.ImportResponse) { m.Err = pb.Err } -func decodeBlockDataRequest(pb *internal.BlockDataRequest, m *pilosa.BlockDataRequest) { +func (s Serializer) decodeBlockDataRequest(pb *internal.BlockDataRequest, m *pilosa.BlockDataRequest) { m.Index = pb.Index m.Field = pb.Field m.View = pb.View @@ -1185,59 +1191,59 @@ func decodeBlockDataRequest(pb *internal.BlockDataRequest, m *pilosa.BlockDataRe m.Block = pb.Block } -func decodeBlockDataResponse(pb *internal.BlockDataResponse, m *pilosa.BlockDataResponse) { +func (s Serializer) decodeBlockDataResponse(pb *internal.BlockDataResponse, m *pilosa.BlockDataResponse) { m.RowIDs = pb.RowIDs m.ColumnIDs = pb.ColumnIDs } -func decodeQueryResponse(pb *internal.QueryResponse, m *pilosa.QueryResponse) { +func (s Serializer) decodeQueryResponse(pb *internal.QueryResponse, m *pilosa.QueryResponse) { m.ColumnAttrSets = make([]*pilosa.ColumnAttrSet, len(pb.ColumnAttrSets)) - decodeColumnAttrSets(pb.ColumnAttrSets, m.ColumnAttrSets) + s.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) + s.decodeQueryResults(pb.Results, m.Results) } -func decodeColumnAttrSets(pb []*internal.ColumnAttrSet, m []*pilosa.ColumnAttrSet) { +func (s Serializer) decodeColumnAttrSets(pb []*internal.ColumnAttrSet, m []*pilosa.ColumnAttrSet) { for i := range pb { m[i] = &pilosa.ColumnAttrSet{} - decodeColumnAttrSet(pb[i], m[i]) + s.decodeColumnAttrSet(pb[i], m[i]) } } -func decodeColumnAttrSet(pb *internal.ColumnAttrSet, m *pilosa.ColumnAttrSet) { +func (s Serializer) decodeColumnAttrSet(pb *internal.ColumnAttrSet, m *pilosa.ColumnAttrSet) { m.ID = pb.ID m.Key = pb.Key - m.Attrs = decodeAttrs(pb.Attrs) + m.Attrs = s.decodeAttrs(pb.Attrs) } -func decodeQueryResults(pb []*internal.QueryResult, m []interface{}) { +func (s Serializer) decodeQueryResults(pb []*internal.QueryResult, m []interface{}) { for i := range pb { - m[i] = decodeQueryResult(pb[i]) + m[i] = s.decodeQueryResult(pb[i]) } } -func decodeTranslateKeysRequest(pb *internal.TranslateKeysRequest, m *pilosa.TranslateKeysRequest) { +func (s Serializer) 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) { +func (s Serializer) decodeTranslateKeysResponse(pb *internal.TranslateKeysResponse, m *pilosa.TranslateKeysResponse) { m.IDs = pb.IDs } -func decodeTranslateIDsRequest(pb *internal.TranslateIDsRequest, m *pilosa.TranslateIDsRequest) { +func (s Serializer) decodeTranslateIDsRequest(pb *internal.TranslateIDsRequest, m *pilosa.TranslateIDsRequest) { m.Index = pb.Index m.Field = pb.Field m.IDs = pb.IDs } -func decodeTranslateIDsResponse(pb *internal.TranslateIDsResponse, m *pilosa.TranslateIDsResponse) { +func (s Serializer) decodeTranslateIDsResponse(pb *internal.TranslateIDsResponse, m *pilosa.TranslateIDsResponse) { m.Keys = pb.Keys } @@ -1278,18 +1284,18 @@ const ( queryResultTypeSignedRow ) -func decodeQueryResult(pb *internal.QueryResult) interface{} { +func (s Serializer) decodeQueryResult(pb *internal.QueryResult) interface{} { switch pb.Type { case queryResultTypeSignedRow: - return decodeSignedRow(pb.SignedRow) + return s.decodeSignedRow(pb.SignedRow) case queryResultTypeRow: - return decodeRow(pb.Row) + return s.decodeRow(pb.Row) case queryResultTypePairs: - return decodePairs(pb.Pairs) + return s.decodePairs(pb.Pairs) case queryResultTypePairsField: - return decodePairsField(pb.PairsField) + return s.decodePairsField(pb.PairsField) case queryResultTypeValCount: - return decodeValCount(pb.ValCount) + return s.decodeValCount(pb.ValCount) case queryResultTypeUint64: return pb.N case queryResultTypeBool: @@ -1299,19 +1305,19 @@ func decodeQueryResult(pb *internal.QueryResult) interface{} { case queryResultTypeRowIDs: return pilosa.RowIDs(pb.RowIDs) case queryResultTypeRowIdentifiers: - return decodeRowIdentifiers(pb.RowIdentifiers) + return s.decodeRowIdentifiers(pb.RowIdentifiers) case queryResultTypeGroupCounts: - return decodeGroupCounts(pb.GroupCounts) + return s.decodeGroupCounts(pb.GroupCounts) case queryResultTypePair: - return decodePair(pb.Pairs[0]) + return s.decodePair(pb.Pairs[0]) case queryResultTypePairField: - return decodePairField(pb.Pairs[0]) + return s.decodePairField(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 { +// s.decodeRow converts r from its internal representation. +func (s Serializer) decodeRow(pr *internal.Row) *pilosa.Row { if pr == nil { return pilosa.NewRow() } @@ -1325,27 +1331,27 @@ func decodeRow(pr *internal.Row) *pilosa.Row { r.SetBit(v) } } - r.Attrs = decodeAttrs(pr.Attrs) + r.Attrs = s.decodeAttrs(pr.Attrs) r.Keys = pr.Keys return r } -func decodeSignedRow(pr *internal.SignedRow) pilosa.SignedRow { +func (s Serializer) decodeSignedRow(pr *internal.SignedRow) pilosa.SignedRow { if pr == nil { return pilosa.SignedRow{} } r := pilosa.SignedRow{ - Pos: decodeRow(pr.Pos), - Neg: decodeRow(pr.Neg), + Pos: s.decodeRow(pr.Pos), + Neg: s.decodeRow(pr.Neg), } return r } -func decodeAttrs(pb []*internal.Attr) map[string]interface{} { +func (s Serializer) decodeAttrs(pb []*internal.Attr) map[string]interface{} { m := make(map[string]interface{}, len(pb)) for i := range pb { - key, value := decodeAttr(pb[i]) + key, value := s.decodeAttr(pb[i]) m[key] = value } return m @@ -1358,7 +1364,7 @@ const ( attrTypeFloat = 4 ) -func decodeAttr(attr *internal.Attr) (key string, value interface{}) { +func (s Serializer) decodeAttr(attr *internal.Attr) (key string, value interface{}) { switch attr.Type { case attrTypeString: return attr.Key, attr.StringValue @@ -1373,18 +1379,18 @@ func decodeAttr(attr *internal.Attr) (key string, value interface{}) { } } -func decodeRowIdentifiers(a *internal.RowIdentifiers) *pilosa.RowIdentifiers { +func (s Serializer) decodeRowIdentifiers(a *internal.RowIdentifiers) *pilosa.RowIdentifiers { return &pilosa.RowIdentifiers{ Rows: a.Rows, Keys: a.Keys, } } -func decodeGroupCounts(a []*internal.GroupCount) []pilosa.GroupCount { +func (s Serializer) 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), + Group: s.decodeFieldRows(a[i].Group), Count: a[i].Count, Sum: a[i].Sum, } @@ -1392,7 +1398,7 @@ func decodeGroupCounts(a []*internal.GroupCount) []pilosa.GroupCount { return other } -func decodeFieldRows(a []*internal.FieldRow) []pilosa.FieldRow { +func (s Serializer) decodeFieldRows(a []*internal.FieldRow) []pilosa.FieldRow { other := make([]pilosa.FieldRow, len(a)) for i := range a { fr := a[i] @@ -1408,26 +1414,26 @@ func decodeFieldRows(a []*internal.FieldRow) []pilosa.FieldRow { return other } -func decodePairs(a []*internal.Pair) []pilosa.Pair { +func (s Serializer) decodePairs(a []*internal.Pair) []pilosa.Pair { other := make([]pilosa.Pair, len(a)) for i := range a { - other[i] = decodePair(a[i]) + other[i] = s.decodePair(a[i]) } return other } -func decodePairsField(a *internal.PairsField) *pilosa.PairsField { +func (s Serializer) decodePairsField(a *internal.PairsField) *pilosa.PairsField { other := &pilosa.PairsField{ Pairs: make([]pilosa.Pair, len(a.Pairs)), } for i := range a.Pairs { - other.Pairs[i] = decodePair(a.Pairs[i]) + other.Pairs[i] = s.decodePair(a.Pairs[i]) } other.Field = a.Field return other } -func decodePair(pb *internal.Pair) pilosa.Pair { +func (s Serializer) decodePair(pb *internal.Pair) pilosa.Pair { return pilosa.Pair{ ID: pb.ID, Key: pb.Key, @@ -1435,7 +1441,7 @@ func decodePair(pb *internal.Pair) pilosa.Pair { } } -func decodePairField(pb *internal.Pair) pilosa.PairField { +func (s Serializer) decodePairField(pb *internal.Pair) pilosa.PairField { return pilosa.PairField{ Pair: pilosa.Pair{ ID: pb.ID, @@ -1446,16 +1452,16 @@ func decodePairField(pb *internal.Pair) pilosa.PairField { } } -func decodeValCount(pb *internal.ValCount) pilosa.ValCount { +func (s Serializer) decodeValCount(pb *internal.ValCount) pilosa.ValCount { return pilosa.ValCount{ Val: pb.Val, FloatVal: pb.FloatVal, - DecimalVal: decodeDecimalStruct(pb.DecimalVal), + DecimalVal: s.decodeDecimalStruct(pb.DecimalVal), Count: pb.Count, } } -func decodeDecimalStruct(pb *internal.Decimal) *pql.Decimal { +func (s Serializer) decodeDecimalStruct(pb *internal.Decimal) *pql.Decimal { if pb == nil { return nil } @@ -1465,60 +1471,60 @@ func decodeDecimalStruct(pb *internal.Decimal) *pql.Decimal { } } -func encodeColumnAttrSets(a []*pilosa.ColumnAttrSet) []*internal.ColumnAttrSet { +func (s Serializer) encodeColumnAttrSets(a []*pilosa.ColumnAttrSet) []*internal.ColumnAttrSet { other := make([]*internal.ColumnAttrSet, len(a)) for i := range a { - other[i] = encodeColumnAttrSet(a[i]) + other[i] = s.encodeColumnAttrSet(a[i]) } return other } -func encodeColumnAttrSet(set *pilosa.ColumnAttrSet) *internal.ColumnAttrSet { +func (s Serializer) encodeColumnAttrSet(set *pilosa.ColumnAttrSet) *internal.ColumnAttrSet { return &internal.ColumnAttrSet{ ID: set.ID, Key: set.Key, - Attrs: encodeAttrs(set.Attrs), + Attrs: s.encodeAttrs(set.Attrs), } } -func encodeSignedRow(r pilosa.SignedRow) *internal.SignedRow { +func (s Serializer) encodeSignedRow(r pilosa.SignedRow) *internal.SignedRow { ir := &internal.SignedRow{ - Pos: encodeRow(r.Pos), - Neg: encodeRow(r.Neg), + Pos: s.encodeRow(r.Pos), + Neg: s.encodeRow(r.Neg), } return ir } -func encodeRow(r *pilosa.Row) *internal.Row { +func (s Serializer) encodeRow(r *pilosa.Row) *internal.Row { if r == nil { return nil } ir := &internal.Row{ Keys: r.Keys, - Attrs: encodeAttrs(r.Attrs), + Attrs: s.encodeAttrs(r.Attrs), } - if true { - ir.Columns = r.Columns() - } else { + if s.RoaringRows { ir.Roaring = r.Roaring() + } else { + ir.Columns = r.Columns() } return ir } -func encodeRowIdentifiers(r pilosa.RowIdentifiers) *internal.RowIdentifiers { +func (s Serializer) encodeRowIdentifiers(r pilosa.RowIdentifiers) *internal.RowIdentifiers { return &internal.RowIdentifiers{ Rows: r.Rows, Keys: r.Keys, - //Attrs: encodeAttrs(r.Attrs), + //Attrs: s.encodeAttrs(r.Attrs), } } -func encodeGroupCounts(counts []pilosa.GroupCount) []*internal.GroupCount { +func (s Serializer) 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), + Group: s.encodeFieldRows(counts[i].Group), Count: counts[i].Count, Sum: counts[i].Sum, } @@ -1526,7 +1532,7 @@ func encodeGroupCounts(counts []pilosa.GroupCount) []*internal.GroupCount { return result } -func encodeFieldRows(a []pilosa.FieldRow) []*internal.FieldRow { +func (s Serializer) encodeFieldRows(a []pilosa.FieldRow) []*internal.FieldRow { other := make([]*internal.FieldRow, len(a)) for i := range a { fr := a[i] @@ -1543,26 +1549,26 @@ func encodeFieldRows(a []pilosa.FieldRow) []*internal.FieldRow { return other } -func encodePairs(a pilosa.Pairs) []*internal.Pair { +func (s Serializer) encodePairs(a pilosa.Pairs) []*internal.Pair { other := make([]*internal.Pair, len(a)) for i := range a { - other[i] = encodePair(a[i]) + other[i] = s.encodePair(a[i]) } return other } -func encodePairsField(a *pilosa.PairsField) *internal.PairsField { +func (s Serializer) encodePairsField(a *pilosa.PairsField) *internal.PairsField { other := &internal.PairsField{ Pairs: make([]*internal.Pair, len(a.Pairs)), } for i := range a.Pairs { - other.Pairs[i] = encodePair(a.Pairs[i]) + other.Pairs[i] = s.encodePair(a.Pairs[i]) } other.Field = a.Field return other } -func encodePair(p pilosa.Pair) *internal.Pair { +func (s Serializer) encodePair(p pilosa.Pair) *internal.Pair { return &internal.Pair{ ID: p.ID, Key: p.Key, @@ -1570,27 +1576,27 @@ func encodePair(p pilosa.Pair) *internal.Pair { } } -func encodePairField(p pilosa.PairField) *internal.Pair { +func (s Serializer) encodePairField(p pilosa.PairField) *internal.Pair { /* // TODO: in order to have this, we need PairField in QueryResponse. return &internal.Pair{ - Pair: encodePair(p.Pair), + Pair: s.encodePair(p.Pair), Field: p.Field, } */ - return encodePair(p.Pair) + return s.encodePair(p.Pair) } -func encodeValCount(vc pilosa.ValCount) *internal.ValCount { +func (s Serializer) encodeValCount(vc pilosa.ValCount) *internal.ValCount { return &internal.ValCount{ Val: vc.Val, FloatVal: vc.FloatVal, - DecimalVal: encodeDecimal(vc.DecimalVal), + DecimalVal: s.encodeDecimal(vc.DecimalVal), Count: vc.Count, } } -func encodeDecimal(p *pql.Decimal) *internal.Decimal { +func (s Serializer) encodeDecimal(p *pql.Decimal) *internal.Decimal { if p == nil { return nil } @@ -1600,7 +1606,7 @@ func encodeDecimal(p *pql.Decimal) *internal.Decimal { } } -func encodeAttrs(m map[string]interface{}) []*internal.Attr { +func (s Serializer) encodeAttrs(m map[string]interface{}) []*internal.Attr { keys := make([]string, 0, len(m)) for k := range m { keys = append(keys, k) @@ -1609,13 +1615,13 @@ func encodeAttrs(m map[string]interface{}) []*internal.Attr { a := make([]*internal.Attr, len(keys)) for i := range keys { - a[i] = encodeAttr(keys[i], m[keys[i]]) + a[i] = s.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 { +// s.encodeAttr converts a key/value pair into an Attr internal representation. +func (s Serializer) encodeAttr(key string, value interface{}) *internal.Attr { pb := &internal.Attr{Key: key} switch value := value.(type) { case string: diff --git a/http/client.go b/http/client.go index 2b112e9c5..d8d9541e2 100644 --- a/http/client.go +++ b/http/client.go @@ -291,6 +291,7 @@ func (c *InternalClient) QueryNode(ctx context.Context, uri *pilosa.URI, index s req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/x-protobuf") req.Header.Set("Accept", "application/x-protobuf") + req.Header.Set("X-Pilosa-Row", "roaring") req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Execute request against the host. @@ -495,6 +496,7 @@ func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, inde req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/x-protobuf") req.Header.Set("Accept", "application/x-protobuf") + req.Header.Set("X-Pilosa-Row", "roaring") req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Execute request against the host. @@ -682,6 +684,7 @@ func (c *InternalClient) ImportRoaring(ctx context.Context, uri *pilosa.URI, ind } httpReq.Header.Set("Content-Type", "application/x-protobuf") httpReq.Header.Set("Accept", "application/x-protobuf") + httpReq.Header.Set("X-Pilosa-Row", "roaring") httpReq.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Execute request against the host. @@ -731,6 +734,7 @@ func (c *InternalClient) ImportColumnAttrs(ctx context.Context, uri *pilosa.URI, } httpReq.Header.Set("Content-Type", "application/x-protobuf") httpReq.Header.Set("Accept", "application/x-protobuf") + httpReq.Header.Set("X-Pilosa-Row", "roaring") httpReq.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Execute request against the host. @@ -1006,6 +1010,7 @@ func (c *InternalClient) BlockData(ctx context.Context, uri *pilosa.URI, index, req.Header.Set("Content-Type", "application/protobuf") req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Accept", "application/protobuf") + req.Header.Set("X-Pilosa-Row", "roaring") req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) resp, err := c.executeRequest(req.WithContext(ctx)) @@ -1163,6 +1168,7 @@ func (c *InternalClient) TranslateKeysNode(ctx context.Context, uri *pilosa.URI, req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/x-protobuf") req.Header.Set("Accept", "application/x-protobuf") + req.Header.Set("X-Pilosa-Row", "roaring") req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Execute request against the host. @@ -1213,6 +1219,7 @@ func (c *InternalClient) TranslateIDsNode(ctx context.Context, uri *pilosa.URI, req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/x-protobuf") req.Header.Set("Accept", "application/x-protobuf") + req.Header.Set("X-Pilosa-Row", "roaring") req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Execute request against the host. diff --git a/http/handler.go b/http/handler.go index 0e80a2b3c..d5cbadef2 100644 --- a/http/handler.go +++ b/http/handler.go @@ -36,6 +36,7 @@ import ( "github.com/gorilla/handlers" "github.com/gorilla/mux" "github.com/pilosa/pilosa/v2" + "github.com/pilosa/pilosa/v2/encoding/proto" "github.com/pilosa/pilosa/v2/logger" "github.com/pilosa/pilosa/v2/pql" "github.com/pilosa/pilosa/v2/tracing" @@ -472,6 +473,17 @@ func validHeaderAcceptJSON(header http.Header) bool { return true } +// headerAcceptRoaringRow tells us that the request should accept roaring +// rows in response. +func headerAcceptRoaringRow(header http.Header) bool { + for _, v := range header["X-Pilosa-Row"] { + if v == "roaring" { + return true + } + } + return false +} + // handleGetSchema handles GET /schema requests. func (h *Handler) handleGetSchema(w http.ResponseWriter, r *http.Request) { if !validHeaderAcceptJSON(r.Header) { @@ -1223,7 +1235,7 @@ func (h *Handler) readProtobufQueryRequest(r *http.Request) (*pilosa.QueryReques } qreq := &pilosa.QueryRequest{} - err = h.api.Serializer.Unmarshal(body, qreq) + err = proto.DefaultSerializer.Unmarshal(body, qreq) if err != nil { return nil, errors.Wrap(err, "unmarshalling query request") } @@ -1271,15 +1283,19 @@ func (h *Handler) readURLQueryRequest(r *http.Request) (*pilosa.QueryRequest, er func (h *Handler) writeQueryResponse(w http.ResponseWriter, r *http.Request, resp *pilosa.QueryResponse) error { if !validHeaderAcceptJSON(r.Header) { w.Header().Set("Content-Type", "application/protobuf") - return h.writeProtobufQueryResponse(w, resp) + return h.writeProtobufQueryResponse(w, resp, headerAcceptRoaringRow(r.Header)) } w.Header().Set("Content-Type", "application/json") return h.writeJSONQueryResponse(w, resp) } // writeProtobufQueryResponse writes the response from the executor to w as protobuf. -func (h *Handler) writeProtobufQueryResponse(w io.Writer, resp *pilosa.QueryResponse) error { - if buf, err := h.api.Serializer.Marshal(resp); err != nil { +func (h *Handler) writeProtobufQueryResponse(w io.Writer, resp *pilosa.QueryResponse, writeRoaring bool) error { + serializer := proto.DefaultSerializer + if writeRoaring { + serializer = proto.RoaringSerializer + } + if buf, err := serializer.Marshal(resp); err != nil { return errors.Wrap(err, "marshalling") } else if _, err := w.Write(buf); err != nil { return errors.Wrap(err, "writing") @@ -1342,7 +1358,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { // Field type: Int // Marshal into request object. req := &pilosa.ImportValueRequest{} - if err := h.api.Serializer.Unmarshal(body, req); err != nil { + if err := proto.DefaultSerializer.Unmarshal(body, req); err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } @@ -1360,7 +1376,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { // Field type: set, time, mutex // Marshal into request object. req := &pilosa.ImportRequest{} - if err := h.api.Serializer.Unmarshal(body, req); err != nil { + if err := proto.DefaultSerializer.Unmarshal(body, req); err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } @@ -1377,7 +1393,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { } // Marshal response object. - buf, e := h.api.Serializer.Marshal(&pilosa.ImportResponse{Err: ""}) + buf, e := proto.DefaultSerializer.Marshal(&pilosa.ImportResponse{Err: ""}) if e != nil { http.Error(w, "marshal import response", http.StatusInternalServerError) return @@ -1885,7 +1901,7 @@ func (h *Handler) handlePostImportColumnAttrs(w http.ResponseWriter, r *http.Req } req := &pilosa.ImportColumnAttrsRequest{} - if err := h.api.Serializer.Unmarshal(body, req); err != nil { + if err := proto.DefaultSerializer.Unmarshal(body, req); err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } @@ -1896,7 +1912,7 @@ func (h *Handler) handlePostImportColumnAttrs(w http.ResponseWriter, r *http.Req } // Marshal response object. - buf, e := h.api.Serializer.Marshal(&pilosa.ImportResponse{Err: ""}) + buf, e := proto.DefaultSerializer.Marshal(&pilosa.ImportResponse{Err: ""}) if e != nil { http.Error(w, "marshal import-column-attrs response", http.StatusInternalServerError) return @@ -1943,7 +1959,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request req := &pilosa.ImportRoaringRequest{} span, _ = tracing.StartSpanFromContext(ctx, "Unmarshal") - err = h.api.Serializer.Unmarshal(body, req) + err = proto.DefaultSerializer.Unmarshal(body, req) span.Finish() if err != nil { http.Error(w, err.Error(), http.StatusBadRequest) @@ -1969,7 +1985,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request } // Marshal response object. - buf, err := h.api.Serializer.Marshal(resp) + buf, err := proto.DefaultSerializer.Marshal(resp) if err != nil { http.Error(w, fmt.Sprintf("marshal import-roaring response: %v", err), http.StatusInternalServerError) return