Fix gcs actor manager destroy actor crash bug (#9329)

This commit is contained in:
fangfengbin
2020-07-07 21:12:30 +08:00
committed by GitHub
parent 079c1eaa5c
commit 8391f66086
2 changed files with 66 additions and 56 deletions
+14 -3
View File
@@ -525,8 +525,11 @@ void GcsActorManager::DestroyActor(const ActorID &actor_id) {
[actor_id](const std::shared_ptr<GcsActor> &actor) {
return actor->GetActorID() == actor_id;
});
RAY_CHECK(pending_it != pending_actors_.end());
pending_actors_.erase(pending_it);
// TODO(rkooo567): The actor may be in the state of leasing worker. We need to add
// the processing logic of this in https://github.com/ray-project/ray/pull/9215.
if (pending_it != pending_actors_.end()) {
pending_actors_.erase(pending_it);
}
}
}
@@ -695,7 +698,15 @@ void GcsActorManager::OnActorCreationFailed(std::shared_ptr<GcsActor> actor) {
void GcsActorManager::OnActorCreationSuccess(const std::shared_ptr<GcsActor> &actor) {
auto actor_id = actor->GetActorID();
RAY_LOG(DEBUG) << "Actor created successfully, actor id = " << actor_id;
RAY_CHECK(registered_actors_.count(actor_id) > 0);
// NOTE: If an actor is deleted immediately after the user creates the actor, reference
// counter may return a reply to the request of WaitForActorOutOfScope to GCS server,
// and GCS server will destroy the actor. The actor creation is asynchronous, it may be
// destroyed before the actor creation is completed.
if (registered_actors_.count(actor_id) == 0) {
RAY_LOG(WARNING) << "Actor is destroyed before the creation is completed, actor id = "
<< actor_id;
return;
}
actor->UpdateState(rpc::ActorTableData::ALIVE);
auto actor_table_data = actor->GetActorTableData();
// The backend storage is reliable in the future, so the status must be ok.
@@ -111,6 +111,15 @@ class GcsActorManagerTest : public ::testing::Test {
EXPECT_TRUE(WaitForCondition(condition, timeout_ms_.count()));
}
rpc::Address RandomAddress() const {
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
return address;
}
boost::asio::io_service io_service_;
std::unique_ptr<std::thread> thread_io_service_;
std::shared_ptr<gcs::StoreClient> store_client_;
@@ -141,12 +150,7 @@ TEST_F(GcsActorManagerTest, TestBasic) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
actor->UpdateAddress(address);
actor->UpdateAddress(RandomAddress());
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
ASSERT_EQ(finished_actors.size(), 1);
@@ -176,12 +180,7 @@ TEST_F(GcsActorManagerTest, TestSchedulingFailed) {
ASSERT_EQ(finished_actors.size(), 0);
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
actor->UpdateAddress(address);
actor->UpdateAddress(RandomAddress());
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
ASSERT_EQ(finished_actors.size(), 1);
@@ -202,11 +201,9 @@ TEST_F(GcsActorManagerTest, TestWorkerFailure) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
auto worker_id = WorkerID::FromBinary(address.worker_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -243,11 +240,8 @@ TEST_F(GcsActorManagerTest, TestNodeFailure) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -286,11 +280,8 @@ TEST_F(GcsActorManagerTest, TestActorReconstruction) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -348,11 +339,8 @@ TEST_F(GcsActorManagerTest, TestActorRestartWhenOwnerDead) {
const auto owner_node_id = actor->GetOwnerNodeID();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -392,12 +380,7 @@ TEST_F(GcsActorManagerTest, TestDetachedActorRestartWhenCreatorDead) {
const auto owner_node_id = actor->GetOwnerNodeID();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
actor->UpdateAddress(address);
actor->UpdateAddress(RandomAddress());
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
ASSERT_EQ(finished_actors.size(), 1);
@@ -468,11 +451,9 @@ TEST_F(GcsActorManagerTest, TestNamedActorDeletionWorkerFailure) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
auto worker_id = WorkerID::FromBinary(address.worker_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -508,11 +489,8 @@ TEST_F(GcsActorManagerTest, TestNamedActorDeletionNodeFailure) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -549,11 +527,9 @@ TEST_F(GcsActorManagerTest, TestNamedActorDeletionNotHappendWhenReconstructed) {
mock_actor_scheduler_->actors.pop_back();
// Check that the actor is in state `ALIVE`.
rpc::Address address;
auto node_id = ClientID::FromRandom();
auto worker_id = WorkerID::FromRandom();
address.set_raylet_id(node_id.Binary());
address.set_worker_id(worker_id.Binary());
auto address = RandomAddress();
auto node_id = ClientID::FromBinary(address.raylet_id());
auto worker_id = WorkerID::FromBinary(address.worker_id());
actor->UpdateAddress(address);
gcs_actor_manager_->OnActorCreationSuccess(actor);
WaitActorCreated(actor->GetActorID());
@@ -576,6 +552,29 @@ TEST_F(GcsActorManagerTest, TestNamedActorDeletionNotHappendWhenReconstructed) {
request1.task_spec().actor_creation_task_spec().actor_id());
}
TEST_F(GcsActorManagerTest, TestDestroyActorBeforeActorCreationCompletes) {
auto job_id = JobID::FromInt(1);
auto create_actor_request = Mocker::GenCreateActorRequest(job_id);
std::vector<std::shared_ptr<gcs::GcsActor>> finished_actors;
RAY_CHECK_OK(gcs_actor_manager_->RegisterActor(
create_actor_request, [&finished_actors](std::shared_ptr<gcs::GcsActor> actor) {
finished_actors.emplace_back(actor);
}));
ASSERT_EQ(finished_actors.size(), 0);
ASSERT_EQ(mock_actor_scheduler_->actors.size(), 1);
auto actor = mock_actor_scheduler_->actors.back();
mock_actor_scheduler_->actors.clear();
// Simulate the reply of WaitForActorOutOfScope request to trigger actor destruction.
ASSERT_TRUE(worker_client_->Reply());
// Check that the actor is in state `DEAD`.
actor->UpdateAddress(RandomAddress());
gcs_actor_manager_->OnActorCreationSuccess(actor);
ASSERT_EQ(actor->GetState(), rpc::ActorTableData::DEAD);
}
} // namespace ray
int main(int argc, char **argv) {