From 2bd0677df9ab0fb238b1ecbcc43ca55332b6c588 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Thu, 26 Oct 2017 23:33:34 -0500 Subject: [PATCH] Remove MaxSlice polling. Add Schema to proto. Update MaxSlices in proto to include both standard and inverse. Update LocalStatus (shared in gossip) to include MaxSlices and Schema. Encode InputDefinitions with Index. TODO: - make sure InputDefinitions are considered in LocalStatus merge. - decide what to do about `/slices/max` endpoint. client.backup is using it. --- cluster.go | 39 +- cmd/server_test.go | 4 - config.go | 2 - ctl/server.go | 1 - gossip/gossip.go | 3 +- handler.go | 4 +- holder.go | 17 + index.go | 23 +- input_definition.go | 59 ++- internal/private.pb.go | 905 +++++++++++++++++++++++------------------ internal/private.proto | 22 +- server.go | 169 +++----- server/server.go | 13 - 13 files changed, 657 insertions(+), 604 deletions(-) diff --git a/cluster.go b/cluster.go index 4f389fb4b..8bb9b5ff2 100644 --- a/cluster.go +++ b/cluster.go @@ -53,24 +53,7 @@ const ( // Node represents a node in the cluster. type Node struct { - //Scheme string `json:"scheme"` - //Host string `json:"host"` // HostPort URI URI `json:"uri"` - - status *internal.NodeStatus `json:"status"` -} - -// SetStatus sets the NodeStatus. -func (n *Node) SetStatus(s *internal.NodeStatus) { - n.status = s -} - -// SetState sets the Node.status.state. -func (n *Node) SetState(s string) { - if n.status == nil { - n.status = &internal.NodeStatus{} - } - n.status.State = s } // Nodes represents a list of nodes. @@ -242,13 +225,12 @@ func (c *Cluster) URISet() []URI { func (c *Cluster) setState(state string) { c.State = state - localNode := c.localNode() - localNode.SetState(state) } -func (c *Cluster) localNode() *Node { - return c.NodeByURI(c.URI) -} +// localNode is not being used. +//func (c *Cluster) localNode() *Node { +// return c.NodeByURI(c.URI) +//} // Status returns the internal ClusterStatus representation. func (c *Cluster) Status() *internal.ClusterStatus { @@ -258,17 +240,6 @@ func (c *Cluster) Status() *internal.ClusterStatus { } } -/* -// encodeNodeStatuses converts a into its internal representation. -func encodeNodeStatuses(a []*Node) []*internal.NodeStatus { - other := make([]*internal.NodeStatus, len(a)) - for i := range a { - other[i] = a[i].status - } - return other -} -*/ - // NodeByURI returns a node reference by uri. func (c *Cluster) NodeByURI(uri URI) *Node { for _, n := range c.Nodes { @@ -609,7 +580,7 @@ func (c *Cluster) handleJoiningHost(uri URI) error { func (c *Cluster) setStateAndBroadcast(state string) error { c.setState(state) - // Broadcast status changes to the cluster. + // Broadcast cluster status changes to the cluster. return c.Broadcaster.SendSync(c.Status()) } diff --git a/cmd/server_test.go b/cmd/server_test.go index 53841a5d9..b09962663 100644 --- a/cmd/server_test.go +++ b/cmd/server_test.go @@ -51,7 +51,6 @@ func TestServerConfig(t *testing.T) { bind = "localhost:0" [cluster] - poll-interval = "45s" type = "static" replicas = 2 hosts = [ @@ -64,7 +63,6 @@ func TestServerConfig(t *testing.T) { v.Check(cmd.Server.Config.Bind, "localhost:10111") v.Check(cmd.Server.Config.Cluster.ReplicaN, 2) v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:10111", "localhost:10110"}) - v.Check(cmd.Server.Config.Cluster.PollInterval, pilosa.Duration(time.Second*182)) return v.Error() }, }, @@ -99,7 +97,6 @@ func TestServerConfig(t *testing.T) { bind = "localhost:19444" data-dir = "` + actualDataDir + `" [cluster] - poll-interval = "2m0s" hosts = [ "localhost:19444", ] @@ -115,7 +112,6 @@ func TestServerConfig(t *testing.T) { validation: func() error { v := validator{} v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:19444"}) - v.Check(cmd.Server.Config.Cluster.PollInterval, pilosa.Duration(time.Minute*2)) v.Check(cmd.Server.Config.AntiEntropy.Interval, pilosa.Duration(time.Minute*11)) v.Check(cmd.Server.CPUProfile, profFile.Name()) v.Check(cmd.Server.CPUTime, time.Minute) diff --git a/config.go b/config.go index 918c0415a..25ae85cd8 100644 --- a/config.go +++ b/config.go @@ -78,7 +78,6 @@ type Config struct { ReplicaN int `toml:"replicas"` Type string `toml:"type"` Hosts []string `toml:"hosts"` - PollInterval Duration `toml:"poll-interval"` LongQueryTime Duration `toml:"long-query-time"` } `toml:"cluster"` @@ -113,7 +112,6 @@ func NewConfig() *Config { } c.Cluster.ReplicaN = DefaultReplicaN c.Cluster.Type = DefaultClusterType - c.Cluster.PollInterval = Duration(DefaultPollingInterval) c.Cluster.Hosts = []string{} c.AntiEntropy.Interval = Duration(DefaultAntiEntropyInterval) c.Metric.Service = DefaultMetrics diff --git a/ctl/server.go b/ctl/server.go index 5c9fd245e..a2531682d 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -35,7 +35,6 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.IntVarP(&srv.Config.MaxWritesPerRequest, "max-writes-per-request", "", srv.Config.MaxWritesPerRequest, "Number of write commands per request.") flags.IntVarP(&srv.Config.Cluster.ReplicaN, "cluster.replicas", "", 1, "Number of hosts each piece of data should be stored on.") flags.StringSliceVarP(&srv.Config.Cluster.Hosts, "cluster.hosts", "", []string{}, "Comma separated list of hosts in cluster.") - flags.DurationVarP((*time.Duration)(&srv.Config.Cluster.PollInterval), "cluster.poll-interval", "", time.Minute, "Polling interval for cluster.") // TODO what actually is this? flags.DurationVarP((*time.Duration)(&srv.Config.Cluster.LongQueryTime), "cluster.long-query-time", "", time.Minute, "Long Query Time.") flags.StringVarP(&srv.Config.Plugins.Path, "plugins.path", "", "", "Path to plugin directory.") flags.StringVar(&srv.Config.LogPath, "log-path", "", "Log path") diff --git a/gossip/gossip.go b/gossip/gossip.go index e17588e51..2cd2e09e0 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -148,8 +148,7 @@ func NewGossipNodeSet(name string, gossipHost string, gossipPort int, gossipSeed g.config.memberlistConfig.BindPort = gossipPort g.config.memberlistConfig.AdvertiseAddr = pilosa.HostToIP(gossipHost) g.config.memberlistConfig.AdvertisePort = gossipPort - // TODO travis: pause node status (remove this next line) - g.config.memberlistConfig.PushPullInterval = 0 * time.Millisecond + //g.config.memberlistConfig.PushPullInterval = 15 * time.Second // Default is 15s in DefaultLocalConfig. g.config.memberlistConfig.Delegate = g g.config.memberlistConfig.SecretKey = secretKey g.config.memberlistConfig.Events = server.Cluster.EventReceiver.(memberlist.EventDelegate) diff --git a/handler.go b/handler.go index ef5655f43..ffe189e35 100644 --- a/handler.go +++ b/handler.go @@ -132,7 +132,7 @@ func NewRouter(handler *Handler) *mux.Router { router.HandleFunc("/index/{index}/time-quantum", handler.handlePatchIndexTimeQuantum).Methods("PATCH") router.HandleFunc("/hosts", handler.handleGetHosts).Methods("GET") router.HandleFunc("/schema", handler.handleGetSchema).Methods("GET") - router.HandleFunc("/slices/max", handler.handleGetSliceMax).Methods("GET") + //router.HandleFunc("/slices/max", handler.handleGetSliceMax).Methods("GET") // TODO: this is being used by the client (for backups) router.HandleFunc("/status", handler.handleGetStatus).Methods("GET") router.HandleFunc("/version", handler.handleGetVersion).Methods("GET") router.HandleFunc("/recalculate-caches", handler.handleRecalculateCaches).Methods("POST") @@ -305,6 +305,7 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { } } +/* func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) { var ms map[string]uint64 if inverse, _ := strconv.ParseBool(r.URL.Query().Get("inverse")); inverse { @@ -327,6 +328,7 @@ func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) { MaxSlices: ms, }) } +*/ type sliceMaxResponse struct { MaxSlices map[string]uint64 `json:"maxSlices"` diff --git a/holder.go b/holder.go index 909233159..ae314a274 100644 --- a/holder.go +++ b/holder.go @@ -26,6 +26,8 @@ import ( "sync" "syscall" "time" + + "github.com/pilosa/pilosa/internal" ) const ( @@ -185,6 +187,21 @@ func (h *Holder) Schema() []*IndexInfo { return a } +// EncodeMaxSlices creates and internal representation of max slices. +func (h *Holder) EncodeMaxSlices() *internal.MaxSlices { + return &internal.MaxSlices{ + Standard: h.MaxSlices(), + Inverse: h.MaxInverseSlices(), + } +} + +// EncodeSchema creates and internal representation of schema. +func (h *Holder) EncodeSchema() *internal.Schema { + return &internal.Schema{ + Indexes: EncodeIndexes(h.Indexes()), + } +} + // IndexPath returns the path where a given index is stored. func (h *Holder) IndexPath(name string) string { return filepath.Join(h.Path, name) } diff --git a/index.go b/index.go index 00e08002f..734ea7ed5 100644 --- a/index.go +++ b/index.go @@ -392,6 +392,20 @@ func (i *Index) Frames() []*Frame { return a } +// InputDefinitions returns a list of all inputDefinitions in the index. +func (i *Index) InputDefinitions() []*InputDefinition { + i.mu.RLock() + defer i.mu.RUnlock() + + a := make([]*InputDefinition, 0, len(i.inputDefinitions)) + for _, d := range i.inputDefinitions { + a = append(a, d) + } + //sort.Sort(inputDefintionSlice(a)) // TODO + + return a +} + // RecalculateCaches recalculates caches on every frame in the index. func (i *Index) RecalculateCaches() { for _, frame := range i.Frames() { @@ -617,12 +631,10 @@ func EncodeIndexes(a []*Index) []*internal.Index { // encodeIndex converts d into its internal representation. func encodeIndex(d *Index) *internal.Index { - io := d.options() return &internal.Index{ - Name: d.name, - Meta: io.Encode(), - MaxSlice: d.MaxSlice(), - Frames: encodeFrames(d.Frames()), + Name: d.name, + Frames: encodeFrames(d.Frames()), + InputDefinitions: encodeInputDefinitions(d.InputDefinitions()), } } @@ -716,7 +728,6 @@ func (i *Index) newInputDefinition(name string) (*InputDefinition, error) { if err != nil { return nil, err } - inputDef.broadcaster = i.broadcaster return inputDef, nil } diff --git a/input_definition.go b/input_definition.go index 13b27bd17..98ca8ca09 100644 --- a/input_definition.go +++ b/input_definition.go @@ -36,12 +36,11 @@ var validValueDestination = []string{InputMapping, InputValueToRow, InputSingleR // InputDefinition represents a container for the data input definition. type InputDefinition struct { - name string - path string - index string - broadcaster Broadcaster - frames []InputFrame - fields []InputDefinitionField + name string + path string + index string + frames []InputFrame + fields []InputDefinitionField } // NewInputDefinition returns a new instance of InputDefinition. @@ -86,16 +85,9 @@ func (i *InputDefinition) LoadDefinition(pb *internal.InputDefinition) error { // Copy metadata fields. i.name = pb.Name for _, fr := range pb.Frames { - frameMeta := fr.Meta inputFrame := InputFrame{ - Name: fr.Name, - Options: FrameOptions{ - RowLabel: frameMeta.RowLabel, - InverseEnabled: frameMeta.InverseEnabled, - CacheSize: frameMeta.CacheSize, - CacheType: frameMeta.CacheType, - TimeQuantum: TimeQuantum(frameMeta.TimeQuantum), - }, + Name: fr.Name, + Options: *decodeFrameOptions(fr.Meta), } i.frames = append(i.frames, inputFrame) } @@ -327,6 +319,43 @@ func (i *InputDefinitionInfo) Encode() *internal.InputDefinition { return &def } +// encodeInputDefinitions converts a into its internal representation. +func encodeInputDefinitions(a []*InputDefinition) []*internal.InputDefinition { + other := make([]*internal.InputDefinition, len(a)) + for i := range a { + other[i] = encodeInputDefinition(a[i]) + } + return other +} + +// encodeInputDefinition converts i into its internal representation. +func encodeInputDefinition(i *InputDefinition) *internal.InputDefinition { + //fo := f.options() + return &internal.InputDefinition{ + Name: i.name, + Frames: encodeInputFrames(i.frames), + Fields: encodeInputDefinitionFields(i.fields), + } +} + +// encodeInputFrames converts a into its internal representation. +func encodeInputFrames(a []InputFrame) []*internal.Frame { + other := make([]*internal.Frame, len(a)) + for i := range a { + other[i] = a[i].Encode() + } + return other +} + +// encodeInputDefinitionFields converts a into its internal representation. +func encodeInputDefinitionFields(a []InputDefinitionField) []*internal.InputDefinitionField { + other := make([]*internal.InputDefinitionField, len(a)) + for i := range a { + other[i] = a[i].Encode() + } + return other +} + // AddFrame manually add frame to input definition. func (i *InputDefinition) AddFrame(frame InputFrame) error { i.frames = append(i.frames, frame) diff --git a/internal/private.pb.go b/internal/private.pb.go index bab2ad1ba..6993a7070 100644 --- a/internal/private.pb.go +++ b/internal/private.pb.go @@ -14,13 +14,14 @@ BlockDataRequest BlockDataResponse Cache - MaxSlicesResponse + MaxSlices CreateSliceMessage DeleteIndexMessage CreateIndexMessage CreateFrameMessage DeleteFrameMessage Frame + Schema Index InputDefinition InputDefinitionField @@ -248,18 +249,26 @@ func (m *Cache) GetIDs() []uint64 { return nil } -type MaxSlicesResponse struct { - MaxSlices map[string]uint64 `protobuf:"bytes,1,rep,name=MaxSlices" json:"MaxSlices,omitempty" protobuf_key:"bytes,1,opt,name=key,proto3" protobuf_val:"varint,2,opt,name=value,proto3"` +type MaxSlices struct { + Standard map[string]uint64 `protobuf:"bytes,1,rep,name=Standard" json:"Standard,omitempty" protobuf_key:"bytes,1,opt,name=key,proto3" protobuf_val:"varint,2,opt,name=value,proto3"` + Inverse map[string]uint64 `protobuf:"bytes,2,rep,name=Inverse" json:"Inverse,omitempty" protobuf_key:"bytes,1,opt,name=key,proto3" protobuf_val:"varint,2,opt,name=value,proto3"` } -func (m *MaxSlicesResponse) Reset() { *m = MaxSlicesResponse{} } -func (m *MaxSlicesResponse) String() string { return proto.CompactTextString(m) } -func (*MaxSlicesResponse) ProtoMessage() {} -func (*MaxSlicesResponse) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{6} } +func (m *MaxSlices) Reset() { *m = MaxSlices{} } +func (m *MaxSlices) String() string { return proto.CompactTextString(m) } +func (*MaxSlices) ProtoMessage() {} +func (*MaxSlices) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{6} } -func (m *MaxSlicesResponse) GetMaxSlices() map[string]uint64 { +func (m *MaxSlices) GetStandard() map[string]uint64 { if m != nil { - return m.MaxSlices + return m.Standard + } + return nil +} + +func (m *MaxSlices) GetInverse() map[string]uint64 { + if m != nil { + return m.Inverse } return nil } @@ -393,8 +402,9 @@ func (m *DeleteFrameMessage) GetFrame() string { } type Frame struct { - Name string `protobuf:"bytes,1,opt,name=Name,proto3" json:"Name,omitempty"` - Meta *FrameMeta `protobuf:"bytes,2,opt,name=Meta" json:"Meta,omitempty"` + Name string `protobuf:"bytes,1,opt,name=Name,proto3" json:"Name,omitempty"` + Meta *FrameMeta `protobuf:"bytes,2,opt,name=Meta" json:"Meta,omitempty"` + Views []string `protobuf:"bytes,3,rep,name=Views" json:"Views,omitempty"` } func (m *Frame) Reset() { *m = Frame{} } @@ -416,19 +426,42 @@ func (m *Frame) GetMeta() *FrameMeta { return nil } +func (m *Frame) GetViews() []string { + if m != nil { + return m.Views + } + return nil +} + +type Schema struct { + Indexes []*Index `protobuf:"bytes,1,rep,name=Indexes" json:"Indexes,omitempty"` +} + +func (m *Schema) Reset() { *m = Schema{} } +func (m *Schema) String() string { return proto.CompactTextString(m) } +func (*Schema) ProtoMessage() {} +func (*Schema) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{13} } + +func (m *Schema) GetIndexes() []*Index { + if m != nil { + return m.Indexes + } + return nil +} + type Index struct { - Name string `protobuf:"bytes,1,opt,name=Name,proto3" json:"Name,omitempty"` - Meta *IndexMeta `protobuf:"bytes,2,opt,name=Meta" json:"Meta,omitempty"` - MaxSlice uint64 `protobuf:"varint,3,opt,name=MaxSlice,proto3" json:"MaxSlice,omitempty"` - Frames []*Frame `protobuf:"bytes,4,rep,name=Frames" json:"Frames,omitempty"` - Slices []uint64 `protobuf:"varint,5,rep,packed,name=Slices" json:"Slices,omitempty"` + Name string `protobuf:"bytes,1,opt,name=Name,proto3" json:"Name,omitempty"` + // IndexMeta Meta = 2; + // uint64 MaxSlice = 3; + Frames []*Frame `protobuf:"bytes,4,rep,name=Frames" json:"Frames,omitempty"` + // repeated uint64 Slices = 5; InputDefinitions []*InputDefinition `protobuf:"bytes,6,rep,name=InputDefinitions" json:"InputDefinitions,omitempty"` } func (m *Index) Reset() { *m = Index{} } func (m *Index) String() string { return proto.CompactTextString(m) } func (*Index) ProtoMessage() {} -func (*Index) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{13} } +func (*Index) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{14} } func (m *Index) GetName() string { if m != nil { @@ -437,20 +470,6 @@ func (m *Index) GetName() string { return "" } -func (m *Index) GetMeta() *IndexMeta { - if m != nil { - return m.Meta - } - return nil -} - -func (m *Index) GetMaxSlice() uint64 { - if m != nil { - return m.MaxSlice - } - return 0 -} - func (m *Index) GetFrames() []*Frame { if m != nil { return m.Frames @@ -458,13 +477,6 @@ func (m *Index) GetFrames() []*Frame { return nil } -func (m *Index) GetSlices() []uint64 { - if m != nil { - return m.Slices - } - return nil -} - func (m *Index) GetInputDefinitions() []*InputDefinition { if m != nil { return m.InputDefinitions @@ -481,7 +493,7 @@ type InputDefinition struct { func (m *InputDefinition) Reset() { *m = InputDefinition{} } func (m *InputDefinition) String() string { return proto.CompactTextString(m) } func (*InputDefinition) ProtoMessage() {} -func (*InputDefinition) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{14} } +func (*InputDefinition) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{15} } func (m *InputDefinition) GetName() string { if m != nil { @@ -513,7 +525,7 @@ type InputDefinitionField struct { func (m *InputDefinitionField) Reset() { *m = InputDefinitionField{} } func (m *InputDefinitionField) String() string { return proto.CompactTextString(m) } func (*InputDefinitionField) ProtoMessage() {} -func (*InputDefinitionField) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{15} } +func (*InputDefinitionField) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{16} } func (m *InputDefinitionField) GetName() string { if m != nil { @@ -546,7 +558,7 @@ type InputDefinitionAction struct { func (m *InputDefinitionAction) Reset() { *m = InputDefinitionAction{} } func (m *InputDefinitionAction) String() string { return proto.CompactTextString(m) } func (*InputDefinitionAction) ProtoMessage() {} -func (*InputDefinitionAction) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{16} } +func (*InputDefinitionAction) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{17} } func (m *InputDefinitionAction) GetFrame() string { if m != nil { @@ -585,7 +597,7 @@ func (m *CreateInputDefinitionMessage) Reset() { *m = CreateInputDefinit func (m *CreateInputDefinitionMessage) String() string { return proto.CompactTextString(m) } func (*CreateInputDefinitionMessage) ProtoMessage() {} func (*CreateInputDefinitionMessage) Descriptor() ([]byte, []int) { - return fileDescriptorPrivate, []int{17} + return fileDescriptorPrivate, []int{18} } func (m *CreateInputDefinitionMessage) GetIndex() string { @@ -611,7 +623,7 @@ func (m *DeleteInputDefinitionMessage) Reset() { *m = DeleteInputDefinit func (m *DeleteInputDefinitionMessage) String() string { return proto.CompactTextString(m) } func (*DeleteInputDefinitionMessage) ProtoMessage() {} func (*DeleteInputDefinitionMessage) Descriptor() ([]byte, []int) { - return fileDescriptorPrivate, []int{18} + return fileDescriptorPrivate, []int{19} } func (m *DeleteInputDefinitionMessage) GetIndex() string { @@ -637,7 +649,7 @@ type URI struct { func (m *URI) Reset() { *m = URI{} } func (m *URI) String() string { return proto.CompactTextString(m) } func (*URI) ProtoMessage() {} -func (*URI) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{19} } +func (*URI) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{20} } func (m *URI) GetScheme() string { if m != nil { @@ -661,16 +673,15 @@ func (m *URI) GetPort() uint32 { } type NodeStatus struct { - URI *URI `protobuf:"bytes,1,opt,name=URI" json:"URI,omitempty"` - State string `protobuf:"bytes,2,opt,name=State,proto3" json:"State,omitempty"` - Indexes []*Index `protobuf:"bytes,3,rep,name=Indexes" json:"Indexes,omitempty"` - URISet []*URI `protobuf:"bytes,4,rep,name=URISet" json:"URISet,omitempty"` + URI *URI `protobuf:"bytes,1,opt,name=URI" json:"URI,omitempty"` + MaxSlices *MaxSlices `protobuf:"bytes,2,opt,name=MaxSlices" json:"MaxSlices,omitempty"` + Schema *Schema `protobuf:"bytes,3,opt,name=Schema" json:"Schema,omitempty"` } func (m *NodeStatus) Reset() { *m = NodeStatus{} } func (m *NodeStatus) String() string { return proto.CompactTextString(m) } func (*NodeStatus) ProtoMessage() {} -func (*NodeStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{20} } +func (*NodeStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{21} } func (m *NodeStatus) GetURI() *URI { if m != nil { @@ -679,23 +690,16 @@ func (m *NodeStatus) GetURI() *URI { return nil } -func (m *NodeStatus) GetState() string { +func (m *NodeStatus) GetMaxSlices() *MaxSlices { if m != nil { - return m.State - } - return "" -} - -func (m *NodeStatus) GetIndexes() []*Index { - if m != nil { - return m.Indexes + return m.MaxSlices } return nil } -func (m *NodeStatus) GetURISet() []*URI { +func (m *NodeStatus) GetSchema() *Schema { if m != nil { - return m.URISet + return m.Schema } return nil } @@ -708,7 +712,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{21} } +func (*ClusterStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{22} } func (m *ClusterStatus) GetState() string { if m != nil { @@ -734,7 +738,7 @@ type Field struct { func (m *Field) Reset() { *m = Field{} } func (m *Field) String() string { return proto.CompactTextString(m) } func (*Field) ProtoMessage() {} -func (*Field) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{22} } +func (*Field) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{23} } func (m *Field) GetName() string { if m != nil { @@ -773,7 +777,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{24} } func (m *DeleteViewMessage) GetIndex() string { if m != nil { @@ -806,7 +810,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{25} } func (m *ResizeInstruction) GetJobID() int64 { if m != nil { @@ -847,7 +851,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{26} } func (m *ResizeSource) GetURI() *URI { if m != nil { @@ -893,7 +897,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{27} } func (m *ResizeInstructionComplete) GetJobID() int64 { @@ -917,7 +921,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{27} } +func (*Topology) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{28} } func (m *Topology) GetURISet() []*URI { if m != nil { @@ -933,13 +937,14 @@ func init() { proto.RegisterType((*BlockDataRequest)(nil), "internal.BlockDataRequest") proto.RegisterType((*BlockDataResponse)(nil), "internal.BlockDataResponse") proto.RegisterType((*Cache)(nil), "internal.Cache") - proto.RegisterType((*MaxSlicesResponse)(nil), "internal.MaxSlicesResponse") + proto.RegisterType((*MaxSlices)(nil), "internal.MaxSlices") proto.RegisterType((*CreateSliceMessage)(nil), "internal.CreateSliceMessage") proto.RegisterType((*DeleteIndexMessage)(nil), "internal.DeleteIndexMessage") proto.RegisterType((*CreateIndexMessage)(nil), "internal.CreateIndexMessage") proto.RegisterType((*CreateFrameMessage)(nil), "internal.CreateFrameMessage") proto.RegisterType((*DeleteFrameMessage)(nil), "internal.DeleteFrameMessage") proto.RegisterType((*Frame)(nil), "internal.Frame") + proto.RegisterType((*Schema)(nil), "internal.Schema") proto.RegisterType((*Index)(nil), "internal.Index") proto.RegisterType((*InputDefinition)(nil), "internal.InputDefinition") proto.RegisterType((*InputDefinitionField)(nil), "internal.InputDefinitionField") @@ -1216,7 +1221,7 @@ func (m *Cache) MarshalTo(dAtA []byte) (int, error) { return i, nil } -func (m *MaxSlicesResponse) Marshal() (dAtA []byte, err error) { +func (m *MaxSlices) Marshal() (dAtA []byte, err error) { size := m.Size() dAtA = make([]byte, size) n, err := m.MarshalTo(dAtA) @@ -1226,16 +1231,32 @@ func (m *MaxSlicesResponse) Marshal() (dAtA []byte, err error) { return dAtA[:n], nil } -func (m *MaxSlicesResponse) MarshalTo(dAtA []byte) (int, error) { +func (m *MaxSlices) MarshalTo(dAtA []byte) (int, error) { var i int _ = i var l int _ = l - if len(m.MaxSlices) > 0 { - for k, _ := range m.MaxSlices { + if len(m.Standard) > 0 { + for k, _ := range m.Standard { dAtA[i] = 0xa i++ - v := m.MaxSlices[k] + v := m.Standard[k] + mapSize := 1 + len(k) + sovPrivate(uint64(len(k))) + 1 + sovPrivate(uint64(v)) + i = encodeVarintPrivate(dAtA, i, uint64(mapSize)) + dAtA[i] = 0xa + i++ + i = encodeVarintPrivate(dAtA, i, uint64(len(k))) + i += copy(dAtA[i:], k) + dAtA[i] = 0x10 + i++ + i = encodeVarintPrivate(dAtA, i, uint64(v)) + } + } + if len(m.Inverse) > 0 { + for k, _ := range m.Inverse { + dAtA[i] = 0x12 + i++ + v := m.Inverse[k] mapSize := 1 + len(k) + sovPrivate(uint64(len(k))) + 1 + sovPrivate(uint64(v)) i = encodeVarintPrivate(dAtA, i, uint64(mapSize)) dAtA[i] = 0xa @@ -1448,6 +1469,51 @@ func (m *Frame) MarshalTo(dAtA []byte) (int, error) { } i += n9 } + if len(m.Views) > 0 { + for _, s := range m.Views { + dAtA[i] = 0x1a + i++ + l = len(s) + for l >= 1<<7 { + dAtA[i] = uint8(uint64(l)&0x7f | 0x80) + l >>= 7 + i++ + } + dAtA[i] = uint8(l) + i++ + i += copy(dAtA[i:], s) + } + } + return i, nil +} + +func (m *Schema) 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 *Schema) MarshalTo(dAtA []byte) (int, error) { + var i int + _ = i + var l int + _ = l + if len(m.Indexes) > 0 { + for _, msg := range m.Indexes { + dAtA[i] = 0xa + 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 } @@ -1472,21 +1538,6 @@ func (m *Index) MarshalTo(dAtA []byte) (int, error) { i = encodeVarintPrivate(dAtA, i, uint64(len(m.Name))) i += copy(dAtA[i:], m.Name) } - if m.Meta != nil { - dAtA[i] = 0x12 - i++ - i = encodeVarintPrivate(dAtA, i, uint64(m.Meta.Size())) - n10, err := m.Meta.MarshalTo(dAtA[i:]) - if err != nil { - return 0, err - } - i += n10 - } - if m.MaxSlice != 0 { - dAtA[i] = 0x18 - i++ - i = encodeVarintPrivate(dAtA, i, uint64(m.MaxSlice)) - } if len(m.Frames) > 0 { for _, msg := range m.Frames { dAtA[i] = 0x22 @@ -1499,23 +1550,6 @@ func (m *Index) MarshalTo(dAtA []byte) (int, error) { i += n } } - if len(m.Slices) > 0 { - dAtA12 := make([]byte, len(m.Slices)*10) - var j11 int - for _, num := range m.Slices { - for num >= 1<<7 { - dAtA12[j11] = uint8(uint64(num)&0x7f | 0x80) - num >>= 7 - j11++ - } - dAtA12[j11] = uint8(num) - j11++ - } - dAtA[i] = 0x2a - i++ - i = encodeVarintPrivate(dAtA, i, uint64(j11)) - i += copy(dAtA[i:], dAtA12[:j11]) - } if len(m.InputDefinitions) > 0 { for _, msg := range m.InputDefinitions { dAtA[i] = 0x32 @@ -1701,11 +1735,11 @@ func (m *CreateInputDefinitionMessage) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x1a i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Definition.Size())) - n13, err := m.Definition.MarshalTo(dAtA[i:]) + n10, err := m.Definition.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n13 + i += n10 } return i, nil } @@ -1794,41 +1828,31 @@ func (m *NodeStatus) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n14, err := m.URI.MarshalTo(dAtA[i:]) + n11, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n14 + i += n11 } - if len(m.State) > 0 { + if m.MaxSlices != nil { dAtA[i] = 0x12 i++ - i = encodeVarintPrivate(dAtA, i, uint64(len(m.State))) - i += copy(dAtA[i:], m.State) - } - if len(m.Indexes) > 0 { - for _, msg := range m.Indexes { - dAtA[i] = 0x1a - i++ - i = encodeVarintPrivate(dAtA, i, uint64(msg.Size())) - n, err := msg.MarshalTo(dAtA[i:]) - if err != nil { - return 0, err - } - i += n + i = encodeVarintPrivate(dAtA, i, uint64(m.MaxSlices.Size())) + n12, err := m.MaxSlices.MarshalTo(dAtA[i:]) + if err != nil { + return 0, err } + i += n12 } - if len(m.URISet) > 0 { - for _, msg := range m.URISet { - 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 + if m.Schema != nil { + dAtA[i] = 0x1a + i++ + i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) + n13, err := m.Schema.MarshalTo(dAtA[i:]) + if err != nil { + return 0, err } + i += n13 } return i, nil } @@ -1969,21 +1993,21 @@ func (m *ResizeInstruction) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n15, err := m.URI.MarshalTo(dAtA[i:]) + n14, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n15 + i += n14 } if m.Coordinator != nil { dAtA[i] = 0x1a i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Coordinator.Size())) - n16, err := m.Coordinator.MarshalTo(dAtA[i:]) + n15, err := m.Coordinator.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n16 + i += n15 } if len(m.Sources) > 0 { for _, msg := range m.Sources { @@ -2019,11 +2043,11 @@ func (m *ResizeSource) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n17, err := m.URI.MarshalTo(dAtA[i:]) + n16, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n17 + i += n16 } if len(m.Index) > 0 { dAtA[i] = 0x12 @@ -2075,11 +2099,11 @@ func (m *ResizeInstructionComplete) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n18, err := m.URI.MarshalTo(dAtA[i:]) + n17, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n18 + i += n17 } return i, nil } @@ -2237,11 +2261,19 @@ func (m *Cache) Size() (n int) { return n } -func (m *MaxSlicesResponse) Size() (n int) { +func (m *MaxSlices) Size() (n int) { var l int _ = l - if len(m.MaxSlices) > 0 { - for k, v := range m.MaxSlices { + if len(m.Standard) > 0 { + for k, v := range m.Standard { + _ = k + _ = v + mapEntrySize := 1 + len(k) + sovPrivate(uint64(len(k))) + 1 + sovPrivate(uint64(v)) + n += mapEntrySize + 1 + sovPrivate(uint64(mapEntrySize)) + } + } + if len(m.Inverse) > 0 { + for k, v := range m.Inverse { _ = k _ = v mapEntrySize := 1 + len(k) + sovPrivate(uint64(len(k))) + 1 + sovPrivate(uint64(v)) @@ -2334,6 +2366,24 @@ func (m *Frame) Size() (n int) { l = m.Meta.Size() n += 1 + l + sovPrivate(uint64(l)) } + if len(m.Views) > 0 { + for _, s := range m.Views { + l = len(s) + n += 1 + l + sovPrivate(uint64(l)) + } + } + return n +} + +func (m *Schema) Size() (n int) { + var l int + _ = l + if len(m.Indexes) > 0 { + for _, e := range m.Indexes { + l = e.Size() + n += 1 + l + sovPrivate(uint64(l)) + } + } return n } @@ -2344,26 +2394,12 @@ func (m *Index) Size() (n int) { if l > 0 { n += 1 + l + sovPrivate(uint64(l)) } - if m.Meta != nil { - l = m.Meta.Size() - n += 1 + l + sovPrivate(uint64(l)) - } - if m.MaxSlice != 0 { - n += 1 + sovPrivate(uint64(m.MaxSlice)) - } if len(m.Frames) > 0 { for _, e := range m.Frames { l = e.Size() n += 1 + l + sovPrivate(uint64(l)) } } - if len(m.Slices) > 0 { - l = 0 - for _, e := range m.Slices { - l += sovPrivate(uint64(e)) - } - n += 1 + sovPrivate(uint64(l)) + l - } if len(m.InputDefinitions) > 0 { for _, e := range m.InputDefinitions { l = e.Size() @@ -2491,21 +2527,13 @@ func (m *NodeStatus) Size() (n int) { l = m.URI.Size() n += 1 + l + sovPrivate(uint64(l)) } - l = len(m.State) - if l > 0 { + if m.MaxSlices != nil { + l = m.MaxSlices.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)) - } - } - if len(m.URISet) > 0 { - for _, e := range m.URISet { - l = e.Size() - n += 1 + l + sovPrivate(uint64(l)) - } + if m.Schema != nil { + l = m.Schema.Size() + n += 1 + l + sovPrivate(uint64(l)) } return n } @@ -3525,7 +3553,7 @@ func (m *Cache) Unmarshal(dAtA []byte) error { } return nil } -func (m *MaxSlicesResponse) Unmarshal(dAtA []byte) error { +func (m *MaxSlices) Unmarshal(dAtA []byte) error { l := len(dAtA) iNdEx := 0 for iNdEx < l { @@ -3548,15 +3576,15 @@ func (m *MaxSlicesResponse) Unmarshal(dAtA []byte) error { fieldNum := int32(wire >> 3) wireType := int(wire & 0x7) if wireType == 4 { - return fmt.Errorf("proto: MaxSlicesResponse: wiretype end group for non-group") + return fmt.Errorf("proto: MaxSlices: wiretype end group for non-group") } if fieldNum <= 0 { - return fmt.Errorf("proto: MaxSlicesResponse: illegal tag %d (wire type %d)", fieldNum, wire) + return fmt.Errorf("proto: MaxSlices: illegal tag %d (wire type %d)", fieldNum, wire) } switch fieldNum { case 1: if wireType != 2 { - return fmt.Errorf("proto: wrong wireType = %d for field MaxSlices", wireType) + return fmt.Errorf("proto: wrong wireType = %d for field Standard", wireType) } var msglen int for shift := uint(0); ; shift += 7 { @@ -3580,8 +3608,8 @@ func (m *MaxSlicesResponse) Unmarshal(dAtA []byte) error { if postIndex > l { return io.ErrUnexpectedEOF } - if m.MaxSlices == nil { - m.MaxSlices = make(map[string]uint64) + if m.Standard == nil { + m.Standard = make(map[string]uint64) } var mapkey string var mapvalue uint64 @@ -3659,7 +3687,114 @@ func (m *MaxSlicesResponse) Unmarshal(dAtA []byte) error { iNdEx += skippy } } - m.MaxSlices[mapkey] = mapvalue + m.Standard[mapkey] = mapvalue + iNdEx = postIndex + case 2: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Inverse", 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.Inverse == nil { + m.Inverse = make(map[string]uint64) + } + var mapkey string + var mapvalue uint64 + for iNdEx < postIndex { + entryPreIndex := 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) + if fieldNum == 1 { + var stringLenmapkey uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLenmapkey |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLenmapkey := int(stringLenmapkey) + if intStringLenmapkey < 0 { + return ErrInvalidLengthPrivate + } + postStringIndexmapkey := iNdEx + intStringLenmapkey + if postStringIndexmapkey > l { + return io.ErrUnexpectedEOF + } + mapkey = string(dAtA[iNdEx:postStringIndexmapkey]) + iNdEx = postStringIndexmapkey + } else if fieldNum == 2 { + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + mapvalue |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + } else { + iNdEx = entryPreIndex + skippy, err := skipPrivate(dAtA[iNdEx:]) + if err != nil { + return err + } + if skippy < 0 { + return ErrInvalidLengthPrivate + } + if (iNdEx + skippy) > postIndex { + return io.ErrUnexpectedEOF + } + iNdEx += skippy + } + } + m.Inverse[mapkey] = mapvalue iNdEx = postIndex default: iNdEx = preIndex @@ -4331,6 +4466,116 @@ func (m *Frame) Unmarshal(dAtA []byte) error { return err } iNdEx = postIndex + case 3: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field Views", 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.Views = append(m.Views, string(dAtA[iNdEx:postIndex])) + 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 *Schema) 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: Schema: wiretype end group for non-group") + } + if fieldNum <= 0 { + return fmt.Errorf("proto: Schema: illegal tag %d (wire type %d)", fieldNum, wire) + } + switch fieldNum { + case 1: + 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, &Index{}) + 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:]) @@ -4410,58 +4655,6 @@ func (m *Index) Unmarshal(dAtA []byte) error { } m.Name = string(dAtA[iNdEx:postIndex]) iNdEx = postIndex - case 2: - if wireType != 2 { - return fmt.Errorf("proto: wrong wireType = %d for field Meta", 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.Meta == nil { - m.Meta = &IndexMeta{} - } - if err := m.Meta.Unmarshal(dAtA[iNdEx:postIndex]); err != nil { - return err - } - iNdEx = postIndex - case 3: - if wireType != 0 { - return fmt.Errorf("proto: wrong wireType = %d for field MaxSlice", wireType) - } - m.MaxSlice = 0 - for shift := uint(0); ; shift += 7 { - if shift >= 64 { - return ErrIntOverflowPrivate - } - if iNdEx >= l { - return io.ErrUnexpectedEOF - } - b := dAtA[iNdEx] - iNdEx++ - m.MaxSlice |= (uint64(b) & 0x7F) << shift - if b < 0x80 { - break - } - } case 4: if wireType != 2 { return fmt.Errorf("proto: wrong wireType = %d for field Frames", wireType) @@ -4493,68 +4686,6 @@ func (m *Index) Unmarshal(dAtA []byte) error { return err } iNdEx = postIndex - case 5: - 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.Slices = append(m.Slices, 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.Slices = append(m.Slices, v) - } - } else { - return fmt.Errorf("proto: wrong wireType = %d for field Slices", wireType) - } case 6: if wireType != 2 { return fmt.Errorf("proto: wrong wireType = %d for field InputDefinitions", wireType) @@ -5523,36 +5654,7 @@ func (m *NodeStatus) Unmarshal(dAtA []byte) error { iNdEx = postIndex case 2: if wireType != 2 { - return fmt.Errorf("proto: wrong wireType = %d for field State", 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.State = string(dAtA[iNdEx:postIndex]) - iNdEx = postIndex - case 3: - if wireType != 2 { - return fmt.Errorf("proto: wrong wireType = %d for field Indexes", wireType) + return fmt.Errorf("proto: wrong wireType = %d for field MaxSlices", wireType) } var msglen int for shift := uint(0); ; shift += 7 { @@ -5576,14 +5678,16 @@ func (m *NodeStatus) Unmarshal(dAtA []byte) error { if postIndex > l { return io.ErrUnexpectedEOF } - m.Indexes = append(m.Indexes, &Index{}) - if err := m.Indexes[len(m.Indexes)-1].Unmarshal(dAtA[iNdEx:postIndex]); err != nil { + if m.MaxSlices == nil { + m.MaxSlices = &MaxSlices{} + } + if err := m.MaxSlices.Unmarshal(dAtA[iNdEx:postIndex]); err != nil { return err } iNdEx = postIndex - case 4: + case 3: if wireType != 2 { - return fmt.Errorf("proto: wrong wireType = %d for field URISet", wireType) + return fmt.Errorf("proto: wrong wireType = %d for field Schema", wireType) } var msglen int for shift := uint(0); ; shift += 7 { @@ -5607,8 +5711,10 @@ func (m *NodeStatus) Unmarshal(dAtA []byte) error { if postIndex > l { return io.ErrUnexpectedEOF } - m.URISet = append(m.URISet, &URI{}) - if err := m.URISet[len(m.URISet)-1].Unmarshal(dAtA[iNdEx:postIndex]); err != nil { + if m.Schema == nil { + m.Schema = &Schema{} + } + if err := m.Schema.Unmarshal(dAtA[iNdEx:postIndex]); err != nil { return err } iNdEx = postIndex @@ -6672,74 +6778,77 @@ var ( func init() { proto.RegisterFile("private.proto", fileDescriptorPrivate) } var fileDescriptorPrivate = []byte{ - // 1096 bytes of a gzipped FileDescriptorProto - 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x9c, 0x57, 0xcd, 0x6e, 0x23, 0x45, - 0x10, 0x66, 0x3c, 0xb6, 0x63, 0x57, 0xd6, 0x89, 0xd3, 0x2c, 0x91, 0x13, 0x45, 0xde, 0xa8, 0x25, - 0xd8, 0x10, 0x89, 0x00, 0x41, 0x42, 0xfc, 0x1d, 0x60, 0xe3, 0xac, 0x32, 0xb0, 0x59, 0x96, 0x76, - 0xb2, 0xdc, 0x90, 0x3a, 0x4e, 0x93, 0x1d, 0x65, 0x3c, 0x6d, 0x66, 0xda, 0x49, 0xbc, 0x07, 0x6e, - 0xf0, 0x06, 0x48, 0x48, 0x1c, 0x39, 0xf3, 0x1e, 0x1c, 0x79, 0x04, 0x14, 0x2e, 0xbc, 0x01, 0x12, - 0x27, 0xd4, 0xd5, 0x3d, 0x3f, 0x1e, 0xff, 0x84, 0xe4, 0x36, 0xf5, 0x75, 0x75, 0xd5, 0xd7, 0x5f, - 0x57, 0x95, 0xdb, 0xd0, 0x18, 0x44, 0xfe, 0x05, 0x57, 0x62, 0x67, 0x10, 0x49, 0x25, 0x49, 0xcd, - 0x0f, 0x95, 0x88, 0x42, 0x1e, 0xd0, 0x2f, 0xa1, 0xee, 0x85, 0xa7, 0xe2, 0xea, 0x50, 0x28, 0x4e, - 0x36, 0x61, 0x71, 0x4f, 0x06, 0xc3, 0x7e, 0xf8, 0x84, 0x9f, 0x88, 0xa0, 0xe5, 0x6c, 0x3a, 0x5b, - 0x75, 0x96, 0x87, 0xb4, 0xc7, 0x91, 0xdf, 0x17, 0x5f, 0x0d, 0x79, 0xa8, 0x86, 0xfd, 0x56, 0xc9, - 0x78, 0xe4, 0x20, 0xfa, 0xaf, 0x03, 0xf5, 0xc7, 0x11, 0xef, 0x0b, 0x8c, 0xb8, 0x0e, 0x35, 0x26, - 0x2f, 0xf3, 0xe1, 0x52, 0x9b, 0xbc, 0x01, 0x4b, 0x5e, 0x78, 0x21, 0xa2, 0x58, 0xec, 0x87, 0xfc, - 0x24, 0x10, 0xa7, 0x18, 0xae, 0xc6, 0x0a, 0x28, 0xd9, 0x80, 0xfa, 0x1e, 0xef, 0xbd, 0x10, 0x47, - 0xa3, 0x81, 0x68, 0xb9, 0x18, 0x24, 0x03, 0xd2, 0xd5, 0xae, 0xff, 0x52, 0xb4, 0xca, 0x9b, 0xce, - 0x56, 0x83, 0x65, 0x40, 0x91, 0x6f, 0x65, 0x82, 0x2f, 0xa1, 0x70, 0x8f, 0xf1, 0xf0, 0x2c, 0xe5, - 0x50, 0x45, 0x0e, 0x63, 0x18, 0x79, 0x08, 0xd5, 0xc7, 0xbe, 0x08, 0x4e, 0xe3, 0xd6, 0xc2, 0xa6, - 0xbb, 0xb5, 0xb8, 0xbb, 0xbc, 0x93, 0xe8, 0xb7, 0x83, 0x38, 0xb3, 0xcb, 0x94, 0xc2, 0x92, 0xd7, - 0x1f, 0xc8, 0x48, 0x31, 0x11, 0x0f, 0x64, 0x18, 0x0b, 0xd2, 0x04, 0x77, 0x3f, 0x8a, 0xec, 0xd9, - 0xf5, 0x27, 0xfd, 0x1e, 0x9a, 0x8f, 0x02, 0xd9, 0x3b, 0xef, 0x70, 0xc5, 0x99, 0xf8, 0x6e, 0x28, - 0x62, 0x45, 0xee, 0x43, 0x05, 0x6f, 0xc1, 0xfa, 0x19, 0x43, 0xa3, 0xa8, 0xa4, 0x95, 0xd9, 0x18, - 0x1a, 0xc5, 0xfd, 0x28, 0x45, 0x99, 0x19, 0x43, 0xa3, 0xdd, 0xc0, 0xef, 0x19, 0x09, 0xca, 0xcc, - 0x18, 0x84, 0x40, 0xf9, 0xb9, 0x2f, 0x2e, 0xed, 0xb9, 0xf1, 0x9b, 0x7a, 0xb0, 0x92, 0xcb, 0x6f, - 0x69, 0xae, 0x42, 0x95, 0xc9, 0x4b, 0xaf, 0x13, 0xb7, 0x9c, 0x4d, 0x77, 0xab, 0xcc, 0xac, 0x85, - 0xea, 0xe2, 0xf5, 0xeb, 0xa5, 0x12, 0x2e, 0x65, 0x00, 0x5d, 0x83, 0x0a, 0x4a, 0xad, 0x4f, 0x99, - 0xed, 0xd5, 0x9f, 0xf4, 0x17, 0x07, 0x56, 0x0e, 0xf9, 0x15, 0xd2, 0x88, 0xd3, 0x34, 0x07, 0x50, - 0x4f, 0x41, 0xf4, 0x5e, 0xdc, 0xdd, 0xce, 0xb4, 0x9c, 0xf0, 0xcf, 0x90, 0xfd, 0x50, 0x45, 0x23, - 0x96, 0x6d, 0x5e, 0xff, 0x04, 0x96, 0xc6, 0x17, 0x35, 0x87, 0x73, 0x31, 0x4a, 0x94, 0x3e, 0x17, - 0x23, 0xad, 0xc9, 0x05, 0x0f, 0x86, 0x46, 0xbf, 0x32, 0x33, 0xc6, 0x47, 0xa5, 0x0f, 0x1c, 0xfa, - 0x0d, 0x90, 0xbd, 0x48, 0x70, 0x25, 0x30, 0xc0, 0xa1, 0x88, 0x63, 0x7e, 0x26, 0x66, 0xdf, 0x82, - 0x51, 0xb6, 0x94, 0x57, 0x76, 0x03, 0xea, 0x5e, 0x6c, 0x0b, 0x15, 0x6f, 0xa2, 0xc6, 0x32, 0x80, - 0x6e, 0x03, 0xe9, 0x88, 0x40, 0x28, 0x61, 0x7b, 0x6b, 0x4e, 0x7c, 0xda, 0x4d, 0xb8, 0xdc, 0xec, - 0x4b, 0x1e, 0x42, 0x59, 0xb7, 0x15, 0x52, 0x59, 0xdc, 0x7d, 0x35, 0x93, 0x2e, 0xed, 0x61, 0x86, - 0x0e, 0xd4, 0x4f, 0x82, 0xda, 0x56, 0xbc, 0xe1, 0x80, 0x53, 0xca, 0x2c, 0x49, 0xe5, 0x16, 0x53, - 0xa5, 0xcd, 0x6d, 0x53, 0x7d, 0x9a, 0x9c, 0xf5, 0xae, 0xa9, 0x68, 0xc7, 0xa2, 0xba, 0x5c, 0x9f, - 0xea, 0x55, 0xb3, 0x07, 0xbf, 0x67, 0x1f, 0xb9, 0xc8, 0xe3, 0x6f, 0xc7, 0xa6, 0xbc, 0x5d, 0x98, - 0x82, 0x72, 0x7a, 0x62, 0x25, 0x85, 0x65, 0x3b, 0x2c, 0xb5, 0x71, 0x0e, 0xe8, 0xac, 0x71, 0xab, - 0x3c, 0x31, 0x07, 0x34, 0xce, 0xec, 0xb2, 0x6e, 0x27, 0x5b, 0xe4, 0x15, 0xd3, 0x4e, 0xc6, 0x22, - 0xfb, 0xd0, 0xf4, 0xc2, 0xc1, 0x50, 0x75, 0xc4, 0xb7, 0x7e, 0xe8, 0x2b, 0x5f, 0x86, 0x71, 0xab, - 0x8a, 0xa1, 0xd6, 0xf2, 0x8c, 0xc6, 0x3c, 0xd8, 0xc4, 0x16, 0xfa, 0xa3, 0x03, 0xcb, 0x05, 0x70, - 0xc6, 0xa1, 0x13, 0xbe, 0xa5, 0xf9, 0x7c, 0xdf, 0x4f, 0x07, 0x9c, 0x8b, 0x8e, 0xed, 0x99, 0x6c, - 0xc6, 0xe7, 0xdd, 0xaf, 0x0e, 0xdc, 0x9f, 0xe6, 0x30, 0x95, 0x4d, 0x1b, 0xe0, 0x59, 0xe4, 0xf7, - 0x79, 0x34, 0xfa, 0x42, 0x8c, 0xec, 0xac, 0xcf, 0x21, 0xe4, 0x6b, 0x58, 0x2d, 0xc4, 0xfa, 0xac, - 0x67, 0x24, 0x32, 0xa4, 0x1e, 0xcc, 0x24, 0x65, 0xfc, 0xd8, 0x8c, 0xed, 0xf4, 0x1f, 0x07, 0x5e, - 0x9b, 0xba, 0x94, 0xd5, 0xa3, 0x93, 0x2f, 0xfd, 0x6d, 0x68, 0x3e, 0xd7, 0xa3, 0xa2, 0x23, 0x62, - 0xe5, 0x87, 0x5c, 0x7b, 0xda, 0x82, 0x9d, 0xc0, 0x89, 0x07, 0x35, 0xc4, 0x0e, 0xf9, 0xc0, 0xd2, - 0x7c, 0xeb, 0x06, 0x9a, 0x3b, 0x89, 0xbf, 0x99, 0x69, 0xe9, 0x76, 0x4d, 0x06, 0xa7, 0x6e, 0x32, - 0xc2, 0xd1, 0x58, 0xff, 0x18, 0x1a, 0x63, 0x1b, 0x6e, 0x35, 0xe7, 0x24, 0x6c, 0x24, 0xb3, 0x65, - 0x8c, 0xc9, 0xfc, 0x2e, 0xfd, 0x10, 0x20, 0x73, 0xb5, 0x03, 0x60, 0x4e, 0x7d, 0xe6, 0x9c, 0xe9, - 0x01, 0x6c, 0x24, 0x83, 0xef, 0x16, 0x09, 0x93, 0x6a, 0x29, 0x65, 0xd5, 0x42, 0xf7, 0xc1, 0x3d, - 0x66, 0x1e, 0x76, 0x52, 0xef, 0x85, 0x48, 0xaf, 0xc8, 0x5a, 0x7a, 0xcb, 0x81, 0x8c, 0x55, 0xb2, - 0x45, 0x7f, 0x6b, 0xec, 0x99, 0x8c, 0x14, 0x32, 0x6e, 0x30, 0xfc, 0xa6, 0x3f, 0x39, 0x00, 0x4f, - 0xe5, 0xa9, 0xe8, 0x2a, 0xae, 0x86, 0x31, 0x79, 0x80, 0x51, 0x31, 0xd6, 0xe2, 0x6e, 0x23, 0x3b, - 0xd3, 0x31, 0xf3, 0x18, 0xe6, 0xd3, 0xd3, 0x5e, 0x71, 0x95, 0x4e, 0x28, 0x34, 0xc8, 0x9b, 0xb0, - 0x80, 0x4c, 0x45, 0x52, 0x8b, 0xcb, 0x85, 0x01, 0xc2, 0x92, 0x75, 0xf2, 0x3a, 0x54, 0x8f, 0x99, - 0xd7, 0x15, 0xca, 0xce, 0x88, 0x42, 0x12, 0xbb, 0x48, 0x9f, 0x40, 0x63, 0x2f, 0x18, 0xc6, 0x4a, - 0x44, 0x96, 0x59, 0x9a, 0xd8, 0xc9, 0x27, 0xce, 0xa2, 0x95, 0xe6, 0x45, 0xeb, 0x42, 0x65, 0x76, - 0xdf, 0x11, 0x28, 0xe3, 0xd3, 0xc9, 0x4a, 0x85, 0xaf, 0xa6, 0x26, 0xb8, 0x87, 0xbe, 0xb9, 0x5b, - 0x97, 0xe9, 0x4f, 0x44, 0xf8, 0x15, 0xd6, 0x9e, 0x46, 0xb8, 0xfe, 0x61, 0x5a, 0x31, 0x77, 0xa9, - 0x9f, 0x0d, 0x77, 0xf9, 0x09, 0x49, 0x5e, 0x1f, 0x6e, 0xee, 0xf5, 0xf1, 0x9b, 0x03, 0x2b, 0x4c, - 0xc4, 0xfe, 0x4b, 0xe1, 0x85, 0xb1, 0x8a, 0x86, 0x69, 0x1f, 0x7e, 0x2e, 0x4f, 0xbc, 0x0e, 0x46, - 0x75, 0x99, 0x31, 0x92, 0xcb, 0x2a, 0xcd, 0xbc, 0xac, 0xb7, 0xf5, 0x7b, 0x55, 0x46, 0xa7, 0xba, - 0x19, 0x65, 0x64, 0x2b, 0xb5, 0xe0, 0x98, 0xf7, 0x20, 0xef, 0xc0, 0x42, 0x57, 0x0e, 0xa3, 0x5e, - 0x3a, 0xc1, 0x57, 0x33, 0x67, 0xc3, 0xca, 0x2c, 0xb3, 0xc4, 0x8d, 0xfe, 0xe0, 0xc0, 0xbd, 0xfc, - 0xca, 0xff, 0xaa, 0x20, 0xa3, 0x50, 0x69, 0xaa, 0x42, 0xee, 0x34, 0x85, 0xca, 0x99, 0x42, 0xd9, - 0x7b, 0xa3, 0x92, 0x7b, 0x6f, 0x50, 0x06, 0x6b, 0x13, 0xb2, 0xed, 0xc9, 0xfe, 0x40, 0xdf, 0xcf, - 0x1d, 0xe5, 0xa3, 0xef, 0x42, 0xed, 0x48, 0x0e, 0x64, 0x20, 0xcf, 0x46, 0xb9, 0x42, 0x73, 0xe6, - 0x14, 0xda, 0xa3, 0xe6, 0xef, 0xd7, 0x6d, 0xe7, 0x8f, 0xeb, 0xb6, 0xf3, 0xe7, 0x75, 0xdb, 0xf9, - 0xf9, 0xaf, 0xf6, 0x2b, 0x27, 0x55, 0xfc, 0x47, 0xf1, 0xde, 0x7f, 0x01, 0x00, 0x00, 0xff, 0xff, - 0xd7, 0x59, 0xea, 0x61, 0x62, 0x0c, 0x00, 0x00, + // 1137 bytes of a gzipped FileDescriptorProto + 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x9c, 0x57, 0x4f, 0x6f, 0x1b, 0x45, + 0x14, 0x67, 0xbd, 0x8e, 0x6b, 0x3f, 0xc7, 0x8d, 0x33, 0x94, 0xc8, 0x89, 0x22, 0xd7, 0x8c, 0x04, + 0x0d, 0x95, 0x08, 0x34, 0x95, 0x10, 0x04, 0x21, 0x41, 0xe3, 0x54, 0x5d, 0x68, 0x4a, 0x19, 0x27, + 0x45, 0xe2, 0x80, 0x34, 0xb1, 0x87, 0x74, 0x95, 0xf5, 0x8e, 0xd9, 0x1d, 0x27, 0x71, 0x0f, 0xdc, + 0xe0, 0x00, 0x5f, 0x80, 0x3b, 0x67, 0xbe, 0x07, 0x47, 0x3e, 0x02, 0x0a, 0x1f, 0x02, 0x89, 0x0b, + 0x68, 0xde, 0xce, 0xec, 0xae, 0xff, 0xc5, 0x4a, 0x6e, 0xfb, 0xde, 0xbc, 0xf7, 0xe6, 0x37, 0xbf, + 0xf7, 0x67, 0x66, 0xa1, 0x36, 0x88, 0xfc, 0x33, 0xae, 0xc4, 0xf6, 0x20, 0x92, 0x4a, 0x92, 0xb2, + 0x1f, 0x2a, 0x11, 0x85, 0x3c, 0xa0, 0x5f, 0x42, 0xc5, 0x0b, 0x7b, 0xe2, 0xe2, 0x40, 0x28, 0x4e, + 0x5a, 0x50, 0xdd, 0x93, 0xc1, 0xb0, 0x1f, 0x3e, 0xe5, 0xc7, 0x22, 0x68, 0x38, 0x2d, 0x67, 0xab, + 0xc2, 0xf2, 0x2a, 0x6d, 0x71, 0xe8, 0xf7, 0xc5, 0x57, 0x43, 0x1e, 0xaa, 0x61, 0xbf, 0x51, 0x48, + 0x2c, 0x72, 0x2a, 0xfa, 0xaf, 0x03, 0x95, 0xc7, 0x11, 0xef, 0x0b, 0x8c, 0xb8, 0x01, 0x65, 0x26, + 0xcf, 0xf3, 0xe1, 0x52, 0x99, 0xbc, 0x0d, 0xb7, 0xbd, 0xf0, 0x4c, 0x44, 0xb1, 0xd8, 0x0f, 0xf9, + 0x71, 0x20, 0x7a, 0x18, 0xae, 0xcc, 0x26, 0xb4, 0x64, 0x13, 0x2a, 0x7b, 0xbc, 0xfb, 0x52, 0x1c, + 0x8e, 0x06, 0xa2, 0xe1, 0x62, 0x90, 0x4c, 0x91, 0xae, 0x76, 0xfc, 0x57, 0xa2, 0x51, 0x6c, 0x39, + 0x5b, 0x35, 0x96, 0x29, 0x26, 0xf1, 0x2e, 0x4d, 0xe1, 0x25, 0x14, 0x96, 0x19, 0x0f, 0x4f, 0x52, + 0x0c, 0x25, 0xc4, 0x30, 0xa6, 0x23, 0xf7, 0xa0, 0xf4, 0xd8, 0x17, 0x41, 0x2f, 0x6e, 0xdc, 0x6a, + 0xb9, 0x5b, 0xd5, 0x9d, 0x95, 0x6d, 0xcb, 0xdf, 0x36, 0xea, 0x99, 0x59, 0xa6, 0x14, 0x6e, 0x7b, + 0xfd, 0x81, 0x8c, 0x14, 0x13, 0xf1, 0x40, 0x86, 0xb1, 0x20, 0x75, 0x70, 0xf7, 0xa3, 0xc8, 0x9c, + 0x5d, 0x7f, 0xd2, 0x1f, 0xa0, 0xfe, 0x28, 0x90, 0xdd, 0xd3, 0x36, 0x57, 0x9c, 0x89, 0xef, 0x87, + 0x22, 0x56, 0xe4, 0x0e, 0x2c, 0x61, 0x16, 0x8c, 0x5d, 0x22, 0x68, 0x2d, 0x32, 0x69, 0x68, 0x4e, + 0x04, 0xad, 0x45, 0x7f, 0xa4, 0xa2, 0xc8, 0x12, 0x41, 0x6b, 0x3b, 0x81, 0xdf, 0x4d, 0x28, 0x28, + 0xb2, 0x44, 0x20, 0x04, 0x8a, 0x2f, 0x7c, 0x71, 0x6e, 0xce, 0x8d, 0xdf, 0xd4, 0x83, 0xd5, 0xdc, + 0xfe, 0x06, 0xe6, 0x1a, 0x94, 0x98, 0x3c, 0xf7, 0xda, 0x71, 0xc3, 0x69, 0xb9, 0x5b, 0x45, 0x66, + 0x24, 0x64, 0x17, 0xd3, 0xaf, 0x97, 0x0a, 0xb8, 0x94, 0x29, 0xe8, 0x3a, 0x2c, 0x21, 0xd5, 0xfa, + 0x94, 0x99, 0xaf, 0xfe, 0xa4, 0xff, 0x39, 0x50, 0x39, 0xe0, 0x17, 0x08, 0x23, 0x26, 0x9f, 0x40, + 0xb9, 0xa3, 0x78, 0xd8, 0xe3, 0x51, 0x0f, 0x8d, 0xaa, 0x3b, 0x6f, 0x66, 0x14, 0xa6, 0x66, 0xdb, + 0xd6, 0x66, 0x3f, 0x54, 0xd1, 0x88, 0xa5, 0x2e, 0x64, 0x17, 0x6e, 0x99, 0x9a, 0x40, 0x0c, 0xd5, + 0x9d, 0xd6, 0x2c, 0xef, 0xb4, 0x6c, 0xb4, 0xb3, 0x75, 0xd8, 0xf8, 0x18, 0x6a, 0x63, 0x61, 0x35, + 0xd6, 0x53, 0x31, 0xb2, 0x19, 0x39, 0x15, 0x23, 0xcd, 0xdd, 0x19, 0x0f, 0x86, 0x09, 0xcf, 0x45, + 0x96, 0x08, 0xbb, 0x85, 0x0f, 0x9d, 0x8d, 0x5d, 0x58, 0xce, 0x47, 0xbd, 0x8e, 0x2f, 0xfd, 0x16, + 0xc8, 0x5e, 0x24, 0xb8, 0x12, 0x08, 0xef, 0x40, 0xc4, 0x31, 0x3f, 0x11, 0xf3, 0x33, 0x9d, 0x64, + 0xaf, 0x90, 0xcf, 0xde, 0x26, 0x54, 0xbc, 0xd8, 0x1e, 0xdc, 0xc5, 0xba, 0xcc, 0x14, 0xf4, 0x3e, + 0x90, 0xb6, 0x08, 0x84, 0x12, 0xa6, 0x7f, 0xaf, 0x88, 0x4f, 0x3b, 0x16, 0xcb, 0x62, 0x5b, 0x72, + 0x0f, 0x8a, 0xba, 0x75, 0x11, 0x4a, 0x75, 0xe7, 0xf5, 0x8c, 0xe9, 0x74, 0x4e, 0x30, 0x34, 0xa0, + 0xbe, 0x0d, 0x6a, 0xda, 0x7d, 0xc1, 0x01, 0x67, 0x94, 0xb2, 0xdd, 0xca, 0x9d, 0xdc, 0x2a, 0x1d, + 0x20, 0x66, 0xab, 0x4f, 0xed, 0x59, 0x6f, 0xba, 0x15, 0xfd, 0xc6, 0x68, 0x75, 0x4b, 0x3c, 0xd3, + 0xab, 0x89, 0x0f, 0x7e, 0xcf, 0x3f, 0xf2, 0x04, 0x0e, 0x1d, 0x5b, 0xf7, 0x50, 0xdc, 0x70, 0x5b, + 0xae, 0x8e, 0x8d, 0x02, 0x7d, 0x08, 0xa5, 0x4e, 0xf7, 0xa5, 0xe8, 0x73, 0xf2, 0x8e, 0x2e, 0xd4, + 0x9e, 0xb8, 0x10, 0xb1, 0x29, 0xf3, 0x95, 0x09, 0xfa, 0x98, 0x5d, 0xa7, 0xbf, 0x38, 0x06, 0xfd, + 0x1c, 0x44, 0x25, 0xdc, 0x3b, 0x6e, 0x14, 0xa7, 0x26, 0x8e, 0xd6, 0x33, 0xb3, 0x4c, 0xf6, 0xa1, + 0xee, 0x85, 0x83, 0xa1, 0x6a, 0x8b, 0xef, 0xfc, 0xd0, 0x57, 0xbe, 0x0c, 0xe3, 0x46, 0x09, 0x5d, + 0xd6, 0xf3, 0x5b, 0x8f, 0x59, 0xb0, 0x29, 0x17, 0xfa, 0x93, 0x03, 0x2b, 0x13, 0xca, 0x05, 0xb8, + 0x0a, 0x57, 0xe3, 0xfa, 0x20, 0x1d, 0x99, 0x2e, 0x1a, 0x36, 0xe7, 0xa2, 0x19, 0x9f, 0xa0, 0xbf, + 0x39, 0x70, 0x67, 0x96, 0xc1, 0x4c, 0x34, 0x4d, 0x80, 0xe7, 0x91, 0xdf, 0xe7, 0xd1, 0xe8, 0x0b, + 0x31, 0x32, 0xb7, 0x47, 0x4e, 0x43, 0xbe, 0x86, 0xb5, 0x89, 0x58, 0x9f, 0x75, 0x13, 0x8a, 0x12, + 0x50, 0x77, 0xe7, 0x82, 0x4a, 0xec, 0xd8, 0x1c, 0x77, 0xfa, 0x8f, 0x03, 0x6f, 0xcc, 0x5c, 0xca, + 0xaa, 0xcf, 0xc9, 0x17, 0xfa, 0x7d, 0xa8, 0xbf, 0xd0, 0x83, 0xa1, 0x2d, 0x62, 0xe5, 0x87, 0x5c, + 0x5b, 0x9a, 0xf2, 0x9c, 0xd2, 0x13, 0x0f, 0xca, 0xa8, 0x3b, 0xe0, 0x03, 0x03, 0xf3, 0xdd, 0x05, + 0x30, 0xb7, 0xad, 0xbd, 0x99, 0x9b, 0x56, 0xd4, 0x60, 0x70, 0x8e, 0xdb, 0x4b, 0x01, 0x05, 0x3d, + 0x11, 0xc7, 0x1c, 0xae, 0x35, 0xd5, 0x24, 0x6c, 0xda, 0x49, 0x32, 0x86, 0xe4, 0xea, 0x9e, 0xfc, + 0x08, 0x20, 0x33, 0x35, 0xed, 0x7e, 0x45, 0x7d, 0xe6, 0x8c, 0xe9, 0x13, 0xd8, 0xb4, 0x63, 0xee, + 0x1a, 0x1b, 0xda, 0x6a, 0x29, 0x64, 0xd5, 0x42, 0xf7, 0xc1, 0x3d, 0x62, 0x9e, 0xbe, 0xea, 0xb0, + 0x5b, 0x6d, 0x8a, 0x8c, 0xa4, 0x5d, 0x9e, 0xc8, 0x58, 0x59, 0x17, 0xfd, 0xad, 0x75, 0xcf, 0x65, + 0xa4, 0x10, 0x71, 0x8d, 0xe1, 0x37, 0xfd, 0xd9, 0x01, 0x78, 0x26, 0x7b, 0xa2, 0xa3, 0xb8, 0x1a, + 0xc6, 0xe4, 0x2e, 0x46, 0xc5, 0x58, 0xd5, 0x9d, 0x5a, 0x76, 0xa6, 0x23, 0xe6, 0x31, 0xdc, 0xef, + 0x41, 0xee, 0x22, 0x9c, 0x9e, 0x30, 0xe9, 0x12, 0xcb, 0x5d, 0x97, 0x5b, 0x76, 0xa0, 0x18, 0xaa, + 0xea, 0x99, 0x7d, 0xa2, 0x37, 0xa0, 0x39, 0x7d, 0x0a, 0xb5, 0xbd, 0x60, 0x18, 0x2b, 0x11, 0x19, + 0x38, 0xfa, 0x26, 0x51, 0x5c, 0xa5, 0xf5, 0x87, 0x02, 0x79, 0x0b, 0x4a, 0x47, 0xcc, 0xeb, 0x08, + 0x65, 0xda, 0x76, 0x02, 0xa7, 0x59, 0xa4, 0x1d, 0x58, 0x9a, 0xdf, 0x6c, 0x04, 0x8a, 0xf8, 0x02, + 0x33, 0xfc, 0xe0, 0xe3, 0xab, 0x0e, 0xee, 0x81, 0x9f, 0x24, 0xd4, 0x65, 0xfa, 0x13, 0x35, 0xfc, + 0x02, 0x0b, 0x4e, 0x6b, 0xb8, 0xbe, 0x7b, 0x56, 0x93, 0x04, 0xea, 0x61, 0x79, 0x93, 0x5b, 0xc2, + 0x3e, 0x62, 0xdc, 0xdc, 0x23, 0xe6, 0x77, 0x07, 0x56, 0x99, 0x88, 0xfd, 0x57, 0xc2, 0x0b, 0x63, + 0x15, 0x0d, 0xd3, 0xe6, 0xfb, 0x5c, 0x1e, 0x7b, 0x6d, 0x8c, 0xea, 0xb2, 0x44, 0xb0, 0x19, 0x2a, + 0xcc, 0xcd, 0xd0, 0x7b, 0xfa, 0xd9, 0x2b, 0xa3, 0x9e, 0xee, 0x40, 0x19, 0x19, 0xce, 0x27, 0x0c, + 0xf3, 0x16, 0xe4, 0x7d, 0xb8, 0xd5, 0x91, 0xc3, 0xa8, 0x9b, 0x8e, 0xe7, 0xb5, 0xcc, 0x38, 0x41, + 0x95, 0x2c, 0x33, 0x6b, 0x46, 0x7f, 0x74, 0x60, 0x39, 0xbf, 0xb2, 0xb8, 0x6c, 0x52, 0x86, 0x0a, + 0x33, 0x19, 0x72, 0x67, 0x31, 0x54, 0xcc, 0x18, 0xca, 0x9e, 0x14, 0x4b, 0xb9, 0x27, 0x05, 0x65, + 0xb0, 0x3e, 0x45, 0xdb, 0x9e, 0xec, 0x0f, 0x74, 0x7e, 0x6e, 0x48, 0x1f, 0x7d, 0x00, 0xe5, 0x43, + 0x39, 0x90, 0x81, 0x3c, 0x19, 0xe5, 0x0a, 0xcd, 0xb9, 0xa2, 0xd0, 0x1e, 0xd5, 0xff, 0xb8, 0x6c, + 0x3a, 0x7f, 0x5e, 0x36, 0x9d, 0xbf, 0x2e, 0x9b, 0xce, 0xaf, 0x7f, 0x37, 0x5f, 0x3b, 0x2e, 0xe1, + 0x8f, 0xc9, 0xc3, 0xff, 0x03, 0x00, 0x00, 0xff, 0xff, 0xb2, 0x75, 0xc6, 0x9e, 0xa9, 0x0c, 0x00, + 0x00, } diff --git a/internal/private.proto b/internal/private.proto index fdb5fe082..5101a6c35 100644 --- a/internal/private.proto +++ b/internal/private.proto @@ -38,8 +38,13 @@ message Cache { repeated uint64 IDs = 1; } -message MaxSlicesResponse { - map MaxSlices = 1; +//message MaxSlicesResponse { +// map MaxSlices = 1; +//} + +message MaxSlices { + map Standard = 1; + map Inverse = 2; } message CreateSliceMessage { @@ -71,14 +76,16 @@ message DeleteFrameMessage { message Frame { string Name = 1; FrameMeta Meta = 2; + repeated string Views = 3; +} + +message Schema { + repeated Index Indexes = 1; } message Index { string Name = 1; - IndexMeta Meta = 2; - uint64 MaxSlice = 3; repeated Frame Frames = 4; - repeated uint64 Slices = 5; repeated InputDefinition InputDefinitions = 6; } @@ -120,9 +127,8 @@ message URI { message NodeStatus { URI URI = 1; - string State = 2; - repeated Index Indexes = 3; - repeated URI URISet = 4; + MaxSlices MaxSlices = 2; + Schema Schema = 3; } message ClusterStatus { diff --git a/server.go b/server.go index 4e4a507fd..3139300c6 100644 --- a/server.go +++ b/server.go @@ -19,11 +19,9 @@ import ( "errors" "fmt" "io" - "io/ioutil" "log" "net" "net/http" - "net/url" "os" "os/exec" "runtime" @@ -40,7 +38,6 @@ import ( // Default server settings. const ( DefaultAntiEntropyInterval = 10 * time.Minute - DefaultPollingInterval = 60 * time.Second ) // Server represents a holder wrapped by a running HTTP server. @@ -65,7 +62,6 @@ type Server struct { // Background monitoring intervals. AntiEntropyInterval time.Duration - PollingInterval time.Duration MetricInterval time.Duration // TLS configuration @@ -92,7 +88,6 @@ func NewServer() *Server { Network: "tcp", AntiEntropyInterval: DefaultAntiEntropyInterval, - PollingInterval: DefaultPollingInterval, MetricInterval: 0, LogOutput: os.Stderr, @@ -198,9 +193,8 @@ func (s *Server) Open() error { /* // Start background monitoring. - s.wg.Add(3) + s.wg.Add(2) go func() { defer s.wg.Done(); s.monitorAntiEntropy() }() - go func() { defer s.wg.Done(); s.monitorMaxSlices() }() go func() { defer s.wg.Done(); s.monitorRuntime() }() */ @@ -275,39 +269,6 @@ func (s *Server) monitorAntiEntropy() { s.Holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0) } -// monitorMaxSlices periodically pulls the highest slice from each node in the cluster. -func (s *Server) monitorMaxSlices() { - ticker := time.NewTicker(s.PollingInterval) - defer ticker.Stop() - - for { - select { - case <-s.closing: - return - case <-ticker.C: - } - - oldmaxslices := s.Holder.MaxSlices() - for _, node := range s.Cluster.Nodes { - if s.URI != node.URI { - maxSlices, _ := s.checkMaxSlices(node.URI) - for index, newmax := range maxSlices { - // if we don't know about an index locally, log an error because - // indexes should be created and synced prior to slice creation - if localIndex := s.Holder.Index(index); localIndex != nil { - if newmax > oldmaxslices[index] { - oldmaxslices[index] = newmax - localIndex.SetRemoteMaxSlice(newmax) - } - } else { - s.Logger().Printf("Local Index not found: %s", index) - } - } - } - } - } -} - // ReceiveMessage represents an implementation of BroadcastHandler. func (s *Server) ReceiveMessage(pb proto.Message) error { switch obj := pb.(type) { @@ -392,10 +353,16 @@ func (s *Server) State() string { return s.Cluster.State } -// LocalStatus returns the state of the local node as well as the -// holder (indexes/frames) according to the local node. -// In a gossip implementation, memberlist.Delegate.LocalState() uses this. // Server implements StatusHandler. +// LocalStatus is used to periodically sync information +// between nodes. Under normal conditions, nodes should +// remain in sync through Broadcast messages. For cases +// where a node fails to receive a Broadcast message, or +// when a new (empty) node needs to get in sync with the +// rest of the cluster, two things are shared via gossip: +// - MaxSlice/MaxInverseSlice by Index +// - Schema +// In a gossip implementation, memberlist.Delegate.LocalState() uses this. func (s *Server) LocalStatus() (proto.Message, error) { if s.Cluster == nil { return nil, errors.New("Server.Cluster is nil") @@ -405,61 +372,66 @@ func (s *Server) LocalStatus() (proto.Message, error) { } ns := internal.NodeStatus{ - URI: encodeURI(s.URI), - State: s.State(), - Indexes: EncodeIndexes(s.Holder.Indexes()), - URISet: encodeURIs(s.Cluster.URISet()), - } - - // TODO: get rid of this - // Append Slice list per this Node's indexes - for _, index := range ns.Indexes { - index.Slices = s.Cluster.OwnsSlices(index.Name, index.MaxSlice, s.URI) + URI: encodeURI(s.URI), + MaxSlices: s.Holder.EncodeMaxSlices(), + Schema: s.Holder.EncodeSchema(), } return &ns, nil } -// ClusterStatus returns the NodeState for all nodes in the cluster. +// ClusterStatus returns the ClusterState and URISet for the cluster. func (s *Server) ClusterStatus() (proto.Message, error) { - // Update local Node.state. - ns, err := s.LocalStatus() - if err != nil { - return nil, err - } - localNode := s.Cluster.localNode() - localNode.SetStatus(ns.(*internal.NodeStatus)) - return s.Cluster.Status(), nil } -// HandleRemoteStatus receives incoming NodeState from remote nodes. +// HandleRemoteStatus receives incoming NodeStatus from remote nodes. func (s *Server) HandleRemoteStatus(pb proto.Message) error { return s.mergeRemoteStatus(pb.(*internal.NodeStatus)) } func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { - // Ignore status updates from self. if s.URI == decodeURI(ns.URI) { return nil } - fmt.Printf("mergeRemoteStatus on (%s) from (%s)\n", s.URI, ns.URI) - // Update Node.state. - // Node can be nil if a merge occurs (via gossip) before the coordinator has - // a chance to broadcast the existence of the node. - uri := decodeURI(ns.URI) - if node := s.Cluster.NodeByURI(uri); node != nil { - node.SetStatus(ns) + // Sync maxSlices (standard). + oldmaxslices := s.Holder.MaxSlices() + for index, newMax := range ns.MaxSlices.Standard { + 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 slice creation + if localIndex == nil { + s.Logger().Printf("Local Index not found: %s", index) + continue + } + if newMax > oldmaxslices[index] { + oldmaxslices[index] = newMax + localIndex.SetRemoteMaxSlice(newMax) + } } - // Create indexes that don't exist. - for _, index := range ns.Indexes { - opt := IndexOptions{ - ColumnLabel: index.Meta.ColumnLabel, - TimeQuantum: TimeQuantum(index.Meta.TimeQuantum), + // Sync maxSlices (inverse). + oldMaxInverseSlices := s.Holder.MaxInverseSlices() + for index, newMaxInverse := range ns.MaxSlices.Inverse { + 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 slice creation + if localIndex == nil { + s.Logger().Printf("Local Index not found: %s", index) + continue } + if newMaxInverse > oldMaxInverseSlices[index] { + oldMaxInverseSlices[index] = newMaxInverse + localIndex.SetRemoteMaxSlice(newMaxInverse) + } + } + + // Sync schema. + // Create indexes that don't exist. + for _, index := range ns.Schema.Indexes { + opt := IndexOptions{} idx, err := s.Holder.CreateIndexIfNotExists(index.Name, opt) if err != nil { return err @@ -472,55 +444,12 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { return err } } + // TODO: Create inputDefinitions that don't exist. } return nil } -func (s *Server) checkMaxSlices(uri URI) (map[string]uint64, error) { - // Create HTTP request. - req, err := http.NewRequest("GET", (&url.URL{ - Scheme: uri.Scheme(), - Host: uri.HostPort(), - Path: "/slices/max", - }).String(), nil) - - if err != nil { - return nil, err - } - - // Require protobuf encoding. - req.Header.Set("Accept", "application/x-protobuf") - req.Header.Set("Content-Type", "application/x-protobuf") - req.Header.Set("User-Agent", "pilosa/"+Version) - - resp, err := s.defaultClient.Do(req) - if err != nil { - return nil, err - } - defer resp.Body.Close() - - // Read response into buffer. - body, err := ioutil.ReadAll(resp.Body) - if err != nil { - return nil, err - } - - // Check status code. - if resp.StatusCode != http.StatusOK { - return nil, fmt.Errorf("invalid status checkMaxSlices: code=%d, err=%s, req=%v", resp.StatusCode, body, req) - } - - // Decode response object. - pb := internal.MaxSlicesResponse{} - - if err = proto.Unmarshal(body, &pb); err != nil { - return nil, err - } - - return pb.MaxSlices, nil -} - // monitorRuntime periodically polls the Go runtime metrics. func (s *Server) monitorRuntime() { // Disable metrics when poll interval is zero diff --git a/server/server.go b/server/server.go index 10220c1d4..34d397fcf 100644 --- a/server/server.go +++ b/server/server.go @@ -124,19 +124,6 @@ func (m *Command) SetupServer() error { cluster.ReplicaN = m.Config.Cluster.ReplicaN cluster.IndexReporter = m.Server.Holder - /* - // TODO travis: get rid of this URI code - for _, address := range m.Config.Cluster.Hosts { - uri, err := pilosa.NewURIFromAddress(address) - if err != nil { - return err - } - cluster.Nodes = append(cluster.Nodes, &pilosa.Node{ - Scheme: uri.Scheme(), - Host: uri.HostPort(), - }) - } - */ m.Server.Cluster = cluster // Setup logging output.