Merge pull request #1600 from benbjohnson/available-shards

Maintain available shards set
This commit is contained in:
Travis Turner 2018-08-23 07:32:37 -05:00 • committed by GitHub
commit 10eea2db4c
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
17 changed files with 833 additions and 277 deletions

12
api.go
View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

36
view.go
View file

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