From 8391f66086882b505e8f5805cedfc17a456d11cf Mon Sep 17 00:00:00 2001 From: fangfengbin <869218239a@zju.edu.cn> Date: Tue, 7 Jul 2020 21:12:30 +0800 Subject: [PATCH] Fix gcs actor manager destroy actor crash bug (#9329) --- src/ray/gcs/gcs_server/gcs_actor_manager.cc | 17 ++- .../gcs_server/test/gcs_actor_manager_test.cc | 105 +++++++++--------- 2 files changed, 66 insertions(+), 56 deletions(-) diff --git a/src/ray/gcs/gcs_server/gcs_actor_manager.cc b/src/ray/gcs/gcs_server/gcs_actor_manager.cc index f18563f25..06034069b 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_actor_manager.cc @@ -525,8 +525,11 @@ void GcsActorManager::DestroyActor(const ActorID &actor_id) { [actor_id](const std::shared_ptr &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 actor) { void GcsActorManager::OnActorCreationSuccess(const std::shared_ptr &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. diff --git a/src/ray/gcs/gcs_server/test/gcs_actor_manager_test.cc b/src/ray/gcs/gcs_server/test/gcs_actor_manager_test.cc index 68f9cd6de..b542b414e 100644 --- a/src/ray/gcs/gcs_server/test/gcs_actor_manager_test.cc +++ b/src/ray/gcs/gcs_server/test/gcs_actor_manager_test.cc @@ -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 thread_io_service_; std::shared_ptr 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> finished_actors; + RAY_CHECK_OK(gcs_actor_manager_->RegisterActor( + create_actor_request, [&finished_actors](std::shared_ptr 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) {