From f4c9c0fed340e5b55f4772fbeee485467f8f374f Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Wed, 15 Aug 2018 07:50:47 -0600 Subject: [PATCH] Maintain available shards set. This commit removes the previous `MaxShard` tracking and replaces it with an `Available Shards` set tracking. This allows sparse shard tracking without implicitly tracking all shards in between. --- api.go | 12 +- apimethod_string.go | 4 +- cluster.go | 36 +- cluster_internal_test.go | 27 +- diagnostics.go | 2 +- encoding/proto/proto.go | 68 +++- executor.go | 15 +- executor_test.go | 1 - field.go | 25 +- gossip/gossip.go | 20 +- holder.go | 17 +- holder_internal_test.go | 6 +- index.go | 30 +- internal/private.pb.go | 763 +++++++++++++++++++++++++++++++-------- internal/private.proto | 13 +- server.go | 35 +- view.go | 36 +- 17 files changed, 833 insertions(+), 277 deletions(-) diff --git a/api.go b/api.go index 88b759903..fb64b5bdc 100644 --- a/api.go +++ b/api.go @@ -27,6 +27,7 @@ import ( "time" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" ) @@ -727,7 +728,16 @@ func (api *API) ImportValue(_ context.Context, req *ImportValueRequest) error { // MaxShards returns the maximum shard number for each index in a map. func (api *API) MaxShards(_ context.Context) map[string]uint64 { - return api.holder.maxShards() + m := make(map[string]uint64) + for k, v := range api.holder.availableShardsByIndex() { + m[k] = v.Max() + } + return m +} + +// AvailableShardsByIndex returns bitmaps of shards with available by index name. +func (api *API) AvailableShardsByIndex(_ context.Context) map[string]*roaring.Bitmap { + return api.holder.availableShardsByIndex() } // StatsWithTags returns an instance of whatever implementation of StatsClient diff --git a/apimethod_string.go b/apimethod_string.go index 01217092f..f3af365f0 100644 --- a/apimethod_string.go +++ b/apimethod_string.go @@ -2,7 +2,7 @@ package pilosa -import "strconv" +import "fmt" const _apiMethod_name = "apiClusterMessageapiCreateFieldapiCreateIndexapiDeleteFieldapiDeleteIndexapiDeleteViewapiExportCSVapiFragmentBlockDataapiFragmentBlocksapiFieldapiFieldAttrDiffapiImportapiImportValueapiIndexapiIndexAttrDiffapiQueryapiRecalculateCachesapiRemoveNodeapiResizeAbortapiSetCoordinatorapiShardNodesapiViews" @@ -10,7 +10,7 @@ var _apiMethod_index = [...]uint16{0, 17, 31, 45, 59, 73, 86, 98, 118, 135, 143, func (i apiMethod) String() string { if i < 0 || i >= apiMethod(len(_apiMethod_index)-1) { - return "apiMethod(" + strconv.FormatInt(int64(i), 10) + ")" + return fmt.Sprintf("apiMethod(%d)", i) } return _apiMethod_name[_apiMethod_index[i]:_apiMethod_index[i+1]] } diff --git a/cluster.go b/cluster.go index 3c2f8070d..868feceee 100644 --- a/cluster.go +++ b/cluster.go @@ -31,6 +31,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" uuid "github.com/satori/go.uuid" ) @@ -640,17 +641,17 @@ func (c *cluster) fragsByHost(idx *Index) fragsByHost { for _, field := range idx.Fields() { for _, view := range field.views() { fieldViews.addView(field.Name(), view.name) - } } - return c.fragCombos(idx.Name(), idx.maxShard(), fieldViews) + return c.fragCombos(idx.Name(), idx.AvailableShards(), fieldViews) } // fragCombos returns a map (by uri) of lists of fragments for a given index -// by creating every combination of field/view specified in `fieldViews` up to maxShard. -func (c *cluster) fragCombos(idx string, maxShard uint64, fieldViews viewsByField) fragsByHost { +// by creating every combination of field/view specified in `fieldViews` up +// for the given set of shards with data. +func (c *cluster) fragCombos(idx string, availableShards *roaring.Bitmap, fieldViews viewsByField) fragsByHost { t := make(fragsByHost) - for i := uint64(0); i <= maxShard; i++ { + availableShards.ForEach(func(i uint64) { nodes := c.shardNodes(idx, i) for _, n := range nodes { // for each field/view combination: @@ -660,7 +661,7 @@ func (c *cluster) fragCombos(idx string, maxShard uint64, fieldViews viewsByFiel } } } - } + }) return t } @@ -838,9 +839,9 @@ func (c *cluster) partitionNodes(partitionID int) []*Node { } // containsShards is like OwnsShards, but it includes replicas. -func (c *cluster) containsShards(index string, maxShard uint64, node *Node) []uint64 { +func (c *cluster) containsShards(index string, availableShards *roaring.Bitmap, node *Node) []uint64 { var shards []uint64 - for i := uint64(0); i <= maxShard; i++ { + availableShards.ForEach(func(i uint64) { p := c.partition(index, i) // Determine the nodes for partition. nodes := c.partitionNodes(p) @@ -849,7 +850,7 @@ func (c *cluster) containsShards(index string, maxShard uint64, node *Node) []ui shards = append(shards, i) } } - } + }) return shards } @@ -1921,6 +1922,7 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { type CreateShardMessage struct { Index string + Field string Shard uint64 } @@ -1975,9 +1977,19 @@ type NodeStateMessage struct { } type NodeStatus struct { - Node *Node - MaxShards map[string]uint64 - Schema *Schema + Node *Node + Indexes []*IndexStatus + Schema *Schema +} + +type IndexStatus struct { + Name string + Fields []*FieldStatus +} + +type FieldStatus struct { + Name string + AvailableShards *roaring.Bitmap } type RecalculateCaches struct{} diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 879ad2f35..a2b529fba 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -25,12 +25,12 @@ import ( "time" "github.com/davecgh/go-spew/spew" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" ) // Ensure that fragCombos creates the correct fragment mapping. func TestFragCombos(t *testing.T) { - uri0, err := NewURIFromAddress("host0") if err != nil { t.Fatal(err) @@ -48,24 +48,24 @@ func TestFragCombos(t *testing.T) { c.addNodeBasicSorted(node1) tests := []struct { - idx string - maxShard uint64 - fieldViews viewsByField - expected fragsByHost + idx string + availableShards *roaring.Bitmap + fieldViews viewsByField + expected fragsByHost }{ { - idx: "i", - maxShard: uint64(2), - fieldViews: viewsByField{"f": []string{"v1", "v2"}}, + idx: "i", + availableShards: roaring.NewBitmap(0, 1, 2), + fieldViews: viewsByField{"f": []string{"v1", "v2"}}, expected: fragsByHost{ "node0": []frag{{"f", "v1", uint64(0)}, {"f", "v2", uint64(0)}}, "node1": []frag{{"f", "v1", uint64(1)}, {"f", "v2", uint64(1)}, {"f", "v1", uint64(2)}, {"f", "v2", uint64(2)}}, }, }, { - idx: "foo", - maxShard: uint64(3), - fieldViews: viewsByField{"f": []string{"v0"}}, + idx: "foo", + availableShards: roaring.NewBitmap(0, 1, 2, 3), + fieldViews: viewsByField{"f": []string{"v0"}}, expected: fragsByHost{ "node0": []frag{{"f", "v0", uint64(1)}, {"f", "v0", uint64(2)}}, "node1": []frag{{"f", "v0", uint64(0)}, {"f", "v0", uint64(3)}}, @@ -73,8 +73,7 @@ func TestFragCombos(t *testing.T) { }, } for _, test := range tests { - - actual := c.fragCombos(test.idx, test.maxShard, test.fieldViews) + actual := c.fragCombos(test.idx, test.availableShards, test.fieldViews) if !reflect.DeepEqual(actual, test.expected) { t.Errorf("expected: %v, but got: %v", test.expected, actual) } @@ -385,7 +384,7 @@ func TestHasher(t *testing.T) { func TestCluster_ContainsShards(t *testing.T) { c := NewTestCluster(5) c.ReplicaN = 3 - shards := c.containsShards("test", 10, c.nodes[2]) + shards := c.containsShards("test", roaring.NewBitmap(0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10), c.nodes[2]) if !reflect.DeepEqual(shards, []uint64{0, 2, 3, 5, 6, 9, 10}) { t.Fatalf("unexpected shars for node's index: %v", shards) diff --git a/diagnostics.go b/diagnostics.go index eb30a5916..cad7725d7 100644 --- a/diagnostics.go +++ b/diagnostics.go @@ -223,7 +223,7 @@ func (d *diagnosticsCollector) EnrichWithSchemaProperties() { timeQuantumEnabled := false for _, index := range d.server.holder.Indexes() { - numShards += index.maxShard() + 1 + numShards += index.AvailableShards().Count() numIndexes += 1 for _, field := range index.Fields() { numFields += 1 diff --git a/encoding/proto/proto.go b/encoding/proto/proto.go index 298683e92..a19f1e548 100644 --- a/encoding/proto/proto.go +++ b/encoding/proto/proto.go @@ -7,6 +7,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" ) @@ -494,6 +495,7 @@ func encodeClusterStatus(m *pilosa.ClusterStatus) *internal.ClusterStatus { func encodeCreateShardMessage(m *pilosa.CreateShardMessage) *internal.CreateShardMessage { return &internal.CreateShardMessage{ Index: m.Index, + Field: m.Field, Shard: m.Shard, } } @@ -584,12 +586,42 @@ func encodeNodeEventMessage(m *pilosa.NodeEvent) *internal.NodeEventMessage { func encodeNodeStatus(m *pilosa.NodeStatus) *internal.NodeStatus { return &internal.NodeStatus{ - Node: encodeNode(m.Node), - MaxShards: &internal.MaxShards{Standard: m.MaxShards}, - Schema: encodeSchema(m.Schema), + Node: encodeNode(m.Node), + Indexes: encodeIndexStatuses(m.Indexes), + Schema: encodeSchema(m.Schema), } } +func encodeIndexStatus(m *pilosa.IndexStatus) *internal.IndexStatus { + return &internal.IndexStatus{ + Name: m.Name, + Fields: encodeFieldStatuses(m.Fields), + } +} + +func encodeIndexStatuses(a []*pilosa.IndexStatus) []*internal.IndexStatus { + other := make([]*internal.IndexStatus, len(a)) + for i := range a { + other[i] = encodeIndexStatus(a[i]) + } + return other +} + +func encodeFieldStatus(m *pilosa.FieldStatus) *internal.FieldStatus { + return &internal.FieldStatus{ + Name: m.Name, + AvailableShards: m.AvailableShards.Slice(), + } +} + +func encodeFieldStatuses(a []*pilosa.FieldStatus) []*internal.FieldStatus { + other := make([]*internal.FieldStatus, len(a)) + for i := range a { + other[i] = encodeFieldStatus(a[i]) + } + return other +} + func encodeRecalculateCaches(*pilosa.RecalculateCaches) *internal.RecalculateCaches { return &internal.RecalculateCaches{} } @@ -697,6 +729,7 @@ func decodeURI(i *internal.URI, m *pilosa.URI) { func decodeCreateShardMessage(pb *internal.CreateShardMessage, m *pilosa.CreateShardMessage) { m.Index = pb.Index + m.Field = pb.Field m.Shard = pb.Shard } @@ -768,12 +801,37 @@ func decodeNodeEventMessage(pb *internal.NodeEventMessage, m *pilosa.NodeEvent) func decodeNodeStatus(pb *internal.NodeStatus, m *pilosa.NodeStatus) { m.Node = &pilosa.Node{} - decodeNode(pb.Node, m.Node) - m.MaxShards = pb.MaxShards.Standard + decodeIndexStatuses(pb.Indexes, m.Indexes) m.Schema = &pilosa.Schema{} decodeSchema(pb.Schema, m.Schema) } +func decodeIndexStatuses(a []*internal.IndexStatus, m []*pilosa.IndexStatus) { + m = m[:0] + for i := range a { + m = append(m, &pilosa.IndexStatus{}) + decodeIndexStatus(a[i], m[i]) + } +} + +func decodeIndexStatus(pb *internal.IndexStatus, m *pilosa.IndexStatus) { + m.Name = pb.Name + decodeFieldStatuses(pb.Fields, m.Fields) +} + +func decodeFieldStatuses(a []*internal.FieldStatus, m []*pilosa.FieldStatus) { + m = m[:0] + for i := range a { + m = append(m, &pilosa.FieldStatus{}) + decodeFieldStatus(a[i], m[i]) + } +} + +func decodeFieldStatus(pb *internal.FieldStatus, m *pilosa.FieldStatus) { + m.Name = pb.Name + m.AvailableShards = roaring.NewBitmap(pb.AvailableShards...) +} + func decodeRecalculateCaches(pb *internal.RecalculateCaches, m *pilosa.RecalculateCaches) {} func decodeQueryRequest(pb *internal.QueryRequest, m *pilosa.QueryRequest) { diff --git a/executor.go b/executor.go index d7e70ce81..6f17cd6a7 100644 --- a/executor.go +++ b/executor.go @@ -140,12 +140,9 @@ func (e *executor) execute(ctx context.Context, index string, q *pql.Query, shar if idx == nil { return nil, ErrIndexNotFound } - maxShard := idx.maxShard() - - // Generate a slice of all shards. - shards = make([]uint64, maxShard+1) - for i := range shards { - shards[i] = uint64(i) + shards = idx.AvailableShards().Slice() + if len(shards) == 0 { + shards = []uint64{0} } } @@ -1443,7 +1440,7 @@ func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64, // Iterate over all map responses and reduce. var result interface{} - var maxShard int + var shardN int for { select { case <-ctx.Done(): @@ -1469,8 +1466,8 @@ func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64, result = reduceFn(result, resp.result) // If all shards have been processed then return. - maxShard += len(resp.shards) - if maxShard >= len(shards) { + shardN += len(resp.shards) + if shardN >= len(shards) { return result, nil } } diff --git a/executor_test.go b/executor_test.go index 1bcb17524..a04b5dba2 100644 --- a/executor_test.go +++ b/executor_test.go @@ -474,7 +474,6 @@ func TestExecutor_Execute_SetRowAttrs(t *testing.T) { } else if _, err := index.CreateFieldIfNotExists("kf", pilosa.OptFieldTypeDefault(), pilosa.OptFieldKeys()); err != nil { t.Fatal(err) } - t.Run("rowID", func(t *testing.T) { // Set two attrs on f/10. // Also set attrs on other bitmaps and fields to test isolation. diff --git a/field.go b/field.go index 04c857ff9..9fb352fc1 100644 --- a/field.go +++ b/field.go @@ -28,6 +28,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" ) @@ -73,6 +74,9 @@ type Field struct { bsiGroups []*bsiGroup + // Shards with data on any node in the cluster, according to this node. + remoteAvailableShards *roaring.Bitmap + logger Logger } @@ -179,6 +183,8 @@ func NewField(path, index, name string, opts FieldOption) (*Field, error) { options: applyDefaultOptions(fo), + remoteAvailableShards: roaring.NewBitmap(), + logger: NopLogger, } return f, nil @@ -196,18 +202,23 @@ func (f *Field) Path() string { return f.path } // RowAttrStore returns the attribute storage. func (f *Field) RowAttrStore() AttrStore { return f.rowAttrStore } -// maxShard returns the max shard in the field. -func (f *Field) maxShard() uint64 { +// AvailableShards returns a bitmap of shards that contain data. +func (f *Field) AvailableShards() *roaring.Bitmap { f.mu.RLock() defer f.mu.RUnlock() - var max uint64 + b := f.remoteAvailableShards.Clone() for _, view := range f.viewMap { - if viewMaxShard := view.calculateMaxShard(); viewMaxShard > max { - max = viewMaxShard - } + b = b.Union(view.availableShards()) } - return max + return b +} + +// addRemoteAvailableShards merges the set of available shards into the current known set. +func (f *Field) addRemoteAvailableShards(b *roaring.Bitmap) { + f.mu.Lock() + defer f.mu.Unlock() + f.remoteAvailableShards = f.remoteAvailableShards.Union(b) } // Type returns the field type. diff --git a/gossip/gossip.go b/gossip/gossip.go index 0fff2e52d..789ed50b0 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -28,6 +28,7 @@ import ( "github.com/hashicorp/memberlist" "github.com/pilosa/pilosa" + "github.com/pilosa/pilosa/roaring" "github.com/pilosa/pilosa/toml" "github.com/pkg/errors" ) @@ -258,9 +259,22 @@ func (g *memberSet) GetBroadcasts(overhead, limit int) [][]byte { // sends this Node's state data. func (g *memberSet) LocalState(join bool) []byte { m := &pilosa.NodeStatus{ - Node: g.papi.Node(), - MaxShards: g.papi.MaxShards(context.Background()), - Schema: &pilosa.Schema{Indexes: g.papi.Schema(context.Background())}, + Node: g.papi.Node(), + Schema: &pilosa.Schema{Indexes: g.papi.Schema(context.Background())}, + } + for _, idx := range m.Schema.Indexes { + is := &pilosa.IndexStatus{Name: idx.Name} + for _, f := range idx.Fields { + availableShards := roaring.NewBitmap() + if field, _ := g.papi.Field(context.Background(), idx.Name, f.Name); field != nil { + availableShards = field.AvailableShards() + } + is.Fields = append(is.Fields, &pilosa.FieldStatus{ + Name: f.Name, + AvailableShards: availableShards, + }) + } + m.Indexes = append(m.Indexes, is) } // Marshal nodestate data to bytes. diff --git a/holder.go b/holder.go index 15c8524b8..957765511 100644 --- a/holder.go +++ b/holder.go @@ -27,6 +27,7 @@ import ( "syscall" "time" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" uuid "github.com/satori/go.uuid" ) @@ -213,13 +214,13 @@ func (h *Holder) HasData() (bool, error) { return false, nil } -// maxShards returns MaxShard map for all indexes. -func (h *Holder) maxShards() map[string]uint64 { - a := make(map[string]uint64) +// availableShardsByIndex returns a bitmap of all shards by indexes. +func (h *Holder) availableShardsByIndex() map[string]*roaring.Bitmap { + m := make(map[string]*roaring.Bitmap) for _, index := range h.Indexes() { - a[index.Name()] = index.maxShard() + m[index.Name()] = index.AvailableShards() } - return a + return m } // Schema returns schema information for all indexes, fields, and views. @@ -647,7 +648,9 @@ func (s *holderSyncer) SyncHolder() error { return nil } - for shard := uint64(0); shard <= s.Holder.Index(di.Name).maxShard(); shard++ { + itr := s.Holder.Index(di.Name).AvailableShards().Iterator() + itr.Seek(0) + for shard, eof := itr.Next(); !eof; shard, eof = itr.Next() { // Ignore shards that this host doesn't own. if !s.Cluster.ownsShard(s.Node.ID, di.Name, shard) { continue @@ -828,7 +831,7 @@ func (c *holderCleaner) CleanHolder() error { } // Get the fragments that node is responsible for (based on hash(index, node)). - containedShards := c.Cluster.containsShards(index.Name(), index.maxShard(), c.Node) + containedShards := c.Cluster.containsShards(index.Name(), index.AvailableShards(), c.Node) // Get the fragments registered in memory. for _, field := range index.Fields() { diff --git a/holder_internal_test.go b/holder_internal_test.go index e9ab1977e..c4230e9aa 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -21,6 +21,8 @@ import ( "reflect" "strings" "testing" + + "github.com/pilosa/pilosa/roaring" ) type tHolder struct { @@ -207,8 +209,8 @@ func TestHolderCleaner_CleanHolder(t *testing.T) { hldr0.SetBit("y", "z", 10, (2*ShardWidth)+7) // Set highest shard. - hldr0.Index("i").setRemoteMaxShard(1) - hldr0.Index("y").setRemoteMaxShard(2) + hldr0.Field("i", "f").addRemoteAvailableShards(roaring.NewBitmap(0, 1)) + hldr0.Field("y", "z").addRemoteAvailableShards(roaring.NewBitmap(0, 1, 2)) // Keep replication the same and ensure we get the expected results. cluster.ReplicaN = 2 diff --git a/index.go b/index.go index d9338cceb..0547d07d2 100644 --- a/index.go +++ b/index.go @@ -25,6 +25,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/internal" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" ) @@ -38,9 +39,6 @@ type Index struct { // Fields by name. fields map[string]*Field - // Max shard on any node in the cluster, according to this node. - remoteMaxShard uint64 - newAttrStore func(string) AttrStore // Column attribute storage and cache. @@ -64,8 +62,6 @@ func NewIndex(path, name string) (*Index, error) { name: name, fields: make(map[string]*Field), - remoteMaxShard: 0, - newAttrStore: newNopAttrStore, columnAttrs: nopStore, @@ -210,30 +206,22 @@ func (i *Index) Close() error { return nil } -// maxShard returns the max shard in the index according to this node. -func (i *Index) maxShard() uint64 { +// AvailableShards returns a bitmap of all shards with data in the index. +func (i *Index) AvailableShards() *roaring.Bitmap { if i == nil { - return 0 + return roaring.NewBitmap() } + i.mu.RLock() defer i.mu.RUnlock() - max := i.remoteMaxShard + b := roaring.NewBitmap() for _, f := range i.fields { - if shard := f.maxShard(); shard > max { - max = shard - } + b = b.Union(f.AvailableShards()) } - i.Stats.Gauge("maxShard", float64(max), 1.0) - return max -} - -// setRemoteMaxShard sets the remote max shard value received from another node. -func (i *Index) setRemoteMaxShard(newmax uint64) { - i.mu.Lock() - defer i.mu.Unlock() - i.remoteMaxShard = newmax + i.Stats.Gauge("maxShard", float64(b.Max()), 1.0) + return b } // fieldPath returns the path to a field in the index. diff --git a/internal/private.pb.go b/internal/private.pb.go index b6c7ff5ca..6247c58d4 100644 --- a/internal/private.pb.go +++ b/internal/private.pb.go @@ -28,6 +28,8 @@ NodeStateMessage NodeEventMessage NodeStatus + IndexStatus + FieldStatus ClusterStatus BSIGroup CreateViewMessage @@ -261,6 +263,7 @@ func (m *MaxShards) GetStandard() map[string]uint64 { type CreateShardMessage struct { Index string `protobuf:"bytes,1,opt,name=Index,proto3" json:"Index,omitempty"` + Field string `protobuf:"bytes,3,opt,name=Field,proto3" json:"Field,omitempty"` Shard uint64 `protobuf:"varint,2,opt,name=Shard,proto3" json:"Shard,omitempty"` } @@ -276,6 +279,13 @@ func (m *CreateShardMessage) GetIndex() string { return "" } +func (m *CreateShardMessage) GetField() string { + if m != nil { + return m.Field + } + return "" +} + func (m *CreateShardMessage) GetShard() uint64 { if m != nil { return m.Shard @@ -564,9 +574,9 @@ func (m *NodeEventMessage) GetNode() *Node { } type NodeStatus struct { - Node *Node `protobuf:"bytes,1,opt,name=Node" json:"Node,omitempty"` - MaxShards *MaxShards `protobuf:"bytes,2,opt,name=MaxShards" json:"MaxShards,omitempty"` - Schema *Schema `protobuf:"bytes,3,opt,name=Schema" json:"Schema,omitempty"` + Node *Node `protobuf:"bytes,1,opt,name=Node" json:"Node,omitempty"` + Schema *Schema `protobuf:"bytes,3,opt,name=Schema" json:"Schema,omitempty"` + Indexes []*IndexStatus `protobuf:"bytes,4,rep,name=Indexes" json:"Indexes,omitempty"` } func (m *NodeStatus) Reset() { *m = NodeStatus{} } @@ -581,16 +591,64 @@ func (m *NodeStatus) GetNode() *Node { return nil } -func (m *NodeStatus) GetMaxShards() *MaxShards { +func (m *NodeStatus) GetSchema() *Schema { if m != nil { - return m.MaxShards + return m.Schema } return nil } -func (m *NodeStatus) GetSchema() *Schema { +func (m *NodeStatus) GetIndexes() []*IndexStatus { if m != nil { - return m.Schema + return m.Indexes + } + return nil +} + +type IndexStatus struct { + Name string `protobuf:"bytes,1,opt,name=Name,proto3" json:"Name,omitempty"` + Fields []*FieldStatus `protobuf:"bytes,2,rep,name=Fields" json:"Fields,omitempty"` +} + +func (m *IndexStatus) Reset() { *m = IndexStatus{} } +func (m *IndexStatus) String() string { return proto.CompactTextString(m) } +func (*IndexStatus) ProtoMessage() {} +func (*IndexStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{20} } + +func (m *IndexStatus) GetName() string { + if m != nil { + return m.Name + } + return "" +} + +func (m *IndexStatus) GetFields() []*FieldStatus { + if m != nil { + return m.Fields + } + return nil +} + +type FieldStatus struct { + Name string `protobuf:"bytes,1,opt,name=Name,proto3" json:"Name,omitempty"` + AvailableShards []uint64 `protobuf:"varint,2,rep,packed,name=AvailableShards" json:"AvailableShards,omitempty"` +} + +func (m *FieldStatus) Reset() { *m = FieldStatus{} } +func (m *FieldStatus) String() string { return proto.CompactTextString(m) } +func (*FieldStatus) ProtoMessage() {} +func (*FieldStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{21} } + +func (m *FieldStatus) GetName() string { + if m != nil { + return m.Name + } + return "" +} + +func (m *FieldStatus) GetAvailableShards() []uint64 { + if m != nil { + return m.AvailableShards } return nil } @@ -604,7 +662,7 @@ type ClusterStatus struct { func (m *ClusterStatus) Reset() { *m = ClusterStatus{} } func (m *ClusterStatus) String() string { return proto.CompactTextString(m) } func (*ClusterStatus) ProtoMessage() {} -func (*ClusterStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{20} } +func (*ClusterStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{22} } func (m *ClusterStatus) GetClusterID() string { if m != nil { @@ -637,7 +695,7 @@ type BSIGroup struct { func (m *BSIGroup) Reset() { *m = BSIGroup{} } func (m *BSIGroup) String() string { return proto.CompactTextString(m) } func (*BSIGroup) ProtoMessage() {} -func (*BSIGroup) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{21} } +func (*BSIGroup) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{23} } func (m *BSIGroup) GetName() string { if m != nil { @@ -676,7 +734,7 @@ type CreateViewMessage struct { func (m *CreateViewMessage) Reset() { *m = CreateViewMessage{} } func (m *CreateViewMessage) String() string { return proto.CompactTextString(m) } func (*CreateViewMessage) ProtoMessage() {} -func (*CreateViewMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{22} } +func (*CreateViewMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{24} } func (m *CreateViewMessage) GetIndex() string { if m != nil { @@ -708,7 +766,7 @@ type DeleteViewMessage struct { func (m *DeleteViewMessage) Reset() { *m = DeleteViewMessage{} } func (m *DeleteViewMessage) String() string { return proto.CompactTextString(m) } func (*DeleteViewMessage) ProtoMessage() {} -func (*DeleteViewMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{23} } +func (*DeleteViewMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{25} } func (m *DeleteViewMessage) GetIndex() string { if m != nil { @@ -743,7 +801,7 @@ type ResizeInstruction struct { func (m *ResizeInstruction) Reset() { *m = ResizeInstruction{} } func (m *ResizeInstruction) String() string { return proto.CompactTextString(m) } func (*ResizeInstruction) ProtoMessage() {} -func (*ResizeInstruction) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{24} } +func (*ResizeInstruction) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{26} } func (m *ResizeInstruction) GetJobID() int64 { if m != nil { @@ -798,7 +856,7 @@ type ResizeSource struct { func (m *ResizeSource) Reset() { *m = ResizeSource{} } func (m *ResizeSource) String() string { return proto.CompactTextString(m) } func (*ResizeSource) ProtoMessage() {} -func (*ResizeSource) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{25} } +func (*ResizeSource) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{27} } func (m *ResizeSource) GetNode() *Node { if m != nil { @@ -845,7 +903,7 @@ func (m *ResizeInstructionComplete) Reset() { *m = ResizeInstructionComp func (m *ResizeInstructionComplete) String() string { return proto.CompactTextString(m) } func (*ResizeInstructionComplete) ProtoMessage() {} func (*ResizeInstructionComplete) Descriptor() ([]byte, []int) { - return fileDescriptorPrivate, []int{26} + return fileDescriptorPrivate, []int{28} } func (m *ResizeInstructionComplete) GetJobID() int64 { @@ -876,7 +934,7 @@ type SetCoordinatorMessage struct { func (m *SetCoordinatorMessage) Reset() { *m = SetCoordinatorMessage{} } func (m *SetCoordinatorMessage) String() string { return proto.CompactTextString(m) } func (*SetCoordinatorMessage) ProtoMessage() {} -func (*SetCoordinatorMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{27} } +func (*SetCoordinatorMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{29} } func (m *SetCoordinatorMessage) GetNew() *Node { if m != nil { @@ -892,7 +950,7 @@ type UpdateCoordinatorMessage struct { func (m *UpdateCoordinatorMessage) Reset() { *m = UpdateCoordinatorMessage{} } func (m *UpdateCoordinatorMessage) String() string { return proto.CompactTextString(m) } func (*UpdateCoordinatorMessage) ProtoMessage() {} -func (*UpdateCoordinatorMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{28} } +func (*UpdateCoordinatorMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{30} } func (m *UpdateCoordinatorMessage) GetNew() *Node { if m != nil { @@ -909,7 +967,7 @@ type Topology struct { func (m *Topology) Reset() { *m = Topology{} } func (m *Topology) String() string { return proto.CompactTextString(m) } func (*Topology) ProtoMessage() {} -func (*Topology) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{29} } +func (*Topology) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{31} } func (m *Topology) GetClusterID() string { if m != nil { @@ -931,7 +989,7 @@ type RecalculateCaches struct { func (m *RecalculateCaches) Reset() { *m = RecalculateCaches{} } func (m *RecalculateCaches) String() string { return proto.CompactTextString(m) } func (*RecalculateCaches) ProtoMessage() {} -func (*RecalculateCaches) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{30} } +func (*RecalculateCaches) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{32} } func init() { proto.RegisterType((*IndexMeta)(nil), "internal.IndexMeta") @@ -954,6 +1012,8 @@ func init() { proto.RegisterType((*NodeStateMessage)(nil), "internal.NodeStateMessage") proto.RegisterType((*NodeEventMessage)(nil), "internal.NodeEventMessage") proto.RegisterType((*NodeStatus)(nil), "internal.NodeStatus") + proto.RegisterType((*IndexStatus)(nil), "internal.IndexStatus") + proto.RegisterType((*FieldStatus)(nil), "internal.FieldStatus") proto.RegisterType((*ClusterStatus)(nil), "internal.ClusterStatus") proto.RegisterType((*BSIGroup)(nil), "internal.BSIGroup") proto.RegisterType((*CreateViewMessage)(nil), "internal.CreateViewMessage") @@ -1272,6 +1332,12 @@ func (m *CreateShardMessage) MarshalTo(dAtA []byte) (int, error) { i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Shard)) } + if len(m.Field) > 0 { + dAtA[i] = 0x1a + i++ + i = encodeVarintPrivate(dAtA, i, uint64(len(m.Field))) + i += copy(dAtA[i:], m.Field) + } return i, nil } @@ -1685,25 +1751,104 @@ func (m *NodeStatus) MarshalTo(dAtA []byte) (int, error) { } i += n12 } - if m.MaxShards != nil { - dAtA[i] = 0x12 + if m.Schema != nil { + dAtA[i] = 0x1a i++ - i = encodeVarintPrivate(dAtA, i, uint64(m.MaxShards.Size())) - n13, err := m.MaxShards.MarshalTo(dAtA[i:]) + i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) + n13, err := m.Schema.MarshalTo(dAtA[i:]) if err != nil { return 0, err } i += n13 } - if m.Schema != nil { - dAtA[i] = 0x1a - i++ - i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) - n14, err := m.Schema.MarshalTo(dAtA[i:]) - if err != nil { - return 0, err + if len(m.Indexes) > 0 { + for _, msg := range m.Indexes { + dAtA[i] = 0x22 + i++ + i = encodeVarintPrivate(dAtA, i, uint64(msg.Size())) + n, err := msg.MarshalTo(dAtA[i:]) + if err != nil { + return 0, err + } + i += n } - i += n14 + } + return i, nil +} + +func (m *IndexStatus) Marshal() (dAtA []byte, err error) { + size := m.Size() + dAtA = make([]byte, size) + n, err := m.MarshalTo(dAtA) + if err != nil { + return nil, err + } + return dAtA[:n], nil +} + +func (m *IndexStatus) MarshalTo(dAtA []byte) (int, error) { + var i int + _ = i + var l int + _ = l + if len(m.Name) > 0 { + dAtA[i] = 0xa + i++ + i = encodeVarintPrivate(dAtA, i, uint64(len(m.Name))) + i += copy(dAtA[i:], m.Name) + } + if len(m.Fields) > 0 { + for _, msg := range m.Fields { + dAtA[i] = 0x12 + i++ + i = encodeVarintPrivate(dAtA, i, uint64(msg.Size())) + n, err := msg.MarshalTo(dAtA[i:]) + if err != nil { + return 0, err + } + i += n + } + } + return i, nil +} + +func (m *FieldStatus) Marshal() (dAtA []byte, err error) { + size := m.Size() + dAtA = make([]byte, size) + n, err := m.MarshalTo(dAtA) + if err != nil { + return nil, err + } + return dAtA[:n], nil +} + +func (m *FieldStatus) MarshalTo(dAtA []byte) (int, error) { + var i int + _ = i + var l int + _ = l + if len(m.Name) > 0 { + dAtA[i] = 0xa + i++ + i = encodeVarintPrivate(dAtA, i, uint64(len(m.Name))) + i += copy(dAtA[i:], m.Name) + } + if len(m.AvailableShards) > 0 { + dAtA15 := make([]byte, len(m.AvailableShards)*10) + var j14 int + for _, num := range m.AvailableShards { + for num >= 1<<7 { + dAtA15[j14] = uint8(uint64(num)&0x7f | 0x80) + num >>= 7 + j14++ + } + dAtA15[j14] = uint8(num) + j14++ + } + dAtA[i] = 0x12 + i++ + i = encodeVarintPrivate(dAtA, i, uint64(j14)) + i += copy(dAtA[i:], dAtA15[:j14]) } return i, nil } @@ -1886,21 +2031,21 @@ func (m *ResizeInstruction) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Node.Size())) - n15, err := m.Node.MarshalTo(dAtA[i:]) + n16, err := m.Node.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n15 + i += n16 } if m.Coordinator != nil { dAtA[i] = 0x1a i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Coordinator.Size())) - n16, err := m.Coordinator.MarshalTo(dAtA[i:]) + n17, err := m.Coordinator.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n16 + i += n17 } if len(m.Sources) > 0 { for _, msg := range m.Sources { @@ -1918,21 +2063,21 @@ func (m *ResizeInstruction) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x2a i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) - n17, err := m.Schema.MarshalTo(dAtA[i:]) + n18, err := m.Schema.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n17 + i += n18 } if m.ClusterStatus != nil { dAtA[i] = 0x32 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.ClusterStatus.Size())) - n18, err := m.ClusterStatus.MarshalTo(dAtA[i:]) + n19, err := m.ClusterStatus.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n18 + i += n19 } return i, nil } @@ -1956,11 +2101,11 @@ func (m *ResizeSource) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Node.Size())) - n19, err := m.Node.MarshalTo(dAtA[i:]) + n20, err := m.Node.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n19 + i += n20 } if len(m.Index) > 0 { dAtA[i] = 0x12 @@ -2012,11 +2157,11 @@ func (m *ResizeInstructionComplete) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Node.Size())) - n20, err := m.Node.MarshalTo(dAtA[i:]) + n21, err := m.Node.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n20 + i += n21 } if len(m.Error) > 0 { dAtA[i] = 0x1a @@ -2046,11 +2191,11 @@ func (m *SetCoordinatorMessage) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.New.Size())) - n21, err := m.New.MarshalTo(dAtA[i:]) + n22, err := m.New.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n21 + i += n22 } return i, nil } @@ -2074,11 +2219,11 @@ func (m *UpdateCoordinatorMessage) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.New.Size())) - n22, err := m.New.MarshalTo(dAtA[i:]) + n23, err := m.New.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n22 + i += n23 } return i, nil } @@ -2279,6 +2424,10 @@ func (m *CreateShardMessage) Size() (n int) { if m.Shard != 0 { n += 1 + sovPrivate(uint64(m.Shard)) } + l = len(m.Field) + if l > 0 { + n += 1 + l + sovPrivate(uint64(l)) + } return n } @@ -2454,14 +2603,49 @@ func (m *NodeStatus) Size() (n int) { l = m.Node.Size() n += 1 + l + sovPrivate(uint64(l)) } - if m.MaxShards != nil { - l = m.MaxShards.Size() - n += 1 + l + sovPrivate(uint64(l)) - } if m.Schema != nil { l = m.Schema.Size() n += 1 + l + sovPrivate(uint64(l)) } + if len(m.Indexes) > 0 { + for _, e := range m.Indexes { + l = e.Size() + n += 1 + l + sovPrivate(uint64(l)) + } + } + return n +} + +func (m *IndexStatus) Size() (n int) { + var l int + _ = l + l = len(m.Name) + if l > 0 { + n += 1 + l + sovPrivate(uint64(l)) + } + if len(m.Fields) > 0 { + for _, e := range m.Fields { + l = e.Size() + n += 1 + l + sovPrivate(uint64(l)) + } + } + return n +} + +func (m *FieldStatus) Size() (n int) { + var l int + _ = l + l = len(m.Name) + if l > 0 { + n += 1 + l + sovPrivate(uint64(l)) + } + if len(m.AvailableShards) > 0 { + l = 0 + for _, e := range m.AvailableShards { + l += sovPrivate(uint64(e)) + } + n += 1 + sovPrivate(uint64(l)) + l + } return n } @@ -3727,6 +3911,35 @@ func (m *CreateShardMessage) Unmarshal(dAtA []byte) error { break } } + case 3: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Field", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + 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 ErrInvalidLengthPrivate + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.Field = string(dAtA[iNdEx:postIndex]) + iNdEx = postIndex default: iNdEx = preIndex skippy, err := skipPrivate(dAtA[iNdEx:]) @@ -5051,39 +5264,6 @@ func (m *NodeStatus) Unmarshal(dAtA []byte) error { return err } iNdEx = postIndex - case 2: - if wireType != 2 { - return fmt.Errorf("proto: wrong wireType = %d for field MaxShards", wireType) - } - var msglen int - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - msglen |= (int(b) & 0x7F) << shift - if b < 0x80 { - break - } - } - if msglen < 0 { - return ErrInvalidLengthPrivate - } - postIndex := iNdEx + msglen - if postIndex > l { - return io.ErrUnexpectedEOF - } - if m.MaxShards == nil { - m.MaxShards = &MaxShards{} - } - if err := m.MaxShards.Unmarshal(dAtA[iNdEx:postIndex]); err != nil { - return err - } - iNdEx = postIndex case 3: if wireType != 2 { return fmt.Errorf("proto: wrong wireType = %d for field Schema", wireType) @@ -5117,6 +5297,288 @@ func (m *NodeStatus) Unmarshal(dAtA []byte) error { return err } iNdEx = postIndex + case 4: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Indexes", wireType) + } + var msglen int + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + msglen |= (int(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + if msglen < 0 { + return ErrInvalidLengthPrivate + } + postIndex := iNdEx + msglen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.Indexes = append(m.Indexes, &IndexStatus{}) + if err := m.Indexes[len(m.Indexes)-1].Unmarshal(dAtA[iNdEx:postIndex]); err != nil { + return err + } + iNdEx = postIndex + default: + iNdEx = preIndex + skippy, err := skipPrivate(dAtA[iNdEx:]) + if err != nil { + return err + } + if skippy < 0 { + return ErrInvalidLengthPrivate + } + if (iNdEx + skippy) > l { + return io.ErrUnexpectedEOF + } + iNdEx += skippy + } + } + + if iNdEx > l { + return io.ErrUnexpectedEOF + } + return nil +} +func (m *IndexStatus) Unmarshal(dAtA []byte) error { + l := len(dAtA) + iNdEx := 0 + for iNdEx < l { + preIndex := iNdEx + var wire uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + wire |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + fieldNum := int32(wire >> 3) + wireType := int(wire & 0x7) + if wireType == 4 { + return fmt.Errorf("proto: IndexStatus: wiretype end group for non-group") + } + if fieldNum <= 0 { + return fmt.Errorf("proto: IndexStatus: illegal tag %d (wire type %d)", fieldNum, wire) + } + switch fieldNum { + case 1: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Name", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + 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 ErrInvalidLengthPrivate + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.Name = string(dAtA[iNdEx:postIndex]) + iNdEx = postIndex + case 2: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Fields", wireType) + } + var msglen int + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + msglen |= (int(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + if msglen < 0 { + return ErrInvalidLengthPrivate + } + postIndex := iNdEx + msglen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.Fields = append(m.Fields, &FieldStatus{}) + if err := m.Fields[len(m.Fields)-1].Unmarshal(dAtA[iNdEx:postIndex]); err != nil { + return err + } + iNdEx = postIndex + default: + iNdEx = preIndex + skippy, err := skipPrivate(dAtA[iNdEx:]) + if err != nil { + return err + } + if skippy < 0 { + return ErrInvalidLengthPrivate + } + if (iNdEx + skippy) > l { + return io.ErrUnexpectedEOF + } + iNdEx += skippy + } + } + + if iNdEx > l { + return io.ErrUnexpectedEOF + } + return nil +} +func (m *FieldStatus) Unmarshal(dAtA []byte) error { + l := len(dAtA) + iNdEx := 0 + for iNdEx < l { + preIndex := iNdEx + var wire uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + wire |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + fieldNum := int32(wire >> 3) + wireType := int(wire & 0x7) + if wireType == 4 { + return fmt.Errorf("proto: FieldStatus: wiretype end group for non-group") + } + if fieldNum <= 0 { + return fmt.Errorf("proto: FieldStatus: illegal tag %d (wire type %d)", fieldNum, wire) + } + switch fieldNum { + case 1: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Name", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + 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 ErrInvalidLengthPrivate + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.Name = string(dAtA[iNdEx:postIndex]) + iNdEx = postIndex + case 2: + if wireType == 0 { + var v uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + v |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + m.AvailableShards = append(m.AvailableShards, v) + } else if wireType == 2 { + var packedLen int + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + packedLen |= (int(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + if packedLen < 0 { + return ErrInvalidLengthPrivate + } + postIndex := iNdEx + packedLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + for iNdEx < postIndex { + var v uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + v |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + m.AvailableShards = append(m.AvailableShards, v) + } + } else { + return fmt.Errorf("proto: wrong wireType = %d for field AvailableShards", wireType) + } default: iNdEx = preIndex skippy, err := skipPrivate(dAtA[iNdEx:]) @@ -6681,70 +7143,73 @@ var ( func init() { proto.RegisterFile("private.proto", fileDescriptorPrivate) } var fileDescriptorPrivate = []byte{ - // 1027 bytes of a gzipped FileDescriptorProto - 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xac, 0x56, 0x4d, 0x6f, 0x1b, 0x45, - 0x18, 0x66, 0xbd, 0x6b, 0xc7, 0x7e, 0x53, 0x87, 0x64, 0x0a, 0x61, 0x8b, 0x50, 0x6a, 0x46, 0x95, - 0x1a, 0x7a, 0x88, 0x4a, 0x7b, 0xe1, 0xab, 0x52, 0x14, 0x3b, 0xc0, 0x02, 0x09, 0x30, 0x9b, 0xf4, - 0xd6, 0xc3, 0xd4, 0x1e, 0x35, 0xab, 0xac, 0x77, 0x96, 0xdd, 0xd9, 0x24, 0xee, 0x81, 0x2b, 0x5c, - 0xb8, 0x23, 0x7e, 0x09, 0x3f, 0x81, 0x23, 0x3f, 0x01, 0x85, 0x3f, 0x82, 0xe6, 0x9d, 0xd9, 0x8f, - 0xc4, 0x4e, 0x53, 0x85, 0xde, 0xe6, 0xfd, 0x7e, 0xe6, 0xfd, 0x9a, 0x81, 0x7e, 0x9a, 0x45, 0x27, - 0x5c, 0x89, 0xad, 0x34, 0x93, 0x4a, 0x92, 0x6e, 0x94, 0x28, 0x91, 0x25, 0x3c, 0xa6, 0x77, 0xa1, - 0x17, 0x24, 0x13, 0x71, 0xb6, 0x27, 0x14, 0x27, 0x04, 0xbc, 0x6f, 0xc5, 0x2c, 0xf7, 0xdd, 0x81, - 0xb3, 0xd9, 0x65, 0x78, 0xa6, 0x7f, 0x3a, 0x70, 0xeb, 0xcb, 0x48, 0xc4, 0x93, 0xef, 0x53, 0x15, - 0xc9, 0x24, 0x27, 0x1f, 0x40, 0x6f, 0xc8, 0xc7, 0x47, 0xe2, 0x60, 0x96, 0x0a, 0xd4, 0xec, 0xb1, - 0x9a, 0x51, 0x49, 0xc3, 0xe8, 0xa5, 0xf0, 0xbd, 0x81, 0xb3, 0xd9, 0x67, 0x35, 0x83, 0x0c, 0x60, - 0xf9, 0x20, 0x9a, 0x8a, 0x1f, 0x0b, 0x9e, 0xa8, 0x62, 0xea, 0xb7, 0xd1, 0xba, 0xc9, 0xd2, 0x10, - 0xd0, 0x71, 0x17, 0x45, 0x78, 0x26, 0xab, 0xe0, 0xee, 0x45, 0x89, 0xdf, 0x1b, 0x38, 0x9b, 0x2e, - 0xd3, 0x47, 0xe4, 0xf0, 0x33, 0x1f, 0x2c, 0x87, 0x9f, 0x55, 0xd0, 0x97, 0x1b, 0xd0, 0x29, 0xac, - 0x04, 0xd3, 0x54, 0x66, 0x8a, 0x89, 0x3c, 0x95, 0x49, 0x8e, 0x9e, 0x76, 0xb3, 0xcc, 0x77, 0xd0, - 0xb9, 0x3e, 0xd2, 0x9f, 0x61, 0x75, 0x27, 0x96, 0xe3, 0xe3, 0x11, 0x57, 0x9c, 0x89, 0x9f, 0x0a, - 0x91, 0x2b, 0xf2, 0x0e, 0xb4, 0x31, 0x27, 0x56, 0xcf, 0x10, 0x9a, 0x8b, 0x79, 0xf0, 0x5b, 0x86, - 0x8b, 0x84, 0xe6, 0xa2, 0x3d, 0x66, 0xc2, 0x63, 0x86, 0xd0, 0xdc, 0xf0, 0x88, 0x67, 0x13, 0xcc, - 0x80, 0xc7, 0x0c, 0xa1, 0x31, 0x3e, 0x8d, 0xc4, 0xa9, 0xbd, 0x36, 0x9e, 0x69, 0x00, 0x6b, 0x8d, - 0xf8, 0x16, 0xe6, 0x3a, 0x74, 0x98, 0x3c, 0x0d, 0x46, 0xb9, 0xef, 0x0c, 0xdc, 0x4d, 0x8f, 0x59, - 0x0a, 0x93, 0x2b, 0xe3, 0x62, 0x9a, 0x68, 0x51, 0x0b, 0x45, 0x35, 0x83, 0xde, 0x81, 0x36, 0x66, - 0x5a, 0xdf, 0xb2, 0xb6, 0xd5, 0x47, 0xfa, 0x8b, 0x03, 0xbd, 0x3d, 0x7e, 0x86, 0x30, 0x72, 0xf2, - 0x04, 0xba, 0xa1, 0xe2, 0xc9, 0x44, 0x03, 0xd4, 0x4a, 0xcb, 0x8f, 0x3e, 0xdc, 0x2a, 0x1b, 0x62, - 0xab, 0x52, 0xdb, 0x2a, 0x75, 0x76, 0x13, 0x95, 0xcd, 0x58, 0x65, 0xf2, 0xfe, 0xe7, 0xd0, 0xbf, - 0x20, 0xd2, 0xf1, 0x8e, 0xc5, 0xac, 0xcc, 0xea, 0xb1, 0x98, 0xe9, 0xfb, 0x9f, 0xf0, 0xb8, 0x10, - 0x98, 0x2b, 0x8f, 0x19, 0xe2, 0xb3, 0xd6, 0x27, 0x0e, 0xdd, 0x06, 0x32, 0xcc, 0x04, 0x57, 0x02, - 0x83, 0xec, 0x89, 0x3c, 0xe7, 0x2f, 0xc4, 0xd5, 0x19, 0x37, 0x59, 0x6c, 0x35, 0xb2, 0x48, 0x1f, - 0x00, 0x19, 0x89, 0x58, 0x28, 0x61, 0xfb, 0xf6, 0x15, 0x1e, 0x68, 0x58, 0x46, 0xbb, 0x5e, 0x97, - 0xdc, 0x07, 0x4f, 0x0f, 0x01, 0x06, 0x5b, 0x7e, 0x74, 0xbb, 0xce, 0x48, 0x35, 0x1f, 0x0c, 0x15, - 0x68, 0x5c, 0x3a, 0xc5, 0x0e, 0xb8, 0xf6, 0x0a, 0x0b, 0x9a, 0xe6, 0x81, 0x0d, 0xe5, 0x62, 0xa8, - 0xf5, 0x3a, 0x54, 0x73, 0xd0, 0x6c, 0xb4, 0xed, 0xf2, 0xba, 0x37, 0x8d, 0x46, 0x9f, 0x59, 0xae, - 0xee, 0xbf, 0x7d, 0x3e, 0x15, 0xd6, 0x06, 0xcf, 0x15, 0x94, 0xd6, 0xf5, 0x50, 0xb4, 0x7b, 0xdd, - 0xb3, 0x7a, 0x3f, 0xb8, 0xda, 0x3d, 0x12, 0xf4, 0x31, 0x74, 0xc2, 0xf1, 0x91, 0x98, 0x72, 0xf2, - 0x11, 0x2c, 0x21, 0x0e, 0x91, 0xdb, 0xb6, 0x7a, 0xfb, 0x52, 0x12, 0x59, 0x29, 0xa7, 0x23, 0x8b, - 0x7f, 0x21, 0xa6, 0xfb, 0xd0, 0xc1, 0xe8, 0xb9, 0xef, 0x5d, 0x76, 0x83, 0x7c, 0x66, 0xc5, 0x74, - 0x17, 0xdc, 0x43, 0x16, 0xe8, 0x71, 0x41, 0x04, 0xa5, 0x17, 0x4b, 0x69, 0xdf, 0x5f, 0xcb, 0x5c, - 0xd9, 0x6c, 0xe0, 0x59, 0xf3, 0x7e, 0x90, 0x99, 0xc2, 0xd4, 0xf7, 0x19, 0x9e, 0xe9, 0x33, 0xf0, - 0xf6, 0xe5, 0x44, 0x90, 0x15, 0x68, 0x05, 0x23, 0xeb, 0xa3, 0x15, 0x8c, 0xc8, 0x5d, 0x74, 0x6f, - 0x53, 0xd3, 0xaf, 0x41, 0x1c, 0xb2, 0x80, 0x61, 0xe0, 0x7b, 0xd0, 0x0f, 0xf2, 0xa1, 0x94, 0xd9, - 0x24, 0x4a, 0xb8, 0x92, 0x99, 0x5d, 0x9c, 0x17, 0x99, 0x74, 0x1b, 0x56, 0xb5, 0xfb, 0x50, 0x71, - 0x25, 0xca, 0xfa, 0xad, 0x43, 0x47, 0xf3, 0xaa, 0x70, 0x96, 0xc2, 0x96, 0xd7, 0x7a, 0x65, 0x05, - 0x91, 0xa0, 0xdf, 0x19, 0x0f, 0xbb, 0x27, 0x22, 0x51, 0x8d, 0x0e, 0x40, 0x1a, 0x1d, 0xf4, 0x99, - 0x21, 0x08, 0x35, 0x57, 0xb1, 0x98, 0x57, 0x6a, 0xcc, 0x9a, 0xcb, 0x50, 0x46, 0x7f, 0x73, 0x00, - 0x4a, 0x40, 0x45, 0x5e, 0x99, 0x38, 0x57, 0x9b, 0x90, 0x8f, 0x1b, 0xeb, 0x63, 0x7e, 0x40, 0x2a, - 0x11, 0x6b, 0x2c, 0x99, 0xcd, 0xb2, 0x2d, 0x6c, 0x97, 0xaf, 0xd6, 0xfa, 0x86, 0x6f, 0xcb, 0xc4, - 0x69, 0x04, 0xfd, 0x61, 0x5c, 0xe4, 0x4a, 0x64, 0x16, 0x91, 0x5e, 0x73, 0x86, 0x51, 0xe5, 0xa7, - 0x66, 0x2c, 0x4e, 0x11, 0xb9, 0x07, 0x6d, 0x8d, 0xd4, 0xf4, 0xe6, 0xfc, 0x35, 0x8c, 0x90, 0x3e, - 0x85, 0xee, 0x4e, 0x18, 0x7c, 0x95, 0xc9, 0x22, 0x5d, 0xd8, 0x79, 0xe5, 0xeb, 0xd3, 0x9a, 0x7f, - 0x7d, 0xdc, 0xb9, 0xd7, 0xc7, 0xab, 0x5e, 0x1f, 0x1a, 0xc2, 0x9a, 0x59, 0x09, 0x7a, 0x24, 0x6e, - 0xb2, 0x11, 0xca, 0xa7, 0xc1, 0x6d, 0x3c, 0x0d, 0x21, 0xac, 0x99, 0xc9, 0x7f, 0x93, 0x4e, 0xff, - 0x68, 0xc1, 0x1a, 0x13, 0x79, 0xf4, 0x52, 0x04, 0x49, 0xae, 0xb2, 0x62, 0xac, 0x07, 0x5c, 0xdb, - 0x7f, 0x23, 0x9f, 0xdb, 0x6c, 0xbb, 0xcc, 0x10, 0xaf, 0xd3, 0x4c, 0xe4, 0x21, 0x2c, 0x5f, 0x1e, - 0x80, 0x79, 0xd5, 0xa6, 0x0a, 0x79, 0x08, 0x4b, 0xa1, 0x2c, 0xb2, 0xb1, 0x28, 0xc7, 0xbb, 0xb1, - 0x74, 0x0c, 0x32, 0x23, 0x66, 0xa5, 0x5a, 0xa3, 0x95, 0xda, 0xaf, 0x6e, 0x25, 0xf2, 0xe4, 0x52, - 0x2b, 0xf9, 0x1d, 0x34, 0x78, 0xaf, 0x36, 0xb8, 0x20, 0x66, 0x17, 0xb5, 0xe9, 0xaf, 0x0e, 0xdc, - 0x6a, 0x42, 0x78, 0xad, 0xd9, 0xa8, 0x2a, 0xd2, 0x5a, 0x58, 0x11, 0x77, 0x51, 0x45, 0xbc, 0xba, - 0x22, 0xf5, 0x2b, 0xd7, 0x6e, 0xbe, 0x72, 0xc7, 0x70, 0x67, 0xae, 0x4c, 0x43, 0x39, 0x4d, 0x75, - 0x3f, 0xfc, 0x8f, 0x72, 0xe9, 0xad, 0x91, 0x65, 0xb6, 0x50, 0x3d, 0x66, 0x08, 0xfa, 0x29, 0xbc, - 0x1b, 0x0a, 0xd5, 0x28, 0x52, 0xd9, 0x6d, 0x03, 0x70, 0xf7, 0xc5, 0xe9, 0x15, 0xd7, 0xd7, 0x22, - 0xfa, 0x05, 0xf8, 0x87, 0xe9, 0x84, 0x2b, 0x71, 0x23, 0xeb, 0x1d, 0xe8, 0x1e, 0xc8, 0x54, 0xc6, - 0xf2, 0xc5, 0xec, 0x9a, 0xa9, 0xf7, 0x61, 0xc9, 0xac, 0x48, 0xf3, 0xf1, 0xe9, 0xb1, 0x92, 0xa4, - 0xb7, 0x75, 0x43, 0x8f, 0x79, 0x3c, 0x2e, 0x62, 0x0d, 0x43, 0xff, 0x80, 0xf2, 0x9d, 0xd5, 0xbf, - 0xce, 0x37, 0x9c, 0xbf, 0xcf, 0x37, 0x9c, 0x7f, 0xce, 0x37, 0x9c, 0xdf, 0xff, 0xdd, 0x78, 0xeb, - 0x79, 0x07, 0x7f, 0xbe, 0x8f, 0xff, 0x0b, 0x00, 0x00, 0xff, 0xff, 0x39, 0x2f, 0x93, 0x68, 0x0a, - 0x0b, 0x00, 0x00, + // 1077 bytes of a gzipped FileDescriptorProto + 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xac, 0x56, 0xdd, 0x6e, 0x1b, 0x45, + 0x14, 0x66, 0x7f, 0xec, 0xda, 0xc7, 0x75, 0x9a, 0x6c, 0x69, 0xd8, 0x22, 0x94, 0x9a, 0x51, 0xa5, + 0x9a, 0x4a, 0x84, 0xaa, 0xbd, 0xe1, 0xaf, 0x52, 0x49, 0x1c, 0x60, 0x29, 0x09, 0x65, 0x36, 0xc9, + 0x5d, 0x2f, 0x26, 0xf6, 0xa8, 0x59, 0x65, 0xbd, 0xb3, 0xec, 0xce, 0x26, 0x71, 0x2f, 0xb8, 0x05, + 0x89, 0x17, 0x40, 0x3c, 0x09, 0x8f, 0xc0, 0x25, 0x8f, 0x80, 0xc2, 0x8b, 0xa0, 0x39, 0x33, 0xfb, + 0x13, 0xc7, 0x21, 0x55, 0xe0, 0x6e, 0xce, 0x77, 0xce, 0x9c, 0xf3, 0xed, 0xf9, 0x9b, 0x85, 0x7e, + 0x9a, 0x45, 0xc7, 0x4c, 0xf2, 0xf5, 0x34, 0x13, 0x52, 0x78, 0x9d, 0x28, 0x91, 0x3c, 0x4b, 0x58, + 0x4c, 0xee, 0x41, 0x37, 0x48, 0x26, 0xfc, 0x74, 0x9b, 0x4b, 0xe6, 0x79, 0xe0, 0x3e, 0xe7, 0xb3, + 0xdc, 0x77, 0x06, 0xd6, 0xb0, 0x43, 0xf1, 0x4c, 0x7e, 0xb7, 0xe0, 0xe6, 0x97, 0x11, 0x8f, 0x27, + 0xdf, 0xa5, 0x32, 0x12, 0x49, 0xee, 0xbd, 0x07, 0xdd, 0x4d, 0x36, 0x3e, 0xe4, 0xbb, 0xb3, 0x94, + 0xa3, 0x65, 0x97, 0xd6, 0x40, 0xa5, 0x0d, 0xa3, 0xd7, 0xdc, 0x77, 0x07, 0xd6, 0xb0, 0x4f, 0x6b, + 0xc0, 0x1b, 0x40, 0x6f, 0x37, 0x9a, 0xf2, 0xef, 0x0b, 0x96, 0xc8, 0x62, 0xea, 0xb7, 0xf0, 0x76, + 0x13, 0x52, 0x14, 0xd0, 0x71, 0x07, 0x55, 0x78, 0xf6, 0x96, 0xc1, 0xd9, 0x8e, 0x12, 0xbf, 0x3b, + 0xb0, 0x86, 0x0e, 0x55, 0x47, 0x44, 0xd8, 0xa9, 0x0f, 0x06, 0x61, 0xa7, 0x15, 0xf5, 0x5e, 0x83, + 0x3a, 0x81, 0xa5, 0x60, 0x9a, 0x8a, 0x4c, 0x52, 0x9e, 0xa7, 0x22, 0xc9, 0xd1, 0xd3, 0x56, 0x96, + 0xf9, 0x16, 0x3a, 0x57, 0x47, 0xf2, 0x23, 0x2c, 0x6f, 0xc4, 0x62, 0x7c, 0x34, 0x62, 0x92, 0x51, + 0xfe, 0x43, 0xc1, 0x73, 0xe9, 0xbd, 0x0d, 0x2d, 0xcc, 0x89, 0xb1, 0xd3, 0x82, 0x42, 0x31, 0x0f, + 0xbe, 0xad, 0x51, 0x14, 0x14, 0x8a, 0xf7, 0x31, 0x13, 0x2e, 0xd5, 0x82, 0x42, 0xc3, 0x43, 0x96, + 0x4d, 0x30, 0x03, 0x2e, 0xd5, 0x82, 0xe2, 0xb8, 0x1f, 0xf1, 0x13, 0xf3, 0xd9, 0x78, 0x26, 0x01, + 0xac, 0x34, 0xe2, 0x1b, 0x9a, 0xab, 0xd0, 0xa6, 0xe2, 0x24, 0x18, 0xe5, 0xbe, 0x35, 0x70, 0x86, + 0x2e, 0x35, 0x12, 0x26, 0x57, 0xc4, 0xc5, 0x34, 0x51, 0x2a, 0x1b, 0x55, 0x35, 0x40, 0xee, 0x42, + 0x0b, 0x33, 0xad, 0xbe, 0xb2, 0xbe, 0xab, 0x8e, 0xe4, 0x27, 0x0b, 0xba, 0xdb, 0xec, 0x14, 0x69, + 0xe4, 0xde, 0x53, 0xe8, 0x84, 0x92, 0x25, 0x13, 0x45, 0x50, 0x19, 0xf5, 0x1e, 0xbf, 0xbf, 0x5e, + 0x36, 0xc4, 0x7a, 0x65, 0xb6, 0x5e, 0xda, 0x6c, 0x25, 0x32, 0x9b, 0xd1, 0xea, 0xca, 0xbb, 0x9f, + 0x41, 0xff, 0x9c, 0x4a, 0xc5, 0x3b, 0xe2, 0xb3, 0x32, 0xab, 0x47, 0x7c, 0xa6, 0xbe, 0xff, 0x98, + 0xc5, 0x05, 0xc7, 0x5c, 0xb9, 0x54, 0x0b, 0x9f, 0xda, 0x1f, 0x5b, 0x64, 0x1f, 0xbc, 0xcd, 0x8c, + 0x33, 0xc9, 0x31, 0xc8, 0x36, 0xcf, 0x73, 0xf6, 0x8a, 0x5f, 0x9e, 0x71, 0x9d, 0x45, 0xbb, 0x99, + 0xc5, 0xaa, 0x0e, 0x4e, 0xa3, 0x0e, 0xe4, 0x21, 0x78, 0x23, 0x1e, 0x73, 0xc9, 0x4d, 0x37, 0xff, + 0x8b, 0x5f, 0x12, 0x96, 0x1c, 0xae, 0xb6, 0xf5, 0x1e, 0x80, 0xab, 0x46, 0x03, 0x29, 0xf4, 0x1e, + 0xdf, 0xae, 0xf3, 0x54, 0x4d, 0x0d, 0x45, 0x03, 0x12, 0x97, 0x4e, 0x91, 0xcf, 0x95, 0x1f, 0xb6, + 0xa0, 0x95, 0x1e, 0x9a, 0x50, 0x0e, 0x86, 0x5a, 0xad, 0x43, 0x35, 0xc7, 0xcf, 0x44, 0x7b, 0x56, + 0x7e, 0xee, 0x75, 0xa3, 0x91, 0x97, 0x06, 0x55, 0x5d, 0xb9, 0xc3, 0xa6, 0xdc, 0xdc, 0xc1, 0x73, + 0x45, 0xc5, 0xbe, 0x9a, 0x8a, 0x72, 0xaf, 0x3a, 0x59, 0x6d, 0x0d, 0x47, 0xb9, 0x47, 0x81, 0x3c, + 0x81, 0x76, 0x38, 0x3e, 0xe4, 0x53, 0xe6, 0x7d, 0x00, 0x37, 0x90, 0x07, 0xcf, 0x4d, 0xb3, 0xdd, + 0x9a, 0x4b, 0x22, 0x2d, 0xf5, 0x64, 0x64, 0xf8, 0x2f, 0xe4, 0xf4, 0x00, 0xda, 0x18, 0x3d, 0xf7, + 0xdd, 0x79, 0x37, 0x88, 0x53, 0xa3, 0x26, 0x5b, 0xe0, 0xec, 0xd1, 0x40, 0x0d, 0x11, 0x32, 0x28, + 0xbd, 0x18, 0x49, 0xf9, 0xfe, 0x5a, 0xe4, 0xd2, 0x64, 0x03, 0xcf, 0x0a, 0x7b, 0x21, 0x32, 0x89, + 0xa9, 0xef, 0x53, 0x3c, 0x93, 0x97, 0xe0, 0xee, 0x88, 0x09, 0xf7, 0x96, 0xc0, 0x0e, 0x46, 0xc6, + 0x87, 0x1d, 0x8c, 0xbc, 0x7b, 0xe8, 0xde, 0xa4, 0xa6, 0x5f, 0x93, 0xd8, 0xa3, 0x01, 0xc5, 0xc0, + 0xf7, 0xa1, 0x1f, 0xe4, 0x9b, 0x42, 0x64, 0x93, 0x28, 0x61, 0x52, 0x64, 0x66, 0x9d, 0x9e, 0x07, + 0xc9, 0x33, 0x58, 0x56, 0xee, 0x43, 0xc9, 0x24, 0x2f, 0xeb, 0xb7, 0x0a, 0x6d, 0x85, 0x55, 0xe1, + 0x8c, 0x84, 0x83, 0xa0, 0xec, 0xca, 0x0a, 0xa2, 0x40, 0xbe, 0xd5, 0x1e, 0xb6, 0x8e, 0x79, 0x22, + 0x1b, 0x1d, 0x80, 0x32, 0x3a, 0xe8, 0x53, 0x2d, 0x78, 0x44, 0x7f, 0x8a, 0xe1, 0xbc, 0x54, 0x73, + 0x56, 0x28, 0x45, 0x1d, 0xf9, 0xc5, 0x02, 0x28, 0x09, 0x15, 0x79, 0x75, 0xc5, 0xba, 0xfc, 0x8a, + 0x37, 0x2c, 0x6b, 0x6c, 0x5a, 0x76, 0xb9, 0xb6, 0xd2, 0x38, 0x2d, 0x7b, 0xe0, 0xa3, 0xba, 0x07, + 0x74, 0xf1, 0xee, 0xcc, 0xf5, 0x80, 0x8e, 0x5a, 0x77, 0xc2, 0x0b, 0xe8, 0x35, 0xf0, 0x85, 0xfd, + 0xf0, 0x61, 0xd5, 0x0f, 0xf6, 0xbc, 0x4b, 0xc4, 0x8d, 0xcb, 0xb2, 0x2b, 0x9e, 0x43, 0xaf, 0x01, + 0x2f, 0xf4, 0x38, 0x84, 0x5b, 0x5f, 0x1c, 0xb3, 0x28, 0x66, 0x07, 0xb1, 0x5e, 0x4f, 0xe5, 0x92, + 0x9d, 0x87, 0x49, 0x04, 0xfd, 0xcd, 0xb8, 0xc8, 0x25, 0xcf, 0x8c, 0x3b, 0xb5, 0x99, 0x35, 0x50, + 0x15, 0xaf, 0x06, 0x16, 0xd7, 0xcf, 0xbb, 0x0f, 0x2d, 0x95, 0x46, 0x3d, 0x38, 0x17, 0x73, 0xac, + 0x95, 0x64, 0x1f, 0x3a, 0x1b, 0x61, 0xf0, 0x55, 0x26, 0x8a, 0x74, 0x21, 0xe9, 0xf2, 0xc1, 0xb4, + 0x2f, 0x3e, 0x98, 0xce, 0x85, 0x07, 0xd3, 0xad, 0x1e, 0x4c, 0x12, 0xc2, 0x8a, 0xde, 0x57, 0x6a, + 0x5e, 0xaf, 0xb3, 0xae, 0xca, 0xd7, 0xcc, 0x69, 0xbc, 0x66, 0x21, 0xac, 0xe8, 0xb5, 0xf4, 0x7f, + 0x3a, 0xfd, 0xcd, 0x86, 0x15, 0xca, 0xf3, 0xe8, 0x35, 0x0f, 0x92, 0x5c, 0x66, 0xc5, 0x58, 0x6d, + 0x1f, 0x75, 0xff, 0x1b, 0x71, 0x60, 0xb2, 0xed, 0x50, 0x2d, 0xbc, 0x49, 0xa7, 0x7b, 0x8f, 0xa0, + 0x37, 0x3f, 0x9d, 0x17, 0x4d, 0x9b, 0x26, 0xde, 0x23, 0xb8, 0x11, 0x8a, 0x22, 0x1b, 0x57, 0xed, + 0xdb, 0xd8, 0x88, 0x9a, 0x99, 0x56, 0xd3, 0xd2, 0xac, 0x31, 0x1a, 0xad, 0x2b, 0x46, 0xe3, 0xe9, + 0x5c, 0x2b, 0xf9, 0x6d, 0xbc, 0xf0, 0x4e, 0x7d, 0xe1, 0x9c, 0x9a, 0x9e, 0xb7, 0x26, 0x3f, 0x5b, + 0x70, 0xb3, 0x49, 0xe1, 0x8d, 0x06, 0xb7, 0xaa, 0x88, 0xbd, 0xb0, 0x22, 0xce, 0xa2, 0x8a, 0xb8, + 0x75, 0x45, 0xea, 0x87, 0xb9, 0xd5, 0x78, 0x98, 0xc9, 0x11, 0xdc, 0xbd, 0x50, 0xa6, 0x4d, 0x31, + 0x4d, 0x55, 0x3f, 0xfc, 0x87, 0x72, 0xa9, 0x95, 0x96, 0x65, 0xa6, 0x50, 0x5d, 0xaa, 0x05, 0xf2, + 0x09, 0xdc, 0x09, 0xb9, 0x6c, 0x14, 0xa9, 0xec, 0xb6, 0x01, 0x38, 0x3b, 0xfc, 0xe4, 0x92, 0xcf, + 0x57, 0x2a, 0xf2, 0x39, 0xf8, 0x7b, 0xe9, 0x84, 0x49, 0x7e, 0xad, 0xdb, 0x1b, 0xd0, 0xd9, 0x15, + 0xa9, 0x88, 0xc5, 0xab, 0xd9, 0x15, 0x53, 0xef, 0xc3, 0x0d, 0xbd, 0xbf, 0xf5, 0x1a, 0xe9, 0xd2, + 0x52, 0x24, 0xb7, 0x55, 0x43, 0x8f, 0x59, 0x3c, 0x2e, 0x62, 0x45, 0x43, 0xfd, 0xb4, 0xe5, 0x1b, + 0xcb, 0x7f, 0x9c, 0xad, 0x59, 0x7f, 0x9e, 0xad, 0x59, 0x7f, 0x9d, 0xad, 0x59, 0xbf, 0xfe, 0xbd, + 0xf6, 0xd6, 0x41, 0x1b, 0x7f, 0xd6, 0x9f, 0xfc, 0x13, 0x00, 0x00, 0xff, 0xff, 0x7e, 0x7c, 0x6d, + 0x22, 0xbd, 0x0b, 0x00, 0x00, } diff --git a/internal/private.proto b/internal/private.proto index 408c065e5..95f211a81 100644 --- a/internal/private.proto +++ b/internal/private.proto @@ -43,6 +43,7 @@ message MaxShards { message CreateShardMessage { string Index = 1; + string Field = 3; uint64 Shard = 2; } @@ -105,8 +106,18 @@ message NodeEventMessage { message NodeStatus { Node Node = 1; - MaxShards MaxShards = 2; Schema Schema = 3; + repeated IndexStatus Indexes = 4; +} + +message IndexStatus { + string Name = 1; + repeated FieldStatus Fields = 2; +} + +message FieldStatus { + string Name = 1; + repeated uint64 AvailableShards = 2; } message ClusterStatus { diff --git a/server.go b/server.go index 70fd1398e..cde1dd406 100644 --- a/server.go +++ b/server.go @@ -27,8 +27,8 @@ import ( "sync" "time" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" - "golang.org/x/sync/errgroup" ) @@ -470,11 +470,11 @@ func (s *Server) monitorAntiEntropy() { func (s *Server) receiveMessage(m Message) error { switch obj := m.(type) { case *CreateShardMessage: - idx := s.holder.Index(obj.Index) - if idx == nil { - return fmt.Errorf("Local Index not found: %s", obj.Index) + f := s.holder.Field(obj.Index, obj.Field) + if f == nil { + return fmt.Errorf("Local field not found: %s/%s", obj.Index, obj.Field) } - idx.setRemoteMaxShard(obj.Shard) + f.addRemoteAvailableShards(roaring.NewBitmap(obj.Shard)) case *CreateIndexMessage: opt := obj.Meta _, err := s.holder.CreateIndex(obj.Index, *opt) @@ -628,19 +628,18 @@ func (s *Server) mergeRemoteStatus(ns *NodeStatus) error { return errors.Wrap(err, "applying schema") } - // Sync maxShards. - oldmaxshards := s.holder.maxShards() - for index, newMax := range ns.MaxShards { - localIndex := s.holder.Index(index) - // if we don't know about an index locally, log an error because - // indexes should be created and synced prior to shard creation - if localIndex == nil { - s.logger.Printf("Local Index not found: %s", index) - continue - } - if newMax > oldmaxshards[index] { - oldmaxshards[index] = newMax - localIndex.setRemoteMaxShard(newMax) + // Sync available shards. + for _, is := range ns.Indexes { + for _, fs := range is.Fields { + f := s.holder.Field(is.Name, fs.Name) + + // if we don't know about an field locally, log a error because + // fields should be created and synced prior to shard creation + if f == nil { + s.logger.Printf("Local Field not found: %s/%s", is.Name, fs.Name) + continue + } + f.addRemoteAvailableShards(fs.AvailableShards) } } diff --git a/view.go b/view.go index d46947cc6..2feac8b44 100644 --- a/view.go +++ b/view.go @@ -23,6 +23,7 @@ import ( "sync" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/roaring" "github.com/pkg/errors" ) @@ -48,10 +49,6 @@ type view struct { // Fragments by shard. fragments map[uint64]*fragment - // maxShard maintains this view's max shard in order to - // prevent sending multiple `CreateShardMessage` messages - maxShard uint64 - broadcaster broadcaster stats StatsClient rowAttrStore AttrStore @@ -160,19 +157,16 @@ func (v *view) close() error { return nil } -// calculateMaxShard returns the max shard in the view. -func (v *view) calculateMaxShard() uint64 { +// availableShards returns a bitmap of shards which contain data. +func (v *view) availableShards() *roaring.Bitmap { v.mu.RLock() defer v.mu.RUnlock() - var max uint64 + b := roaring.NewBitmap() for shard := range v.fragments { - if shard > max { - max = shard - } + b.Add(shard) // ignore error, no writer attached } - - return max + return b } // fragmentPath returns the path to a fragment in the view. @@ -229,18 +223,12 @@ func (v *view) createFragmentIfNotExists(shard uint64) (*fragment, error) { frag.RowAttrStore = v.rowAttrStore // Broadcast a message that a new max shard was just created. - if shard > v.maxShard { - v.maxShard = shard - - // Send the create shard message to all nodes. - err := v.broadcaster.SendSync( - &CreateShardMessage{ - Index: v.index, - Shard: shard, - }) - if err != nil { - return nil, errors.Wrap(err, "sending createshard message") - } + if err := v.broadcaster.SendSync(&CreateShardMessage{ + Index: v.index, + Field: v.field, + Shard: shard, + }); err != nil { + return nil, errors.Wrap(err, "sending createshard message") } // Save to lookup.