Merge pull request #1065 from travisturner/fix-cluster-tests

Fix cluster tests
This commit is contained in:
Travis Turner 2018-01-25 09:34:04 -06:00 committed by GitHub
commit b8150d345e
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
15 changed files with 887 additions and 520 deletions

View file

@ -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

View file

@ -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()

View file

@ -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{

View file

@ -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.")

View file

@ -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.

View file

@ -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
}
}

View file

@ -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))
}

View file

@ -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 {

View file

@ -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,
}

View file

@ -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;
}

View file

@ -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
}

View file

@ -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())
}
})
}

View file

@ -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)

View file

@ -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

View file

@ -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
}