diff --git a/src/ray/gcs/gcs_server/gcs_placement_group_manager.cc b/src/ray/gcs/gcs_server/gcs_placement_group_manager.cc index 6a15cd32f..3521dd225 100644 --- a/src/ray/gcs/gcs_server/gcs_placement_group_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_placement_group_manager.cc @@ -125,7 +125,7 @@ PlacementGroupID GcsPlacementGroupManager::GetPlacementGroupIDByName( void GcsPlacementGroupManager::OnPlacementGroupCreationFailed( std::shared_ptr placement_group) { RAY_LOG(INFO) << "Failed to create placement group " << placement_group->GetName() - << ", try again."; + << ", id: " << placement_group->GetPlacementGroupID() << ", try again."; // We will attempt to schedule this placement_group once an eligible node is // registered. auto state = placement_group->GetState(); @@ -148,7 +148,8 @@ void GcsPlacementGroupManager::OnPlacementGroupCreationFailed( void GcsPlacementGroupManager::OnPlacementGroupCreationSuccess( const std::shared_ptr &placement_group) { - RAY_LOG(INFO) << "Successfully created placement group " << placement_group->GetName(); + RAY_LOG(INFO) << "Successfully created placement group " << placement_group->GetName() + << ", id: " << placement_group->GetPlacementGroupID(); placement_group->UpdateState(rpc::PlacementGroupTableData::CREATED); auto placement_group_id = placement_group->GetPlacementGroupID(); RAY_CHECK_OK(gcs_table_storage_->PlacementGroupTable().Put( diff --git a/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.cc b/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.cc index f5332d0d6..bfeb209f3 100644 --- a/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.cc +++ b/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.cc @@ -202,6 +202,7 @@ void GcsPlacementGroupScheduler::ScheduleUnplacedBundles( auto strategy = placement_group->GetStrategy(); RAY_LOG(INFO) << "Scheduling placement group " << placement_group->GetName() + << ", id: " << placement_group->GetPlacementGroupID() << ", bundles size = " << bundles.size(); auto selected_nodes = scheduler_strategies_[strategy]->Schedule( bundles, GetScheduleContext(placement_group->GetPlacementGroupID())); @@ -209,6 +210,7 @@ void GcsPlacementGroupScheduler::ScheduleUnplacedBundles( // If no nodes are available, scheduling fails. if (selected_nodes.empty()) { RAY_LOG(INFO) << "Failed to schedule placement group " << placement_group->GetName() + << ", id: " << placement_group->GetPlacementGroupID() << ", because no nodes are available."; failure_callback(placement_group); return; @@ -242,9 +244,10 @@ void GcsPlacementGroupScheduler::ScheduleUnplacedBundles( void GcsPlacementGroupScheduler::DestroyPlacementGroupBundleResourcesIfExists( const PlacementGroupID &placement_group_id) { - bool is_committed = false; - bool is_prepared = false; - std::shared_ptr bundle_locations = std::make_shared(); + std::shared_ptr committed_bundle_locations = + std::make_shared(); + std::shared_ptr leasing_bundle_locations = + std::make_shared(); // Check if we can find committed bundle locations. const auto &maybe_bundle_locations = @@ -252,23 +255,25 @@ void GcsPlacementGroupScheduler::DestroyPlacementGroupBundleResourcesIfExists( // If bundle location has been already removed, it means bundles // are already destroyed. Do nothing. if (maybe_bundle_locations.has_value()) { - is_committed = true; - bundle_locations = maybe_bundle_locations.value(); + committed_bundle_locations = maybe_bundle_locations.value(); } + // Now let's see if there are leasing bundles. There could be leasing bundles and + // committed bundles at the same time if plaement groups are reshceduling. auto it = placement_group_leasing_in_progress_.find(placement_group_id); if (it != placement_group_leasing_in_progress_.end()) { const auto &leasing_context = it->second; - is_prepared = true; - bundle_locations = leasing_context->GetPreparedBundleLocations(); + leasing_bundle_locations = leasing_context->GetPreparedBundleLocations(); } - RAY_CHECK(!(is_committed && is_prepared)) - << "Anomaly detected. It shouldn't be possible that placement group is both " - "committing and preparing."; - // Cancel all resource reservation. - for (const auto &iter : *(bundle_locations)) { + for (const auto &iter : *(committed_bundle_locations)) { + auto &bundle_spec = iter.second.second; + auto &node_id = iter.second.first; + CancelResourceReserve(bundle_spec, gcs_node_manager_.GetNode(node_id)); + } + + for (const auto &iter : *(leasing_bundle_locations)) { auto &bundle_spec = iter.second.second; auto &node_id = iter.second.first; CancelResourceReserve(bundle_spec, gcs_node_manager_.GetNode(node_id)); @@ -293,19 +298,19 @@ void GcsPlacementGroupScheduler::PrepareResources( } const auto lease_client = GetLeaseClientFromNode(node.value()); const auto node_id = NodeID::FromBinary(node.value()->node_id()); - RAY_LOG(INFO) << "Preparing resource from node " << node_id - << " for a bundle: " << bundle->DebugString(); + RAY_LOG(DEBUG) << "Preparing resource from node " << node_id + << " for a bundle: " << bundle->DebugString(); lease_client->PrepareBundleResources( *bundle, [node_id, bundle, callback]( const Status &status, const rpc::PrepareBundleResourcesReply &reply) { auto result = reply.success() ? Status::OK() : Status::IOError("Failed to reserve resource"); if (result.ok()) { - RAY_LOG(INFO) << "Finished leasing resource from " << node_id - << " for bundle: " << bundle->DebugString(); + RAY_LOG(DEBUG) << "Finished leasing resource from " << node_id + << " for bundle: " << bundle->DebugString(); } else { - RAY_LOG(INFO) << "Failed to lease resource from " << node_id - << " for bundle: " << bundle->DebugString(); + RAY_LOG(DEBUG) << "Failed to lease resource from " << node_id + << " for bundle: " << bundle->DebugString(); } callback(result); }); @@ -318,17 +323,17 @@ void GcsPlacementGroupScheduler::CommitResources( RAY_CHECK(node.has_value()); const auto lease_client = GetLeaseClientFromNode(node.value()); const auto node_id = NodeID::FromBinary(node.value()->node_id()); - RAY_LOG(INFO) << "Committing resource to a node " << node_id - << " for a bundle: " << bundle->DebugString(); + RAY_LOG(DEBUG) << "Committing resource to a node " << node_id + << " for a bundle: " << bundle->DebugString(); lease_client->CommitBundleResources( *bundle, [bundle, node_id, callback](const Status &status, const rpc::CommitBundleResourcesReply &reply) { if (status.ok()) { - RAY_LOG(INFO) << "Finished committing resource to " << node_id - << " for bundle: " << bundle->DebugString(); + RAY_LOG(DEBUG) << "Finished committing resource to " << node_id + << " for bundle: " << bundle->DebugString(); } else { - RAY_LOG(INFO) << "Failed to commit resource to " << node_id - << " for bundle: " << bundle->DebugString(); + RAY_LOG(DEBUG) << "Failed to commit resource to " << node_id + << " for bundle: " << bundle->DebugString(); } RAY_CHECK(callback); callback(status); @@ -345,14 +350,14 @@ void GcsPlacementGroupScheduler::CancelResourceReserve( return; } auto node_id = NodeID::FromBinary(node.value()->node_id()); - RAY_LOG(INFO) << "Cancelling the resource reserved for bundle: " - << bundle_spec->DebugString() << " at node " << node_id; + RAY_LOG(DEBUG) << "Cancelling the resource reserved for bundle: " + << bundle_spec->DebugString() << " at node " << node_id; const auto return_client = GetLeaseClientFromNode(node.value()); return_client->CancelResourceReserve( *bundle_spec, [bundle_spec, node_id](const Status &status, const rpc::CancelResourceReserveReply &reply) { - RAY_LOG(INFO) << "Finished cancelling the resource reserved for bundle: " - << bundle_spec->DebugString() << " at node " << node_id; + RAY_LOG(DEBUG) << "Finished cancelling the resource reserved for bundle: " + << bundle_spec->DebugString() << " at node " << node_id; }); } @@ -480,22 +485,26 @@ void GcsPlacementGroupScheduler::OnAllBundleCommitRequestReturned( committed_bundle_location_index_.AddBundleLocations(placement_group_id, prepared_bundle_locations); - if (!lease_status_tracker->AllCommitRequestsSuccessful()) { - if (lease_status_tracker->GetLeasingState() == LeasingState::CANCELLED) { - DestroyPlacementGroupBundleResourcesIfExists(placement_group_id); - } else { - // Update the state to be reschedule so that the failure handle will reschedule the - // failed bundles. - const auto &uncommitted_bundle_locations = - lease_status_tracker->GetUnCommittedBundleLocations(); - for (const auto &bundle : *uncommitted_bundle_locations) { - placement_group->GetMutableBundle(bundle.first.second)->clear_node_id(); - } - placement_group->UpdateState(rpc::PlacementGroupTableData::RESCHEDULING); - } + // If the placement group scheduling has been cancelled, destroy them. + if (lease_status_tracker->GetLeasingState() == LeasingState::CANCELLED) { + DestroyPlacementGroupBundleResourcesIfExists(placement_group_id); schedule_failure_handler(placement_group); return; } + + if (!lease_status_tracker->AllCommitRequestsSuccessful()) { + // Update the state to be reschedule so that the failure handle will reschedule the + // failed bundles. + const auto &uncommitted_bundle_locations = + lease_status_tracker->GetUnCommittedBundleLocations(); + for (const auto &bundle : *uncommitted_bundle_locations) { + placement_group->GetMutableBundle(bundle.first.second)->clear_node_id(); + } + placement_group->UpdateState(rpc::PlacementGroupTableData::RESCHEDULING); + schedule_failure_handler(placement_group); + return; + } + schedule_success_handler(placement_group); } diff --git a/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.h b/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.h index 0eb4d7fc7..5baadc6d4 100644 --- a/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.h +++ b/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.h @@ -475,7 +475,6 @@ class GcsPlacementGroupScheduler : public GcsPlacementGroupSchedulerInterface { BundleLocationIndex committed_bundle_location_index_; /// Set of placement group that have lease requests in flight to nodes. - /// It is required to know if placement group has been removed or not. absl::flat_hash_map> placement_group_leasing_in_progress_; }; 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 8cd16ca04..765430c04 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 @@ -59,6 +59,15 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { EXPECT_TRUE(WaitForCondition(condition, timeout_ms_.count())); } + template + void WaitPendingDone(const std::list &data, int expected_count) { + auto condition = [this, &data, expected_count]() { + absl::MutexLock lock(&vector_mutex_); + return (int)data.size() == expected_count; + }; + EXPECT_TRUE(WaitForCondition(condition, timeout_ms_.count())); + } + void AddNode(const std::shared_ptr &node, int cpu_num = 10) { gcs_node_manager_->AddNode(node); rpc::HeartbeatTableData heartbeat; @@ -115,6 +124,9 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { ASSERT_EQ(2, raylet_clients_[0]->lease_callbacks.size()); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); WaitPendingDone(failure_placement_groups_, 0); WaitPendingDone(success_placement_groups_, 1); ASSERT_EQ(placement_group, success_placement_groups_.front()); @@ -147,6 +159,9 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test { success_handler); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); WaitPendingDone(success_placement_groups_, 1); } @@ -314,6 +329,9 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyBalancedScheduling) node_commit_count[node_index] += 2; ASSERT_TRUE(raylet_clients_[node_index]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[node_index]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[node_index]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[node_index]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[node_index]->GrantCommitBundleResources()); auto condition = [this, node_index, node_commit_count]() { return raylet_clients_[node_index]->num_commit_requested == node_commit_count[node_index]; @@ -346,6 +364,9 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyResourceCheck) { scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); WaitPendingDone(success_placement_groups_, 1); // Node1 has less number of bundles, but it doesn't satisfy the resource @@ -359,6 +380,9 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictPackStrategyResourceCheck) { scheduler_->ScheduleUnplacedBundles(placement_group2, failure_handler, success_handler); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); WaitPendingDone(success_placement_groups_, 2); } @@ -385,6 +409,9 @@ TEST_F(GcsPlacementGroupSchedulerTest, DestroyPlacementGroup) { }); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); WaitPendingDone(failure_placement_groups_, 0); WaitPendingDone(success_placement_groups_, 1); const auto &placement_group_id = placement_group->GetPlacementGroupID(); @@ -429,6 +456,41 @@ TEST_F(GcsPlacementGroupSchedulerTest, DestroyCancelledPlacementGroup) { WaitPendingDone(failure_placement_groups_, 1); } +TEST_F(GcsPlacementGroupSchedulerTest, PlacementGroupCancelledDuringCommit) { + auto node = Mocker::GenNodeInfo(); + AddNode(node); + ASSERT_EQ(1, gcs_node_manager_->GetAllAliveNodes().size()); + + auto create_placement_group_request = Mocker::GenCreatePlacementGroupRequest(); + auto placement_group = + std::make_shared(create_placement_group_request); + const auto &placement_group_id = placement_group->GetPlacementGroupID(); + + // Schedule the placement_group with 1 available node, and the lease request should be + // send to the node. + scheduler_->ScheduleUnplacedBundles( + placement_group, + [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + failure_placement_groups_.emplace_back(std::move(placement_group)); + }, + [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + success_placement_groups_.emplace_back(std::move(placement_group)); + }); + + // Now, cancel the schedule request. + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + scheduler_->MarkScheduleCancelled(placement_group_id); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 2); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve()); + ASSERT_TRUE(raylet_clients_[0]->GrantCancelResourceReserve()); + WaitPendingDone(failure_placement_groups_, 1); +} + TEST_F(GcsPlacementGroupSchedulerTest, TestPackStrategyReschedulingWhenNodeAdd) { ReschedulingWhenNodeAddTest(rpc::PlacementStrategy::PACK); } @@ -459,6 +521,17 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestPackStrategyLargeBundlesScheduling) { for (int index = 0; index < raylet_clients_[1]->num_lease_requested; ++index) { ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); } + // Wait until all resources are prepared. + WaitPendingDone(raylet_clients_[0]->commit_callbacks, + raylet_clients_[0]->num_lease_requested); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, + raylet_clients_[1]->num_lease_requested); + for (int index = 0; index < raylet_clients_[0]->num_commit_requested; ++index) { + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + } + for (int index = 0; index < raylet_clients_[1]->num_commit_requested; ++index) { + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources()); + } WaitPendingDone(success_placement_groups_, 1); } @@ -486,6 +559,10 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestRescheduleWhenNodeDead) { scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 1); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources()); WaitPendingDone(success_placement_groups_, 1); auto bundles_on_node0 = @@ -501,8 +578,17 @@ 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. + + auto commit_ready = [this]() { + absl::MutexLock lock(&vector_mutex_); + return raylet_clients_[0]->commit_callbacks.size() == 1 || + raylet_clients_[1]->commit_callbacks.size() == 1; + }; raylet_clients_[0]->GrantPrepareBundleResources(); raylet_clients_[1]->GrantPrepareBundleResources(); + EXPECT_TRUE(WaitForCondition(commit_ready, timeout_ms_.count())); + raylet_clients_[0]->GrantCommitBundleResources(); + raylet_clients_[1]->GrantCommitBundleResources(); WaitPendingDone(success_placement_groups_, 2); } @@ -537,6 +623,10 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestStrictSpreadStrategyResourceCheck) { scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); ASSERT_TRUE(raylet_clients_[2]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[2]->commit_callbacks, 1); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[2]->GrantCommitBundleResources()); WaitPendingDone(success_placement_groups_, 1); } @@ -616,7 +706,43 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestBundleLocationIndex) { ASSERT_TRUE(bundle_location_index.GetBundleLocationsOnNode(node2).value()->size() == 1); } -TEST_F(GcsPlacementGroupSchedulerTest, TestNodeDeadDuringCommitResources) { +TEST_F(GcsPlacementGroupSchedulerTest, TestNodeDeadDuringPreparingResources) { + auto node0 = Mocker::GenNodeInfo(0); + auto node1 = Mocker::GenNodeInfo(1); + AddNode(node0); + AddNode(node1); + ASSERT_EQ(2, gcs_node_manager_->GetAllAliveNodes().size()); + + auto create_placement_group_request = Mocker::GenCreatePlacementGroupRequest(); + auto placement_group = + std::make_shared(create_placement_group_request); + + // Schedule the placement group. + // One node is dead, so one bundle failed to schedule. + auto failure_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + ASSERT_TRUE(placement_group->GetUnplacedBundles().size() == 2); + failure_placement_groups_.emplace_back(std::move(placement_group)); + }; + auto success_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + success_placement_groups_.emplace_back(std::move(placement_group)); + }; + + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + gcs_node_manager_->RemoveNode(NodeID::FromBinary(node1->node_id())); + // This should fail because the node is dead. + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources(false)); + ASSERT_TRUE(raylet_clients_[0]->commit_callbacks.size() == 0); + ASSERT_TRUE(raylet_clients_[1]->commit_callbacks.size() == 0); + WaitPendingDone(failure_placement_groups_, 1); +} + +TEST_F(GcsPlacementGroupSchedulerTest, + TestNodeDeadDuringPreparingResourcesRaceCondition) { + // This covers the scnario where the node is dead right after raylet sends a success + // response. auto node0 = Mocker::GenNodeInfo(0); auto node1 = Mocker::GenNodeInfo(1); AddNode(node0); @@ -642,7 +768,216 @@ TEST_F(GcsPlacementGroupSchedulerTest, TestNodeDeadDuringCommitResources) { scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); gcs_node_manager_->RemoveNode(NodeID::FromBinary(node1->node_id())); + // If node is dead right after raylet succeds to create a bundle, it will reply that + // the request has been succeed. In this case, we should just treating like a committed + // bundle that is just removed. ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + // This won't send a commit request because a node is dead. + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 0); + ASSERT_FALSE(raylet_clients_[1]->GrantCommitBundleResources()); + // In this case, we treated the placement group creation successful. Instead, + // we will reschedule them. + WaitPendingDone(failure_placement_groups_, 1); +} + +TEST_F(GcsPlacementGroupSchedulerTest, TestNodeDeadDuringCommittingResources) { + auto node0 = Mocker::GenNodeInfo(0); + auto node1 = Mocker::GenNodeInfo(1); + AddNode(node0); + AddNode(node1); + ASSERT_EQ(2, gcs_node_manager_->GetAllAliveNodes().size()); + + auto create_placement_group_request = Mocker::GenCreatePlacementGroupRequest(); + auto placement_group = + std::make_shared(create_placement_group_request); + + // Schedule the placement group. + // One node is dead, so one bundle failed to schedule. + auto failure_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + ASSERT_TRUE(placement_group->GetUnplacedBundles().size() == 2); + failure_placement_groups_.emplace_back(std::move(placement_group)); + }; + auto success_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + success_placement_groups_.emplace_back(std::move(placement_group)); + }; + + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 1); + gcs_node_manager_->RemoveNode(NodeID::FromBinary(node1->node_id())); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + // Commit will fail because the node is dead. + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources(false)); + WaitPendingDone(success_placement_groups_, 1); +} + +TEST_F(GcsPlacementGroupSchedulerTest, TestNodeDeadDuringRescheduling) { + auto node0 = Mocker::GenNodeInfo(0); + auto node1 = Mocker::GenNodeInfo(1); + AddNode(node0); + AddNode(node1); + ASSERT_EQ(2, gcs_node_manager_->GetAllAliveNodes().size()); + + auto create_placement_group_request = Mocker::GenCreatePlacementGroupRequest(); + auto placement_group = + std::make_shared(create_placement_group_request); + + // Schedule the placement group successfully. + auto failure_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + failure_placement_groups_.emplace_back(std::move(placement_group)); + }; + auto success_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + success_placement_groups_.emplace_back(std::move(placement_group)); + }; + + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 1); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources()); + WaitPendingDone(success_placement_groups_, 1); + + auto bundles_on_node0 = + scheduler_->GetBundlesOnNode(NodeID::FromBinary(node0->node_id())); + ASSERT_EQ(1, bundles_on_node0.size()); + auto bundles_on_node1 = + scheduler_->GetBundlesOnNode(NodeID::FromBinary(node1->node_id())); + ASSERT_EQ(1, bundles_on_node1.size()); + // all nodes are dead, reschedule the placement group. + placement_group->GetMutableBundle(0)->clear_node_id(); + placement_group->GetMutableBundle(1)->clear_node_id(); + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + // Before prepare requests are done, suppose a node is dead. + gcs_node_manager_->RemoveNode(NodeID::FromBinary(node1->node_id())); + // This should fail since the node is dead. + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources(false)); + // Make sure the commit requests are not sent. + ASSERT_TRUE(raylet_clients_[0]->commit_callbacks.size() == 0); + ASSERT_TRUE(raylet_clients_[1]->commit_callbacks.size() == 0); + + WaitPendingDone(success_placement_groups_, 1); + WaitPendingDone(failure_placement_groups_, 1); +} + +TEST_F(GcsPlacementGroupSchedulerTest, TestPGCancelledDuringReschedulingCommit) { + auto node0 = Mocker::GenNodeInfo(0); + auto node1 = Mocker::GenNodeInfo(1); + AddNode(node0); + AddNode(node1); + ASSERT_EQ(2, gcs_node_manager_->GetAllAliveNodes().size()); + + auto create_placement_group_request = Mocker::GenCreatePlacementGroupRequest(); + auto placement_group = + std::make_shared(create_placement_group_request); + + // Schedule the placement group successfully. + auto failure_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + failure_placement_groups_.emplace_back(std::move(placement_group)); + }; + auto success_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + success_placement_groups_.emplace_back(std::move(placement_group)); + }; + + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 1); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources()); + WaitPendingDone(success_placement_groups_, 1); + + auto bundles_on_node0 = + scheduler_->GetBundlesOnNode(NodeID::FromBinary(node0->node_id())); + ASSERT_EQ(1, bundles_on_node0.size()); + auto bundles_on_node1 = + scheduler_->GetBundlesOnNode(NodeID::FromBinary(node1->node_id())); + ASSERT_EQ(1, bundles_on_node1.size()); + // all nodes are dead, reschedule the placement group. + placement_group->GetMutableBundle(0)->clear_node_id(); + placement_group->GetMutableBundle(1)->clear_node_id(); + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + + // Rescheduling happening. + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + // Make sure the commit requests are not sent. + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 1); + // Cancel the placement group scheduling before commit requests are granted. + scheduler_->MarkScheduleCancelled(placement_group->GetPlacementGroupID()); + // After commits are granted the placement group will be removed. + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources()); + WaitPendingDone(success_placement_groups_, 1); + WaitPendingDone(failure_placement_groups_, 1); +} + +TEST_F(GcsPlacementGroupSchedulerTest, TestPGCancelledDuringReschedulingCommitPrepare) { + auto node0 = Mocker::GenNodeInfo(0); + auto node1 = Mocker::GenNodeInfo(1); + AddNode(node0); + AddNode(node1); + ASSERT_EQ(2, gcs_node_manager_->GetAllAliveNodes().size()); + + auto create_placement_group_request = Mocker::GenCreatePlacementGroupRequest(); + auto placement_group = + std::make_shared(create_placement_group_request); + + // Schedule the placement group successfully. + auto failure_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + failure_placement_groups_.emplace_back(std::move(placement_group)); + }; + auto success_handler = [this](std::shared_ptr placement_group) { + absl::MutexLock lock(&vector_mutex_); + success_placement_groups_.emplace_back(std::move(placement_group)); + }; + + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 1); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 1); + ASSERT_TRUE(raylet_clients_[0]->GrantCommitBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantCommitBundleResources()); + WaitPendingDone(success_placement_groups_, 1); + + auto bundles_on_node0 = + scheduler_->GetBundlesOnNode(NodeID::FromBinary(node0->node_id())); + ASSERT_EQ(1, bundles_on_node0.size()); + auto bundles_on_node1 = + scheduler_->GetBundlesOnNode(NodeID::FromBinary(node1->node_id())); + ASSERT_EQ(1, bundles_on_node1.size()); + // All nodes are dead, reschedule the placement group. + placement_group->GetMutableBundle(0)->clear_node_id(); + placement_group->GetMutableBundle(1)->clear_node_id(); + scheduler_->ScheduleUnplacedBundles(placement_group, failure_handler, success_handler); + + // Rescheduling happening. + // Cancel the placement group scheduling before prepare requests are granted. + scheduler_->MarkScheduleCancelled(placement_group->GetPlacementGroupID()); + ASSERT_TRUE(raylet_clients_[0]->GrantPrepareBundleResources()); + ASSERT_TRUE(raylet_clients_[1]->GrantPrepareBundleResources()); + // Make sure the commit requests are not sent. + WaitPendingDone(raylet_clients_[0]->commit_callbacks, 0); + WaitPendingDone(raylet_clients_[1]->commit_callbacks, 0); + // Make sure the placement group creation has failed. + WaitPendingDone(success_placement_groups_, 1); WaitPendingDone(failure_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 67420ee03..6ef863baf 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 @@ -172,8 +172,7 @@ struct GcsServerMocker { const ray::rpc::ClientCallback &callback) override { num_commit_requested += 1; - rpc::CommitBundleResourcesReply reply; - callback(Status::OK(), reply); + commit_callbacks.push_back(callback); } void CancelResourceReserve( @@ -184,7 +183,7 @@ struct GcsServerMocker { return_callbacks.push_back(callback); } - // Trigger reply to RequestWorkerLease. + // Trigger reply to PrepareBundleResources. bool GrantPrepareBundleResources(bool success = true) { Status status = Status::OK(); rpc::PrepareBundleResourcesReply reply; @@ -199,6 +198,20 @@ struct GcsServerMocker { } } + // Trigger reply to CommitBundleResources. + bool GrantCommitBundleResources(bool success = true) { + Status status = Status::OK(); + rpc::CommitBundleResourcesReply reply; + if (commit_callbacks.size() == 0) { + return false; + } else { + auto callback = commit_callbacks.front(); + callback(status, reply); + commit_callbacks.pop_front(); + return true; + } + } + // Trigger reply to CancelResourceReserve. bool GrantCancelResourceReserve(bool success = true) { Status status = Status::OK(); @@ -220,6 +233,7 @@ struct GcsServerMocker { int num_commit_requested = 0; NodeID node_id = NodeID::FromRandom(); std::list> lease_callbacks = {}; + std::list> commit_callbacks = {}; std::list> return_callbacks = {}; }; class MockedGcsActorScheduler : public gcs::GcsActorScheduler {