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