mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-06 19:07:50 +00:00
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.
This commit is contained in:
parent
b22780ccf1
commit
f4c9c0fed3
17 changed files with 833 additions and 277 deletions
12
api.go
12
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
|
||||
|
|
|
|||
|
|
@ -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]]
|
||||
}
|
||||
|
|
|
|||
36
cluster.go
36
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{}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
15
executor.go
15
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
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
25
field.go
25
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.
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
17
holder.go
17
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() {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
30
index.go
30
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.
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
35
server.go
35
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)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
36
view.go
36
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.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue