diff --git a/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc b/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc index 28ddaf61e..3ab438eb2 100644 --- a/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_placement_group_scheduler_test.cc @@ -29,9 +29,10 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { io_service_.run(); })); - raylet_client_ = std::make_shared(); - raylet_client1_ = std::make_shared(); - raylet_client2_ = std::make_shared(); + for (int index = 0; index < 3; ++index) { + raylet_clients_.push_back( + std::make_shared()); + } gcs_table_storage_ = std::make_shared(io_service_); gcs_pub_sub_ = std::make_shared(redis_client_); gcs_node_manager_ = std::make_shared( @@ -41,15 +42,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { scheduler_ = std::make_shared( io_service_, gcs_table_storage_, *gcs_node_manager_, /*lease_client_fplacement_groupy=*/ - [this](const rpc::Address &address) { - if (0 == address.port()) { - return raylet_client_; - } else if (1 == address.port()) { - return raylet_client1_; - } else { - return raylet_client2_; - } - }); + [this](const rpc::Address &address) { return raylet_clients_[address.port()]; }); } void TearDown() override { @@ -91,7 +84,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { // The lease request should not be send and the scheduling of placement_group should // fail as there are no available nodes. - ASSERT_EQ(raylet_client_->num_lease_requested, 0); + ASSERT_EQ(raylet_clients_[0]->num_lease_requested, 0); ASSERT_EQ(0, success_placement_groups_.size()); ASSERT_EQ(1, failure_placement_groups_.size()); ASSERT_EQ(placement_group, failure_placement_groups_.front()); @@ -118,10 +111,10 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { success_placement_groups_.emplace_back(std::move(placement_group)); }); - ASSERT_EQ(2, raylet_client_->num_lease_requested); - ASSERT_EQ(2, raylet_client_->lease_callbacks.size()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + ASSERT_EQ(2, raylet_clients_[0]->num_lease_requested); + ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); WaitPendingDone(failure_placement_groups_, 0); WaitPendingDone(success_placement_groups_, 1); ASSERT_EQ(placement_group, success_placement_groups_.front()); @@ -152,8 +145,8 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { AddNode(Mocker::GenNodeInfo(0), 2); scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); WaitPendingDone(success_placement_groups_, 1); } @@ -164,9 +157,7 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { boost::asio::io_service io_service_; std::shared_ptr store_client_; - std::shared_ptr raylet_client_; - std::shared_ptr raylet_client1_; - std::shared_ptr raylet_client2_; + std::vector> raylet_clients_; std::shared_ptr gcs_node_manager_; std::shared_ptr scheduler_; std::vector> success_placement_groups_; @@ -225,12 +216,12 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestSchedulePlacementGroupReplyFailure) { success_placement_groups_.emplace_back(std::move(placement_group)); }); - ASSERT_EQ(2, raylet_client_->num_lease_requested); - ASSERT_EQ(2, raylet_client_->lease_callbacks.size()); + ASSERT_EQ(2, raylet_clients_[0]->num_lease_requested); + ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size()); // Reply failure, so the placement group scheduling failed. - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources(false)); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources(false)); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources(false)); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources(false)); WaitPendingDone(failure_placement_groups_, 1); WaitPendingDone(success_placement_groups_, 0); ASSERT_EQ(placement_group, failure_placement_groups_.front()); @@ -282,12 +273,12 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestSchedulePlacementGroupReturnResource) success_placement_groups_.emplace_back(std::move(placement_group)); }); - ASSERT_EQ(2, raylet_client_->num_lease_requested); - ASSERT_EQ(2, raylet_client_->lease_callbacks.size()); + ASSERT_EQ(2, raylet_clients_[0]->num_lease_requested); + ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size()); // One bundle success and the other failed. - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources(false)); - ASSERT_EQ(1, raylet_client_->num_return_requested); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources(false)); + ASSERT_EQ(1, raylet_clients_[0]->num_return_requested); // Reply the placement_group creation request, then the placement_group should be // scheduled successfully. WaitPendingDone(failure_placement_groups_, 1); @@ -308,8 +299,9 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyBalancedScheduling) }; // Schedule placement group, it will be evenly distributed over the two nodes. - int select_node0_count = 0; - int select_node1_count = 0; + int node_select_count[2] = {0, 0}; + int node_commit_count[2] = {0, 0}; + int node_index = 0; for (int index = 0; index < 10; ++index) { auto request = Mocker::GenCreatePlacementGroupRequest("", rpc::PlacementStrategy::STRICT_PACK); @@ -317,19 +309,20 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyBalancedScheduling) scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); - if (!raylet_client_->lease_callbacks.empty()) { - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ++select_node0_count; - } else { - ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources()); - ++select_node1_count; - } + node_index = !raylet_clients_[0]->lease_callbacks.empty() ? 0 : 1; + ++node_select_count[node_index]; + node_commit_count[node_index] += 2; + ASSERT_TRUE(raylet_clients_[node_index]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[node_index]->GrantPrepareBundleResources()); + auto condition = [this, node_index, node_commit_count]() { + return raylet_clients_[node_index]->num_commit_requested == + node_commit_count[node_index]; + }; + EXPECT_TRUE(WaitForCondition(condition, timeout_ms_.count())); } WaitPendingDone(success_placement_groups_, 10); - ASSERT_EQ(select_node0_count, 5); - ASSERT_EQ(select_node1_count, 5); + ASSERT_EQ(node_select_count[0], 5); + ASSERT_EQ(node_select_count[1], 5); } TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyReschedulingWhenNodeAdd) { @@ -351,8 +344,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyResourceCheck) { Mocker::GenCreatePlacementGroupRequest("", rpc::PlacementStrategy::STRICT_PACK); auto placement_group = std::make_shared(request); scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); WaitPendingDone(success_placement_groups_, 1); // Node1 has less number of bundles, but it doesn't satisfy the resource @@ -364,8 +357,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyResourceCheck) { auto placement_group2 = std::make_shared(create_placement_group_request2); scheduler_->ScheduleUnplacedBundles(placement_group2, failure_handler, success_handler); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); WaitPendingDone(success_placement_groups_, 2); } @@ -390,19 +383,18 @@ TEST_F(GcsPlacementGroupSchedulerTest, DestroyPlacementGroup) { absl::MutexLock lock(&vector_mutex_); success_placement_groups_.emplace_back(std::move(placement_group)); }); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); WaitPendingDone(failure_placement_groups_, 0); WaitPendingDone(success_placement_groups_, 1); - RAY_LOG(ERROR) << "sanngbin1"; const auto &placement_group_id = placement_group->GetPlacementGroupID(); scheduler_->DestroyPlacementGroupBundleResourcesIfExists(placement_group_id); - ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve()); - ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve()); + ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve()); + ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve()); // Subsequent destroy request should not do anything. scheduler_->DestroyPlacementGroupBundleResourcesIfExists(placement_group_id); - ASSERT_FALSE(raylet_client_->GrantCancelResourceReserve()); - ASSERT_FALSE(raylet_client_->GrantCancelResourceReserve()); + ASSERT_FALSE(raylet_clients_[0]->GrantCancelResourceReserve()); + ASSERT_FALSE(raylet_clients_[0]->GrantCancelResourceReserve()); } TEST_F(GcsPlacementGroupSchedulerTest, DestroyCancelledPlacementGroup) { @@ -429,11 +421,11 @@ TEST_F(GcsPlacementGroupSchedulerTest, DestroyCancelledPlacementGroup) { }); // Now, cancel the schedule request. - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); scheduler_->MarkScheduleCancelled(placement_group_id); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve()); - ASSERT_TRUE(raylet_client_->GrantCancelResourceReserve()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve()); + ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve()); WaitPendingDone(failure_placement_groups_, 1); } @@ -459,13 +451,13 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestPackStrategyLargeBundlesScheduling) { Mocker::GenCreatePlacementGroupRequest("", rpc::PlacementStrategy::PACK, 15); auto placement_group = std::make_shared(request); scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); - RAY_CHECK(raylet_client_->num_lease_requested > 0); - RAY_CHECK(raylet_client1_->num_lease_requested > 0); - for (int index = 0; index < raylet_client_->num_lease_requested; ++index) { - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); + RAY_CHECK(raylet_clients_[0]->num_lease_requested > 0); + RAY_CHECK(raylet_clients_[1]->num_lease_requested > 0); + for (int index = 0; index < raylet_clients_[0]->num_lease_requested; ++index) { + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); } - for (int index = 0; index < raylet_client1_->num_lease_requested; ++index) { - ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources()); + for (int index = 0; index < raylet_clients_[1]->num_lease_requested; ++index) { + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); } WaitPendingDone(success_placement_groups_, 1); } @@ -492,8 +484,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestRescheduleWhenNodeDead) { }; scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client1_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); WaitPendingDone(success_placement_groups_, 1); auto bundles_on_node0 = @@ -509,8 +501,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestRescheduleWhenNodeDead) { // TODO(ffbin): We need to see which node the other bundles that have been placed are // deployed on, and spread them as far as possible. It will be implemented in the next // pr. - raylet_client_->GrantPrepareBundleResources(); - raylet_client1_->GrantPrepareBundleResources(); + raylet_clients_[0]->GrantPrepareBundleResources(); + raylet_clients_[1]->GrantPrepareBundleResources(); WaitPendingDone(success_placement_groups_, 2); } @@ -543,8 +535,8 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictSpreadStrategyResourceCheck) { auto node2 = Mocker::GenNodeInfo(2); AddNode(node2); scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); - ASSERT_TRUE(raylet_client_->GrantPrepareBundleResources()); - ASSERT_TRUE(raylet_client2_->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[2]->GrantPrepareBundleResources()); WaitPendingDone(success_placement_groups_, 1); } diff --git a/src/ray/gcs/gcs_server/test/gcs_server_test_util.h b/src/ray/gcs/gcs_server/test/gcs_server_test_util.h index 3c6755917..835993fb9 100644 --- a/src/ray/gcs/gcs_server/test/gcs_server_test_util.h +++ b/src/ray/gcs/gcs_server/test/gcs_server_test_util.h @@ -170,6 +170,7 @@ struct GcsServerMocker { const BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) override { + num_commit_requested += 1; rpc::CommitBundleResourcesReply reply; callback(Status::OK(), reply); } @@ -215,6 +216,7 @@ struct GcsServerMocker { int num_lease_requested = 0; int num_return_requested = 0; + int num_commit_requested = 0; ClientID node_id = ClientID::FromRandom(); std::list> lease_callbacks = {}; std::list> return_callbacks = {};