Change coordinator to primary

Signed-off-by: Antonio Navarro Perez <antnavper@gmail.com>
This commit is contained in:
Antonio Navarro Perez 2021-02-02 19:09:53 +01:00 • committed by Travis
parent 28a19cdff7
commit c45e21640c
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
22 changed files with 215 additions and 829 deletions

36
api.go
View file

@ -844,8 +844,8 @@ func (api *API) Node() *topology.Node {
return api.server.node()
}
// CoordinatorNode returns the coordinator node for the cluster.
func (api *API) CoordinatorNode() *topology.Node {
// PrimaryNode returns the coordinator node for the cluster.
func (api *API) PrimaryNode() *topology.Node {
// Create a snapshot of the cluster to use for node/partition calculations.
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
return snap.PrimaryFieldTranslationNode()
@ -1738,38 +1738,6 @@ func (api *API) indexField(indexName string, fieldName string, shard uint64) (*I
return index, field, nil
}
// SetCoordinator makes a new Node the cluster coordinator.
func (api *API) SetCoordinator(ctx context.Context, id string) (oldNode, newNode *topology.Node, err error) {
span, _ := tracing.StartSpanFromContext(ctx, "API.SetCoordinator")
defer span.Finish()
if err := api.validate(apiSetCoordinator); err != nil {
return nil, nil, errors.Wrap(err, "validating api method")
}
oldNode = api.cluster.nodeByID(api.cluster.Coordinator)
newNode = api.cluster.nodeByID(id)
if newNode == nil {
return nil, nil, errors.Wrap(ErrNodeIDNotExists, "getting new node")
}
// If the new coordinator is this node, do the SetCoordinator directly.
if newNode.ID == api.Node().ID {
return oldNode, newNode, api.cluster.setCoordinator(newNode)
}
// Send the set-coordinator message to new node.
err = api.server.SendTo(
newNode,
&SetCoordinatorMessage{
New: newNode,
})
if err != nil {
return nil, nil, fmt.Errorf("problem sending SetCoordinator message: %s", err)
}
return oldNode, newNode, nil
}
// RemoveNode puts the cluster into the "RESIZING" state and begins the job of
// removing the given node.
func (api *API) RemoveNode(id string) (*topology.Node, error) {

View file

@ -106,10 +106,6 @@ func getMessage(typ byte) Message {
return &ResizeInstruction{}
case messageTypeResizeInstructionComplete:
return &ResizeInstructionComplete{}
case messageTypeSetCoordinator:
return &SetCoordinatorMessage{}
case messageTypeUpdateCoordinator:
return &UpdateCoordinatorMessage{}
case messageTypeNodeState:
return &NodeStateMessage{}
case messageTypeRecalculateCaches:
@ -147,10 +143,6 @@ func getMessageType(m Message) byte {
return messageTypeResizeInstruction
case *ResizeInstructionComplete:
return messageTypeResizeInstructionComplete
case *SetCoordinatorMessage:
return messageTypeSetCoordinator
case *UpdateCoordinatorMessage:
return messageTypeUpdateCoordinator
case *NodeStateMessage:
return messageTypeNodeState
case *RecalculateCaches:

View file

@ -108,7 +108,6 @@ type cluster struct { // nolint: maligned
// Required for cluster Resize.
Static bool // Static is primarily used for testing in a non-gossip environment.
state string
Coordinator string
holder *Holder
broadcaster broadcaster
@ -228,18 +227,6 @@ func (c *cluster) setCoordinator(n *topology.Node) error {
return fmt.Errorf("coordinator node does not match this node")
}
// Update IsCoordinator on all nodes (locally).
_ = c.unprotectedUpdateCoordinator(n)
// Send the update coordinator message to all nodes.
err := c.unprotectedSendSync(
&UpdateCoordinatorMessage{
New: n,
})
if err != nil {
return fmt.Errorf("problem sending UpdateCoordinator message: %v", err)
}
// Broadcast cluster status.
return c.unprotectedSendSync(c.unprotectedStatus())
}
@ -262,40 +249,9 @@ func (c *cluster) unprotectedSendSync(m Message) error {
return eg.Wait()
}
// updateCoordinator updates this nodes Coordinator value as well as
// changing the corresponding node's IsCoordinator value
// to true, and sets all other nodes to false. Returns true if the value
// changed.
func (c *cluster) updateCoordinator(n *topology.Node) bool { // nolint: unparam
c.mu.Lock()
defer c.mu.Unlock()
return c.unprotectedUpdateCoordinator(n)
}
func (c *cluster) unprotectedUpdateCoordinator(n *topology.Node) bool {
var changed bool
if c.Coordinator != n.ID {
c.Coordinator = n.ID
changed = true
}
for _, node := range c.noder.Nodes() {
if node.ID == n.ID {
node.IsCoordinator = true
} else {
node.IsCoordinator = false
}
}
return changed
}
// addNode adds a node to the Cluster and updates and saves the
// new topology. unprotected.
func (c *cluster) addNode(node *topology.Node) error {
// If the node being added is the coordinator, set it for this node.
if node.IsCoordinator {
c.Coordinator = node.ID
}
// add to cluster
if !c.addNodeBasicSorted(node) {
return nil
@ -541,9 +497,9 @@ func (c *cluster) addNodeBasicSorted(node *topology.Node) bool {
n.Mu.Lock()
defer n.Mu.Unlock()
if n.State != node.State || n.IsCoordinator != node.IsCoordinator || n.URI != node.URI {
if n.State != node.State || n.IsPrimary != node.IsPrimary || n.URI != node.URI {
n.State = node.State
n.IsCoordinator = node.IsCoordinator
n.IsPrimary = node.IsPrimary
n.URI = node.URI
n.GRPCURI = node.GRPCURI
return true
@ -570,7 +526,7 @@ func (c *cluster) Nodes() []*topology.Node {
// Set node states and IsPrimary.
for _, node := range nodes {
node.IsCoordinator = node.ID == primaryNode.ID
node.IsPrimary = node.ID == primaryNode.ID
s, err := c.stator.NodeState(context.Background(), node.ID)
if err != nil {
@ -1432,7 +1388,7 @@ func (c *cluster) unprotectedGenerateResizeJobByAction(nodeAction nodeAction) (*
instr := &ResizeInstruction{
JobID: j.ID,
Node: toCluster.unprotectedNodeByID(node.ID),
Coordinator: snap.PrimaryFieldTranslationNode(),
Primary: snap.PrimaryFieldTranslationNode(),
Sources: fragmentSourcesByNode[node.ID],
TranslationSources: translationSourcesByNode[node.ID],
NodeStatus: c.nodeStatus(), // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node.
@ -1609,7 +1565,7 @@ func (c *cluster) followResizeInstruction(instr *ResizeInstruction) error {
complete.Error = err.Error()
}
if err := c.sendTo(instr.Coordinator, complete); err != nil {
if err := c.sendTo(instr.Primary, complete); err != nil {
c.logger.Printf("sending resizeInstructionComplete error: err=%s", err)
}
}()
@ -2361,6 +2317,7 @@ func (c *cluster) PrimaryReplicaNode() *topology.Node {
func (c *cluster) unprotectedPrimaryReplicaNode() *topology.Node {
pos := c.nodePositionByID(c.Node.ID)
if pos <= 0 {
fmt.Println("----------------------- PRIMARY NOT FOUND")
return nil
}
cNodes := c.noder.Nodes()
@ -2975,7 +2932,7 @@ type ClusterStatus struct {
type ResizeInstruction struct {
JobID int64
Node *topology.Node
Coordinator *topology.Node
Primary *topology.Node
Sources []*ResizeSource
TranslationSources []*TranslationResizeSource
NodeStatus *NodeStatus
@ -3101,16 +3058,6 @@ type ResizeInstructionComplete struct {
Error string
}
// SetCoordinatorMessage is an internal message instructing nodes to honor a new coordinator.
type SetCoordinatorMessage struct {
New *topology.Node
}
// UpdateCoordinatorMessage is an internal message for reassigning the coordinator.
type UpdateCoordinatorMessage struct {
New *topology.Node
}
// NodeStateMessage is an internal message for broadcasting a node's state.
type NodeStateMessage struct {
NodeID string `protobuf:"bytes,1,opt,name=NodeID,proto3" json:"NodeID,omitempty"`

View file

@ -288,8 +288,8 @@ func TestFragSources(t *testing.T) {
"node0": {},
"node1": {},
"node2": {
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(0)},
{&topology.Node{ID: "node1", URI: pnet.URI{Scheme: "http", Host: "host1", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(2)},
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(0)},
{&topology.Node{ID: "node1", URI: pnet.URI{Scheme: "http", Host: "host1", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(2)},
},
},
err: "",
@ -300,11 +300,11 @@ func TestFragSources(t *testing.T) {
idx: idx,
expected: map[string][]*ResizeSource{
"node0": {
{&topology.Node{ID: "node1", URI: pnet.URI{Scheme: "http", Host: "host1", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(1)},
{&topology.Node{ID: "node1", URI: pnet.URI{Scheme: "http", Host: "host1", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(1)},
},
"node1": {
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(0)},
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(2)},
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(0)},
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(2)},
},
},
err: "",
@ -315,11 +315,11 @@ func TestFragSources(t *testing.T) {
idx: idx,
expected: map[string][]*ResizeSource{
"node0": {
{&topology.Node{ID: "node2", URI: pnet.URI{Scheme: "http", Host: "host2", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(0)},
{&topology.Node{ID: "node2", URI: pnet.URI{Scheme: "http", Host: "host2", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(2)},
{&topology.Node{ID: "node2", URI: pnet.URI{Scheme: "http", Host: "host2", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(0)},
{&topology.Node{ID: "node2", URI: pnet.URI{Scheme: "http", Host: "host2", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(2)},
},
"node1": {
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsCoordinator: false}, "i", "f", "standard", uint64(3)},
{&topology.Node{ID: "node0", URI: pnet.URI{Scheme: "http", Host: "host0", Port: 10101}, IsPrimary: false}, "i", "f", "standard", uint64(3)},
},
"node2": {},
},
@ -617,6 +617,9 @@ func TestCluster_PreviousNode(t *testing.T) {
// NEXT: move this test to internal and unexport IsCoordinator
func TestCluster_Coordinator(t *testing.T) {
// TODO check if this test still makes sense
t.Skip()
const urisCount = 2
var uris []pnet.URI
if err := port.GetPorts(func(ports []int) error {
@ -634,11 +637,11 @@ func TestCluster_Coordinator(t *testing.T) {
c1 := *newCluster()
c1.Node = node1
c1.Coordinator = node1.ID
// c1.Coordinator = node1.ID
c1.noder = noder
c2 := *newCluster()
c2.Node = node2
c2.Coordinator = node1.ID
// c2.Coordinator = node1.ID
c2.noder = noder
t.Run("IsCoordinator", func(t *testing.T) {
@ -1070,32 +1073,6 @@ func TestAE(t *testing.T) {
})
}
// Ensures that coordinator can be changed.
func TestCluster_UpdateCoordinator(t *testing.T) {
t.Run("UpdateCoordinator", func(t *testing.T) {
c := NewTestCluster(t, 2)
cNodes := c.noder.Nodes()
oldNode := cNodes[0]
newNode := cNodes[1]
// Update coordinator to the same value.
if c.updateCoordinator(oldNode) {
t.Errorf("did not expect coordinator to change")
} else if c.Coordinator != oldNode.ID {
t.Errorf("expected coordinator: %s, but got: %s", c.Coordinator, oldNode.URI)
}
// Update coordinator to a new value.
if !c.updateCoordinator(newNode) {
t.Errorf("expected coordinator to change")
} else if c.Coordinator != newNode.ID {
t.Errorf("expected coordinator: %s, but got: %s", c.Coordinator, newNode.URI)
}
})
}
func TestCluster_confirmNodeDownUp(t *testing.T) {
t.Skip("does a listen on :0, skip for now. TODO(jea) restore this.")
r := mux.NewRouter()

View file

@ -47,7 +47,6 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
flags.StringSliceVar(&srv.Config.Handler.AllowedOrigins, "handler.allowed-origins", []string{}, "Comma separated list of allowed origin URIs (for CORS/Web UI).")
// Cluster
flags.BoolVar(&srv.Config.Cluster.Coordinator, "cluster.coordinator", srv.Config.Cluster.Coordinator, "Host that will act as cluster coordinator during startup and resizing.")
flags.IntVar(&srv.Config.Cluster.ReplicaN, "cluster.replicas", 1, "Number of hosts each piece of data should be stored on.")
flags.DurationVar((*time.Duration)(&srv.Config.Cluster.LongQueryTime), "cluster.long-query-time", time.Duration(srv.Config.Cluster.LongQueryTime), "RENAMED TO 'long-query-time': Duration that will trigger log and stat messages for slow queries.") // negative duration indicates invalid value because 0 is meaningful
flags.StringVar(&srv.Config.Cluster.Name, "cluster.name", srv.Config.Cluster.Name, "Human-readable name for the cluster.")

View file

@ -138,22 +138,6 @@ func (s Serializer) Unmarshal(buf []byte, m pilosa.Message) error {
}
s.decodeResizeInstructionComplete(msg, mt)
return nil
case *pilosa.SetCoordinatorMessage:
msg := &internal.SetCoordinatorMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling SetCoordinatorMessage")
}
s.decodeSetCoordinatorMessage(msg, mt)
return nil
case *pilosa.UpdateCoordinatorMessage:
msg := &internal.UpdateCoordinatorMessage{}
err := proto.Unmarshal(buf, msg)
if err != nil {
return errors.Wrap(err, "unmarshaling UpdateCoordinatorMessage")
}
s.decodeUpdateCoordinatorMessage(msg, mt)
return nil
case *pilosa.NodeStateMessage:
msg := &internal.NodeStateMessage{}
err := proto.Unmarshal(buf, msg)
@ -351,10 +335,6 @@ func (s Serializer) encodeToProto(m pilosa.Message) proto.Message {
return s.encodeResizeInstruction(mt)
case *pilosa.ResizeInstructionComplete:
return s.encodeResizeInstructionComplete(mt)
case *pilosa.SetCoordinatorMessage:
return s.encodeSetCoordinatorMessage(mt)
case *pilosa.UpdateCoordinatorMessage:
return s.encodeUpdateCoordinatorMessage(mt)
case *pilosa.NodeStateMessage:
return s.encodeNodeStateMessage(mt)
case *pilosa.RecalculateCaches:
@ -574,7 +554,7 @@ func (s Serializer) encodeResizeInstruction(m *pilosa.ResizeInstruction) *intern
return &internal.ResizeInstruction{
JobID: m.JobID,
Node: s.encodeNode(m.Node),
Coordinator: s.encodeNode(m.Coordinator),
Primary: s.encodeNode(m.Primary),
Sources: s.encodeResizeSources(m.Sources),
TranslationSources: s.encodeTranslationResizeSources(m.TranslationSources),
NodeStatus: s.encodeNodeStatus(m.NodeStatus),
@ -693,11 +673,10 @@ func (s Serializer) encodeNodes(a []*topology.Node) []*internal.Node {
func (s Serializer) encodeNode(m *topology.Node) *internal.Node {
n := m.ProtectedClone()
return &internal.Node{
ID: n.ID,
URI: s.encodeURI(n.URI),
IsCoordinator: n.IsCoordinator,
State: n.State,
GRPCURI: s.encodeURI(n.GRPCURI),
ID: n.ID,
URI: s.encodeURI(n.URI),
State: n.State,
GRPCURI: s.encodeURI(n.GRPCURI),
}
}
@ -795,18 +774,6 @@ func (s Serializer) encodeResizeInstructionComplete(m *pilosa.ResizeInstructionC
}
}
func (s Serializer) encodeSetCoordinatorMessage(m *pilosa.SetCoordinatorMessage) *internal.SetCoordinatorMessage {
return &internal.SetCoordinatorMessage{
New: s.encodeNode(m.New),
}
}
func (s Serializer) encodeUpdateCoordinatorMessage(m *pilosa.UpdateCoordinatorMessage) *internal.UpdateCoordinatorMessage {
return &internal.UpdateCoordinatorMessage{
New: s.encodeNode(m.New),
}
}
func (s Serializer) encodeNodeStateMessage(m *pilosa.NodeStateMessage) *internal.NodeStateMessage {
return &internal.NodeStateMessage{
NodeID: m.NodeID,
@ -953,8 +920,8 @@ func (s Serializer) decodeResizeInstruction(ri *internal.ResizeInstruction, m *p
m.JobID = ri.JobID
m.Node = &topology.Node{}
s.decodeNode(ri.Node, m.Node)
m.Coordinator = &topology.Node{}
s.decodeNode(ri.Coordinator, m.Coordinator)
m.Primary = &topology.Node{}
s.decodeNode(ri.Primary, m.Primary)
m.Sources = make([]*pilosa.ResizeSource, len(ri.Sources))
s.decodeResizeSources(ri.Sources, m.Sources)
m.TranslationSources = make([]*pilosa.TranslationResizeSource, len(ri.TranslationSources))
@ -1073,7 +1040,6 @@ func (s Serializer) decodeNode(node *internal.Node, m *topology.Node) {
m.ID = node.ID
s.decodeURI(node.URI, &m.URI)
s.decodeURI(node.GRPCURI, &m.GRPCURI)
m.IsCoordinator = node.IsCoordinator
m.State = node.State
}
@ -1145,16 +1111,6 @@ func (s Serializer) decodeResizeInstructionComplete(pb *internal.ResizeInstructi
m.Error = pb.Error
}
func (s Serializer) decodeSetCoordinatorMessage(pb *internal.SetCoordinatorMessage, m *pilosa.SetCoordinatorMessage) {
m.New = &topology.Node{}
s.decodeNode(pb.New, m.New)
}
func (s Serializer) decodeUpdateCoordinatorMessage(pb *internal.UpdateCoordinatorMessage, m *pilosa.UpdateCoordinatorMessage) {
m.New = &topology.Node{}
s.decodeNode(pb.New, m.New)
}
func (s Serializer) decodeNodeStateMessage(pb *internal.NodeStateMessage, m *pilosa.NodeStateMessage) {
m.NodeID = pb.NodeID
m.State = pb.State

View file

@ -647,7 +647,7 @@ func (h *Holder) Open() error {
return errors.Wrap(err, "opening index")
}
if h.isCoordinator() {
if h.isPrimary() {
index.createdAt = timestamp()
err = index.OpenWithTimestamp()
} else {
@ -1201,9 +1201,9 @@ func (h *Holder) recalculateCaches() {
}
// TODO: this needs to be removed
func (h *Holder) isCoordinator() bool {
func (h *Holder) isPrimary() bool {
if s, ok := h.broadcaster.(*Server); ok {
return s.isCoordinator
return s.IsPrimary()
}
return false
}

View file

@ -173,7 +173,7 @@ func (c *InternalClient) CreateIndex(ctx context.Context, index string, opt pilo
if err != nil {
return fmt.Errorf("getting nodes: %s", err)
}
coord := getCoordinatorNode(nodes)
coord := getPrimaryNode(nodes)
if coord == nil {
return fmt.Errorf("could not find the coordinator node")
}
@ -370,9 +370,9 @@ func (c *InternalClient) Import(ctx context.Context, index, field string, shard
return nil
}
func getCoordinatorNode(nodes []*topology.Node) *topology.Node {
func getPrimaryNode(nodes []*topology.Node) *topology.Node {
for _, node := range nodes {
if node.IsCoordinator {
if node.IsPrimary {
return node
}
}
@ -417,7 +417,7 @@ func (c *InternalClient) ImportK(ctx context.Context, index, field string, bits
if err != nil {
return fmt.Errorf("getting nodes: %s", err)
}
coord := getCoordinatorNode(nodes)
coord := getPrimaryNode(nodes)
if coord == nil {
return fmt.Errorf("could not find the coordinator node")
}
@ -629,7 +629,7 @@ func (c *InternalClient) ImportValueK(ctx context.Context, index, field string,
if err != nil {
return fmt.Errorf("getting nodes: %s", err)
}
coord := getCoordinatorNode(nodes)
coord := getPrimaryNode(nodes)
if coord == nil {
return fmt.Errorf("could not find the coordinator node")
}
@ -939,7 +939,7 @@ func (c *InternalClient) CreateFieldWithOptions(ctx context.Context, index, fiel
if err != nil {
return fmt.Errorf("getting nodes: %s", err)
}
coord := getCoordinatorNode(nodes)
coord := getPrimaryNode(nodes)
if coord == nil {
return fmt.Errorf("could not find the coordinator node")
}

View file

@ -367,7 +367,6 @@ func newRouter(handler *Handler) http.Handler {
router := mux.NewRouter()
router.HandleFunc("/cluster/resize/abort", handler.handlePostClusterResizeAbort).Methods("POST").Name("PostClusterResizeAbort")
router.HandleFunc("/cluster/resize/remove-node", handler.handlePostClusterResizeRemoveNode).Methods("POST").Name("PostClusterResizeRemoveNode")
router.HandleFunc("/cluster/resize/set-coordinator", handler.handlePostClusterResizeSetCoordinator).Methods("POST").Name("PostClusterResizeSetCoordinator")
router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET")
router.Handle("/debug/vars", expvar.Handler()).Methods("GET")
router.Handle("/metrics", promhttp.Handler())
@ -2029,38 +2028,6 @@ func parseUint64Slice(s string) ([]uint64, error) {
return a, nil
}
func (h *Handler) handlePostClusterResizeSetCoordinator(w http.ResponseWriter, r *http.Request) {
if !validHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
// Decode request.
var req setCoordinatorRequest
err := json.NewDecoder(r.Body).Decode(&req)
if err != nil {
http.Error(w, "decoding request "+err.Error(), http.StatusBadRequest)
return
}
oldNode, newNode, err := h.api.SetCoordinator(r.Context(), req.ID)
if err != nil {
if errors.Cause(err) == pilosa.ErrNodeIDNotExists {
http.Error(w, "setting new coordinator: "+err.Error(), http.StatusNotFound)
} else {
http.Error(w, "setting new coordinator: "+err.Error(), http.StatusInternalServerError)
}
return
}
// Encode response.
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(setCoordinatorResponse{
Old: oldNode,
New: newNode,
}); err != nil {
h.logger.Printf("response encoding error: %s", err)
}
}
type setCoordinatorRequest struct {
ID string `json:"id"`
}

View file

@ -1120,7 +1120,7 @@ func (m *URI) GetPort() uint32 {
type Node struct {
ID string `protobuf:"bytes,1,opt,name=ID,proto3" json:"ID,omitempty"`
URI *URI `protobuf:"bytes,2,opt,name=URI,proto3" json:"URI,omitempty"`
IsCoordinator bool `protobuf:"varint,3,opt,name=IsCoordinator,proto3" json:"IsCoordinator,omitempty"`
IsPrimary bool `protobuf:"varint,3,opt,name=IsPrimary,proto3" json:"IsPrimary,omitempty"`
State string `protobuf:"bytes,4,opt,name=State,proto3" json:"State,omitempty"`
GRPCURI *URI `protobuf:"bytes,5,opt,name=GRPCURI,proto3" json:"GRPCURI,omitempty"`
XXX_NoUnkeyedLiteral struct{} `json:"-"`
@ -1175,9 +1175,9 @@ func (m *Node) GetURI() *URI {
return nil
}
func (m *Node) GetIsCoordinator() bool {
func (m *Node) GetIsPrimary() bool {
if m != nil {
return m.IsCoordinator
return m.IsPrimary
}
return false
}
@ -1766,7 +1766,7 @@ func (m *DeleteViewMessage) GetView() string {
type ResizeInstruction struct {
JobID int64 `protobuf:"varint,1,opt,name=JobID,proto3" json:"JobID,omitempty"`
Node *Node `protobuf:"bytes,2,opt,name=Node,proto3" json:"Node,omitempty"`
Coordinator *Node `protobuf:"bytes,3,opt,name=Coordinator,proto3" json:"Coordinator,omitempty"`
Primary *Node `protobuf:"bytes,3,opt,name=Primary,proto3" json:"Primary,omitempty"`
Sources []*ResizeSource `protobuf:"bytes,4,rep,name=Sources,proto3" json:"Sources,omitempty"`
TranslationSources []*TranslationResizeSource `protobuf:"bytes,8,rep,name=TranslationSources,proto3" json:"TranslationSources,omitempty"`
NodeStatus *NodeStatus `protobuf:"bytes,7,opt,name=NodeStatus,proto3" json:"NodeStatus,omitempty"`
@ -1823,9 +1823,9 @@ func (m *ResizeInstruction) GetNode() *Node {
return nil
}
func (m *ResizeInstruction) GetCoordinator() *Node {
func (m *ResizeInstruction) GetPrimary() *Node {
if m != nil {
return m.Coordinator
return m.Primary
}
return nil
}
@ -2063,100 +2063,6 @@ func (m *ResizeInstructionComplete) GetError() string {
return ""
}
type SetCoordinatorMessage struct {
New *Node `protobuf:"bytes,1,opt,name=New,proto3" json:"New,omitempty"`
XXX_NoUnkeyedLiteral struct{} `json:"-"`
XXX_unrecognized []byte `json:"-"`
XXX_sizecache int32 `json:"-"`
}
func (m *SetCoordinatorMessage) Reset() { *m = SetCoordinatorMessage{} }
func (m *SetCoordinatorMessage) String() string { return proto.CompactTextString(m) }
func (*SetCoordinatorMessage) ProtoMessage() {}
func (*SetCoordinatorMessage) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{31}
}
func (m *SetCoordinatorMessage) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
}
func (m *SetCoordinatorMessage) XXX_Marshal(b []byte, deterministic bool) ([]byte, error) {
if deterministic {
return xxx_messageInfo_SetCoordinatorMessage.Marshal(b, m, deterministic)
} else {
b = b[:cap(b)]
n, err := m.MarshalToSizedBuffer(b)
if err != nil {
return nil, err
}
return b[:n], nil
}
}
func (m *SetCoordinatorMessage) XXX_Merge(src proto.Message) {
xxx_messageInfo_SetCoordinatorMessage.Merge(m, src)
}
func (m *SetCoordinatorMessage) XXX_Size() int {
return m.Size()
}
func (m *SetCoordinatorMessage) XXX_DiscardUnknown() {
xxx_messageInfo_SetCoordinatorMessage.DiscardUnknown(m)
}
var xxx_messageInfo_SetCoordinatorMessage proto.InternalMessageInfo
func (m *SetCoordinatorMessage) GetNew() *Node {
if m != nil {
return m.New
}
return nil
}
type UpdateCoordinatorMessage struct {
New *Node `protobuf:"bytes,1,opt,name=New,proto3" json:"New,omitempty"`
XXX_NoUnkeyedLiteral struct{} `json:"-"`
XXX_unrecognized []byte `json:"-"`
XXX_sizecache int32 `json:"-"`
}
func (m *UpdateCoordinatorMessage) Reset() { *m = UpdateCoordinatorMessage{} }
func (m *UpdateCoordinatorMessage) String() string { return proto.CompactTextString(m) }
func (*UpdateCoordinatorMessage) ProtoMessage() {}
func (*UpdateCoordinatorMessage) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{32}
}
func (m *UpdateCoordinatorMessage) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
}
func (m *UpdateCoordinatorMessage) XXX_Marshal(b []byte, deterministic bool) ([]byte, error) {
if deterministic {
return xxx_messageInfo_UpdateCoordinatorMessage.Marshal(b, m, deterministic)
} else {
b = b[:cap(b)]
n, err := m.MarshalToSizedBuffer(b)
if err != nil {
return nil, err
}
return b[:n], nil
}
}
func (m *UpdateCoordinatorMessage) XXX_Merge(src proto.Message) {
xxx_messageInfo_UpdateCoordinatorMessage.Merge(m, src)
}
func (m *UpdateCoordinatorMessage) XXX_Size() int {
return m.Size()
}
func (m *UpdateCoordinatorMessage) XXX_DiscardUnknown() {
xxx_messageInfo_UpdateCoordinatorMessage.DiscardUnknown(m)
}
var xxx_messageInfo_UpdateCoordinatorMessage proto.InternalMessageInfo
func (m *UpdateCoordinatorMessage) GetNew() *Node {
if m != nil {
return m.New
}
return nil
}
type Topology struct {
ClusterID string `protobuf:"bytes,1,opt,name=ClusterID,proto3" json:"ClusterID,omitempty"`
NodeIDs []string `protobuf:"bytes,2,rep,name=NodeIDs,proto3" json:"NodeIDs,omitempty"`
@ -2169,7 +2075,7 @@ func (m *Topology) Reset() { *m = Topology{} }
func (m *Topology) String() string { return proto.CompactTextString(m) }
func (*Topology) ProtoMessage() {}
func (*Topology) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{33}
return fileDescriptor_d2a91b51c7bdc125, []int{31}
}
func (m *Topology) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
@ -2222,7 +2128,7 @@ func (m *RecalculateCaches) Reset() { *m = RecalculateCaches{} }
func (m *RecalculateCaches) String() string { return proto.CompactTextString(m) }
func (*RecalculateCaches) ProtoMessage() {}
func (*RecalculateCaches) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{34}
return fileDescriptor_d2a91b51c7bdc125, []int{32}
}
func (m *RecalculateCaches) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
@ -2263,7 +2169,7 @@ func (m *TransactionMessage) Reset() { *m = TransactionMessage{} }
func (m *TransactionMessage) String() string { return proto.CompactTextString(m) }
func (*TransactionMessage) ProtoMessage() {}
func (*TransactionMessage) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{35}
return fileDescriptor_d2a91b51c7bdc125, []int{33}
}
func (m *TransactionMessage) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
@ -2322,7 +2228,7 @@ func (m *Transaction) Reset() { *m = Transaction{} }
func (m *Transaction) String() string { return proto.CompactTextString(m) }
func (*Transaction) ProtoMessage() {}
func (*Transaction) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{36}
return fileDescriptor_d2a91b51c7bdc125, []int{34}
}
func (m *Transaction) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
@ -2403,7 +2309,7 @@ func (m *TransactionStats) Reset() { *m = TransactionStats{} }
func (m *TransactionStats) String() string { return proto.CompactTextString(m) }
func (*TransactionStats) ProtoMessage() {}
func (*TransactionStats) Descriptor() ([]byte, []int) {
return fileDescriptor_d2a91b51c7bdc125, []int{37}
return fileDescriptor_d2a91b51c7bdc125, []int{35}
}
func (m *TransactionStats) XXX_Unmarshal(b []byte) error {
return m.Unmarshal(b)
@ -2465,8 +2371,6 @@ func init() {
proto.RegisterType((*ResizeSource)(nil), "internal.ResizeSource")
proto.RegisterType((*TranslationResizeSource)(nil), "internal.TranslationResizeSource")
proto.RegisterType((*ResizeInstructionComplete)(nil), "internal.ResizeInstructionComplete")
proto.RegisterType((*SetCoordinatorMessage)(nil), "internal.SetCoordinatorMessage")
proto.RegisterType((*UpdateCoordinatorMessage)(nil), "internal.UpdateCoordinatorMessage")
proto.RegisterType((*Topology)(nil), "internal.Topology")
proto.RegisterType((*RecalculateCaches)(nil), "internal.RecalculateCaches")
proto.RegisterType((*TransactionMessage)(nil), "internal.TransactionMessage")
@ -2477,99 +2381,96 @@ func init() {
func init() { proto.RegisterFile("private.proto", fileDescriptor_d2a91b51c7bdc125) }
var fileDescriptor_d2a91b51c7bdc125 = []byte{
// 1458 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xac, 0x58, 0xcb, 0x72, 0x1b, 0x45,
0x17, 0xfe, 0x47, 0x23, 0xd9, 0xd2, 0x91, 0xe5, 0xc8, 0x9d, 0xc4, 0x99, 0xf8, 0xff, 0xcb, 0xbf,
0x68, 0x52, 0x44, 0xa4, 0x2a, 0x26, 0x95, 0x50, 0xc5, 0x35, 0x55, 0x89, 0x2d, 0x27, 0x08, 0xb0,
0x93, 0xb4, 0x9c, 0xec, 0xdb, 0xa3, 0xae, 0x78, 0xca, 0xa3, 0x19, 0x65, 0x2e, 0x8e, 0x1c, 0xaa,
0xd8, 0x42, 0xc1, 0x8a, 0x62, 0xc3, 0x82, 0x05, 0xef, 0xc1, 0x0b, 0xb0, 0xe4, 0x11, 0xa8, 0xf0,
0x14, 0xec, 0xa8, 0x3e, 0xdd, 0x3d, 0x17, 0x59, 0x8e, 0x4c, 0xc2, 0x6e, 0xce, 0xfd, 0x3b, 0x97,
0x3e, 0xdd, 0x12, 0xb4, 0xc6, 0x91, 0x77, 0xc4, 0x13, 0xb1, 0x31, 0x8e, 0xc2, 0x24, 0x24, 0x75,
0x2f, 0x48, 0x44, 0x14, 0x70, 0x7f, 0x6d, 0x69, 0x9c, 0xee, 0xfb, 0x9e, 0xab, 0xf8, 0xf4, 0x3e,
0x34, 0xfa, 0xc1, 0x50, 0x4c, 0x76, 0x44, 0xc2, 0x09, 0x81, 0xea, 0x17, 0xe2, 0x38, 0x76, 0xec,
0x8e, 0xd5, 0xad, 0x33, 0xfc, 0x26, 0xef, 0xc0, 0xf2, 0x5e, 0xc4, 0xdd, 0xc3, 0xed, 0x89, 0x17,
0x27, 0x22, 0x70, 0x85, 0x53, 0x45, 0xe9, 0x14, 0x97, 0xfe, 0x62, 0xc3, 0xd2, 0x3d, 0x4f, 0xf8,
0xc3, 0x07, 0xe3, 0xc4, 0x0b, 0x83, 0x58, 0x3a, 0xdb, 0x3b, 0x1e, 0x0b, 0xa7, 0xde, 0xb1, 0xba,
0x0d, 0x86, 0xdf, 0xe4, 0x7f, 0xd0, 0xd8, 0xe2, 0xee, 0x81, 0x40, 0x81, 0x8d, 0x82, 0x9c, 0x91,
0x49, 0x07, 0xde, 0x0b, 0x15, 0xa5, 0xc5, 0x72, 0x06, 0xe9, 0x40, 0x73, 0xcf, 0x1b, 0x89, 0x47,
0x29, 0x0f, 0x92, 0x74, 0xe4, 0xd4, 0xd0, 0xba, 0xc8, 0x22, 0xab, 0xb0, 0xf0, 0xc0, 0x1f, 0xee,
0x78, 0x81, 0xd3, 0xe8, 0x58, 0x5d, 0x9b, 0x69, 0xca, 0xf0, 0xf9, 0xc4, 0x81, 0x9c, 0xcf, 0x27,
0x59, 0xba, 0xcd, 0x72, 0xba, 0xbb, 0xe1, 0x20, 0xe1, 0xc1, 0x90, 0x47, 0xc3, 0x27, 0x9e, 0x78,
0xee, 0x2c, 0xa9, 0x74, 0xcb, 0x5c, 0x69, 0xbb, 0xc9, 0x63, 0xe1, 0xb4, 0xd0, 0x23, 0x7e, 0x93,
0x35, 0xa8, 0x6f, 0x7a, 0x49, 0x4f, 0x8c, 0x93, 0x03, 0x67, 0xb9, 0x63, 0x75, 0xab, 0x2c, 0xa3,
0xc9, 0x05, 0xa8, 0x0d, 0x5c, 0xee, 0x0b, 0xe7, 0x1c, 0x1a, 0x28, 0x82, 0x50, 0x58, 0xba, 0x17,
0x46, 0xc2, 0x7b, 0x1a, 0x60, 0x13, 0x9c, 0x36, 0x26, 0x55, 0xe2, 0x91, 0xb7, 0xc1, 0x96, 0x29,
0xad, 0x74, 0xac, 0x6e, 0xf3, 0xe6, 0xca, 0x86, 0xe9, 0xe3, 0x46, 0x4f, 0xb8, 0xde, 0x88, 0xfb,
0x4c, 0x4a, 0x51, 0x89, 0x4f, 0x1c, 0x72, 0xba, 0x12, 0x9f, 0x50, 0x0a, 0xcb, 0xfd, 0xd1, 0x38,
0x8c, 0x12, 0x26, 0xe2, 0x71, 0x18, 0xc4, 0x82, 0xb4, 0xc1, 0xde, 0x8e, 0x22, 0xc7, 0xc2, 0xb0,
0xf2, 0x93, 0x7e, 0x0d, 0xed, 0x4d, 0x3f, 0x74, 0x0f, 0x7b, 0x3c, 0xe1, 0x4c, 0x3c, 0x4b, 0x45,
0x9c, 0x48, 0xec, 0x0a, 0x9e, 0xd2, 0x53, 0x84, 0xe4, 0x62, 0xbf, 0x9d, 0x8a, 0xe2, 0x22, 0x21,
0xeb, 0x82, 0x55, 0x53, 0xed, 0xc1, 0x6f, 0xcc, 0xfd, 0x80, 0x47, 0x43, 0xec, 0x69, 0x95, 0x29,
0x42, 0x72, 0x31, 0x12, 0xce, 0x41, 0x95, 0x29, 0x82, 0xf6, 0x61, 0xa5, 0x10, 0x5f, 0xc3, 0x5c,
0x85, 0x05, 0x16, 0x3e, 0xef, 0xf7, 0x62, 0xc7, 0xea, 0xd8, 0xdd, 0x2a, 0xd3, 0x14, 0x0e, 0x4c,
0xe8, 0xa7, 0xa3, 0x40, 0x8a, 0x2a, 0x28, 0xca, 0x19, 0xf4, 0x32, 0xd4, 0x70, 0x7a, 0x64, 0x96,
0xb9, 0xad, 0xfc, 0xa4, 0xdf, 0x58, 0xd0, 0xd8, 0xe1, 0x13, 0x04, 0x12, 0x93, 0xdb, 0x50, 0x37,
0xbd, 0x45, 0xa5, 0xe6, 0xcd, 0xb7, 0xf2, 0x0a, 0x66, 0x6a, 0x1b, 0x46, 0x67, 0x3b, 0x48, 0xa2,
0x63, 0x96, 0x99, 0xac, 0x7d, 0x02, 0xad, 0x92, 0x48, 0xc6, 0x3b, 0x14, 0xc7, 0xa6, 0xaa, 0x87,
0xe2, 0x58, 0xe6, 0x7a, 0xc4, 0xfd, 0x54, 0x60, 0xad, 0xaa, 0x4c, 0x11, 0x1f, 0x57, 0x3e, 0xb4,
0xe8, 0x13, 0x20, 0x5b, 0x91, 0xe0, 0x89, 0xc0, 0x20, 0x3b, 0x22, 0x8e, 0xf9, 0x53, 0x31, 0xaf,
0xe2, 0x76, 0xb1, 0xe2, 0x59, 0x75, 0x2b, 0x85, 0xea, 0xd2, 0x6b, 0x40, 0x7a, 0xc2, 0x17, 0x89,
0xd0, 0xa7, 0xfb, 0x15, 0x7e, 0xe9, 0x33, 0x83, 0x61, 0xbe, 0x2e, 0xb9, 0x0a, 0x55, 0xb9, 0x2a,
0x30, 0x58, 0xf3, 0xe6, 0xf9, 0xbc, 0x4e, 0xd9, 0x16, 0x61, 0xa8, 0x80, 0xbd, 0x41, 0xa7, 0xc3,
0xbb, 0x09, 0x02, 0xb6, 0x59, 0xce, 0xa0, 0xdf, 0x59, 0x26, 0x26, 0x26, 0x71, 0xc6, 0xbc, 0x4b,
0x93, 0x76, 0x4d, 0x23, 0xb1, 0x11, 0xc9, 0x6a, 0x8e, 0xa4, 0xb8, 0x85, 0x66, 0x81, 0xa9, 0x4e,
0x83, 0xb9, 0x63, 0x6a, 0xf5, 0xba, 0x58, 0xa8, 0x0b, 0xff, 0x55, 0x1e, 0xee, 0x1e, 0x71, 0xcf,
0xe7, 0xfb, 0xfe, 0x3f, 0x6a, 0x67, 0x29, 0x2d, 0x07, 0x16, 0xd1, 0xb6, 0xdf, 0xd3, 0x07, 0xc3,
0x90, 0xf4, 0x2b, 0xc8, 0xcf, 0xd8, 0x2e, 0x1f, 0x09, 0xed, 0x0d, 0xbf, 0xb3, 0x6a, 0x54, 0xce,
0x50, 0x8d, 0x0b, 0x50, 0x93, 0xe7, 0x52, 0xee, 0x79, 0x5b, 0x06, 0x46, 0x62, 0x4e, 0x8d, 0x6e,
0xc1, 0xc2, 0xc0, 0x3d, 0x10, 0x23, 0x4e, 0xde, 0x85, 0x45, 0xc4, 0x2f, 0x62, 0x7d, 0x58, 0xce,
0x4d, 0x0d, 0x01, 0x33, 0x72, 0xfa, 0x83, 0xa5, 0x13, 0x9f, 0x09, 0xb9, 0x14, 0xb0, 0x32, 0x15,
0x90, 0x5c, 0x87, 0x45, 0x8d, 0x1a, 0x77, 0xc9, 0x29, 0xb3, 0x66, 0x74, 0xc8, 0x55, 0x58, 0xc0,
0x4c, 0x63, 0xa7, 0x3a, 0x0d, 0x0a, 0xf9, 0x4c, 0x8b, 0xe9, 0x36, 0xd8, 0x8f, 0x59, 0x5f, 0xae,
0x14, 0xcc, 0xc7, 0x40, 0xd2, 0x94, 0x04, 0xfa, 0x59, 0x18, 0x27, 0xba, 0x27, 0xf8, 0x2d, 0x79,
0x0f, 0xc3, 0x48, 0x4d, 0x71, 0x8b, 0xe1, 0x37, 0xfd, 0xd9, 0x82, 0xea, 0x6e, 0x38, 0x14, 0x64,
0x19, 0x2a, 0xfd, 0x9e, 0x76, 0x52, 0xe9, 0xf7, 0xc8, 0xff, 0xd1, 0xbf, 0xee, 0x43, 0x2b, 0x47,
0xf1, 0x98, 0xf5, 0x19, 0x46, 0xbe, 0x02, 0xad, 0x7e, 0xbc, 0x15, 0x86, 0xd1, 0xd0, 0x0b, 0x78,
0x12, 0x46, 0xfa, 0xb6, 0x2d, 0x33, 0xf1, 0x54, 0x27, 0x3c, 0x51, 0xf7, 0x60, 0x83, 0x29, 0x82,
0x5c, 0x85, 0xc5, 0xfb, 0xec, 0xe1, 0x96, 0x0c, 0x50, 0x9b, 0x15, 0xc0, 0x48, 0xe9, 0x1d, 0x68,
0x4b, 0x74, 0x68, 0x65, 0xa6, 0x70, 0x15, 0x16, 0x24, 0x2f, 0x43, 0xab, 0xa9, 0x3c, 0x54, 0xa5,
0x10, 0x8a, 0x7e, 0xa9, 0x3c, 0x6c, 0x1f, 0x89, 0x20, 0x29, 0xcc, 0x31, 0xd2, 0xe8, 0xa0, 0xc5,
0x14, 0x41, 0xa8, 0xaa, 0x84, 0x4e, 0x79, 0x39, 0x47, 0x24, 0xb9, 0x0c, 0x65, 0xf4, 0x7b, 0x0b,
0xc0, 0x00, 0x4a, 0xe3, 0xcc, 0xc4, 0x3a, 0xdd, 0x84, 0x74, 0xcd, 0xc4, 0xe9, 0x13, 0xde, 0xce,
0xb5, 0x14, 0x9f, 0x99, 0x89, 0x7c, 0x2f, 0x9f, 0x48, 0xd5, 0xfc, 0x8b, 0x53, 0xa3, 0xa2, 0xa2,
0xe6, 0x73, 0x19, 0x40, 0xb3, 0xc0, 0x9f, 0x39, 0x9c, 0xd7, 0xb3, 0x79, 0xaa, 0x4c, 0xbb, 0x44,
0xbe, 0x76, 0xa9, 0x95, 0xe6, 0x6c, 0x3b, 0x0f, 0x9a, 0x05, 0xa3, 0x99, 0xf1, 0xba, 0x70, 0xae,
0xbc, 0x3b, 0xcc, 0x85, 0x36, 0xcd, 0x9e, 0x13, 0xea, 0x47, 0x0b, 0x5a, 0x5b, 0x7e, 0x1a, 0x27,
0x22, 0xd2, 0xd1, 0xa4, 0xbe, 0x62, 0x64, 0x9d, 0xcf, 0x19, 0xb3, 0x9b, 0x4f, 0xae, 0x40, 0x4d,
0xf6, 0x40, 0x6d, 0x88, 0x93, 0x0d, 0x52, 0xc2, 0x42, 0x87, 0xaa, 0xaf, 0xee, 0x10, 0x7d, 0x02,
0xf5, 0xcd, 0x41, 0xff, 0x7e, 0x14, 0xa6, 0xe3, 0x99, 0xd9, 0x9b, 0xb7, 0x62, 0xa5, 0xf0, 0x56,
0x6c, 0xab, 0x77, 0x8f, 0xca, 0x10, 0x1f, 0x39, 0x6d, 0xf5, 0xc8, 0xa9, 0x6a, 0x0e, 0x9f, 0xd0,
0x01, 0xac, 0xa8, 0xd4, 0xe5, 0x0a, 0x7b, 0x9d, 0x6d, 0x6b, 0x9e, 0x2b, 0x76, 0xfe, 0x5c, 0x91,
0x4e, 0xd5, 0x32, 0xff, 0x37, 0x9d, 0xfe, 0x55, 0x81, 0x15, 0x26, 0x62, 0xef, 0x85, 0xe8, 0x07,
0x71, 0x12, 0xa5, 0xae, 0x5c, 0x5b, 0xd2, 0xfe, 0xf3, 0x70, 0x5f, 0xf7, 0xc5, 0x66, 0x8a, 0x38,
0xcb, 0x81, 0x22, 0x37, 0xa0, 0x39, 0xbd, 0x43, 0x4e, 0xaa, 0x16, 0x55, 0xc8, 0x0d, 0x58, 0x1c,
0x84, 0x69, 0xe4, 0x66, 0xa7, 0xa4, 0x70, 0x49, 0x28, 0x64, 0x4a, 0xcc, 0x8c, 0x1a, 0x79, 0x04,
0x64, 0x2f, 0xe2, 0x41, 0xec, 0x73, 0x09, 0xd6, 0x18, 0xd7, 0xa7, 0x5f, 0x48, 0x05, 0x9d, 0x92,
0x9f, 0x19, 0xc6, 0xe4, 0xfd, 0xe2, 0x1a, 0x70, 0x16, 0x11, 0xf5, 0x85, 0x32, 0x6a, 0x7d, 0xb2,
0x8a, 0xeb, 0xe2, 0xf6, 0xd4, 0x4c, 0x3b, 0x0b, 0x68, 0x78, 0x29, 0x37, 0x2c, 0x89, 0x59, 0x59,
0x9b, 0x7e, 0x6b, 0xc1, 0x52, 0x11, 0xd9, 0x99, 0xd6, 0x4f, 0xd6, 0xf0, 0xca, 0xfc, 0x27, 0x98,
0x69, 0x78, 0x75, 0xd6, 0xa3, 0xb7, 0x56, 0x7c, 0x96, 0xa5, 0x70, 0xe9, 0x94, 0x72, 0xbd, 0x01,
0xa8, 0x0e, 0x34, 0x1f, 0xf2, 0x28, 0xf1, 0xa4, 0x4b, 0xfd, 0x6c, 0xa8, 0xb1, 0x22, 0x8b, 0x1e,
0xc2, 0xe5, 0x13, 0xc3, 0xb7, 0x15, 0x8e, 0xc6, 0x72, 0xca, 0xdf, 0x60, 0x08, 0xe5, 0x7d, 0x10,
0x45, 0x7a, 0xfc, 0x1a, 0x4c, 0x11, 0xf4, 0x23, 0xb8, 0x38, 0x10, 0x49, 0x61, 0xf4, 0xcc, 0x19,
0xea, 0x80, 0xbd, 0x2b, 0x9e, 0x9f, 0x92, 0xa0, 0x14, 0xd1, 0x4f, 0xc1, 0x79, 0x3c, 0x1e, 0xf2,
0x44, 0xbc, 0x96, 0xf5, 0x26, 0xd4, 0xf7, 0xc2, 0x71, 0xe8, 0x87, 0x4f, 0x8f, 0xe7, 0x6c, 0x3d,
0x07, 0x16, 0xd5, 0xe5, 0xa7, 0xb6, 0x6c, 0x83, 0x19, 0x92, 0x9e, 0x97, 0xc7, 0xd4, 0xe5, 0xbe,
0x9b, 0xfa, 0x12, 0x86, 0xfc, 0xfd, 0x10, 0x53, 0xa1, 0x0f, 0x02, 0xc7, 0xc2, 0x15, 0xee, 0xd3,
0xbb, 0xc8, 0x30, 0xf7, 0xa9, 0xa2, 0xc8, 0x07, 0xd0, 0x2c, 0x68, 0xeb, 0x02, 0x5e, 0x9c, 0x3a,
0x2f, 0x4a, 0xc8, 0x8a, 0x9a, 0xf4, 0x57, 0xab, 0x64, 0x79, 0xe2, 0x69, 0xa1, 0x03, 0x1e, 0xa9,
0xa6, 0xd4, 0x99, 0xa6, 0x64, 0xae, 0xdb, 0x13, 0xd7, 0x4f, 0x63, 0x29, 0x52, 0xaf, 0x89, 0x9c,
0x21, 0x73, 0x95, 0x3f, 0x92, 0xc3, 0xd4, 0xbc, 0xea, 0x0c, 0x29, 0x7f, 0xaf, 0xf6, 0x04, 0x1f,
0xfa, 0x5e, 0x20, 0x70, 0x4a, 0x6d, 0x96, 0xd1, 0xe4, 0x86, 0xba, 0x17, 0xcc, 0x51, 0x5b, 0x9b,
0x09, 0x1f, 0x35, 0xd4, 0x9d, 0x11, 0x53, 0x02, 0xed, 0x69, 0xd1, 0x66, 0xfb, 0xb7, 0x97, 0xeb,
0xd6, 0xef, 0x2f, 0xd7, 0xad, 0x3f, 0x5e, 0xae, 0x5b, 0x3f, 0xfd, 0xb9, 0xfe, 0x9f, 0xfd, 0x05,
0xfc, 0xdb, 0xe1, 0xd6, 0xdf, 0x01, 0x00, 0x00, 0xff, 0xff, 0x31, 0xb0, 0x31, 0x3c, 0x9f, 0x10,
0x00, 0x00,
// 1420 bytes of a gzipped FileDescriptorProto
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0xac, 0x57, 0xdd, 0x6e, 0x1b, 0x45,
0x14, 0x66, 0xbd, 0xeb, 0xd8, 0x3e, 0x8e, 0x53, 0x67, 0xda, 0xa6, 0xdb, 0x50, 0x05, 0x33, 0x20,
0x6a, 0x2a, 0x35, 0x54, 0x2d, 0x12, 0x08, 0x54, 0xa9, 0x4d, 0x9c, 0x16, 0x03, 0x69, 0xd3, 0x49,
0xda, 0xfb, 0xc9, 0x7a, 0xd4, 0xac, 0xb2, 0xde, 0x75, 0xf7, 0x27, 0x75, 0x8a, 0xc4, 0x2d, 0x08,
0xae, 0x10, 0x5c, 0x70, 0xc9, 0x7b, 0xf0, 0x02, 0x5c, 0xf2, 0x08, 0xa8, 0x3c, 0x01, 0x6f, 0x80,
0xe6, 0xcc, 0xcc, 0xee, 0xda, 0x71, 0xea, 0xd0, 0x72, 0xb7, 0xe7, 0xff, 0x3b, 0x3f, 0x73, 0x66,
0x16, 0x5a, 0xa3, 0xd8, 0x3f, 0xe2, 0xa9, 0x58, 0x1f, 0xc5, 0x51, 0x1a, 0x91, 0xba, 0x1f, 0xa6,
0x22, 0x0e, 0x79, 0xb0, 0xba, 0x38, 0xca, 0xf6, 0x03, 0xdf, 0x53, 0x7c, 0x7a, 0x1f, 0x1a, 0xfd,
0x70, 0x20, 0xc6, 0xdb, 0x22, 0xe5, 0x84, 0x80, 0xf3, 0x95, 0x38, 0x4e, 0x5c, 0xbb, 0x63, 0x75,
0xeb, 0x0c, 0xbf, 0xc9, 0x07, 0xb0, 0xb4, 0x17, 0x73, 0xef, 0x70, 0x6b, 0xec, 0x27, 0xa9, 0x08,
0x3d, 0xe1, 0x3a, 0x28, 0x9d, 0xe2, 0xd2, 0xdf, 0x6c, 0x58, 0xbc, 0xe7, 0x8b, 0x60, 0xf0, 0x70,
0x94, 0xfa, 0x51, 0x98, 0x48, 0x67, 0x7b, 0xc7, 0x23, 0xe1, 0xd6, 0x3b, 0x56, 0xb7, 0xc1, 0xf0,
0x9b, 0x5c, 0x81, 0xc6, 0x26, 0xf7, 0x0e, 0x04, 0x0a, 0x6c, 0x14, 0x14, 0x8c, 0x5c, 0xba, 0xeb,
0xbf, 0x50, 0x51, 0x5a, 0xac, 0x60, 0x90, 0x0e, 0x34, 0xf7, 0xfc, 0xa1, 0x78, 0x94, 0xf1, 0x30,
0xcd, 0x86, 0x6e, 0x15, 0xad, 0xcb, 0x2c, 0xb2, 0x02, 0x0b, 0x0f, 0x83, 0xc1, 0xb6, 0x1f, 0xba,
0x8d, 0x8e, 0xd5, 0xb5, 0x99, 0xa6, 0x0c, 0x9f, 0x8f, 0x5d, 0x28, 0xf8, 0x7c, 0x9c, 0xa7, 0xdb,
0x9c, 0x4c, 0xf7, 0x41, 0xb4, 0x9b, 0xf2, 0x70, 0xc0, 0xe3, 0xc1, 0x13, 0x5f, 0x3c, 0x77, 0x17,
0x55, 0xba, 0x93, 0x5c, 0x69, 0xbb, 0xc1, 0x13, 0xe1, 0xb6, 0xd0, 0x23, 0x7e, 0x93, 0x55, 0xa8,
0x6f, 0xf8, 0x69, 0x4f, 0x8c, 0xd2, 0x03, 0x77, 0xa9, 0x63, 0x75, 0x1d, 0x96, 0xd3, 0xe4, 0x02,
0x54, 0x77, 0x3d, 0x1e, 0x08, 0xf7, 0x1c, 0x1a, 0x28, 0x82, 0x50, 0x58, 0xbc, 0x17, 0xc5, 0xc2,
0x7f, 0x1a, 0x62, 0x13, 0xdc, 0x36, 0x26, 0x35, 0xc1, 0x23, 0xef, 0x81, 0x2d, 0x53, 0x5a, 0xee,
0x58, 0xdd, 0xe6, 0xcd, 0xe5, 0x75, 0xd3, 0xc7, 0xf5, 0x9e, 0xf0, 0xfc, 0x21, 0x0f, 0x98, 0x94,
0xa2, 0x12, 0x1f, 0xbb, 0xe4, 0x74, 0x25, 0x3e, 0xa6, 0x14, 0x96, 0xfa, 0xc3, 0x51, 0x14, 0xa7,
0x4c, 0x24, 0xa3, 0x28, 0x4c, 0x04, 0x69, 0x83, 0xbd, 0x15, 0xc7, 0xae, 0x85, 0x61, 0xe5, 0x27,
0xfd, 0x16, 0xda, 0x1b, 0x41, 0xe4, 0x1d, 0xf6, 0x78, 0xca, 0x99, 0x78, 0x96, 0x89, 0x24, 0x95,
0xd8, 0x15, 0x3c, 0xa5, 0xa7, 0x08, 0xc9, 0xc5, 0x7e, 0xbb, 0x15, 0xc5, 0x45, 0x42, 0xd6, 0x05,
0xab, 0xa6, 0xda, 0x83, 0xdf, 0x98, 0xfb, 0x01, 0x8f, 0x07, 0xd8, 0x53, 0x87, 0x29, 0x42, 0x72,
0x31, 0x12, 0xce, 0x81, 0xc3, 0x14, 0x41, 0xfb, 0xb0, 0x5c, 0x8a, 0xaf, 0x61, 0xae, 0xc0, 0x02,
0x8b, 0x9e, 0xf7, 0x7b, 0x89, 0x6b, 0x75, 0xec, 0xae, 0xc3, 0x34, 0x85, 0x03, 0x13, 0x05, 0xd9,
0x30, 0x94, 0xa2, 0x0a, 0x8a, 0x0a, 0x06, 0xbd, 0x0c, 0x55, 0x9c, 0x1e, 0x99, 0x65, 0x61, 0x2b,
0x3f, 0xe9, 0x77, 0x16, 0x34, 0xb6, 0xf9, 0x18, 0x81, 0x24, 0xe4, 0x36, 0xd4, 0x4d, 0x6f, 0x51,
0xa9, 0x79, 0xf3, 0xdd, 0xa2, 0x82, 0xb9, 0xda, 0xba, 0xd1, 0xd9, 0x0a, 0xd3, 0xf8, 0x98, 0xe5,
0x26, 0xab, 0x9f, 0x43, 0x6b, 0x42, 0x24, 0xe3, 0x1d, 0x8a, 0x63, 0x53, 0xd5, 0x43, 0x71, 0x2c,
0x73, 0x3d, 0xe2, 0x41, 0x26, 0xb0, 0x56, 0x0e, 0x53, 0xc4, 0x67, 0x95, 0x4f, 0x2d, 0xfa, 0x04,
0xc8, 0x66, 0x2c, 0x78, 0x2a, 0x30, 0xc8, 0xb6, 0x48, 0x12, 0xfe, 0x54, 0xcc, 0xab, 0xb8, 0x5d,
0xae, 0x78, 0x5e, 0xdd, 0x4a, 0xa9, 0xba, 0xf4, 0x1a, 0x90, 0x9e, 0x08, 0x44, 0x2a, 0xf4, 0xe9,
0x7e, 0x85, 0x5f, 0xfa, 0xcc, 0x60, 0x98, 0xaf, 0x4b, 0xae, 0x82, 0x23, 0x57, 0x05, 0x06, 0x6b,
0xde, 0x3c, 0x5f, 0xd4, 0x29, 0xdf, 0x22, 0x0c, 0x15, 0xb0, 0x37, 0xe8, 0x74, 0x70, 0x37, 0x45,
0xc0, 0x36, 0x2b, 0x18, 0xf4, 0x07, 0xcb, 0xc4, 0xc4, 0x24, 0xce, 0x98, 0xf7, 0xc4, 0xa4, 0x5d,
0xd3, 0x48, 0x6c, 0x44, 0xb2, 0x52, 0x20, 0x29, 0x6f, 0xa1, 0x59, 0x60, 0x9c, 0x69, 0x30, 0x77,
0x4c, 0xad, 0x5e, 0x17, 0x0b, 0xf5, 0xe0, 0x6d, 0xe5, 0xe1, 0xee, 0x11, 0xf7, 0x03, 0xbe, 0x1f,
0xfc, 0xa7, 0x76, 0x4e, 0xa4, 0xe5, 0x42, 0x0d, 0x6d, 0xfb, 0x3d, 0x7d, 0x30, 0x0c, 0x49, 0xbf,
0x81, 0xe2, 0x8c, 0x3d, 0xe0, 0x43, 0xa1, 0xbd, 0xe1, 0x77, 0x5e, 0x8d, 0xca, 0x19, 0xaa, 0x71,
0x01, 0xaa, 0xf2, 0x5c, 0xca, 0x3d, 0x6f, 0xcb, 0xc0, 0x48, 0xcc, 0xa9, 0xd1, 0x2d, 0x58, 0xd8,
0xf5, 0x0e, 0xc4, 0x90, 0x93, 0x0f, 0xa1, 0x86, 0xf8, 0x45, 0xa2, 0x0f, 0xcb, 0xb9, 0xa9, 0x21,
0x60, 0x46, 0x4e, 0x7f, 0xb2, 0x74, 0xe2, 0x33, 0x21, 0x4f, 0x04, 0xac, 0x4c, 0x05, 0x24, 0xd7,
0xa1, 0xa6, 0x51, 0xe3, 0x2e, 0x39, 0x65, 0xd6, 0x8c, 0x0e, 0xb9, 0x0a, 0x0b, 0x98, 0x69, 0xe2,
0x3a, 0xd3, 0xa0, 0x90, 0xcf, 0xb4, 0x98, 0x6e, 0x81, 0xfd, 0x98, 0xf5, 0xe5, 0x4a, 0xc1, 0x7c,
0x0c, 0x24, 0x4d, 0x49, 0xa0, 0x5f, 0x44, 0x49, 0xaa, 0x7b, 0x82, 0xdf, 0x92, 0xb7, 0x13, 0xc5,
0x6a, 0x8a, 0x5b, 0x0c, 0xbf, 0xe9, 0x2f, 0x16, 0x38, 0x0f, 0xa2, 0x81, 0x20, 0x4b, 0x50, 0xe9,
0xf7, 0xb4, 0x93, 0x4a, 0xbf, 0x47, 0xde, 0x41, 0xff, 0xba, 0x0f, 0xad, 0x02, 0xc5, 0x63, 0xd6,
0x67, 0x18, 0xf9, 0x0a, 0x34, 0xfa, 0xc9, 0x4e, 0xec, 0x0f, 0x79, 0x7c, 0xac, 0x6f, 0xda, 0x82,
0x81, 0xa7, 0x39, 0xe5, 0xa9, 0xba, 0xff, 0x1a, 0x4c, 0x11, 0xe4, 0x2a, 0xd4, 0xee, 0xb3, 0x9d,
0x4d, 0xe9, 0xb8, 0x3a, 0xcb, 0xb1, 0x91, 0xd2, 0x3b, 0xd0, 0x96, 0xa8, 0xd0, 0xca, 0x4c, 0xdf,
0x0a, 0x2c, 0x48, 0x5e, 0x8e, 0x52, 0x53, 0x45, 0xa8, 0x4a, 0x29, 0x14, 0xfd, 0x5a, 0x79, 0xd8,
0x3a, 0x12, 0x61, 0x5a, 0x9a, 0x5f, 0xa4, 0xd1, 0x41, 0x8b, 0x29, 0x82, 0x50, 0x55, 0x01, 0x9d,
0xea, 0x52, 0x81, 0x48, 0x72, 0x19, 0xca, 0xe8, 0x8f, 0x16, 0x80, 0x01, 0x94, 0x25, 0xb9, 0x89,
0x75, 0xba, 0x09, 0xe9, 0x9a, 0x49, 0xd3, 0x27, 0xbb, 0x5d, 0x68, 0x29, 0x3e, 0x33, 0x93, 0xf8,
0x51, 0x31, 0x89, 0xaa, 0xe9, 0x17, 0xa7, 0x46, 0x44, 0x45, 0x2d, 0xe6, 0x31, 0x84, 0x66, 0x89,
0x3f, 0x73, 0x28, 0xaf, 0xe7, 0x73, 0x54, 0x99, 0x76, 0x89, 0x7c, 0xed, 0x52, 0x2b, 0xcd, 0xd9,
0x72, 0x3e, 0x34, 0x4b, 0x46, 0x33, 0xe3, 0x75, 0xe1, 0xdc, 0xe4, 0xce, 0x30, 0x17, 0xd9, 0x34,
0x7b, 0x4e, 0xa8, 0x9f, 0x2d, 0x68, 0x6d, 0x06, 0x59, 0x92, 0x8a, 0x58, 0x47, 0x93, 0xfa, 0x8a,
0x91, 0x77, 0xbe, 0x60, 0xcc, 0x6e, 0x3e, 0x79, 0x1f, 0xaa, 0xb2, 0x07, 0x6a, 0x33, 0x9c, 0x6c,
0x90, 0x12, 0x96, 0x3a, 0xe4, 0xbc, 0xba, 0x43, 0xf4, 0x09, 0xd4, 0x37, 0x76, 0xfb, 0xf7, 0xe3,
0x28, 0x1b, 0xcd, 0xcc, 0xde, 0xbc, 0x11, 0x2b, 0xa5, 0x37, 0x62, 0x5b, 0xbd, 0x77, 0x54, 0x86,
0xf8, 0xb8, 0x69, 0xab, 0xc7, 0x8d, 0xa3, 0x39, 0x7c, 0x4c, 0x77, 0x61, 0x59, 0xa5, 0x2e, 0x57,
0xd7, 0xeb, 0x6c, 0x59, 0xf3, 0x4c, 0xb1, 0x8b, 0x67, 0x8a, 0x74, 0xaa, 0x96, 0xf8, 0xff, 0xe9,
0xf4, 0x9f, 0x0a, 0x2c, 0x33, 0x91, 0xf8, 0x2f, 0x44, 0x3f, 0x4c, 0xd2, 0x38, 0xf3, 0xe4, 0xba,
0x92, 0xf6, 0x5f, 0x46, 0xfb, 0xba, 0x2f, 0x36, 0x53, 0xc4, 0x59, 0x0e, 0x14, 0xe9, 0x42, 0xad,
0xbc, 0x3b, 0x4e, 0xaa, 0x19, 0x31, 0xb9, 0x01, 0xb5, 0xdd, 0x28, 0x8b, 0xbd, 0xfc, 0x74, 0x94,
0x2e, 0x05, 0x85, 0x48, 0x89, 0x99, 0x51, 0x23, 0x8f, 0x80, 0xec, 0xc5, 0x3c, 0x4c, 0x02, 0x2e,
0x41, 0x1a, 0xe3, 0xfa, 0xf4, 0x8b, 0xa8, 0xa4, 0x33, 0xe1, 0x67, 0x86, 0x31, 0xf9, 0xb8, 0x7c,
0xfc, 0xdd, 0x1a, 0x22, 0xbe, 0x30, 0x89, 0x58, 0x9f, 0xa8, 0xf2, 0x9a, 0xb8, 0x3d, 0x35, 0xcb,
0xee, 0x02, 0x1a, 0x5e, 0x2a, 0x0c, 0x27, 0xc4, 0x6c, 0x52, 0x9b, 0x7e, 0x6f, 0xc1, 0x62, 0x19,
0xd9, 0x99, 0xd6, 0x4e, 0xde, 0xe8, 0xca, 0xfc, 0x27, 0x97, 0x69, 0xb4, 0x33, 0xeb, 0x91, 0x5b,
0x2d, 0x3f, 0xc3, 0x32, 0xb8, 0x74, 0x4a, 0xb9, 0xde, 0x00, 0x54, 0x07, 0x9a, 0x3b, 0x3c, 0x4e,
0x7d, 0xe9, 0x52, 0x3f, 0x13, 0xaa, 0xac, 0xcc, 0xa2, 0x87, 0x70, 0xf9, 0xc4, 0xd0, 0x6d, 0x46,
0xc3, 0x91, 0x9c, 0xee, 0x37, 0x18, 0x3e, 0x79, 0x0f, 0xc4, 0x71, 0x14, 0x9b, 0x6a, 0x20, 0x41,
0x37, 0xa0, 0xbe, 0x17, 0x8d, 0xa2, 0x20, 0x7a, 0x7a, 0x3c, 0x67, 0xe9, 0xb8, 0x50, 0x53, 0x77,
0x8f, 0x5a, 0x72, 0x0d, 0x66, 0x48, 0x7a, 0x5e, 0x9e, 0x12, 0x8f, 0x07, 0x5e, 0x16, 0xf0, 0x54,
0xe0, 0xb3, 0x3d, 0xa1, 0x42, 0xcf, 0x23, 0x47, 0xfc, 0xa5, 0xeb, 0xec, 0x2e, 0x32, 0xcc, 0x75,
0xa6, 0x28, 0xf2, 0x09, 0x34, 0x4b, 0xda, 0x3a, 0x8f, 0x8b, 0x53, 0x63, 0xab, 0x84, 0xac, 0xac,
0x49, 0x7f, 0xb7, 0x26, 0x2c, 0x4f, 0xdc, 0xe8, 0x3a, 0xe0, 0x91, 0xaa, 0x4d, 0x9d, 0x69, 0x4a,
0xe6, 0xba, 0x35, 0xf6, 0x82, 0x2c, 0x91, 0x22, 0x7d, 0x91, 0xe7, 0x0c, 0x99, 0xab, 0xfc, 0x37,
0x8d, 0x32, 0xf3, 0x98, 0x32, 0xa4, 0xfc, 0x4d, 0xec, 0x09, 0x3e, 0x08, 0xfc, 0x50, 0xe0, 0xb0,
0xd8, 0x2c, 0xa7, 0xc9, 0x0d, 0xb5, 0x96, 0xcd, 0xc4, 0xaf, 0xce, 0x84, 0x8f, 0x1a, 0x6a, 0x65,
0x27, 0x94, 0x40, 0x7b, 0x5a, 0xb4, 0xd1, 0xfe, 0xe3, 0xe5, 0x9a, 0xf5, 0xe7, 0xcb, 0x35, 0xeb,
0xaf, 0x97, 0x6b, 0xd6, 0xaf, 0x7f, 0xaf, 0xbd, 0xb5, 0xbf, 0x80, 0x7f, 0xfb, 0xb7, 0xfe, 0x0d,
0x00, 0x00, 0xff, 0xff, 0x63, 0xcb, 0x53, 0xd8, 0x16, 0x10, 0x00, 0x00,
}
func (m *IndexMeta) Marshal() (dAtA []byte, err error) {
@ -3529,9 +3430,9 @@ func (m *Node) MarshalToSizedBuffer(dAtA []byte) (int, error) {
i--
dAtA[i] = 0x22
}
if m.IsCoordinator {
if m.IsPrimary {
i--
if m.IsCoordinator {
if m.IsPrimary {
dAtA[i] = 1
} else {
dAtA[i] = 0
@ -4111,9 +4012,9 @@ func (m *ResizeInstruction) MarshalToSizedBuffer(dAtA []byte) (int, error) {
dAtA[i] = 0x22
}
}
if m.Coordinator != nil {
if m.Primary != nil {
{
size, err := m.Coordinator.MarshalToSizedBuffer(dAtA[:i])
size, err := m.Primary.MarshalToSizedBuffer(dAtA[:i])
if err != nil {
return 0, err
}
@ -4310,84 +4211,6 @@ func (m *ResizeInstructionComplete) MarshalToSizedBuffer(dAtA []byte) (int, erro
return len(dAtA) - i, nil
}
func (m *SetCoordinatorMessage) Marshal() (dAtA []byte, err error) {
size := m.Size()
dAtA = make([]byte, size)
n, err := m.MarshalToSizedBuffer(dAtA[:size])
if err != nil {
return nil, err
}
return dAtA[:n], nil
}
func (m *SetCoordinatorMessage) MarshalTo(dAtA []byte) (int, error) {
size := m.Size()
return m.MarshalToSizedBuffer(dAtA[:size])
}
func (m *SetCoordinatorMessage) MarshalToSizedBuffer(dAtA []byte) (int, error) {
i := len(dAtA)
_ = i
var l int
_ = l
if m.XXX_unrecognized != nil {
i -= len(m.XXX_unrecognized)
copy(dAtA[i:], m.XXX_unrecognized)
}
if m.New != nil {
{
size, err := m.New.MarshalToSizedBuffer(dAtA[:i])
if err != nil {
return 0, err
}
i -= size
i = encodeVarintPrivate(dAtA, i, uint64(size))
}
i--
dAtA[i] = 0xa
}
return len(dAtA) - i, nil
}
func (m *UpdateCoordinatorMessage) Marshal() (dAtA []byte, err error) {
size := m.Size()
dAtA = make([]byte, size)
n, err := m.MarshalToSizedBuffer(dAtA[:size])
if err != nil {
return nil, err
}
return dAtA[:n], nil
}
func (m *UpdateCoordinatorMessage) MarshalTo(dAtA []byte) (int, error) {
size := m.Size()
return m.MarshalToSizedBuffer(dAtA[:size])
}
func (m *UpdateCoordinatorMessage) MarshalToSizedBuffer(dAtA []byte) (int, error) {
i := len(dAtA)
_ = i
var l int
_ = l
if m.XXX_unrecognized != nil {
i -= len(m.XXX_unrecognized)
copy(dAtA[i:], m.XXX_unrecognized)
}
if m.New != nil {
{
size, err := m.New.MarshalToSizedBuffer(dAtA[:i])
if err != nil {
return 0, err
}
i -= size
i = encodeVarintPrivate(dAtA, i, uint64(size))
}
i--
dAtA[i] = 0xa
}
return len(dAtA) - i, nil
}
func (m *Topology) Marshal() (dAtA []byte, err error) {
size := m.Size()
dAtA = make([]byte, size)
@ -5052,7 +4875,7 @@ func (m *Node) Size() (n int) {
l = m.URI.Size()
n += 1 + l + sovPrivate(uint64(l))
}
if m.IsCoordinator {
if m.IsPrimary {
n += 2
}
l = len(m.State)
@ -5302,8 +5125,8 @@ func (m *ResizeInstruction) Size() (n int) {
l = m.Node.Size()
n += 1 + l + sovPrivate(uint64(l))
}
if m.Coordinator != nil {
l = m.Coordinator.Size()
if m.Primary != nil {
l = m.Primary.Size()
n += 1 + l + sovPrivate(uint64(l))
}
if len(m.Sources) > 0 {
@ -5409,38 +5232,6 @@ func (m *ResizeInstructionComplete) Size() (n int) {
return n
}
func (m *SetCoordinatorMessage) Size() (n int) {
if m == nil {
return 0
}
var l int
_ = l
if m.New != nil {
l = m.New.Size()
n += 1 + l + sovPrivate(uint64(l))
}
if m.XXX_unrecognized != nil {
n += len(m.XXX_unrecognized)
}
return n
}
func (m *UpdateCoordinatorMessage) Size() (n int) {
if m == nil {
return 0
}
var l int
_ = l
if m.New != nil {
l = m.New.Size()
n += 1 + l + sovPrivate(uint64(l))
}
if m.XXX_unrecognized != nil {
n += len(m.XXX_unrecognized)
}
return n
}
func (m *Topology) Size() (n int) {
if m == nil {
return 0
@ -8237,7 +8028,7 @@ func (m *Node) Unmarshal(dAtA []byte) error {
iNdEx = postIndex
case 3:
if wireType != 0 {
return fmt.Errorf("proto: wrong wireType = %d for field IsCoordinator", wireType)
return fmt.Errorf("proto: wrong wireType = %d for field IsPrimary", wireType)
}
var v int
for shift := uint(0); ; shift += 7 {
@ -8254,7 +8045,7 @@ func (m *Node) Unmarshal(dAtA []byte) error {
break
}
}
m.IsCoordinator = bool(v != 0)
m.IsPrimary = bool(v != 0)
case 4:
if wireType != 2 {
return fmt.Errorf("proto: wrong wireType = %d for field State", wireType)
@ -9755,7 +9546,7 @@ func (m *ResizeInstruction) Unmarshal(dAtA []byte) error {
iNdEx = postIndex
case 3:
if wireType != 2 {
return fmt.Errorf("proto: wrong wireType = %d for field Coordinator", wireType)
return fmt.Errorf("proto: wrong wireType = %d for field Primary", wireType)
}
var msglen int
for shift := uint(0); ; shift += 7 {
@ -9782,10 +9573,10 @@ func (m *ResizeInstruction) Unmarshal(dAtA []byte) error {
if postIndex > l {
return io.ErrUnexpectedEOF
}
if m.Coordinator == nil {
m.Coordinator = &Node{}
if m.Primary == nil {
m.Primary = &Node{}
}
if err := m.Coordinator.Unmarshal(dAtA[iNdEx:postIndex]); err != nil {
if err := m.Primary.Unmarshal(dAtA[iNdEx:postIndex]); err != nil {
return err
}
iNdEx = postIndex
@ -10429,180 +10220,6 @@ func (m *ResizeInstructionComplete) Unmarshal(dAtA []byte) error {
}
return nil
}
func (m *SetCoordinatorMessage) 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: SetCoordinatorMessage: wiretype end group for non-group")
}
if fieldNum <= 0 {
return fmt.Errorf("proto: SetCoordinatorMessage: illegal tag %d (wire type %d)", fieldNum, wire)
}
switch fieldNum {
case 1:
if wireType != 2 {
return fmt.Errorf("proto: wrong wireType = %d for field New", 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 < 0 {
return ErrInvalidLengthPrivate
}
if postIndex > l {
return io.ErrUnexpectedEOF
}
if m.New == nil {
m.New = &Node{}
}
if err := m.New.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) || (iNdEx+skippy) < 0 {
return ErrInvalidLengthPrivate
}
if (iNdEx + skippy) > l {
return io.ErrUnexpectedEOF
}
m.XXX_unrecognized = append(m.XXX_unrecognized, dAtA[iNdEx:iNdEx+skippy]...)
iNdEx += skippy
}
}
if iNdEx > l {
return io.ErrUnexpectedEOF
}
return nil
}
func (m *UpdateCoordinatorMessage) 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: UpdateCoordinatorMessage: wiretype end group for non-group")
}
if fieldNum <= 0 {
return fmt.Errorf("proto: UpdateCoordinatorMessage: illegal tag %d (wire type %d)", fieldNum, wire)
}
switch fieldNum {
case 1:
if wireType != 2 {
return fmt.Errorf("proto: wrong wireType = %d for field New", 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 < 0 {
return ErrInvalidLengthPrivate
}
if postIndex > l {
return io.ErrUnexpectedEOF
}
if m.New == nil {
m.New = &Node{}
}
if err := m.New.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) || (iNdEx+skippy) < 0 {
return ErrInvalidLengthPrivate
}
if (iNdEx + skippy) > l {
return io.ErrUnexpectedEOF
}
m.XXX_unrecognized = append(m.XXX_unrecognized, dAtA[iNdEx:iNdEx+skippy]...)
iNdEx += skippy
}
}
if iNdEx > l {
return io.ErrUnexpectedEOF
}
return nil
}
func (m *Topology) Unmarshal(dAtA []byte) error {
l := len(dAtA)
iNdEx := 0

View file

@ -112,7 +112,7 @@ message URI {
message Node {
string ID = 1;
URI URI = 2;
bool IsCoordinator = 3;
bool IsPrimary = 3;
string State = 4;
URI GRPCURI = 5;
}
@ -174,7 +174,7 @@ message DeleteViewMessage {
message ResizeInstruction {
int64 JobID = 1;
Node Node = 2;
Node Coordinator = 3;
Node Primary = 3;
repeated ResizeSource Sources = 4;
repeated TranslationResizeSource TranslationSources = 8;
NodeStatus NodeStatus = 7;
@ -201,14 +201,6 @@ message ResizeInstructionComplete {
string Error = 3;
}
message SetCoordinatorMessage {
Node New = 1;
}
message UpdateCoordinatorMessage {
Node New = 1;
}
message Topology {
string ClusterID = 1;
repeated string NodeIDs = 2;

View file

@ -91,7 +91,6 @@ type Server struct { // nolint: maligned
maxWritesPerRequest int
confirmDownSleep time.Duration
confirmDownRetries int
isCoordinator bool
syncer holderSyncer
translationSyncer TranslationSyncer
@ -296,15 +295,6 @@ func OptServerSerializer(ser Serializer) ServerOption {
}
}
// OptServerIsCoordinator is a functional option on Server
// used to specify whether or not this server is the coordinator.
func OptServerIsCoordinator(is bool) ServerOption {
return func(s *Server) error {
s.isCoordinator = is
return nil
}
}
// OptServerNodeID is a functional option on Server
// used to set the server node ID.
func OptServerNodeID(nodeID string) ServerOption {
@ -575,12 +565,14 @@ func (s *Server) Open() error {
// Set node ID.
s.nodeID = s.disCo.ID()
// TODO we cannot set IsPrimary here because we don't have all the needed info
node := &topology.Node{
ID: s.nodeID,
URI: s.uri,
GRPCURI: s.grpcURI,
IsCoordinator: s.isCoordinator,
State: nodeStateDown,
ID: s.nodeID,
URI: s.uri,
GRPCURI: s.grpcURI,
State: nodeStateDown,
// TODO set primary
IsPrimary: false,
}
// Set metadata for this node.
@ -876,7 +868,7 @@ func (s *Server) receiveMessage(m Message) error {
if err != nil {
return err
}
if !s.isCoordinator {
if !s.IsPrimary() {
if obj.Schema != nil {
s.holder.applyCreatedAt(obj.Schema.Indexes)
}
@ -892,10 +884,6 @@ func (s *Server) receiveMessage(m Message) error {
if err != nil {
return err
}
case *SetCoordinatorMessage:
return s.cluster.setCoordinator(obj.New)
case *UpdateCoordinatorMessage:
s.cluster.updateCoordinator(obj.New)
case *NodeStateMessage:
err := s.cluster.receiveNodeState(obj.NodeID, obj.State)
if err != nil {
@ -1058,6 +1046,12 @@ func (s *Server) mergeRemoteStatus(ns *NodeStatus) error {
return nil
}
// IsPrimary returns if this node is primary right now or not.
func (s *Server) IsPrimary() bool {
primary := s.cluster.PrimaryReplicaNode()
return s.nodeID == primary.ID
}
// monitorDiagnostics periodically polls the Pilosa Indexes for cluster info.
func (s *Server) monitorDiagnostics() {
// Do not send more than once a minute
@ -1157,11 +1151,12 @@ func (s *Server) monitorRuntime() {
}
func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, remote bool) (*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !remote && !node.IsCoordinator && len(srv.cluster.Nodes()) > 1 {
if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
}
if remote && (node.IsCoordinator || len(srv.cluster.Nodes()) == 1) {
if remote && (snap.IsPrimaryFieldTranslationNode(node.ID) || len(srv.cluster.Nodes()) == 1) {
return nil, errors.New("unexpected remote start call to coordinator or single node cluster")
}
@ -1203,11 +1198,12 @@ func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time
}
func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool) (*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !remote && !node.IsCoordinator && len(srv.cluster.Nodes()) > 1 {
if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
}
if remote && (node.IsCoordinator || len(srv.cluster.Nodes()) == 1) {
if remote && (snap.IsPrimaryFieldTranslationNode(node.ID) || len(srv.cluster.Nodes()) == 1) {
return nil, errors.New("unexpected remote finish call to coordinator or single node cluster")
}
@ -1232,8 +1228,9 @@ func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool
}
func (srv *Server) Transactions(ctx context.Context) (map[string]*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !node.IsCoordinator && len(srv.cluster.Nodes()) > 1 {
if !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
}
@ -1241,12 +1238,14 @@ func (srv *Server) Transactions(ctx context.Context) (map[string]*Transaction, e
}
func (srv *Server) GetTransaction(ctx context.Context, id string, remote bool) (*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !remote && !node.IsCoordinator && len(srv.cluster.Nodes()) > 1 {
if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
}
if remote && (node.IsCoordinator || len(srv.cluster.Nodes()) == 1) {
if remote && (snap.IsPrimaryFieldTranslationNode(node.ID) || len(srv.cluster.Nodes()) == 1) {
return nil, errors.New("unexpected remote get call to coordinator or single node cluster")
}

View file

@ -186,7 +186,7 @@ func TestClusterResize_AddNode(t *testing.T) {
}
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
@ -247,7 +247,7 @@ func TestClusterResize_AddNode(t *testing.T) {
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
@ -308,7 +308,7 @@ func TestClusterResize_AddNode(t *testing.T) {
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
@ -374,7 +374,7 @@ func TestClusterResize_AddNode(t *testing.T) {
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
@ -435,7 +435,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
}()
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
@ -497,7 +497,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
@ -565,7 +565,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
@ -631,7 +631,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
m0.QueryExpect(t, "i", "", `Row(f=1)`, exp)
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
m1.Config.Gossip.Seeds = []string{seed}
if err := port.GetListeners(func(lsns []*net.TCPListener) error {
portsCfg := test.GenPortsConfig(test.NewPorts(lsns))
@ -677,7 +677,7 @@ func TestCluster_GossipMembership(t *testing.T) {
var eg errgroup.Group
// Configure node1
m1 := test.NewCommandNode(t, false)
m1 := test.NewCommandNode(t)
defer m1.Close()
eg.Go(func() error {
// Pass invalid seed as first in list
@ -693,7 +693,7 @@ func TestCluster_GossipMembership(t *testing.T) {
})
// Configure node1
m2 := test.NewCommandNode(t, false)
m2 := test.NewCommandNode(t)
defer m2.Close()
eg.Go(func() error {
// Pass invalid seed as first in list

View file

@ -123,9 +123,8 @@ type Config struct {
ImportWorkerPoolSize int `toml:"-"`
Cluster struct {
Coordinator bool `toml:"coordinator"`
ReplicaN int `toml:"replicas"`
Name string `toml:"name"`
ReplicaN int `toml:"replicas"`
Name string `toml:"name"`
// This LongQueryTime is deprecated but still exists for backward compatibility
LongQueryTime toml.Duration `toml:"long-query-time"`
} `toml:"cluster"`

View file

@ -1394,7 +1394,7 @@ func TestHandler_Endpoints(t *testing.T) {
func TestCluster_TranslateStore(t *testing.T) {
cluster := test.MustNewCluster(t, 1)
cluster.Nodes[0] = test.NewCommandNode(t, true,
cluster.Nodes[0] = test.NewCommandNode(t,
server.OptCommandServerOptions(
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderWithLockerFunc(nil, &sync.Mutex{})),

View file

@ -387,12 +387,6 @@ func (m *Command) SetupServer() error {
m.logger.Printf("DEPRECATED: Configuration parameter cluster.long-query-time has been renamed to long-query-time")
}
// Set Coordinator.
coordinatorOpt := pilosa.OptServerIsCoordinator(false)
if m.Config.Cluster.Coordinator || len(m.Config.Gossip.Seeds) == 0 {
coordinatorOpt = pilosa.OptServerIsCoordinator(true)
}
// Use other config parameters to set Etcd parameters which we don't want to
// expose in the user-facing config.
//
@ -440,7 +434,6 @@ func (m *Command) SetupServer() error {
pilosa.OptServerRowcacheOn(m.Config.RowcacheOn),
pilosa.OptServerRBFConfig(m.Config.RBFConfig),
pilosa.OptServerQueryHistoryLength(m.Config.QueryHistoryLength),
coordinatorOpt,
discoOpt,
}

View file

@ -137,7 +137,7 @@ func (c *Cluster) GetNode(n int) *Command {
// need to act on the coordinator.
func (c *Cluster) GetCoordinator() *Command {
for _, n := range c.Nodes {
if n.IsCoordinator() {
if n.IsPrimary() {
return n
}
}
@ -147,7 +147,7 @@ func (c *Cluster) GetCoordinator() *Command {
// GetNonCoordinator gets first first non-coordinator node in the list of nodes.
func (c *Cluster) GetNonCoordinator() *Command {
for _, n := range c.Nodes {
if !n.IsCoordinator() {
if !n.IsPrimary() {
return n
}
}
@ -158,7 +158,7 @@ func (c *Cluster) GetNonCoordinator() *Command {
func (c *Cluster) GetNonCoordinators() []*Command {
rtn := make([]*Command, 0)
for _, n := range c.Nodes {
if !n.IsCoordinator() {
if !n.IsPrimary() {
rtn = append(rtn, n)
}
}
@ -453,7 +453,7 @@ func (c *Cluster) Close() error {
func (c *Cluster) CloseAndRemoveNonCoordinator() error {
for i, n := range c.Nodes {
if !n.IsCoordinator() {
if !n.IsPrimary() {
return c.CloseAndRemove(i)
}
}
@ -522,7 +522,7 @@ func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (*Clust
if len(opts) > 0 {
commandOpts = opts[i%len(opts)]
}
m := NewCommandNode(tb, i == 0, commandOpts...)
m := NewCommandNode(tb, commandOpts...)
cluster.Nodes[i] = m
}

View file

@ -90,13 +90,12 @@ func newCommand(tb testing.TB, opts ...server.CommandOption) *Command {
}
// NewCommandNode returns a new instance of Command with clustering enabled.
func NewCommandNode(tb testing.TB, isCoordinator bool, opts ...server.CommandOption) *Command {
func NewCommandNode(tb testing.TB, opts ...server.CommandOption) *Command {
// We want tests to default to using the in-memory translate store, so we
// prepend opts with that functional option. If a different translate store
// has been specified, it will override this one.
opts = prependTestServerOpts(opts)
m := newCommand(tb, opts...)
m.Config.Cluster.Coordinator = isCoordinator
return m
}
@ -191,9 +190,9 @@ func (m *Command) URL() string { return m.API.Node().URI.String() }
// ID returns the node ID used by the running program.
func (m *Command) ID() string { return m.API.Node().ID }
// IsCoordinator returns true if this is the coordinator.
func (m *Command) IsCoordinator() bool {
coord := m.API.CoordinatorNode()
// IsPrimary returns true if this is the primary.
func (m *Command) IsPrimary() bool {
coord := m.API.PrimaryNode()
if coord == nil {
return false
}

View file

@ -85,7 +85,7 @@ func TestNewCluster(t *testing.T) {
func getCoordinator(m *test.Command) string {
hosts := m.API.Hosts(context.Background())
for _, host := range hosts {
if host.IsCoordinator {
if host.IsPrimary {
return host.ID
}
}

View file

@ -25,11 +25,11 @@ import (
type Node struct {
Mu sync.Mutex `json:"-"` // TODO: we really need to get rid of this
ID string `json:"id"`
URI net.URI `json:"uri"`
GRPCURI net.URI `json:"grpc-uri"`
IsCoordinator bool `json:"isCoordinator"`
State string `json:"state"`
ID string `json:"id"`
URI net.URI `json:"uri"`
GRPCURI net.URI `json:"grpc-uri"`
IsPrimary bool `json:"isPrimary"`
State string `json:"state"`
}
func (n *Node) ProtectedClone() *Node {
@ -46,13 +46,13 @@ func (n *Node) Clone() *Node {
other.ID = n.ID
other.URI = n.URI
other.GRPCURI = n.GRPCURI
other.IsCoordinator = n.IsCoordinator
other.IsPrimary = n.IsPrimary
other.State = n.State
return &other
}
func (n *Node) String() string {
return fmt.Sprintf("Node:%s:%s:%s(%v)", n.URI, n.State, n.ID, n.IsCoordinator)
return fmt.Sprintf("Node:%s:%s:%s(%v)", n.URI, n.State, n.ID, n.IsPrimary)
}
// Nodes represents a list of nodes.

View file

@ -204,28 +204,24 @@ func TestTranslation_Reset(t *testing.T) {
c := test.MustRunCluster(t, 4,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(true),
pilosa.OptServerNodeID("2node0"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("4node1"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("3node2"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("1node3"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
@ -304,28 +300,24 @@ func TestTranslation_KeyNotFound(t *testing.T) {
c := test.MustRunCluster(t, 4,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(true),
pilosa.OptServerNodeID("node0"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node1"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node2"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node3"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
@ -462,21 +454,18 @@ func TestTranslation_Replication(t *testing.T) {
c := test.MustRunCluster(t, 3,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(true),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
pilosa.OptServerReplicaN(2),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
pilosa.OptServerReplicaN(2),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
pilosa.OptServerReplicaN(2),
@ -552,14 +541,12 @@ func TestTranslation_Coordinator(t *testing.T) {
c := test.MustRunCluster(t, 2,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(true),
pilosa.OptServerNodeID("node0"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node1"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
@ -624,28 +611,24 @@ func TestTranslation_TranslateIDsOnCluster(t *testing.T) {
c := test.MustRunCluster(t, 4,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(true),
pilosa.OptServerNodeID("node0"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node1"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node2"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node3"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),

View file

@ -84,7 +84,6 @@ func NewTestCluster(tb testing.TB, n int) *cluster {
cNodes := c.noder.Nodes()
c.Node = cNodes[0]
c.Coordinator = cNodes[0].ID
c.SetState(string(ClusterStateNormal))
return c
@ -266,9 +265,8 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*cluster, error)
uri := NewTestURI("http", fmt.Sprintf("host%d", i), uint16(0))
node := &topology.Node{
ID: id,
URI: uri,
IsCoordinator: i == 0,
ID: id,
URI: uri,
}
// add URI to common
@ -296,7 +294,7 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*cluster, error)
c.Topology = NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c)
c.holder = h
c.Node = node
c.Coordinator = t.common.Nodes[0].ID // the first node is the coordinator
// c.Coordinator = t.common.Nodes[0].ID // the first node is the coordinator
c.broadcaster = t.broadcaster(c)
// add nodes
@ -530,7 +528,7 @@ func (t *ClusterCluster) FollowResizeInstruction(instr *ResizeInstruction) error
complete.Error = err.Error()
}
node := instr.Coordinator
node := instr.Primary
return bcast{t: t}.SendTo(node, complete)
}
@ -567,7 +565,7 @@ func NewTestClusterWithReplication(tb testing.TB, nNodes, nReplicas, partitionN
cNodes := c.noder.Nodes()
c.Node = cNodes[0]
c.Coordinator = cNodes[0].ID
// c.Coordinator = cNodes[0].ID
c.SetState(string(ClusterStateNormal))
if err := c.holder.Open(); err != nil {