diff --git a/cluster.go b/cluster.go index 479a5c6f2..ed1d92eb9 100644 --- a/cluster.go +++ b/cluster.go @@ -110,6 +110,17 @@ func (a Nodes) ContainsID(id string) bool { return false } +// NodeByID returns the node for an ID. If the ID is not found, +// it returns nil. +func (a Nodes) NodeByID(id string) *Node { + for _, n := range a { + if n.ID == id { + return n + } + } + return nil +} + // Filter returns a new list of nodes with node removed. func (a Nodes) Filter(n *Node) []*Node { other := make([]*Node, 0, len(a)) @@ -1027,20 +1038,49 @@ func (c *cluster) ownsShard(nodeID string, index string, shard uint64) bool { func (c *cluster) partitionNodes(partitionID int) []*Node { // Default replica count to between one and the number of nodes. // The replica count can be zero if there are no nodes. + + // Assume that c.nodes may be missing a node that is part of the cluster but not currently present. + // The partition calculation must use the full cluster size in BOTH cases: + // - use len(c.Topology.nodeIDs) instead of len(c.nodes), + // - collect nodes from c.Topology.nodeIDs rather than from c.nodes, + // - when the node is missing, it should be considered, found absent from c.nodes, then omitted from the return slice. + + // Use c.Topology to determine cluster membership when it + // exists and contains data. Otherwise, fall back to using + // c.nodes. The only time c.Topology should be nil is in + // tests. + var useTopology bool + if c.Topology != nil && len(c.Topology.nodeIDs) > 0 { + useTopology = true + } + replicaN := c.ReplicaN - if replicaN > len(c.nodes) { - replicaN = len(c.nodes) + var nodeN int + if useTopology { + nodeN = len(c.Topology.nodeIDs) + } else { + nodeN = len(c.nodes) + } + if replicaN > nodeN { + replicaN = nodeN } else if replicaN == 0 { replicaN = 1 } // Determine primary owner node. - nodeIndex := c.Hasher.Hash(uint64(partitionID), len(c.nodes)) + nodeIndex := c.Hasher.Hash(uint64(partitionID), nodeN) // Collect nodes around the ring. - nodes := make([]*Node, replicaN) + nodes := make([]*Node, 0, replicaN) for i := 0; i < replicaN; i++ { - nodes[i] = c.nodes[(nodeIndex+i)%len(c.nodes)] + if useTopology { + maybeNodeID := c.Topology.nodeIDs[(nodeIndex+i)%nodeN] + if node := Nodes(c.nodes).NodeByID(maybeNodeID); node != nil { + nodes = append(nodes, node) + } + } else { + nodes = append(nodes, c.nodes[(nodeIndex+i)%len(c.nodes)]) + } } return nodes @@ -2193,7 +2233,7 @@ func (c *cluster) nodeStatus() *NodeStatus { func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error { c.mu.Lock() defer c.mu.Unlock() - c.logger.Printf("merge cluster status: node=%s cluster=%v", c.Node.ID, cs) + c.logger.Printf("merge cluster status: node=%s cluster=%v, topologySize=%v", c.Node.ID, cs, len(c.Topology.nodeIDs)) // Ignore status updates from self (coordinator). if c.unprotectedIsCoordinator() { return nil diff --git a/server/cluster_test.go b/server/cluster_test.go index d788de803..744b63dcb 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -143,9 +143,9 @@ func TestClusterResize_AddNode(t *testing.T) { clus := test.MustRunCluster(t, 2) defer clus.Close() - if !checkClusterState(clus[0], pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(clus[0], pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", clus[0].API.State()) - } else if !checkClusterState(clus[1], pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(clus[1], pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", clus[1].API.State()) } }) @@ -176,9 +176,9 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } }) @@ -224,9 +224,9 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } @@ -273,9 +273,9 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } @@ -326,9 +326,9 @@ func TestClusterResize_AddNode(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } @@ -373,9 +373,9 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } @@ -431,9 +431,9 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { }() defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } @@ -489,9 +489,9 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } @@ -545,9 +545,9 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } defer m1.Close() - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) @@ -598,11 +598,11 @@ func TestCluster_GossipMembership(t *testing.T) { t.Fatal(err) } - if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + if !test.CheckClusterState(m0, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) - } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) - } else if !checkClusterState(m2, pilosa.ClusterStateNormal, 1000) { + } else if !test.CheckClusterState(m2, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node2 cluster state: %s", m2.API.State()) } @@ -725,15 +725,3 @@ func TestClusterMutualTLS(t *testing.T) { t.Fatal(err) } } - -// checkClusterState polls a given cluster for its state until it -// receives a matching state. It polls up to n times before returning. -func checkClusterState(m *test.Command, state string, n int) bool { - for i := 0; i < n; i++ { - if m.API.State() == state { - return true - } - time.Sleep(10 * time.Millisecond) - } - return false -} diff --git a/test/cluster.go b/test/cluster.go index babb3890f..4fbc2050e 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -193,6 +193,18 @@ func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu return c } +// CheckClusterState polls a given cluster for its state until it +// receives a matching state. It polls up to n times before returning. +func CheckClusterState(m *Command, state string, n int) bool { + for i := 0; i < n; i++ { + if m.API.State() == state { + return true + } + time.Sleep(10 * time.Millisecond) + } + return false +} + // newCluster creates a new cluster func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (Cluster, error) { if size == 0 { diff --git a/translator_test.go b/translator_test.go index 72a4e8ddd..0acce9b66 100644 --- a/translator_test.go +++ b/translator_test.go @@ -296,6 +296,82 @@ func TestTranslation_Reset(t *testing.T) { }) } +// Test index key translation replication under node failure. +func TestTranslation_Replication(t *testing.T) { + t.Run("Replication", func(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), + )}, + ) + + node0 := c[0] + node1 := c[1] + + ctx := context.Background() + idx := "i" + field := "f" + + // Create an index with keys. + if _, err := node0.API.CreateIndex(ctx, idx, + pilosa.IndexOptions{ + Keys: true, + }); err != nil { + t.Fatal(err) + } + + if _, err := node0.API.CreateField(ctx, idx, field); err != nil { + t.Fatal(err) + } + + // Write data on first node. + // these keys are a minimal example to reproduce the problem for the case of a 3-node cluster with replication factor 2 + if _, err := node0.Queryf(t, idx, "", ` + Set("x1", f=1) + Set("x2", f=1) + `); err != nil { + t.Fatal(err) + } + + exp := `{"results":[{"attrs":{},"columns":[],"keys":["x1","x2"]}]}` + + if !test.CheckClusterState(node0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", node0.API.State()) + } else if !test.CheckClusterState(node1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", node1.API.State()) + } + + // Verify the data exists + node0.QueryExpect(t, idx, "", `Row(f=1)`, exp) + + // Kill one node. + if err := node1.Command.Close(); err != nil { + t.Fatal(err) + } + + // Verify the data exists with one node down + node0.QueryExpect(t, idx, "", `Row(f=1)`, exp) + }) +} + // Test key translation with multiple nodes. func TestTranslation_Coordinator(t *testing.T) { // Ensure that field key translations requests sent to