mirror of
https://github.com/wassname/ray.git
synced 2026-07-27 11:26:41 +08:00
[GCS]Fix actor idempotent bug (#10856)
This commit is contained in:
@@ -140,6 +140,12 @@ TEST_F(DirectActorSubmitterTest, TestSubmitTask) {
|
||||
ASSERT_TRUE(worker_client_->ReplyPushTask());
|
||||
}
|
||||
ASSERT_THAT(worker_client_->received_seq_nos, ElementsAre(0, 1));
|
||||
|
||||
// Connect to the actor again.
|
||||
// Because the IP and port of address are not modified, it will skip directly and will
|
||||
// not reset `received_seq_nos`.
|
||||
submitter_.ConnectActor(actor_id, addr, 0);
|
||||
ASSERT_THAT(worker_client_->received_seq_nos, ElementsAre(0, 1));
|
||||
}
|
||||
|
||||
TEST_F(DirectActorSubmitterTest, TestDependencies) {
|
||||
@@ -251,6 +257,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartNoRetry) {
|
||||
ActorID actor_id = ActorID::Of(JobID::FromInt(0), TaskID::Nil(), 0);
|
||||
submitter_.AddActorQueueIfNotExists(actor_id);
|
||||
gcs::ActorTableData actor_data;
|
||||
addr.set_port(0);
|
||||
submitter_.ConnectActor(actor_id, addr, 0);
|
||||
ASSERT_EQ(worker_client_->callbacks.size(), 0);
|
||||
|
||||
@@ -277,6 +284,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartNoRetry) {
|
||||
ASSERT_TRUE(worker_client_->ReplyPushTask(Status::IOError("")));
|
||||
|
||||
// Actor gets restarted.
|
||||
addr.set_port(1);
|
||||
submitter_.ConnectActor(actor_id, addr, 1);
|
||||
ASSERT_TRUE(submitter_.SubmitTask(task4).ok());
|
||||
ASSERT_TRUE(worker_client_->ReplyPushTask(Status::OK()));
|
||||
@@ -292,6 +300,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartRetry) {
|
||||
ActorID actor_id = ActorID::Of(JobID::FromInt(0), TaskID::Nil(), 0);
|
||||
submitter_.AddActorQueueIfNotExists(actor_id);
|
||||
gcs::ActorTableData actor_data;
|
||||
addr.set_port(0);
|
||||
submitter_.ConnectActor(actor_id, addr, 0);
|
||||
ASSERT_EQ(worker_client_->callbacks.size(), 0);
|
||||
|
||||
@@ -321,6 +330,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartRetry) {
|
||||
ASSERT_TRUE(worker_client_->ReplyPushTask(Status::IOError("")));
|
||||
|
||||
// Actor gets restarted.
|
||||
addr.set_port(1);
|
||||
submitter_.ConnectActor(actor_id, addr, 1);
|
||||
// A new task is submitted.
|
||||
ASSERT_TRUE(submitter_.SubmitTask(task4).ok());
|
||||
@@ -342,6 +352,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) {
|
||||
ActorID actor_id = ActorID::Of(JobID::FromInt(0), TaskID::Nil(), 0);
|
||||
submitter_.AddActorQueueIfNotExists(actor_id);
|
||||
gcs::ActorTableData actor_data;
|
||||
addr.set_port(0);
|
||||
submitter_.ConnectActor(actor_id, addr, 0);
|
||||
ASSERT_EQ(worker_client_->callbacks.size(), 0);
|
||||
ASSERT_EQ(num_clients_connected_, 1);
|
||||
@@ -354,6 +365,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) {
|
||||
ASSERT_TRUE(worker_client_->ReplyPushTask(Status::OK()));
|
||||
|
||||
// Actor restarts, but we don't receive the disconnect message until later.
|
||||
addr.set_port(1);
|
||||
submitter_.ConnectActor(actor_id, addr, 1);
|
||||
ASSERT_EQ(num_clients_connected_, 2);
|
||||
// Submit a task.
|
||||
@@ -381,6 +393,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) {
|
||||
ASSERT_FALSE(worker_client_->ReplyPushTask(Status::OK()));
|
||||
|
||||
// We receive the late messages. Nothing happens.
|
||||
addr.set_port(2);
|
||||
submitter_.ConnectActor(actor_id, addr, 2);
|
||||
submitter_.DisconnectActor(actor_id, 2, /*dead=*/false);
|
||||
ASSERT_EQ(num_clients_connected_, 2);
|
||||
@@ -392,6 +405,7 @@ TEST_F(DirectActorSubmitterTest, TestActorRestartOutOfOrderGcs) {
|
||||
|
||||
// We receive more late messages. Nothing happens because the actor is dead.
|
||||
submitter_.DisconnectActor(actor_id, 4, /*dead=*/false);
|
||||
addr.set_port(3);
|
||||
submitter_.ConnectActor(actor_id, addr, 4);
|
||||
ASSERT_EQ(num_clients_connected_, 2);
|
||||
// Submit a task.
|
||||
|
||||
@@ -135,6 +135,13 @@ void CoreWorkerDirectActorTaskSubmitter::ConnectActor(const ActorID &actor_id,
|
||||
return;
|
||||
}
|
||||
|
||||
if (queue->second.rpc_client &&
|
||||
queue->second.rpc_client->Addr().ip_address() == address.ip_address() &&
|
||||
queue->second.rpc_client->Addr().port() == address.port()) {
|
||||
RAY_LOG(DEBUG) << "Skip actor that has already been connected, actor_id=" << actor_id;
|
||||
return;
|
||||
}
|
||||
|
||||
if (queue->second.state == rpc::ActorTableData::DEAD) {
|
||||
// This message is about an old version of the actor and the actor has
|
||||
// already died since then. Skip the connection.
|
||||
|
||||
Reference in New Issue
Block a user