diff --git a/cluster.go b/cluster.go index a2d537b11..ee0912b09 100644 --- a/cluster.go +++ b/cluster.go @@ -739,7 +739,7 @@ func (c *Cluster) fragSources(to *Cluster, idx *Index) (map[string][]*internal.R // the fragment. srcNodeID, ok := srcNodesByFrag[frag] if !ok { - return nil, errors.New("not enough data to perform resize") + return nil, errors.New("not enough data to perform resize (replica factor may need to be increased)") } src := &internal.ResizeSource{ @@ -939,6 +939,10 @@ func (c *Cluster) allNodesReady() bool { func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { j, err := c.generateResizeJob(nodeAction) if err != nil { + c.logger().Printf("generateResizeJob error: err=%s", err) + if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil { + c.logger().Printf("setStateAndBroadcast error: err=%s", err) + } return err } @@ -1000,7 +1004,13 @@ func (c *Cluster) ListenForJoins() { } func (c *Cluster) listenForJoins() { - var uriJoined bool + // When a cluster starts, the state is STARTING. + // We first want to wait for at least one node to join. + // Then we want to clear out the joiningLeavingNodes queue (buffered channel). + // Then we want to set the cluster state to NORMAL and resume processing of joiningLeavingNodes events. + // We use a bool `setNormal` to indicate when at least one node has joined. + + var setNormal bool for { @@ -1012,13 +1022,13 @@ func (c *Cluster) listenForJoins() { c.logger().Printf("handleNodeAction error: err=%s", err) continue } - uriJoined = true + setNormal = true continue default: } // Only change state to NORMAL if we have successfully added at least one host. - if uriJoined { + if setNormal { // Put the cluster back to state NORMAL and broadcast. if err := c.setStateAndBroadcast(ClusterStateNormal); err != nil { c.logger().Printf("setStateAndBroadcast error: err=%s", err) @@ -1035,7 +1045,7 @@ func (c *Cluster) listenForJoins() { c.logger().Printf("handleNodeAction error: err=%s", err) continue } - uriJoined = true + setNormal = true continue } } @@ -1704,6 +1714,13 @@ func (c *Cluster) NodeLeave(node *Node) error { return fmt.Errorf("The coordinator node cannot be removed. First, make a different node the new coordinator.") } + // See if resize job can be generated + _, err := c.generateResizeJobByAction(nodeAction{c.nodeByID(node.ID), ResizeJobActionRemove}) + + if err != nil { + return err + } + return c.nodeLeave(node) } diff --git a/server/cluster_test.go b/server/cluster_test.go index 5ec805f81..299415bc3 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -545,4 +545,37 @@ func TestClusterResize_RemoveNode(t *testing.T) { t.Fatalf("expected Body '%s' but got '%s'", expBody, strings.TrimSpace(resp.Body)) } }) + + t.Run("ErrorRemoveWithoutReplicas", func(t *testing.T) { + client0 := m0.Client() + + // Create indexes and frames on one node. + if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { + t.Fatal(err) + } else if err := client0.CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { + t.Fatal(err) + } + + // This is an attempt to ensure there is data on both nodes, but is not guaranteed. + // TODO: Deterministic node IDs would ensure consistent results + setBits := "" + for i := 0; i < 20; i++ { + setBits += fmt.Sprintf("SetBit(rowID=1, frame=\"f\", columnID=%d) ", i*pilosa.SliceWidth) + } + + if _, err := m0.Query("i", "", setBits); err != nil { + t.Fatal(err) + } + + resp := test.MustDo("GET", m1.URL()+fmt.Sprintf("/id"), "") + nodeID := resp.Body + + resp = test.MustDo("POST", m0.URL()+fmt.Sprintf("/cluster/resize/remove-node"), fmt.Sprintf(`{"id": "%s"}`, nodeID)) + expBody := "not enough data to perform resize" + if resp.StatusCode != http.StatusInternalServerError { + t.Fatalf("expected StatusCode %d but got %d", http.StatusInternalServerError, resp.StatusCode) + } else if !strings.Contains(resp.Body, expBody) { + t.Fatalf("expected to contain '%s' but got '%s'", expBody, strings.TrimSpace(resp.Body)) + } + }) }