diff --git a/api.go b/api.go index c9242a9bd..847cc7ebd 100644 --- a/api.go +++ b/api.go @@ -734,13 +734,18 @@ func (api *API) RemoveNode(id string) (*Node, error) { return nil, errors.Wrap(err, "validating api method") } - removeNode := api.cluster.unprotectedNodeByID(id) + removeNode := api.cluster.nodeByID(id) if removeNode == nil { - return nil, errors.Wrap(ErrNodeIDNotExists, "finding node to remove") + if !api.cluster.topologyContainsNode(id) { + return nil, errors.Wrap(ErrNodeIDNotExists, "finding node to remove") + } + removeNode = &Node{ + ID: id, + } } // Start the resize process (similar to NodeJoin) - err := api.cluster.nodeLeave(removeNode) + err := api.cluster.nodeLeave(id) if err != nil { return removeNode, errors.Wrap(err, "calling node leave") } diff --git a/cluster.go b/cluster.go index 86fdf7712..309846e0b 100644 --- a/cluster.go +++ b/cluster.go @@ -341,17 +341,15 @@ func (c *cluster) addNode(node *Node) error { // removeNode removes a node from the Cluster and updates and saves the // new topology. unprotected. -func (c *cluster) removeNode(node *Node) error { +func (c *cluster) removeNode(nodeID string) error { // remove from cluster - if !c.removeNodeBasicSorted(node) { - return nil - } + c.removeNodeBasicSorted(nodeID) // remove from topology if c.Topology == nil { return fmt.Errorf("Cluster.Topology is nil") } - if !c.Topology.removeID(node.ID) { + if !c.Topology.removeID(nodeID) { return nil } @@ -511,6 +509,17 @@ func (c *cluster) unprotectedNodeByID(id string) *Node { return nil } +func (c *cluster) topologyContainsNode(id string) bool { + c.Topology.mu.RLock() + defer c.Topology.mu.RUnlock() + for _, nid := range c.Topology.nodeIDs { + if id == nid { + return true + } + } + return false +} + // nodePositionByID returns the position of the node in slice c.Nodes. func (c *cluster) nodePositionByID(nodeID string) int { for i, n := range c.Nodes { @@ -539,8 +548,8 @@ func (c *cluster) addNodeBasicSorted(node *Node) bool { // removeNodeBasicSorted removes a node from the cluster, maintaining the sort // order. Returns true if the node was removed. unprotected. -func (c *cluster) removeNodeBasicSorted(node *Node) bool { - i := c.nodePositionByID(node.ID) +func (c *cluster) removeNodeBasicSorted(nodeID string) bool { + i := c.nodePositionByID(nodeID) if i < 0 { return false } @@ -970,7 +979,7 @@ func (c *cluster) handleNodeAction(nodeAction nodeAction) error { if j.action == resizeJobActionRemove { c.mu.Lock() defer c.mu.Unlock() - return c.removeNode(nodeAction.node) + return c.removeNode(nodeAction.node.ID) } else if j.action == resizeJobActionAdd { c.mu.Lock() defer c.mu.Unlock() @@ -1100,7 +1109,7 @@ func (c *cluster) unprotectedGenerateResizeJobByAction(nodeAction nodeAction) (* toCluster.partitionN = c.partitionN toCluster.ReplicaN = c.ReplicaN if nodeAction.action == resizeJobActionRemove { - toCluster.removeNodeBasicSorted(nodeAction.node) + toCluster.removeNodeBasicSorted(nodeAction.node.ID) } else if nodeAction.action == resizeJobActionAdd { toCluster.addNodeBasicSorted(nodeAction.node) } @@ -1598,7 +1607,7 @@ func (c *cluster) ReceiveEvent(e *NodeEvent) (err error) { // not already removed by a removeNode request. We treat this as the // host being temporarily unavailable, and expect it to come back // up. - if c.removeNodeBasicSorted(e.Node) { + if c.removeNodeBasicSorted(e.Node.ID) { c.Topology.nodeStates[e.Node.ID] = nodeStateDown // put the cluster into STARTING if we've lost a number of nodes // equal to or greater than ReplicaN @@ -1685,47 +1694,45 @@ func (c *cluster) nodeJoin(node *Node) error { } // nodeLeave initiates the removal of a node from the cluster. -func (c *cluster) nodeLeave(node *Node) error { +func (c *cluster) nodeLeave(nodeID string) error { c.mu.Lock() defer c.mu.Unlock() // Refuse the request if this is not the coordinator. if !c.unprotectedIsCoordinator() { - return fmt.Errorf("node removal requests are only valid on the coordinator node: %s", c.unprotectedCoordinatorNode().ID) + return fmt.Errorf("node removal requests are only valid on the coordinator node: %s", + c.unprotectedCoordinatorNode().ID) } - if c.state != ClusterStateNormal { - return fmt.Errorf("Cluster must be in state %s to remove a node. Current state: %s", ClusterStateNormal, c.state) + if c.state != ClusterStateNormal && c.state != ClusterStateDegraded { + return fmt.Errorf("Cluster must be in state %s to remove a node. Current state: %s", + ClusterStateNormal, c.state) } // Ensure that node is in the cluster. - if c.unprotectedNodeByID(node.ID) == nil { - return fmt.Errorf("Node is not a member of the cluster: %s", node.ID) + if !c.topologyContainsNode(nodeID) { + return fmt.Errorf("Node is not a member of the cluster: %s", nodeID) } // Prevent removing the coordinator node (this node). - if node.ID == c.Node.ID { + if nodeID == c.Node.ID { return fmt.Errorf("coordinator cannot be removed; first, make a different node the new coordinator.") } // See if resize job can be generated - if _, err := c.unprotectedGenerateResizeJobByAction(nodeAction{c.unprotectedNodeByID(node.ID), resizeJobActionRemove}); err != nil { + if _, err := c.unprotectedGenerateResizeJobByAction( + nodeAction{ + node: &Node{ID: nodeID}, + action: resizeJobActionRemove}, + ); err != nil { return errors.Wrap(err, "generating job") } - // Get the actual node in the local cluster. - n := c.unprotectedNodeByID(node.ID) - - // Don't do anything else if the cluster doesn't contain the node. - if n == nil { - return nil - } - // If the holder does not yet contain data, go ahead and remove the node. if ok, err := c.holder.HasData(); !ok && err == nil { - if err := c.removeNode(n); err != nil { + if err := c.removeNode(nodeID); err != nil { return errors.Wrap(err, "removing node") } - return c.unprotectedSetStateAndBroadcast(ClusterStateNormal) + return c.unprotectedSetStateAndBroadcast(c.determineClusterState()) } else if err != nil { return errors.Wrap(err, "checking if holder has data") } @@ -1735,7 +1742,7 @@ func (c *cluster) nodeLeave(node *Node) error { if err := c.unprotectedSetStateAndBroadcast(ClusterStateResizing); err != nil { return errors.Wrap(err, "broadcasting state") } - c.joiningLeavingNodes <- nodeAction{n, resizeJobActionRemove} + c.joiningLeavingNodes <- nodeAction{node: &Node{ID: nodeID}, action: resizeJobActionRemove} return nil } @@ -1777,7 +1784,7 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error { } for _, nodeID := range nodeIDsToRemove { - if err := c.removeNode(c.unprotectedNodeByID(nodeID)); err != nil { + if err := c.removeNode(nodeID); err != nil { return errors.Wrap(err, "removing node") } } diff --git a/server/server_test.go b/server/server_test.go index 9fa2df34d..1f00ebeb5 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -396,7 +396,7 @@ func TestClusteringNodesReplica1(t *testing.T) { t.Fatalf("closing third node: %v", err) } - // TODO: confirm that cluster stops accepting queries after one node closes + // confirm that cluster stops accepting queries after one node closes if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { t.Fatalf("got unexpected error querying an incomplete cluster: %v", err) } @@ -456,12 +456,12 @@ func TestClusteringNodesReplica2(t *testing.T) { t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) } - // TODO: implement and confirm that cluster keeps accepting queries if replication > 1 + // confirm that cluster keeps accepting queries if replication > 1 if _, err := cluster[0].API.CreateIndex(context.Background(), "anewindex", pilosa.IndexOptions{}); err != nil { t.Fatalf("got unexpected error creating index: %v", err) } - // TODO: confirm that cluster stops accepting queries if 2 nodes fail and replication == 2 + // confirm that cluster stops accepting queries if 2 nodes fail and replication == 2 if err := cluster[1].Command.Close(); err != nil { t.Fatalf("closing 2nd node: %v", err) } @@ -516,3 +516,48 @@ func TestClusteringNodesReplica2(t *testing.T) { time.Sleep(time.Millisecond) } } + +func TestRemoveNodeAfterItDies(t *testing.T) { + cluster := test.MustNewCluster(t, 3) + for _, c := range cluster { + c.Config.Cluster.ReplicaN = 2 + } + err := cluster.Start() + if err != nil { + t.Fatalf("starting cluster: %v", err) + } + + var wait = true + for wait { + wait = false + for _, node := range cluster { + if node.API.State() != pilosa.ClusterStateNormal { + wait = true + } + } + time.Sleep(time.Millisecond * 1) + } + + if err := cluster[2].Command.Close(); err != nil { + t.Fatalf("closing third node: %v", err) + } + + if cluster[0].API.State() != pilosa.ClusterStateDegraded { + t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) + } + + if _, err := cluster[0].API.RemoveNode(cluster[2].API.Node().ID); err != nil { + t.Fatalf("removing failed node: %v", err) + } + + if cluster[0].API.State() != pilosa.ClusterStateNormal { + t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) + } + + hosts := cluster[0].API.Hosts(context.Background()) + if len(hosts) != 2 { + t.Fatalf("unexpected hosts: %v", hosts) + } +} + +// TODO: confirm that things keep working if a node is hard-closed (no nodeLeave event) and immediately restarted with a different address.