[Placement Group] Fix placement group bugs that happen when rescheduling. (#11263)

* Fix placement group bugs while autoscaling.

* Addressed code review.
This commit is contained in:
SangBin Cho
2020-10-08 08:58:59 -07:00
committed by GitHub
parent f9621ce23c
commit 37fa86f9a0
5 changed files with 406 additions and 48 deletions
@@ -125,7 +125,7 @@ PlacementGroupID GcsPlacementGroupManager::GetPlacementGroupIDByName(
void GcsPlacementGroupManager::OnPlacementGroupCreationFailed(
std::shared_ptr<GcsPlacementGroup> 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<GcsPlacementGroup> &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(
@@ -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<BundleLocations> bundle_locations = std::make_shared<BundleLocations>();
std::shared_ptr<BundleLocations> committed_bundle_locations =
std::make_shared<BundleLocations>();
std::shared_ptr<BundleLocations> leasing_bundle_locations =
std::make_shared<BundleLocations>();
// 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);
}
@@ -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<PlacementGroupID, std::shared_ptr<LeaseStatusTracker>>
placement_group_leasing_in_progress_;
};
@@ -59,6 +59,15 @@ class GcsPlacementGroupSchedulerTest : public ::testing::Test {
EXPECT_TRUE(WaitForCondition(condition, timeout_ms_.count()));
}
template <typename Data>
void WaitPendingDone(const std::list<Data> &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<rpc::GcsNodeInfo> &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<gcs::GcsPlacementGroup>(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<gcs::GcsPlacementGroup> placement_group) {
absl::MutexLock lock(&vector_mutex_);
failure_placement_groups_.emplace_back(std::move(placement_group));
},
[this](std::shared_ptr<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup>(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<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup>(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<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup>(create_placement_group_request);
// Schedule the placement group successfully.
auto failure_handler = [this](std::shared_ptr<gcs::GcsPlacementGroup> placement_group) {
absl::MutexLock lock(&vector_mutex_);
failure_placement_groups_.emplace_back(std::move(placement_group));
};
auto success_handler = [this](std::shared_ptr<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup>(create_placement_group_request);
// Schedule the placement group successfully.
auto failure_handler = [this](std::shared_ptr<gcs::GcsPlacementGroup> placement_group) {
absl::MutexLock lock(&vector_mutex_);
failure_placement_groups_.emplace_back(std::move(placement_group));
};
auto success_handler = [this](std::shared_ptr<gcs::GcsPlacementGroup> 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<gcs::GcsPlacementGroup>(create_placement_group_request);
// Schedule the placement group successfully.
auto failure_handler = [this](std::shared_ptr<gcs::GcsPlacementGroup> placement_group) {
absl::MutexLock lock(&vector_mutex_);
failure_placement_groups_.emplace_back(std::move(placement_group));
};
auto success_handler = [this](std::shared_ptr<gcs::GcsPlacementGroup> 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);
}
@@ -172,8 +172,7 @@ struct GcsServerMocker {
const ray::rpc::ClientCallback<ray::rpc::CommitBundleResourcesReply> &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<rpc::ClientCallback<rpc::PrepareBundleResourcesReply>> lease_callbacks = {};
std::list<rpc::ClientCallback<rpc::CommitBundleResourcesReply>> commit_callbacks = {};
std::list<rpc::ClientCallback<rpc::CancelResourceReserveReply>> return_callbacks = {};
};
class MockedGcsActorScheduler : public gcs::GcsActorScheduler {