diff --git a/holder.go b/holder.go index bb4206dfa..3b9631761 100644 --- a/holder.go +++ b/holder.go @@ -1233,6 +1233,15 @@ func (c *holderCleaner) CleanHolder() error { // Get the fragments registered in memory. for _, field := range index.Fields() { + // deletedShards is used to track which shards for the field + // were deleted. Any shards that get deleted from this node + // get added to remoteAvailableShards. This is done because + // the CleanHolder process is cleaning up shards which got + // moved to other nodes. Because those shards still exist + // (just no longer on this particular node), this node still + // needs to consider each of them as an available shard in + // the cluster. + var deletedShards []uint64 for _, view := range field.views() { for _, fragment := range view.allFragments() { fragShard := fragment.shard @@ -1244,6 +1253,12 @@ func (c *holderCleaner) CleanHolder() error { if err := view.deleteFragment(fragShard); err != nil { return errors.Wrap(err, "deleting fragment") } + deletedShards = append(deletedShards, fragShard) + } + } + if len(deletedShards) > 0 { + if err := field.AddRemoteAvailableShards(roaring.NewBitmap(deletedShards...)); err != nil { + return errors.Wrap(err, "adding remote available shards") } } } diff --git a/server/cluster_test.go b/server/cluster_test.go index 9ecb57437..0899b8e97 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -138,7 +138,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) { } // Ensure that adding a node correctly resizes the cluster. -func TestXClusterResize_AddNode(t *testing.T) { +func TestClusterResize_AddNode(t *testing.T) { t.Run("NoData", func(t *testing.T) { clus := test.MustRunCluster(t, 2) defer clus.Close() @@ -387,8 +387,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { t.Fatal(err) } else if res != exp { t.Fatalf("unexpected result: %s", res) - } else { - fmt.Println(res) } // Configure node1 diff --git a/view.go b/view.go index 65965b43e..06d38dce8 100644 --- a/view.go +++ b/view.go @@ -278,11 +278,11 @@ func (v *view) CreateFragmentIfNotExists(shard uint64) (*fragment, error) { frag.RowAttrStore = v.rowAttrStore v.fragments[shard] = frag - v.notifyIfNew(shard) + v.notifyIfNewShard(shard) return frag, nil } -func (v *view) notifyIfNew(shard uint64) { +func (v *view) notifyIfNewShard(shard uint64) { if v.remoteShardPresent(shard) { //checks the fields remoteShards bitmap to see if broadcast needed return }