diff --git a/client.go b/client.go index d3bdb1071..c2e88c77b 100644 --- a/client.go +++ b/client.go @@ -303,6 +303,31 @@ func (c *InternalHTTPClient) Import(ctx context.Context, index, frame string, sl return nil } +// ImportK bulk imports bits to a host. +func (c *InternalHTTPClient) ImportK(ctx context.Context, index, frame string, bits []Bit) error { + if index == "" { + return ErrIndexRequired + } else if frame == "" { + return ErrFrameRequired + } + + buf, err := marshalImportPayloadK(index, frame, bits) + if err != nil { + return fmt.Errorf("Error Creating Payload: %s", err) + } + + node := &Node{ + URI: *c.defaultURI, + } + + // Import to node. + if err := c.importNode(ctx, node, buf); err != nil { + return fmt.Errorf("import node: host=%s, err=%s", node.URI, err) + } + + return nil +} + func (c *InternalHTTPClient) EnsureIndex(ctx context.Context, name string, options IndexOptions) error { err := c.CreateIndex(ctx, name, options) if err == nil || err == ErrIndexExists { @@ -341,6 +366,27 @@ func marshalImportPayload(index, frame string, slice uint64, bits []Bit) ([]byte return buf, nil } +// marshalImportPayloadK marshalls the import parameters into a protobuf byte slice. +func marshalImportPayloadK(index, frame string, bits []Bit) ([]byte, error) { + // Separate row and column IDs to reduce allocations. + rowKeys := Bits(bits).RowKeys() + columnKeys := Bits(bits).ColumnKeys() + timestamps := Bits(bits).Timestamps() + + // Marshal bits to protobufs. + buf, err := proto.Marshal(&internal.ImportRequest{ + Index: index, + Frame: frame, + RowKeys: rowKeys, + ColumnKeys: columnKeys, + Timestamps: timestamps, + }) + if err != nil { + return nil, fmt.Errorf("marshal import request: %s", err) + } + return buf, nil +} + // importNode sends a pre-marshaled import request to a node. func (c *InternalHTTPClient) importNode(ctx context.Context, node *Node, buf []byte) error { // Create URL & HTTP request. @@ -1124,6 +1170,8 @@ func (c *InternalHTTPClient) NodeID(uri *URI) (string, error) { type Bit struct { RowID uint64 ColumnID uint64 + RowKey string + ColumnKey string Timestamp int64 } @@ -1161,6 +1209,24 @@ func (p Bits) ColumnIDs() []uint64 { return other } +// RowKeys returns a slice of all the row keys. +func (p Bits) RowKeys() []string { + other := make([]string, len(p)) + for i := range p { + other[i] = p[i].RowKey + } + return other +} + +// ColumnKeys returns a slice of all the column keys. +func (p Bits) ColumnKeys() []string { + other := make([]string, len(p)) + for i := range p { + other[i] = p[i].ColumnKey + } + return other +} + // Timestamps returns a slice of all the timestamps. func (p Bits) Timestamps() []int64 { other := make([]int64, len(p)) @@ -1280,6 +1346,7 @@ type InternalClient interface { FragmentNodes(ctx context.Context, index string, slice uint64) ([]*Node, error) ExecuteQuery(ctx context.Context, index string, queryRequest *internal.QueryRequest) (*internal.QueryResponse, error) Import(ctx context.Context, index, frame string, slice uint64, bits []Bit) error + ImportK(ctx context.Context, index, frame string, bits []Bit) error EnsureIndex(ctx context.Context, name string, options IndexOptions) error EnsureFrame(ctx context.Context, indexName string, frameName string, options FrameOptions) error ImportValue(ctx context.Context, index, frame, field string, slice uint64, vals []FieldValue) error diff --git a/cluster.go b/cluster.go index 4f29196bc..b2ed0e372 100644 --- a/cluster.go +++ b/cluster.go @@ -175,7 +175,7 @@ type Cluster struct { // Required for cluster Resize. Static bool // Static is primarily used for testing in a non-gossip environment. - State string + state string Coordinator URI Holder *Holder Broadcaster Broadcaster @@ -302,9 +302,21 @@ func (c *Cluster) setID(id string) { c.Topology.ClusterID = c.ID } +func (c *Cluster) State() string { + c.mu.RLock() + defer c.mu.RUnlock() + return c.state +} + +func (c *Cluster) SetState(state string) { + c.mu.Lock() + defer c.mu.Unlock() + c.setState(state) +} + func (c *Cluster) setState(state string) { // Ignore cases where the state hasn't changed. - if state == c.State { + if state == c.state { return } @@ -321,12 +333,12 @@ func (c *Cluster) setState(state string) { // - ClusterStateStarting // If state is RESIZING -> NORMAL then run cleanup. - if c.State == ClusterStateResizing { + if c.state == ClusterStateResizing { doCleanup = true } } - c.State = state + c.state = state // TODO: consider NOT running cleanup on an active node that has // been removed. @@ -373,7 +385,7 @@ func (c *Cluster) ReceiveNodeState(uri URI, state string) error { } // This method is really only useful during initial startup. - if c.State != ClusterStateStarting { + if c.State() != ClusterStateStarting { return nil } @@ -397,7 +409,7 @@ func (c *Cluster) ReceiveNodeState(uri URI, state string) error { func (c *Cluster) Status() *internal.ClusterStatus { return &internal.ClusterStatus{ ClusterID: c.ID, - State: c.State, + State: c.state, NodeSet: encodeURIs(c.NodeSet()), } } @@ -763,7 +775,7 @@ func (h *jmphasher) Hash(key uint64, n int) int { func (c *Cluster) Open() error { // Cluster always comes up in state STARTING until cluster membership is determined. - c.State = ClusterStateStarting + c.state = ClusterStateStarting // Load topology file if it exists. if err := c.loadTopology(); err != nil { @@ -820,7 +832,7 @@ func (c *Cluster) markAsJoined() { } func (c *Cluster) needTopologyAgreement() bool { - return c.State == ClusterStateStarting && !URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) + return c.State() == ClusterStateStarting && !URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } func (c *Cluster) haveTopologyAgreement() bool { @@ -886,7 +898,7 @@ func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { } func (c *Cluster) setStateAndBroadcast(state string) error { - c.setState(state) + c.SetState(state) // Broadcast cluster status changes to the cluster. c.logger().Printf("broadcasting ClusterStatus: %s", state) return c.Broadcaster.SendSync(c.Status()) @@ -1618,7 +1630,7 @@ func (c *Cluster) NodeLeave(uri URI) error { return fmt.Errorf("Node removal requests are only valid on the Coordinator node: %s", c.Coordinator) } - if c.State != ClusterStateNormal { + if c.State() != ClusterStateNormal { return fmt.Errorf("Cluster must be in state %s to remove a node. Current state: %s", ClusterStateNormal, c.State) } @@ -1683,7 +1695,7 @@ func (c *Cluster) MergeClusterStatus(cs *internal.ClusterStatus) error { } } - c.setState(cs.State) + c.SetState(cs.State) c.markAsJoined() diff --git a/cluster_test.go b/cluster_test.go index a7fb77c83..b1b62c815 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -255,8 +255,8 @@ func TestCluster_ResizeStates(t *testing.T) { node := tc.Clusters[0] // Ensure that node comes up in state NORMAL. - if node.State != pilosa.ClusterStateNormal { - t.Errorf("expected state: %v, but got: %v", pilosa.ClusterStateNormal, node.State) + if node.State() != pilosa.ClusterStateNormal { + t.Errorf("expected state: %v, but got: %v", pilosa.ClusterStateNormal, node.State()) } expectedTop := &pilosa.Topology{ @@ -292,8 +292,8 @@ func TestCluster_ResizeStates(t *testing.T) { } // Ensure that node comes up in state NORMAL. - if node.State != pilosa.ClusterStateNormal { - t.Errorf("expected state: %v, but got: %v", pilosa.ClusterStateNormal, node.State) + if node.State() != pilosa.ClusterStateNormal { + t.Errorf("expected state: %v, but got: %v", pilosa.ClusterStateNormal, node.State()) } // Close TestCluster. @@ -344,10 +344,10 @@ func TestCluster_ResizeStates(t *testing.T) { node1 := tc.Clusters[1] // Ensure that nodes comes up in state NORMAL. - if node0.State != pilosa.ClusterStateNormal { - t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State) - } else if node1.State != pilosa.ClusterStateNormal { - t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node1.State) + if node0.State() != pilosa.ClusterStateNormal { + t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State()) + } else if node1.State() != pilosa.ClusterStateNormal { + t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node1.State()) } expectedTop := &pilosa.Topology{ @@ -388,8 +388,8 @@ func TestCluster_ResizeStates(t *testing.T) { } // Ensure that node is in state STARTING before the other node joins. - if node0.State != pilosa.ClusterStateStarting { - t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateStarting, node0.State) + if node0.State() != pilosa.ClusterStateStarting { + t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateStarting, node0.State()) } // Expect an error by adding a node not in the topology. @@ -403,10 +403,10 @@ func TestCluster_ResizeStates(t *testing.T) { node2 := tc.Clusters[2] // Ensure that node comes up in state NORMAL. - if node0.State != pilosa.ClusterStateNormal { - t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State) - } else if node2.State != pilosa.ClusterStateNormal { - t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node2.State) + if node0.State() != pilosa.ClusterStateNormal { + t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State()) + } else if node2.State() != pilosa.ClusterStateNormal { + t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node2.State()) } // Close TestCluster. @@ -470,10 +470,10 @@ func TestCluster_ResizeStates(t *testing.T) { node1 := tc.Clusters[1] // Ensure that nodes come up in state NORMAL. - if node0.State != pilosa.ClusterStateNormal { - t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State) - } else if node1.State != pilosa.ClusterStateNormal { - t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node1.State) + if node0.State() != pilosa.ClusterStateNormal { + t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State()) + } else if node1.State() != pilosa.ClusterStateNormal { + t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node1.State()) } expectedTop := &pilosa.Topology{ diff --git a/cmd/import.go b/cmd/import.go index 250dccfd3..8692fa318 100644 --- a/cmd/import.go +++ b/cmd/import.go @@ -56,6 +56,7 @@ omitted. If it is present then its format should be YYYY-MM-DDTHH:MM. flags.StringVarP(&Importer.Index, "index", "i", "", "Pilosa index to import into.") flags.StringVarP(&Importer.Frame, "frame", "f", "", "Frame to import into.") flags.StringVarP(&Importer.Field, "field", "", "", "Field to import into.") + flags.BoolVar(&Importer.StringKeys, "string-keys", false, "Treat payload as string keys.") flags.IntVarP(&Importer.BufferSize, "buffer-size", "s", 10000000, "Number of bits to buffer/sort before importing.") flags.BoolVarP(&Importer.Sort, "sort", "", false, "Enables sorting before import.") flags.BoolVarP(&Importer.CreateSchema, "create", "e", false, "Create the schema if it does not exist before import.") diff --git a/ctl/import.go b/ctl/import.go index 179385db5..8103bf673 100644 --- a/ctl/import.go +++ b/ctl/import.go @@ -48,6 +48,9 @@ type ImportCommand struct { // For Range-Encoded fields, name of the Field to import into. Field string `json:"field"` + // Indicates that the payload should be treated as string keys. + StringKeys bool `json:"StringKeys"` + // Filenames to import from. Paths []string `json:"paths"` @@ -131,7 +134,11 @@ func (cmd *ImportCommand) importPath(ctx context.Context, path string) error { if cmd.Field != "" { return cmd.bufferFieldValues(ctx, path) } else { - return cmd.bufferBits(ctx, path) + if cmd.StringKeys { + return cmd.bufferBitsK(ctx, path) + } else { + return cmd.bufferBits(ctx, path) + } } } @@ -240,7 +247,102 @@ func (cmd *ImportCommand) importBits(ctx context.Context, bits []pilosa.Bit) err } return nil +} +// bufferBitsK buffers slices of keys to be imported as a batch. +func (cmd *ImportCommand) bufferBitsK(ctx context.Context, path string) error { + a := make([]pilosa.Bit, 0, cmd.BufferSize) + + var r *csv.Reader + + if path != "-" { + // Open file for reading. + f, err := os.Open(path) + if err != nil { + return err + } + defer f.Close() + + // Read rows as bits. + r = csv.NewReader(f) + } else { + r = csv.NewReader(cmd.Stdin) + } + + r.FieldsPerRecord = -1 + rnum := 0 + for { + rnum++ + + // Read CSV row. + record, err := r.Read() + if err == io.EOF { + break + } else if err != nil { + return err + } + + // Ignore blank rows. + if record[0] == "" { + continue + } else if len(record) < 2 { + return fmt.Errorf("bad column count on row %d: col=%d", rnum, len(record)) + } + + var bit pilosa.Bit + + // Parse row key. + if record[0] == "" { + return fmt.Errorf("invalid row key on row %d: %q", rnum, record[0]) + } + bit.RowKey = record[0] + + // Parse column key. + if record[1] == "" { + return fmt.Errorf("invalid column id on row %d: %q", rnum, record[1]) + } + bit.ColumnKey = record[1] + + // Parse time, if exists. + if len(record) > 2 && record[2] != "" { + t, err := time.Parse(pilosa.TimeFormat, record[2]) + if err != nil { + return fmt.Errorf("invalid timestamp on row %d: %q", rnum, record[2]) + } + bit.Timestamp = t.UnixNano() + } + + a = append(a, bit) + + // If we've reached the buffer size then import bits. + if len(a) == cmd.BufferSize { + if err := cmd.importBitsK(ctx, a); err != nil { + return err + } + a = a[:0] + } + } + + // If there are still bitKs in the buffer then flush them. + if err := cmd.importBitsK(ctx, a); err != nil { + return err + } + + return nil +} + +// importBitsK sends batches of bitKs to the server. +func (cmd *ImportCommand) importBitsK(ctx context.Context, bits []pilosa.Bit) error { + logger := log.New(cmd.Stderr, "", log.LstdFlags) + + // TODO: does it help to sort the rowKeys? + + logger.Printf("importing keys: n=%d", len(bits)) + if err := cmd.Client.ImportK(ctx, cmd.Index, cmd.Frame, bits); err != nil { + return err + } + + return nil } // bufferFieldValues buffers slices of fieldValues to be imported as a batch. diff --git a/handler.go b/handler.go index d3e4dc5fa..c89abc3bf 100644 --- a/handler.go +++ b/handler.go @@ -1658,6 +1658,16 @@ func (h *Handler) logger() *log.Logger { return log.New(h.LogOutput, "", log.LstdFlags) } +// QueryResult types. +const ( + QueryResultTypeNil uint32 = iota + QueryResultTypeBitmap + QueryResultTypePairs + QueryResultTypeSumCount + QueryResultTypeUint64 + QueryResultTypeBool +) + // QueryRequest represent a request to process a query. type QueryRequest struct { // Index to execute query against. @@ -1737,15 +1747,22 @@ func encodeQueryResponse(resp *QueryResponse) *internal.QueryResponse { switch result := resp.Results[i].(type) { case *Bitmap: + pb.Results[i].Type = QueryResultTypeBitmap pb.Results[i].Bitmap = encodeBitmap(result) case []Pair: + pb.Results[i].Type = QueryResultTypePairs pb.Results[i].Pairs = encodePairs(result) case SumCount: + pb.Results[i].Type = QueryResultTypeSumCount pb.Results[i].SumCount = encodeSumCount(result) case uint64: + pb.Results[i].Type = QueryResultTypeUint64 pb.Results[i].N = result case bool: + pb.Results[i].Type = QueryResultTypeBool pb.Results[i].Changed = result + case nil: + pb.Results[i].Type = QueryResultTypeNil } } diff --git a/handler_test.go b/handler_test.go index 7413cb3ea..b89b7f59c 100644 --- a/handler_test.go +++ b/handler_test.go @@ -376,6 +376,8 @@ func TestHandler_Query_Uint64_Protobuf(t *testing.T) { var resp internal.QueryResponse if err := proto.Unmarshal(w.Body.Bytes(), &resp); err != nil { t.Fatal(err) + } else if rt := resp.Results[0].Type; rt != pilosa.QueryResultTypeUint64 { + t.Fatalf("unexpected response type: %s", resp.Results[0].Type) } else if n := resp.Results[0].N; n != 100 { t.Fatalf("unexpected n: %d", n) } @@ -462,6 +464,8 @@ func TestHandler_Query_Bitmap_Protobuf(t *testing.T) { var resp internal.QueryResponse if err := proto.Unmarshal(w.Body.Bytes(), &resp); err != nil { t.Fatal(err) + } else if rt := resp.Results[0].Type; rt != pilosa.QueryResultTypeBitmap { + t.Fatalf("unexpected response type: %s", resp.Results[0].Type) } else if bits := resp.Results[0].Bitmap.Bits; !reflect.DeepEqual(bits, []uint64{1, SliceWidth + 1}) { t.Fatalf("unexpected bits: %+v", bits) } else if attrs := resp.Results[0].Bitmap.Attrs; len(attrs) != 3 { @@ -521,6 +525,8 @@ func TestHandler_Query_Bitmap_ColumnAttrs_Protobuf(t *testing.T) { } if bits := resp.Results[0].Bitmap.Bits; !reflect.DeepEqual(bits, []uint64{1, SliceWidth + 1}) { t.Fatalf("unexpected bits: %+v", bits) + } else if rt := resp.Results[0].Type; rt != pilosa.QueryResultTypeBitmap { + t.Fatalf("unexpected response type: %s", resp.Results[0].Type) } else if attrs := resp.Results[0].Bitmap.Attrs; len(attrs) != 3 { t.Fatalf("unexpected attr length: %d", len(attrs)) } else if k, v := attrs[0].Key, attrs[0].StringValue; k != "a" || v != "b" { @@ -592,6 +598,8 @@ func TestHandler_Query_Pairs_Protobuf(t *testing.T) { var resp internal.QueryResponse if err := proto.Unmarshal(w.Body.Bytes(), &resp); err != nil { t.Fatal(err) + } else if rt := resp.Results[0].Type; rt != pilosa.QueryResultTypePairs { + t.Fatalf("unexpected response type: %s", resp.Results[0].Type) } else if a := resp.Results[0].GetPairs(); len(a) != 2 { t.Fatalf("unexpected pair length: %d", len(a)) } diff --git a/internal/private.pb.go b/internal/private.pb.go index a0184a96f..7e4fca2ea 100644 --- a/internal/private.pb.go +++ b/internal/private.pb.go @@ -1,5 +1,6 @@ -// Code generated by protoc-gen-gogo. DO NOT EDIT. +// Code generated by protoc-gen-gogo. // source: private.proto +// DO NOT EDIT! /* Package internal is a generated protocol buffer package. @@ -2337,6 +2338,24 @@ func (m *Topology) MarshalTo(dAtA []byte) (int, error) { return i, nil } +func encodeFixed64Private(dAtA []byte, offset int, v uint64) int { + dAtA[offset] = uint8(v) + dAtA[offset+1] = uint8(v >> 8) + dAtA[offset+2] = uint8(v >> 16) + dAtA[offset+3] = uint8(v >> 24) + dAtA[offset+4] = uint8(v >> 32) + dAtA[offset+5] = uint8(v >> 40) + dAtA[offset+6] = uint8(v >> 48) + dAtA[offset+7] = uint8(v >> 56) + return offset + 8 +} +func encodeFixed32Private(dAtA []byte, offset int, v uint32) int { + dAtA[offset] = uint8(v) + dAtA[offset+1] = uint8(v >> 8) + dAtA[offset+2] = uint8(v >> 16) + dAtA[offset+3] = uint8(v >> 24) + return offset + 4 +} func encodeVarintPrivate(dAtA []byte, offset int, v uint64) int { for v >= 1<<7 { dAtA[offset] = uint8(v&0x7f | 0x80) @@ -3855,14 +3874,51 @@ func (m *MaxSlices) Unmarshal(dAtA []byte) error { if postIndex > l { return io.ErrUnexpectedEOF } + var keykey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + keykey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + var stringLenmapkey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLenmapkey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLenmapkey := int(stringLenmapkey) + if intStringLenmapkey < 0 { + return ErrInvalidLengthPrivate + } + postStringIndexmapkey := iNdEx + intStringLenmapkey + if postStringIndexmapkey > l { + return io.ErrUnexpectedEOF + } + mapkey := string(dAtA[iNdEx:postStringIndexmapkey]) + iNdEx = postStringIndexmapkey if m.Standard == nil { m.Standard = make(map[string]uint64) } - var mapkey string - var mapvalue uint64 - for iNdEx < postIndex { - entryPreIndex := iNdEx - var wire uint64 + if iNdEx < postIndex { + var valuekey uint64 for shift := uint(0); ; shift += 7 { if shift >= 64 { return ErrIntOverflowPrivate @@ -3872,69 +3928,31 @@ func (m *MaxSlices) Unmarshal(dAtA []byte) error { } b := dAtA[iNdEx] iNdEx++ - wire |= (uint64(b) & 0x7F) << shift + valuekey |= (uint64(b) & 0x7F) << shift if b < 0x80 { break } } - fieldNum := int32(wire >> 3) - if fieldNum == 1 { - var stringLenmapkey uint64 - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - stringLenmapkey |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } + var mapvalue uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate } - intStringLenmapkey := int(stringLenmapkey) - if intStringLenmapkey < 0 { - return ErrInvalidLengthPrivate - } - postStringIndexmapkey := iNdEx + intStringLenmapkey - if postStringIndexmapkey > l { + if iNdEx >= l { return io.ErrUnexpectedEOF } - mapkey = string(dAtA[iNdEx:postStringIndexmapkey]) - iNdEx = postStringIndexmapkey - } else if fieldNum == 2 { - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - mapvalue |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } + b := dAtA[iNdEx] + iNdEx++ + mapvalue |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break } - } else { - iNdEx = entryPreIndex - skippy, err := skipPrivate(dAtA[iNdEx:]) - if err != nil { - return err - } - if skippy < 0 { - return ErrInvalidLengthPrivate - } - if (iNdEx + skippy) > postIndex { - return io.ErrUnexpectedEOF - } - iNdEx += skippy } + m.Standard[mapkey] = mapvalue + } else { + var mapvalue uint64 + m.Standard[mapkey] = mapvalue } - m.Standard[mapkey] = mapvalue iNdEx = postIndex case 2: if wireType != 2 { @@ -3962,14 +3980,51 @@ func (m *MaxSlices) Unmarshal(dAtA []byte) error { if postIndex > l { return io.ErrUnexpectedEOF } + var keykey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + keykey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + var stringLenmapkey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLenmapkey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLenmapkey := int(stringLenmapkey) + if intStringLenmapkey < 0 { + return ErrInvalidLengthPrivate + } + postStringIndexmapkey := iNdEx + intStringLenmapkey + if postStringIndexmapkey > l { + return io.ErrUnexpectedEOF + } + mapkey := string(dAtA[iNdEx:postStringIndexmapkey]) + iNdEx = postStringIndexmapkey if m.Inverse == nil { m.Inverse = make(map[string]uint64) } - var mapkey string - var mapvalue uint64 - for iNdEx < postIndex { - entryPreIndex := iNdEx - var wire uint64 + if iNdEx < postIndex { + var valuekey uint64 for shift := uint(0); ; shift += 7 { if shift >= 64 { return ErrIntOverflowPrivate @@ -3979,69 +4034,31 @@ func (m *MaxSlices) Unmarshal(dAtA []byte) error { } b := dAtA[iNdEx] iNdEx++ - wire |= (uint64(b) & 0x7F) << shift + valuekey |= (uint64(b) & 0x7F) << shift if b < 0x80 { break } } - fieldNum := int32(wire >> 3) - if fieldNum == 1 { - var stringLenmapkey uint64 - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - stringLenmapkey |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } + var mapvalue uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate } - intStringLenmapkey := int(stringLenmapkey) - if intStringLenmapkey < 0 { - return ErrInvalidLengthPrivate - } - postStringIndexmapkey := iNdEx + intStringLenmapkey - if postStringIndexmapkey > l { + if iNdEx >= l { return io.ErrUnexpectedEOF } - mapkey = string(dAtA[iNdEx:postStringIndexmapkey]) - iNdEx = postStringIndexmapkey - } else if fieldNum == 2 { - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - mapvalue |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } + b := dAtA[iNdEx] + iNdEx++ + mapvalue |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break } - } else { - iNdEx = entryPreIndex - skippy, err := skipPrivate(dAtA[iNdEx:]) - if err != nil { - return err - } - if skippy < 0 { - return ErrInvalidLengthPrivate - } - if (iNdEx + skippy) > postIndex { - return io.ErrUnexpectedEOF - } - iNdEx += skippy } + m.Inverse[mapkey] = mapvalue + } else { + var mapvalue uint64 + m.Inverse[mapkey] = mapvalue } - m.Inverse[mapkey] = mapvalue iNdEx = postIndex default: iNdEx = preIndex @@ -5369,14 +5386,51 @@ func (m *InputDefinitionAction) Unmarshal(dAtA []byte) error { if postIndex > l { return io.ErrUnexpectedEOF } + var keykey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + keykey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + var stringLenmapkey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLenmapkey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLenmapkey := int(stringLenmapkey) + if intStringLenmapkey < 0 { + return ErrInvalidLengthPrivate + } + postStringIndexmapkey := iNdEx + intStringLenmapkey + if postStringIndexmapkey > l { + return io.ErrUnexpectedEOF + } + mapkey := string(dAtA[iNdEx:postStringIndexmapkey]) + iNdEx = postStringIndexmapkey if m.ValueMap == nil { m.ValueMap = make(map[string]uint64) } - var mapkey string - var mapvalue uint64 - for iNdEx < postIndex { - entryPreIndex := iNdEx - var wire uint64 + if iNdEx < postIndex { + var valuekey uint64 for shift := uint(0); ; shift += 7 { if shift >= 64 { return ErrIntOverflowPrivate @@ -5386,69 +5440,31 @@ func (m *InputDefinitionAction) Unmarshal(dAtA []byte) error { } b := dAtA[iNdEx] iNdEx++ - wire |= (uint64(b) & 0x7F) << shift + valuekey |= (uint64(b) & 0x7F) << shift if b < 0x80 { break } } - fieldNum := int32(wire >> 3) - if fieldNum == 1 { - var stringLenmapkey uint64 - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - stringLenmapkey |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } + var mapvalue uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate } - intStringLenmapkey := int(stringLenmapkey) - if intStringLenmapkey < 0 { - return ErrInvalidLengthPrivate - } - postStringIndexmapkey := iNdEx + intStringLenmapkey - if postStringIndexmapkey > l { + if iNdEx >= l { return io.ErrUnexpectedEOF } - mapkey = string(dAtA[iNdEx:postStringIndexmapkey]) - iNdEx = postStringIndexmapkey - } else if fieldNum == 2 { - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - mapvalue |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } + b := dAtA[iNdEx] + iNdEx++ + mapvalue |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break } - } else { - iNdEx = entryPreIndex - skippy, err := skipPrivate(dAtA[iNdEx:]) - if err != nil { - return err - } - if skippy < 0 { - return ErrInvalidLengthPrivate - } - if (iNdEx + skippy) > postIndex { - return io.ErrUnexpectedEOF - } - iNdEx += skippy } + m.ValueMap[mapkey] = mapvalue + } else { + var mapvalue uint64 + m.ValueMap[mapkey] = mapvalue } - m.ValueMap[mapkey] = mapvalue iNdEx = postIndex case 4: if wireType != 0 { diff --git a/internal/public.pb.go b/internal/public.pb.go index 22afdc0f9..dfb9a819e 100644 --- a/internal/public.pb.go +++ b/internal/public.pb.go @@ -1,5 +1,6 @@ -// Code generated by protoc-gen-gogo. DO NOT EDIT. +// Code generated by protoc-gen-gogo. // source: public.proto +// DO NOT EDIT! /* Package internal is a generated protocol buffer package. @@ -27,8 +28,6 @@ import proto "github.com/golang/protobuf/proto" import fmt "fmt" import math "math" -import encoding_binary "encoding/binary" - import io "io" // Reference imports to suppress errors if they are not otherwise used. @@ -355,6 +354,7 @@ func (m *QueryResponse) GetColumnAttrSets() []*ColumnAttrSet { } type QueryResult struct { + Type uint32 `protobuf:"varint,6,opt,name=Type,proto3" json:"Type,omitempty"` Bitmap *Bitmap `protobuf:"bytes,1,opt,name=Bitmap" json:"Bitmap,omitempty"` N uint64 `protobuf:"varint,2,opt,name=N,proto3" json:"N,omitempty"` Pairs []*Pair `protobuf:"bytes,3,rep,name=Pairs" json:"Pairs,omitempty"` @@ -367,6 +367,13 @@ func (m *QueryResult) String() string { return proto.CompactTextStrin func (*QueryResult) ProtoMessage() {} func (*QueryResult) Descriptor() ([]byte, []int) { return fileDescriptorPublic, []int{9} } +func (m *QueryResult) GetType() uint32 { + if m != nil { + return m.Type + } + return 0 +} + func (m *QueryResult) GetBitmap() *Bitmap { if m != nil { return m.Bitmap @@ -408,6 +415,8 @@ type ImportRequest struct { Slice uint64 `protobuf:"varint,3,opt,name=Slice,proto3" json:"Slice,omitempty"` RowIDs []uint64 `protobuf:"varint,4,rep,packed,name=RowIDs" json:"RowIDs,omitempty"` ColumnIDs []uint64 `protobuf:"varint,5,rep,packed,name=ColumnIDs" json:"ColumnIDs,omitempty"` + RowKeys []string `protobuf:"bytes,7,rep,name=RowKeys" json:"RowKeys,omitempty"` + ColumnKeys []string `protobuf:"bytes,8,rep,name=ColumnKeys" json:"ColumnKeys,omitempty"` Timestamps []int64 `protobuf:"varint,6,rep,packed,name=Timestamps" json:"Timestamps,omitempty"` } @@ -451,6 +460,20 @@ func (m *ImportRequest) GetColumnIDs() []uint64 { return nil } +func (m *ImportRequest) GetRowKeys() []string { + if m != nil { + return m.RowKeys + } + return nil +} + +func (m *ImportRequest) GetColumnKeys() []string { + if m != nil { + return m.ColumnKeys + } + return nil +} + func (m *ImportRequest) GetTimestamps() []int64 { if m != nil { return m.Timestamps @@ -459,12 +482,13 @@ func (m *ImportRequest) GetTimestamps() []int64 { } type ImportValueRequest struct { - Index string `protobuf:"bytes,1,opt,name=Index,proto3" json:"Index,omitempty"` - Frame string `protobuf:"bytes,2,opt,name=Frame,proto3" json:"Frame,omitempty"` - Slice uint64 `protobuf:"varint,3,opt,name=Slice,proto3" json:"Slice,omitempty"` - Field string `protobuf:"bytes,4,opt,name=Field,proto3" json:"Field,omitempty"` - ColumnIDs []uint64 `protobuf:"varint,5,rep,packed,name=ColumnIDs" json:"ColumnIDs,omitempty"` - Values []int64 `protobuf:"varint,6,rep,packed,name=Values" json:"Values,omitempty"` + Index string `protobuf:"bytes,1,opt,name=Index,proto3" json:"Index,omitempty"` + Frame string `protobuf:"bytes,2,opt,name=Frame,proto3" json:"Frame,omitempty"` + Slice uint64 `protobuf:"varint,3,opt,name=Slice,proto3" json:"Slice,omitempty"` + Field string `protobuf:"bytes,4,opt,name=Field,proto3" json:"Field,omitempty"` + ColumnIDs []uint64 `protobuf:"varint,5,rep,packed,name=ColumnIDs" json:"ColumnIDs,omitempty"` + ColumnKeys []string `protobuf:"bytes,7,rep,name=ColumnKeys" json:"ColumnKeys,omitempty"` + Values []int64 `protobuf:"varint,6,rep,packed,name=Values" json:"Values,omitempty"` } func (m *ImportValueRequest) Reset() { *m = ImportValueRequest{} } @@ -507,6 +531,13 @@ func (m *ImportValueRequest) GetColumnIDs() []uint64 { return nil } +func (m *ImportValueRequest) GetColumnKeys() []string { + if m != nil { + return m.ColumnKeys + } + return nil +} + func (m *ImportValueRequest) GetValues() []int64 { if m != nil { return m.Values @@ -776,8 +807,7 @@ func (m *Attr) MarshalTo(dAtA []byte) (int, error) { if m.FloatValue != 0 { dAtA[i] = 0x31 i++ - encoding_binary.LittleEndian.PutUint64(dAtA[i:], uint64(math.Float64bits(float64(m.FloatValue)))) - i += 8 + i = encodeFixed64Public(dAtA, i, uint64(math.Float64bits(float64(m.FloatValue)))) } return i, nil } @@ -1003,6 +1033,11 @@ func (m *QueryResult) MarshalTo(dAtA []byte) (int, error) { } i += n6 } + if m.Type != 0 { + dAtA[i] = 0x30 + i++ + i = encodeVarintPublic(dAtA, i, uint64(m.Type)) + } return i, nil } @@ -1090,6 +1125,36 @@ func (m *ImportRequest) MarshalTo(dAtA []byte) (int, error) { i = encodeVarintPublic(dAtA, i, uint64(j11)) i += copy(dAtA[i:], dAtA12[:j11]) } + if len(m.RowKeys) > 0 { + for _, s := range m.RowKeys { + dAtA[i] = 0x3a + i++ + l = len(s) + for l >= 1<<7 { + dAtA[i] = uint8(uint64(l)&0x7f | 0x80) + l >>= 7 + i++ + } + dAtA[i] = uint8(l) + i++ + i += copy(dAtA[i:], s) + } + } + if len(m.ColumnKeys) > 0 { + for _, s := range m.ColumnKeys { + dAtA[i] = 0x42 + i++ + l = len(s) + for l >= 1<<7 { + dAtA[i] = uint8(uint64(l)&0x7f | 0x80) + l >>= 7 + i++ + } + dAtA[i] = uint8(l) + i++ + i += copy(dAtA[i:], s) + } + } return i, nil } @@ -1166,9 +1231,42 @@ func (m *ImportValueRequest) MarshalTo(dAtA []byte) (int, error) { i = encodeVarintPublic(dAtA, i, uint64(j15)) i += copy(dAtA[i:], dAtA16[:j15]) } + if len(m.ColumnKeys) > 0 { + for _, s := range m.ColumnKeys { + dAtA[i] = 0x3a + i++ + l = len(s) + for l >= 1<<7 { + dAtA[i] = uint8(uint64(l)&0x7f | 0x80) + l >>= 7 + i++ + } + dAtA[i] = uint8(l) + i++ + i += copy(dAtA[i:], s) + } + } return i, nil } +func encodeFixed64Public(dAtA []byte, offset int, v uint64) int { + dAtA[offset] = uint8(v) + dAtA[offset+1] = uint8(v >> 8) + dAtA[offset+2] = uint8(v >> 16) + dAtA[offset+3] = uint8(v >> 24) + dAtA[offset+4] = uint8(v >> 32) + dAtA[offset+5] = uint8(v >> 40) + dAtA[offset+6] = uint8(v >> 48) + dAtA[offset+7] = uint8(v >> 56) + return offset + 8 +} +func encodeFixed32Public(dAtA []byte, offset int, v uint32) int { + dAtA[offset] = uint8(v) + dAtA[offset+1] = uint8(v >> 8) + dAtA[offset+2] = uint8(v >> 16) + dAtA[offset+3] = uint8(v >> 24) + return offset + 4 +} func encodeVarintPublic(dAtA []byte, offset int, v uint64) int { for v >= 1<<7 { dAtA[offset] = uint8(v&0x7f | 0x80) @@ -1377,6 +1475,9 @@ func (m *QueryResult) Size() (n int) { l = m.SumCount.Size() n += 1 + l + sovPublic(uint64(l)) } + if m.Type != 0 { + n += 1 + sovPublic(uint64(m.Type)) + } return n } @@ -1415,6 +1516,18 @@ func (m *ImportRequest) Size() (n int) { } n += 1 + sovPublic(uint64(l)) + l } + if len(m.RowKeys) > 0 { + for _, s := range m.RowKeys { + l = len(s) + n += 1 + l + sovPublic(uint64(l)) + } + } + if len(m.ColumnKeys) > 0 { + for _, s := range m.ColumnKeys { + l = len(s) + n += 1 + l + sovPublic(uint64(l)) + } + } return n } @@ -1450,6 +1563,12 @@ func (m *ImportValueRequest) Size() (n int) { } n += 1 + sovPublic(uint64(l)) + l } + if len(m.ColumnKeys) > 0 { + for _, s := range m.ColumnKeys { + l = len(s) + n += 1 + l + sovPublic(uint64(l)) + } + } return n } @@ -2232,8 +2351,15 @@ func (m *Attr) Unmarshal(dAtA []byte) error { if (iNdEx + 8) > l { return io.ErrUnexpectedEOF } - v = uint64(encoding_binary.LittleEndian.Uint64(dAtA[iNdEx:])) iNdEx += 8 + v = uint64(dAtA[iNdEx-8]) + v |= uint64(dAtA[iNdEx-7]) << 8 + v |= uint64(dAtA[iNdEx-6]) << 16 + v |= uint64(dAtA[iNdEx-5]) << 24 + v |= uint64(dAtA[iNdEx-4]) << 32 + v |= uint64(dAtA[iNdEx-3]) << 40 + v |= uint64(dAtA[iNdEx-2]) << 48 + v |= uint64(dAtA[iNdEx-1]) << 56 m.FloatValue = float64(math.Float64frombits(v)) default: iNdEx = preIndex @@ -2864,6 +2990,25 @@ func (m *QueryResult) Unmarshal(dAtA []byte) error { return err } iNdEx = postIndex + case 6: + if wireType != 0 { + return fmt.Errorf("proto: wrong wireType = %d for field Type", wireType) + } + m.Type = 0 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPublic + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + m.Type |= (uint32(b) & 0x7F) << shift + if b < 0x80 { + break + } + } default: iNdEx = preIndex skippy, err := skipPublic(dAtA[iNdEx:]) @@ -3177,6 +3322,64 @@ func (m *ImportRequest) Unmarshal(dAtA []byte) error { } else { return fmt.Errorf("proto: wrong wireType = %d for field Timestamps", wireType) } + case 7: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field RowKeys", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPublic + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLen |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLen := int(stringLen) + if intStringLen < 0 { + return ErrInvalidLengthPublic + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.RowKeys = append(m.RowKeys, string(dAtA[iNdEx:postIndex])) + iNdEx = postIndex + case 8: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field ColumnKeys", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPublic + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLen |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLen := int(stringLen) + if intStringLen < 0 { + return ErrInvalidLengthPublic + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.ColumnKeys = append(m.ColumnKeys, string(dAtA[iNdEx:postIndex])) + iNdEx = postIndex default: iNdEx = preIndex skippy, err := skipPublic(dAtA[iNdEx:]) @@ -3457,6 +3660,35 @@ func (m *ImportValueRequest) Unmarshal(dAtA []byte) error { } else { return fmt.Errorf("proto: wrong wireType = %d for field Values", wireType) } + case 7: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field ColumnKeys", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPublic + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLen |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLen := int(stringLen) + if intStringLen < 0 { + return ErrInvalidLengthPublic + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.ColumnKeys = append(m.ColumnKeys, string(dAtA[iNdEx:postIndex])) + iNdEx = postIndex default: iNdEx = preIndex skippy, err := skipPublic(dAtA[iNdEx:]) @@ -3586,47 +3818,50 @@ var ( func init() { proto.RegisterFile("public.proto", fileDescriptorPublic) } var fileDescriptorPublic = []byte{ - // 671 bytes of a gzipped FileDescriptorProto - 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xac, 0x54, 0xbb, 0x6e, 0xd4, 0x40, - 0x14, 0x65, 0xd6, 0xde, 0xd7, 0xdd, 0x4d, 0x14, 0x8d, 0x20, 0x58, 0x08, 0xad, 0x2c, 0x8b, 0xc2, - 0xd5, 0x46, 0x5a, 0x7a, 0x10, 0x9b, 0x87, 0x64, 0x45, 0x44, 0x70, 0x37, 0x84, 0xda, 0x49, 0x46, - 0xc1, 0x92, 0x5f, 0xd8, 0x63, 0x91, 0xfd, 0x0e, 0x1a, 0x6a, 0x1a, 0xf8, 0x01, 0x3a, 0x3e, 0x80, - 0x92, 0x4f, 0x40, 0xe1, 0x47, 0xd0, 0x9d, 0xf1, 0xd8, 0x5e, 0x22, 0x01, 0x05, 0xdd, 0x9c, 0x73, - 0x66, 0xae, 0xef, 0xe3, 0x5c, 0xc3, 0x34, 0xaf, 0xce, 0xe3, 0xe8, 0x62, 0x9e, 0x17, 0x99, 0xcc, - 0xf8, 0x28, 0x4a, 0xa5, 0x28, 0xd2, 0x30, 0xf6, 0xce, 0x60, 0xb0, 0x8c, 0x64, 0x12, 0xe6, 0x9c, - 0x83, 0xbd, 0x8c, 0x64, 0xe9, 0x30, 0xd7, 0xf2, 0x6d, 0x54, 0x67, 0xfe, 0x08, 0xfa, 0xcf, 0xa4, - 0x2c, 0x4a, 0xa7, 0xe7, 0x5a, 0xfe, 0x64, 0xb1, 0x3d, 0x37, 0xef, 0xe6, 0x44, 0xa3, 0x16, 0xe9, - 0xe5, 0xb1, 0x58, 0x97, 0x8e, 0xe5, 0x5a, 0xfe, 0x18, 0xd5, 0xd9, 0x7b, 0x02, 0xf6, 0x8b, 0x30, - 0x2a, 0xf8, 0x36, 0xf4, 0x82, 0x03, 0x87, 0xb9, 0xcc, 0xb7, 0xb1, 0x17, 0x1c, 0xf0, 0xbb, 0xd0, - 0xdf, 0xcf, 0xaa, 0x54, 0x3a, 0x3d, 0x45, 0x69, 0xc0, 0x77, 0xc0, 0x3a, 0x16, 0x6b, 0xc7, 0x72, - 0x99, 0x3f, 0x46, 0x3a, 0x7a, 0x0b, 0x18, 0xad, 0xaa, 0xa4, 0x51, 0x57, 0x55, 0xa2, 0x82, 0x58, - 0x48, 0xc7, 0xcd, 0x28, 0x56, 0x1d, 0xc5, 0x7b, 0x05, 0xd6, 0x32, 0x92, 0x24, 0x62, 0xf6, 0xae, - 0xf9, 0xaa, 0x06, 0xfc, 0x01, 0x8c, 0xf6, 0xb3, 0xb8, 0x4a, 0xd2, 0xe0, 0xa0, 0xfe, 0x76, 0x83, - 0xf9, 0x43, 0x18, 0x9f, 0x46, 0x89, 0x28, 0x65, 0x98, 0xe4, 0x2a, 0x09, 0x0b, 0x5b, 0xc2, 0x7b, - 0x0d, 0x5b, 0xfa, 0x26, 0x55, 0xbb, 0x12, 0xf2, 0x56, 0x4d, 0xff, 0xd6, 0xa5, 0xdb, 0x35, 0x7e, - 0x66, 0x60, 0x93, 0x66, 0x24, 0xd6, 0x48, 0xd4, 0xd2, 0xd3, 0x75, 0x2e, 0xea, 0x4c, 0xd5, 0x99, - 0xbb, 0x30, 0x59, 0xc9, 0x22, 0x4a, 0xaf, 0xce, 0xc2, 0xb8, 0x12, 0x75, 0xa0, 0x2e, 0x45, 0x35, - 0x06, 0xa9, 0xd4, 0xb2, 0xad, 0xca, 0x68, 0x30, 0xd5, 0xb8, 0xcc, 0xb2, 0x58, 0x8b, 0x7d, 0x97, - 0xf9, 0x23, 0x6c, 0x09, 0x3e, 0x03, 0x38, 0x8a, 0xb3, 0xb0, 0x7e, 0x3b, 0x70, 0x99, 0xcf, 0xb0, - 0xc3, 0x78, 0x7b, 0x30, 0xa4, 0x4c, 0x9f, 0x87, 0x79, 0x5b, 0x2d, 0xfb, 0x43, 0xb5, 0xde, 0x57, - 0x06, 0xd3, 0x97, 0x95, 0x28, 0xd6, 0x28, 0xde, 0x56, 0xa2, 0x54, 0x53, 0x51, 0xb8, 0xae, 0x52, - 0x03, 0xbe, 0x0b, 0x83, 0x55, 0x1c, 0x5d, 0x08, 0xdd, 0x3b, 0x1b, 0x6b, 0x44, 0xb5, 0xb6, 0x3d, - 0x2f, 0x55, 0xad, 0x23, 0xec, 0x52, 0xf4, 0x12, 0x45, 0x92, 0x49, 0x53, 0x4c, 0x8d, 0xb8, 0x07, - 0xd3, 0xc3, 0xeb, 0x8b, 0xb8, 0xba, 0x14, 0xfa, 0xe9, 0x40, 0xa9, 0x1b, 0x1c, 0x45, 0xaf, 0xb1, - 0x72, 0xfc, 0x50, 0x47, 0xef, 0x50, 0xde, 0x7b, 0x06, 0x5b, 0x75, 0xfa, 0x65, 0x9e, 0xa5, 0xa5, - 0xa0, 0x19, 0x1d, 0x16, 0x85, 0x99, 0xd1, 0x61, 0x51, 0xf0, 0x3d, 0x18, 0xa2, 0x28, 0xab, 0x58, - 0x9a, 0xc1, 0xdf, 0x6b, 0x5b, 0x61, 0xde, 0x56, 0xb1, 0x44, 0x73, 0x8b, 0x3f, 0x85, 0xed, 0x0d, - 0x23, 0xe9, 0x8d, 0x99, 0x2c, 0xee, 0xb7, 0xef, 0x36, 0x74, 0xfc, 0xed, 0xba, 0xf7, 0x85, 0xc1, - 0xa4, 0x13, 0x99, 0xfb, 0x66, 0x79, 0x55, 0x5a, 0x93, 0xc5, 0x4e, 0x1b, 0x48, 0xf3, 0x68, 0x96, - 0x7b, 0x0a, 0xec, 0xa4, 0x36, 0x13, 0x3b, 0xa1, 0x11, 0xd2, 0x72, 0x9a, 0xef, 0x77, 0x46, 0x48, - 0x34, 0x6a, 0x91, 0x3b, 0x30, 0xdc, 0x7f, 0x13, 0xa6, 0x57, 0xe2, 0x52, 0x99, 0x69, 0x84, 0x06, - 0xf2, 0x79, 0xbb, 0x9c, 0xaa, 0xfb, 0x93, 0x05, 0x6f, 0x43, 0x18, 0x05, 0x9b, 0x3b, 0xde, 0x27, - 0x06, 0x5b, 0x41, 0x92, 0x67, 0x85, 0xec, 0xb8, 0x21, 0x48, 0x2f, 0xc5, 0xb5, 0x71, 0x83, 0x02, - 0xc4, 0x1e, 0x15, 0x61, 0xa2, 0x6d, 0x3f, 0x46, 0x0d, 0x88, 0x55, 0xae, 0x50, 0x2e, 0xb0, 0x51, - 0x03, 0x35, 0x7f, 0x5a, 0xec, 0xd2, 0xb1, 0xb5, 0x73, 0x34, 0x22, 0x9f, 0x9b, 0xbd, 0x2e, 0x9d, - 0xbe, 0x92, 0x5a, 0x82, 0x7c, 0xde, 0x2c, 0x36, 0x79, 0xc3, 0xf2, 0x2d, 0xec, 0x30, 0xde, 0x47, - 0x06, 0x5c, 0x67, 0xaa, 0x7c, 0xff, 0xff, 0xd2, 0xa5, 0xbb, 0x91, 0x88, 0x75, 0x2b, 0xe9, 0x2e, - 0x81, 0xbf, 0x24, 0xbb, 0x0b, 0x03, 0x95, 0x85, 0x49, 0xb4, 0x46, 0xcb, 0x9d, 0x6f, 0x37, 0x33, - 0xf6, 0xfd, 0x66, 0xc6, 0x7e, 0xdc, 0xcc, 0xd8, 0x87, 0x9f, 0xb3, 0x3b, 0xe7, 0x03, 0xf5, 0x5b, - 0x7f, 0xfc, 0x2b, 0x00, 0x00, 0xff, 0xff, 0x0f, 0xf2, 0x1f, 0x86, 0xe6, 0x05, 0x00, 0x00, + // 705 bytes of a gzipped FileDescriptorProto + 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xac, 0x55, 0xcb, 0x6e, 0xd3, 0x40, + 0x14, 0x65, 0x62, 0x27, 0x71, 0x6e, 0x92, 0xaa, 0x1a, 0x41, 0xb1, 0x10, 0x8a, 0x2c, 0x8b, 0x85, + 0x57, 0xa9, 0x14, 0xf6, 0x20, 0xd2, 0x87, 0x14, 0x55, 0x54, 0x30, 0x29, 0x65, 0xed, 0xb6, 0xa3, + 0x62, 0xc9, 0x2f, 0xec, 0xb1, 0xda, 0x7c, 0x07, 0x1b, 0x3e, 0x81, 0x8f, 0x60, 0xc5, 0x0a, 0x76, + 0x7c, 0x02, 0x94, 0x1f, 0x41, 0xf7, 0x8e, 0x27, 0x76, 0x5a, 0x09, 0x58, 0xb0, 0x9b, 0x73, 0xce, + 0xcc, 0xf5, 0x9c, 0xb9, 0xe7, 0x26, 0x30, 0xca, 0xab, 0xb3, 0x38, 0x3a, 0x9f, 0xe6, 0x45, 0xa6, + 0x32, 0xee, 0x44, 0xa9, 0x92, 0x45, 0x1a, 0xc6, 0xfe, 0x29, 0xf4, 0xe6, 0x91, 0x4a, 0xc2, 0x9c, + 0x73, 0xb0, 0xe7, 0x91, 0x2a, 0x5d, 0xe6, 0x59, 0x81, 0x2d, 0x68, 0xcd, 0x9f, 0x40, 0xf7, 0x85, + 0x52, 0x45, 0xe9, 0x76, 0x3c, 0x2b, 0x18, 0xce, 0xb6, 0xa6, 0xe6, 0xdc, 0x14, 0x69, 0xa1, 0x45, + 0x3c, 0x79, 0x24, 0x57, 0xa5, 0x6b, 0x79, 0x56, 0x30, 0x10, 0xb4, 0xf6, 0x9f, 0x81, 0xfd, 0x2a, + 0x8c, 0x0a, 0xbe, 0x05, 0x9d, 0xc5, 0xbe, 0xcb, 0x3c, 0x16, 0xd8, 0xa2, 0xb3, 0xd8, 0xe7, 0xf7, + 0xa1, 0xbb, 0x97, 0x55, 0xa9, 0x72, 0x3b, 0x44, 0x69, 0xc0, 0xb7, 0xc1, 0x3a, 0x92, 0x2b, 0xd7, + 0xf2, 0x58, 0x30, 0x10, 0xb8, 0xf4, 0x67, 0xe0, 0x2c, 0xab, 0x64, 0xad, 0x2e, 0xab, 0x84, 0x8a, + 0x58, 0x02, 0x97, 0x9b, 0x55, 0xac, 0xba, 0x8a, 0xff, 0x06, 0xac, 0x79, 0xa4, 0x50, 0x14, 0xd9, + 0xd5, 0xfa, 0xab, 0x1a, 0xf0, 0x47, 0xe0, 0xec, 0x65, 0x71, 0x95, 0xa4, 0x8b, 0xfd, 0xfa, 0xdb, + 0x6b, 0xcc, 0x1f, 0xc3, 0xe0, 0x24, 0x4a, 0x64, 0xa9, 0xc2, 0x24, 0xa7, 0x4b, 0x58, 0xa2, 0x21, + 0xfc, 0xb7, 0x30, 0xd6, 0x3b, 0xd1, 0xed, 0x52, 0xaa, 0x3b, 0x9e, 0xfe, 0xed, 0x95, 0xee, 0x7a, + 0xfc, 0xc4, 0xc0, 0x46, 0xcd, 0x48, 0x6c, 0x2d, 0xe1, 0x93, 0x9e, 0xac, 0x72, 0x59, 0xdf, 0x94, + 0xd6, 0xdc, 0x83, 0xe1, 0x52, 0x15, 0x51, 0x7a, 0x79, 0x1a, 0xc6, 0x95, 0xac, 0x0b, 0xb5, 0x29, + 0xf4, 0xb8, 0x48, 0x95, 0x96, 0x6d, 0xb2, 0xb1, 0xc6, 0xe8, 0x71, 0x9e, 0x65, 0xb1, 0x16, 0xbb, + 0x1e, 0x0b, 0x1c, 0xd1, 0x10, 0x7c, 0x02, 0x70, 0x18, 0x67, 0x61, 0x7d, 0xb6, 0xe7, 0xb1, 0x80, + 0x89, 0x16, 0xe3, 0xef, 0x42, 0x1f, 0x6f, 0xfa, 0x32, 0xcc, 0x1b, 0xb7, 0xec, 0x0f, 0x6e, 0xfd, + 0xcf, 0x0c, 0x46, 0xaf, 0x2b, 0x59, 0xac, 0x84, 0x7c, 0x5f, 0xc9, 0x92, 0xba, 0x42, 0xb8, 0x76, + 0xa9, 0x01, 0xdf, 0x81, 0xde, 0x32, 0x8e, 0xce, 0xa5, 0x7e, 0x3b, 0x5b, 0xd4, 0x08, 0xbd, 0x36, + 0x6f, 0x5e, 0x92, 0x57, 0x47, 0xb4, 0x29, 0x3c, 0x29, 0x64, 0x92, 0x29, 0x63, 0xa6, 0x46, 0xdc, + 0x87, 0xd1, 0xc1, 0xf5, 0x79, 0x5c, 0x5d, 0x48, 0x7d, 0xb4, 0x47, 0xea, 0x06, 0x87, 0xd5, 0x6b, + 0x4c, 0x89, 0xef, 0xeb, 0xea, 0x2d, 0xca, 0xff, 0xc0, 0x60, 0x5c, 0x5f, 0xbf, 0xcc, 0xb3, 0xb4, + 0x94, 0xd8, 0xa3, 0x83, 0xa2, 0x30, 0x3d, 0x3a, 0x28, 0x0a, 0xbe, 0x0b, 0x7d, 0x21, 0xcb, 0x2a, + 0x56, 0xa6, 0xf1, 0x0f, 0x9a, 0xa7, 0x30, 0x67, 0xab, 0x58, 0x09, 0xb3, 0x8b, 0x3f, 0x87, 0xad, + 0x8d, 0x20, 0xe9, 0x89, 0x19, 0xce, 0x1e, 0x36, 0xe7, 0x36, 0x74, 0x71, 0x6b, 0xbb, 0xff, 0x8d, + 0xc1, 0xb0, 0x55, 0x99, 0x07, 0x66, 0x78, 0xe9, 0x5a, 0xc3, 0xd9, 0x76, 0x53, 0x48, 0xf3, 0xc2, + 0x0c, 0xf7, 0x08, 0xd8, 0x71, 0x1d, 0x26, 0x76, 0x8c, 0x2d, 0xc4, 0xe1, 0x34, 0xdf, 0x6f, 0xb5, + 0x10, 0x69, 0xa1, 0x45, 0xee, 0x42, 0x7f, 0xef, 0x5d, 0x98, 0x5e, 0xca, 0x0b, 0x0a, 0x93, 0x23, + 0x0c, 0xe4, 0xd3, 0x66, 0x38, 0xe9, 0xf5, 0x87, 0x33, 0xde, 0x94, 0x30, 0x8a, 0x68, 0x06, 0xd8, + 0xa4, 0x19, 0x7b, 0x31, 0xd6, 0x69, 0xf6, 0x7f, 0x32, 0x18, 0x2f, 0x92, 0x3c, 0x2b, 0x54, 0x2b, + 0x21, 0x8b, 0xf4, 0x42, 0x5e, 0x9b, 0x84, 0x10, 0x40, 0xf6, 0xb0, 0x08, 0x13, 0x3d, 0x0a, 0x03, + 0xa1, 0x01, 0xb2, 0x94, 0x14, 0x4a, 0x86, 0x2d, 0x34, 0xa0, 0x4c, 0xe0, 0xb0, 0x97, 0xae, 0xad, + 0xd3, 0xa4, 0x11, 0x66, 0xdf, 0xcc, 0x7a, 0xe9, 0x76, 0x49, 0x6a, 0x08, 0xcc, 0xfe, 0x7a, 0xd8, + 0x31, 0x2f, 0x56, 0x60, 0x89, 0x16, 0x83, 0xef, 0x20, 0xb2, 0x2b, 0xfa, 0x85, 0xeb, 0xd3, 0x2f, + 0x9c, 0x81, 0x78, 0x52, 0x97, 0x21, 0xd1, 0x21, 0xb1, 0xc5, 0xf8, 0x5f, 0x18, 0x70, 0xed, 0x91, + 0xa6, 0xe8, 0xff, 0x19, 0xc5, 0xbd, 0x91, 0x8c, 0x75, 0x63, 0x70, 0x2f, 0x82, 0xbf, 0xd8, 0xdc, + 0x81, 0x1e, 0xdd, 0xc2, 0x58, 0xac, 0xd1, 0x2d, 0x13, 0xfd, 0xdb, 0x26, 0xe6, 0xdb, 0x5f, 0x6f, + 0x26, 0xec, 0xfb, 0xcd, 0x84, 0xfd, 0xb8, 0x99, 0xb0, 0x8f, 0xbf, 0x26, 0xf7, 0xce, 0x7a, 0xf4, + 0x27, 0xf2, 0xf4, 0x77, 0x00, 0x00, 0x00, 0xff, 0xff, 0xa3, 0xa0, 0xd2, 0x51, 0x54, 0x06, 0x00, + 0x00, } diff --git a/internal/public.proto b/internal/public.proto index 025b50baa..5fa40cfab 100644 --- a/internal/public.proto +++ b/internal/public.proto @@ -41,7 +41,7 @@ message Attr { } message AttrMap { - repeated Attr Attrs = 1; + repeated Attr Attrs = 1; } message QueryRequest { @@ -60,6 +60,7 @@ message QueryResponse { } message QueryResult { + uint32 Type = 6; Bitmap Bitmap = 1; uint64 N = 2; repeated Pair Pairs = 3; @@ -73,6 +74,8 @@ message ImportRequest { uint64 Slice = 3; repeated uint64 RowIDs = 4; repeated uint64 ColumnIDs = 5; + repeated string RowKeys = 7; + repeated string ColumnKeys = 8; repeated int64 Timestamps = 6; } @@ -82,5 +85,6 @@ message ImportValueRequest { uint64 Slice = 3; string Field = 4; repeated uint64 ColumnIDs = 5; + repeated string ColumnKeys = 7; repeated int64 Values = 6; } diff --git a/server.go b/server.go index 94f823345..aeb1e9519 100644 --- a/server.go +++ b/server.go @@ -510,7 +510,7 @@ func (s *Server) ClusterStatus() (proto.Message, error) { // HandleRemoteStatus receives incoming NodeStatus from remote nodes. func (s *Server) HandleRemoteStatus(pb proto.Message) error { // Ignore NodeStatus messages until the cluster is in a Normal state. - if s.Cluster.State != ClusterStateNormal { + if s.Cluster.State() != ClusterStateNormal { return nil } diff --git a/server/cluster_test.go b/server/cluster_test.go index fdf6fdc88..98b723335 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -24,15 +24,16 @@ import ( "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/gossip" + "github.com/pilosa/pilosa/test" ) // Ensure program can send/receive broadcast messages. func TestMain_SendReceiveMessage(t *testing.T) { - m0 := MustRunMain() + m0 := test.MustRunMain() defer m0.Close() - m1 := MustRunMain() + m1 := test.MustRunMain() defer m1.Close() // Update cluster config @@ -205,18 +206,18 @@ func TestMain_SendReceiveMessage(t *testing.T) { // Ensure that an empty node comes up in a NORMAL state. func TestClusterResize_EmptyNode(t *testing.T) { - m0 := MustRunMain() + m0 := test.MustRunMain() defer m0.Close() - if m0.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected cluster state: %s", m0.Server.Cluster.State) + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected cluster state: %s", m0.Server.Cluster.State()) } } // Ensure that a cluster of empty nodes comes up in a NORMAL state. func TestClusterResize_EmptyNodes(t *testing.T) { // Configure node0 - m0 := NewMain() + m0 := test.NewMain() defer m0.Close() gossipHost := "localhost" @@ -227,7 +228,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) { } // Configure node1 - m1 := NewMain() + m1 := test.NewMain() defer m1.Close() seed, coord, err = m1.RunWithTransport(gossipHost, gossipPort, seed, &coord) @@ -235,10 +236,10 @@ func TestClusterResize_EmptyNodes(t *testing.T) { t.Fatal(err) } - if m0.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State) - } else if m1.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State) + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) + } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) } } @@ -246,7 +247,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) { func TestClusterResize_AddNode(t *testing.T) { t.Run("NoData", func(t *testing.T) { // Configure node0 - m0 := NewMain() + m0 := test.NewMain() defer m0.Close() seed, coord, err := m0.RunWithTransport("localhost", 0, "", nil) @@ -255,7 +256,7 @@ func TestClusterResize_AddNode(t *testing.T) { } // Configure node1 - m1 := NewMain() + m1 := test.NewMain() defer m1.Close() var eg errgroup.Group @@ -272,15 +273,15 @@ func TestClusterResize_AddNode(t *testing.T) { time.Sleep(1 * time.Second) - if m0.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State) - } else if m1.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State) + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) + } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) } }) t.Run("WithIndex", func(t *testing.T) { // Configure node0 - m0 := NewMain() + m0 := test.NewMain() defer m0.Close() seed, coord, err := m0.RunWithTransport("localhost", 0, "", nil) @@ -299,7 +300,7 @@ func TestClusterResize_AddNode(t *testing.T) { } // Configure node1 - m1 := NewMain() + m1 := test.NewMain() defer m1.Close() var eg errgroup.Group @@ -317,16 +318,16 @@ func TestClusterResize_AddNode(t *testing.T) { // Give the cluster time to settle. time.Sleep(1 * time.Second) - if m0.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State) - } else if m1.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State) + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) + } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) } }) t.Run("ContinuousSlices", func(t *testing.T) { // Configure node0 - m0 := NewMain() + m0 := test.NewMain() defer m0.Close() seed, coord, err := m0.RunWithTransport("localhost", 0, "", nil) @@ -354,7 +355,7 @@ func TestClusterResize_AddNode(t *testing.T) { } // Configure node1 - m1 := NewMain() + m1 := test.NewMain() defer m1.Close() var eg errgroup.Group @@ -372,16 +373,16 @@ func TestClusterResize_AddNode(t *testing.T) { // Give the cluster time to settle. time.Sleep(1 * time.Second) - if m0.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State) - } else if m1.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State) + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) + } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) } }) t.Run("SkippedSlice", func(t *testing.T) { // Configure node0 - m0 := NewMain() + m0 := test.NewMain() defer m0.Close() seed, coord, err := m0.RunWithTransport("localhost", 0, "", nil) @@ -409,7 +410,7 @@ func TestClusterResize_AddNode(t *testing.T) { } // Configure node1 - m1 := NewMain() + m1 := test.NewMain() defer m1.Close() var eg errgroup.Group @@ -427,10 +428,10 @@ func TestClusterResize_AddNode(t *testing.T) { // Give the cluster time to settle. time.Sleep(1 * time.Second) - if m0.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State) - } else if m1.Server.Cluster.State != pilosa.ClusterStateNormal { - t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State) + if m0.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node0 cluster state: %s", m0.Server.Cluster.State()) + } else if m1.Server.Cluster.State() != pilosa.ClusterStateNormal { + t.Fatalf("unexpected node1 cluster state: %s", m1.Server.Cluster.State()) } }) } diff --git a/server/server_test.go b/server/server_test.go index c7b5038f7..379175c75 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -15,26 +15,19 @@ package server_test import ( - "bytes" "context" "encoding/json" "fmt" - "io" "io/ioutil" "math/rand" - "net/http" - "os" "reflect" "runtime" "sort" - "strings" "testing" "testing/quick" "github.com/BurntSushi/toml" "github.com/pilosa/pilosa" - "github.com/pilosa/pilosa/gossip" - "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/test" ) @@ -45,7 +38,7 @@ func TestMain_Set_Quick(t *testing.T) { } if err := quick.Check(func(cmds []SetCommand) bool { - m := MustRunMain() + m := test.MustRunMain() defer m.Close() // Create client. @@ -121,7 +114,7 @@ func TestMain_Set_Quick(t *testing.T) { // Ensure program can set row attributes and retrieve them. func TestMain_SetRowAttrs(t *testing.T) { - m := MustRunMain() + m := test.MustRunMain() defer m.Close() // Create frames. @@ -198,7 +191,7 @@ func TestMain_SetRowAttrs(t *testing.T) { // Ensure program can set column attributes and retrieve them. func TestMain_SetColumnAttrs(t *testing.T) { - m := MustRunMain() + m := test.MustRunMain() defer m.Close() // Create frames. @@ -242,7 +235,7 @@ func TestMain_SetColumnAttrs(t *testing.T) { // Ensure program can set column attributes with columnLabel option. func TestMain_SetColumnAttrsWithColumnOption(t *testing.T) { - m := MustRunMain() + m := test.MustRunMain() defer m.Close() // Create frames. @@ -276,17 +269,15 @@ func TestMain_SetColumnAttrsWithColumnOption(t *testing.T) { // Ensure program can set bits on one cluster and then restore to a second cluster. func TestMain_FrameRestore(t *testing.T) { - m0 := MustRunMain() - // TODO: this test used to start a two node cluster, but there was a race - // condition with anti-entropy. We need some general code for starting up - // arbitrarily sized Pilosa clusters for testing, and then we should - // re-instate the multi-node nature of this test. + mains1 := test.NewMainArrayWithCluster(2) + m0 := mains1[0] // Create frames. client := m0.Client() if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { t.Fatal("create index:", err) - } else if err := client.CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { + } + if err := client.CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { t.Fatal("create frame:", err) } @@ -311,18 +302,22 @@ func TestMain_FrameRestore(t *testing.T) { } // Start second cluster. - m2 := MustRunMain() + mains2 := test.NewMainArrayWithCluster(2) + m2 := mains2[0] defer m2.Close() // Import from first cluster. client, err := pilosa.NewInternalHTTPClient(m2.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) if err != nil { t.Fatal("new client:", err) - } else if err := m2.Client().CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { + } + if err := m2.Client().CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { t.Fatal("create new index:", err) - } else if err := m2.Client().CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { + } + if err := m2.Client().CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { t.Fatal("create new frame:", err) - } else if err := client.RestoreFrame(context.Background(), m0.Server.URI.HostPort(), "i", "f"); err != nil { + } + if err := client.RestoreFrame(context.Background(), m0.Server.URI.HostPort(), "i", "f"); err != nil { t.Fatal("restore frame:", err) } @@ -376,179 +371,6 @@ func TestCountOpenFiles(t *testing.T) { } } -// Main represents a test wrapper for main.Main. -type Main struct { - *server.Command - - Stdin bytes.Buffer - Stdout bytes.Buffer - Stderr bytes.Buffer -} - -// NewMain returns a new instance of Main with a temporary data directory and random port. -func NewMain() *Main { - path, err := ioutil.TempDir("", "pilosa-") - if err != nil { - panic(err) - } - - m := &Main{Command: server.NewCommand(os.Stdin, os.Stdout, os.Stderr)} - m.Server.Network = *test.Network - m.Config.DataDir = path - m.Config.Bind = "localhost:0" - m.Config.Cluster.Type = "static" - m.Command.Stdin = &m.Stdin - m.Command.Stdout = &m.Stdout - m.Command.Stderr = &m.Stderr - - if testing.Verbose() { - m.Command.Stdout = io.MultiWriter(os.Stdout, m.Command.Stdout) - m.Command.Stderr = io.MultiWriter(os.Stderr, m.Command.Stderr) - } - - return m -} - -// MustRunMain returns a new, running Main. Panic on error. -func MustRunMain() *Main { - m := NewMain() - m.Config.Metric.Diagnostics = false // Disable diagnostics. - if err := m.Run(); err != nil { - panic(err) - } - return m -} - -// Close closes the program and removes the underlying data directory. -func (m *Main) Close() error { - defer os.RemoveAll(m.Config.DataDir) - return m.Command.Close() -} - -// Reopen closes the program and reopens it. -func (m *Main) Reopen() error { - if err := m.Command.Close(); err != nil { - return err - } - - // Create new main with the same config. - config := m.Config - m.Command = server.NewCommand(os.Stdin, os.Stdout, os.Stderr) - m.Server.Network = *test.Network - m.Config = config - - // Run new program. - if err := m.Run(); err != nil { - return err - } - return nil -} - -// RunWithTransport runs Main and returns the dynamically allocated gossip port. -func (m *Main) RunWithTransport(host string, bindPort int, joinSeed string, coordinator *pilosa.URI) (seed string, coord pilosa.URI, err error) { - defer close(m.Started) - - m.Config.Cluster.Type = "gossip" - - /* - TEST: - - SetupServer (just static settings from config) - - OpenListener (sets Server.Name to use in gossip) - - NewTransport (gossip) - - SetupNetworking (does the gossip or static stuff) - uses Server.Name - - Open server - - PRODUCTION: - - SetupServer (just static settings from config) - - SetupNetworking (does the gossip or static stuff) - calls NewTransport - - Open server - calls OpenListener - */ - - // SetupServer - err = m.SetupServer() - if err != nil { - return seed, coord, err - } - - // Open server listener. - // This is used to set Server.Name, which is used as the node - // name for identifying a memberlist node. - err = m.Server.OpenListener() - if err != nil { - return seed, coord, err - } - - // Open gossip transport to use in SetupServer. - transport, err := gossip.NewTransport(host, bindPort) - if err != nil { - return seed, coord, err - } - m.GossipTransport = transport - - if joinSeed != "" { - m.Config.Gossip.Seed = joinSeed - } else { - m.Config.Gossip.Seed = transport.URI.String() - } - seed = m.Config.Gossip.Seed - - // SetupNetworking - err = m.SetupNetworking() - if err != nil { - return seed, coord, err - } - - if err = m.Server.BroadcastReceiver.Start(m.Server); err != nil { - return seed, coord, err - } - - if coordinator != nil { - coord = *coordinator - } else { - coord = m.Server.URI - } - m.Server.Cluster.Coordinator = coord - m.Server.Cluster.Static = false - - // Initialize server. - err = m.Server.Open() - if err != nil { - return seed, coord, err - } - - return seed, coord, nil -} - -// URL returns the base URL string for accessing the running program. -func (m *Main) URL() string { return "http://" + m.Server.Addr().String() } - -// Client returns a client to connect to the program. -func (m *Main) Client() *pilosa.InternalHTTPClient { - client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) - if err != nil { - panic(err) - } - return client -} - -// Query executes a query against the program through the HTTP API. -func (m *Main) Query(index, rawQuery, query string) (string, error) { - resp := MustDo("POST", m.URL()+fmt.Sprintf("/index/%s/query?", index)+rawQuery, query) - if resp.StatusCode != http.StatusOK { - return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body) - } - return resp.Body, nil -} - -// CreateDefinition. -func (m *Main) CreateDefinition(index, def, query string) (string, error) { - resp := MustDo("POST", m.URL()+fmt.Sprintf("/index/%s/input-definition/%s", index, def), query) - if resp.StatusCode != http.StatusOK { - return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body) - } - return resp.Body, nil -} - // SetCommand represents a command to set a bit. type SetCommand struct { ID uint64 @@ -605,32 +427,6 @@ func ParseConfig(s string) (pilosa.Config, error) { return c, err } -// MustDo executes http.Do() with an http.NewRequest(). Panic on error. -func MustDo(method, urlStr string, body string) *httpResponse { - req, err := http.NewRequest(method, urlStr, strings.NewReader(body)) - if err != nil { - panic(err) - } - resp, err := http.DefaultClient.Do(req) - if err != nil { - panic(err) - } - defer resp.Body.Close() - - buf, err := ioutil.ReadAll(resp.Body) - if err != nil { - panic(err) - } - - return &httpResponse{Response: resp, Body: string(buf)} -} - -// httpResponse is a wrapper for http.Response that holds the Body as a string. -type httpResponse struct { - *http.Response - Body string -} - // MustMarshalJSON marshals v into a string. Panic on error. func MustMarshalJSON(v interface{}) string { buf, err := json.Marshal(v) diff --git a/test/cluster.go b/test/cluster.go index 5c95c962e..ebe48fd35 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -7,6 +7,7 @@ import ( "io/ioutil" "path/filepath" "sort" + "sync" "time" "github.com/gogo/protobuf/proto" @@ -77,6 +78,8 @@ type TestCluster struct { common *commonClusterSettings + mu sync.RWMutex + resizing bool resizeDone chan struct{} } @@ -184,8 +187,11 @@ func (t *TestCluster) AddNode(saveTopology bool) error { } // Wait for the AddNode job to finish. - if c.State != pilosa.ClusterStateNormal { + if c.State() != pilosa.ClusterStateNormal { t.resizeDone = make(chan struct{}) + t.mu.Lock() + t.resizing = true + t.mu.Unlock() <-t.resizeDone } } @@ -266,7 +272,7 @@ func NewTestCluster(n int) *TestCluster { // SetState sets the state of the cluster on each node. func (t *TestCluster) SetState(state string) { for _, c := range t.Clusters { - c.State = state + c.SetState(state) } } @@ -314,9 +320,11 @@ func (t *TestCluster) SendSync(pb proto.Message) error { for _, c := range t.Clusters { c.MergeClusterStatus(obj) } - if obj.State == pilosa.ClusterStateNormal && t.resizeDone != nil { + t.mu.RLock() + if obj.State == pilosa.ClusterStateNormal && t.resizing { close(t.resizeDone) } + t.mu.RUnlock() } return nil diff --git a/test/pilosa.go b/test/pilosa.go index d733e99c0..1278dc17a 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -2,9 +2,12 @@ package test import ( "bytes" + "fmt" "io" "io/ioutil" + "net/http" "os" + "strings" "testing" "github.com/pilosa/pilosa" @@ -47,12 +50,53 @@ func NewMain() *Main { return m } +func NewMainArrayWithCluster(size int) []*Main { + cluster, err := NewServerCluster(size) + if err != nil { + panic(err) + } + mainArray := make([]*Main, size) + for i := 0; i < size; i++ { + mainArray[i] = cluster.Servers[i] + } + return mainArray +} + +// MustRunMain returns a new, running Main. Panic on error. +func MustRunMain() *Main { + m := NewMain() + m.Config.Metric.Diagnostics = false // Disable diagnostics. + if err := m.Run(); err != nil { + panic(err) + } + return m +} + // Close closes the program and removes the underlying data directory. func (m *Main) Close() error { defer os.RemoveAll(m.Config.DataDir) return m.Command.Close() } +// Reopen closes the program and reopens it. +func (m *Main) Reopen() error { + if err := m.Command.Close(); err != nil { + return err + } + + // Create new main with the same config. + config := m.Config + m.Command = server.NewCommand(os.Stdin, os.Stdout, os.Stderr) + m.Server.Network = *Network + m.Config = config + + // Run new program. + if err := m.Run(); err != nil { + return err + } + return nil +} + // RunWithTransport runs Main and returns the dynamically allocated gossip port. func (m *Main) RunWithTransport(host string, bindPort int, joinSeed string, coordinator *pilosa.URI) (seed string, coord pilosa.URI, err error) { defer close(m.Started) @@ -128,6 +172,36 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeed string, coor return seed, coord, nil } +// URL returns the base URL string for accessing the running program. +func (m *Main) URL() string { return "http://" + m.Server.Addr().String() } + +// Client returns a client to connect to the program. +func (m *Main) Client() *pilosa.InternalHTTPClient { + client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) + if err != nil { + panic(err) + } + return client +} + +// Query executes a query against the program through the HTTP API. +func (m *Main) Query(index, rawQuery, query string) (string, error) { + resp := MustDo("POST", m.URL()+fmt.Sprintf("/index/%s/query?", index)+rawQuery, query) + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body) + } + return resp.Body, nil +} + +// CreateDefinition. +func (m *Main) CreateDefinition(index, def, query string) (string, error) { + resp := MustDo("POST", m.URL()+fmt.Sprintf("/index/%s/input-definition/%s", index, def), query) + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body) + } + return resp.Body, nil +} + //////////////////////////////////////////////////////////////////////////////////// type Cluster struct { @@ -169,3 +243,29 @@ func NewServerCluster(size int) (cluster *Cluster, err error) { return cluster, nil } + +// MustDo executes http.Do() with an http.NewRequest(). Panic on error. +func MustDo(method, urlStr string, body string) *httpResponse { + req, err := http.NewRequest(method, urlStr, strings.NewReader(body)) + if err != nil { + panic(err) + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + panic(err) + } + defer resp.Body.Close() + + buf, err := ioutil.ReadAll(resp.Body) + if err != nil { + panic(err) + } + + return &httpResponse{Response: resp, Body: string(buf)} +} + +// httpResponse is a wrapper for http.Response that holds the Body as a string. +type httpResponse struct { + *http.Response + Body string +}