From 4f2e4f31dda0d11b01e761f282114a890a6a622f Mon Sep 17 00:00:00 2001 From: Zhijun Fu <37800433+zhijunfu@users.noreply.github.com> Date: Tue, 4 Aug 2020 03:44:02 +0800 Subject: [PATCH] async grpc calls should always return void (#9533) --- src/ray/core_worker/core_worker.cc | 17 +-- src/ray/core_worker/future_resolver.cc | 5 +- .../core_worker/object_recovery_manager.cc | 32 ++--- src/ray/core_worker/reference_count.cc | 4 +- src/ray/core_worker/reference_count_test.cc | 3 +- .../test/direct_actor_transport_test.cc | 6 +- .../test/direct_task_transport_test.cc | 24 ++-- .../test/object_recovery_manager_test.cc | 3 +- .../transport/direct_actor_transport.cc | 6 +- .../transport/direct_task_transport.cc | 86 +++++------ src/ray/gcs/gcs_server/gcs_actor_manager.cc | 4 +- src/ray/gcs/gcs_server/gcs_actor_scheduler.cc | 20 +-- .../gcs_placement_group_scheduler.cc | 9 +- .../gcs_server/test/gcs_actor_manager_test.cc | 9 +- .../gcs_server/test/gcs_server_test_util.h | 18 +-- src/ray/raylet/node_manager.cc | 60 ++++---- src/ray/raylet/worker.cc | 9 +- src/ray/raylet_client/raylet_client.cc | 31 ++-- src/ray/raylet_client/raylet_client.h | 26 ++-- src/ray/rpc/grpc_client.h | 19 +-- .../rpc/node_manager/node_manager_client.h | 16 +-- src/ray/rpc/worker/core_worker_client.h | 133 +++++++----------- 22 files changed, 225 insertions(+), 315 deletions(-) diff --git a/src/ray/core_worker/core_worker.cc b/src/ray/core_worker/core_worker.cc index 42c174bd9..7b89836d7 100644 --- a/src/ray/core_worker/core_worker.cc +++ b/src/ray/core_worker/core_worker.cc @@ -745,7 +745,7 @@ void CoreWorker::PutObjectIntoPlasma(const RayObject &object, const ObjectID &ob if (!object_exists) { // Tell the raylet to pin the object **after** it is created. RAY_LOG(DEBUG) << "Pinning put object " << object_id; - RAY_CHECK_OK(local_raylet_client_->PinObjectIDs( + local_raylet_client_->PinObjectIDs( rpc_address_, {object_id}, [this, object_id](const Status &status, const rpc::PinObjectIDsReply &reply) { // Only release the object once the raylet has responded to avoid the race @@ -754,7 +754,7 @@ void CoreWorker::PutObjectIntoPlasma(const RayObject &object, const ObjectID &ob RAY_LOG(ERROR) << "Failed to release ObjectID (" << object_id << "), might cause a leak in plasma."; } - })); + }); } RAY_CHECK(memory_store_->Put(RayObject(rpc::ErrorType::OBJECT_IN_PLASMA), object_id)); } @@ -837,7 +837,7 @@ Status CoreWorker::Put(const RayObject &object, if (pin_object) { // Tell the raylet to pin the object **after** it is created. RAY_LOG(DEBUG) << "Pinning put object " << object_id; - RAY_CHECK_OK(local_raylet_client_->PinObjectIDs( + local_raylet_client_->PinObjectIDs( rpc_address_, {object_id}, [this, object_id](const Status &status, const rpc::PinObjectIDsReply &reply) { // Only release the object once the raylet has responded to avoid the race @@ -846,7 +846,7 @@ Status CoreWorker::Put(const RayObject &object, RAY_LOG(ERROR) << "Failed to release ObjectID (" << object_id << "), might cause a leak in plasma."; } - })); + }); } else { RAY_RETURN_NOT_OK(plasma_store_provider_->Release(object_id)); } @@ -895,7 +895,7 @@ Status CoreWorker::Seal(const ObjectID &object_id, bool pin_object, if (pin_object) { // Tell the raylet to pin the object **after** it is created. RAY_LOG(DEBUG) << "Pinning sealed object " << object_id; - RAY_CHECK_OK(local_raylet_client_->PinObjectIDs( + local_raylet_client_->PinObjectIDs( owner_address.has_value() ? *owner_address : rpc_address_, {object_id}, [this, object_id](const Status &status, const rpc::PinObjectIDsReply &reply) { // Only release the object once the raylet has responded to avoid the race @@ -904,7 +904,7 @@ Status CoreWorker::Seal(const ObjectID &object_id, bool pin_object, RAY_LOG(ERROR) << "Failed to release ObjectID (" << object_id << "), might cause a leak in plasma."; } - })); + }); } else { RAY_RETURN_NOT_OK(plasma_store_provider_->Release(object_id)); reference_counter_->FreePlasmaObjects({object_id}); @@ -1119,15 +1119,12 @@ Status CoreWorker::Delete(const std::vector &object_ids, bool local_on } void CoreWorker::TriggerGlobalGC() { - auto status = local_raylet_client_->GlobalGC( + local_raylet_client_->GlobalGC( [](const Status &status, const rpc::GlobalGCReply &reply) { if (!status.ok()) { RAY_LOG(ERROR) << "Failed to send global GC request: " << status.ToString(); } }); - if (!status.ok()) { - RAY_LOG(ERROR) << "Failed to send global GC request: " << status.ToString(); - } } std::string CoreWorker::MemoryUsageString() { diff --git a/src/ray/core_worker/future_resolver.cc b/src/ray/core_worker/future_resolver.cc index 6d1317415..79aa458e3 100644 --- a/src/ray/core_worker/future_resolver.cc +++ b/src/ray/core_worker/future_resolver.cc @@ -35,13 +35,14 @@ void FutureResolver::ResolveFutureAsync(const ObjectID &object_id, rpc::GetObjectStatusRequest request; request.set_object_id(object_id.Binary()); request.set_owner_worker_id(owner_worker_id.Binary()); - RAY_CHECK_OK(it->second->GetObjectStatus( + it->second->GetObjectStatus( request, [this, object_id](const Status &status, const rpc::GetObjectStatusReply &reply) { if (!status.ok()) { RAY_LOG(WARNING) << "Error retrieving the value of object ID " << object_id << " that was deserialized: " << status.ToString(); } + if (!status.ok() || reply.status() == rpc::GetObjectStatusReply::OUT_OF_SCOPE) { // The owner is gone or the owner replied that the object has gone // out of scope (this is an edge case in the distributed ref counting @@ -57,7 +58,7 @@ void FutureResolver::ResolveFutureAsync(const ObjectID &object_id, RAY_UNUSED(in_memory_store_->Put(RayObject(rpc::ErrorType::OBJECT_IN_PLASMA), object_id)); } - })); + }); } } // namespace ray diff --git a/src/ray/core_worker/object_recovery_manager.cc b/src/ray/core_worker/object_recovery_manager.cc index 9f06919e3..92beaa21f 100644 --- a/src/ray/core_worker/object_recovery_manager.cc +++ b/src/ray/core_worker/object_recovery_manager.cc @@ -100,22 +100,22 @@ void ObjectRecoveryManager::PinExistingObjectCopy( client = client_it->second; } - RAY_UNUSED(client->PinObjectIDs( - rpc_address_, {object_id}, - [this, object_id, other_locations, node_id](const Status &status, - const rpc::PinObjectIDsReply &reply) { - if (status.ok()) { - // TODO(swang): Make sure that the node is still alive when - // marking the object as pinned. - RAY_CHECK(in_memory_store_->Put(RayObject(rpc::ErrorType::OBJECT_IN_PLASMA), - object_id)); - reference_counter_->UpdateObjectPinnedAtRaylet(object_id, node_id); - } else { - RAY_LOG(INFO) << "Error pinning new copy of lost object " << object_id - << ", trying again"; - PinOrReconstructObject(object_id, other_locations); - } - })); + client->PinObjectIDs(rpc_address_, {object_id}, + [this, object_id, other_locations, node_id]( + const Status &status, const rpc::PinObjectIDsReply &reply) { + if (status.ok()) { + // TODO(swang): Make sure that the node is still alive when + // marking the object as pinned. + RAY_CHECK(in_memory_store_->Put( + RayObject(rpc::ErrorType::OBJECT_IN_PLASMA), object_id)); + reference_counter_->UpdateObjectPinnedAtRaylet(object_id, + node_id); + } else { + RAY_LOG(INFO) << "Error pinning new copy of lost object " + << object_id << ", trying again"; + PinOrReconstructObject(object_id, other_locations); + } + }); } void ObjectRecoveryManager::ReconstructObject(const ObjectID &object_id) { diff --git a/src/ray/core_worker/reference_count.cc b/src/ray/core_worker/reference_count.cc index cfcd255ca..8be25de67 100644 --- a/src/ray/core_worker/reference_count.cc +++ b/src/ray/core_worker/reference_count.cc @@ -742,7 +742,7 @@ void ReferenceCounter::WaitForRefRemoved(const ReferenceTable::iterator &ref_it, << addr.port << " for object " << object_id; // Send the borrower a message about this object. The borrower responds once // it is no longer using the object ID. - RAY_CHECK_OK(it->second->WaitForRefRemoved( + it->second->WaitForRefRemoved( request, [this, object_id, addr](const Status &status, const rpc::WaitForRefRemovedReply &reply) { RAY_LOG(DEBUG) << "Received reply from borrower " << addr.ip_address << ":" @@ -759,7 +759,7 @@ void ReferenceCounter::WaitForRefRemoved(const ReferenceTable::iterator &ref_it, RAY_CHECK(it != object_id_refs_.end()); RAY_CHECK(it->second.borrowers.erase(addr)); DeleteReferenceInternal(it, nullptr); - })); + }); } void ReferenceCounter::AddNestedObjectIds(const ObjectID &object_id, diff --git a/src/ray/core_worker/reference_count_test.cc b/src/ray/core_worker/reference_count_test.cc index 8f29456a1..ea17419b5 100644 --- a/src/ray/core_worker/reference_count_test.cc +++ b/src/ray/core_worker/reference_count_test.cc @@ -67,7 +67,7 @@ class MockWorkerClient : public rpc::CoreWorkerClientInterface { /*distributed_ref_counting_enabled=*/true, /*lineage_pinning_enabled=*/false, client_factory) {} - ray::Status WaitForRefRemoved( + void WaitForRefRemoved( const rpc::WaitForRefRemovedRequest &request, const rpc::ClientCallback &callback) override { auto r = num_requests_; @@ -93,7 +93,6 @@ class MockWorkerClient : public rpc::CoreWorkerClientInterface { borrower_callbacks_[r] = borrower_callback; num_requests_++; - return Status::OK(); } bool FlushBorrowerCallbacks() { diff --git a/src/ray/core_worker/test/direct_actor_transport_test.cc b/src/ray/core_worker/test/direct_actor_transport_test.cc index b6320b877..f6962b1a8 100644 --- a/src/ray/core_worker/test/direct_actor_transport_test.cc +++ b/src/ray/core_worker/test/direct_actor_transport_test.cc @@ -59,12 +59,10 @@ class MockWorkerClient : public rpc::CoreWorkerClientInterface { public: const rpc::Address &Addr() const override { return addr; } - ray::Status PushActorTask( - std::unique_ptr request, bool skip_queue, - const rpc::ClientCallback &callback) override { + void PushActorTask(std::unique_ptr request, bool skip_queue, + const rpc::ClientCallback &callback) override { received_seq_nos.push_back(request->sequence_number()); callbacks.push_back(callback); - return Status::OK(); } bool ReplyPushTask(Status status = Status::OK()) { diff --git a/src/ray/core_worker/test/direct_task_transport_test.cc b/src/ray/core_worker/test/direct_task_transport_test.cc index 6fef0fabd..80126d872 100644 --- a/src/ray/core_worker/test/direct_task_transport_test.cc +++ b/src/ray/core_worker/test/direct_task_transport_test.cc @@ -31,11 +31,9 @@ int64_t kLongTimeout = 1024 * 1024 * 1024; class MockWorkerClient : public rpc::CoreWorkerClientInterface { public: - ray::Status PushNormalTask( - std::unique_ptr request, - const rpc::ClientCallback &callback) override { + void PushNormalTask(std::unique_ptr request, + const rpc::ClientCallback &callback) override { callbacks.push_back(callback); - return Status::OK(); } bool ReplyPushTask(Status status = Status::OK(), bool exit = false) { @@ -52,11 +50,9 @@ class MockWorkerClient : public rpc::CoreWorkerClientInterface { return true; } - ray::Status CancelTask( - const rpc::CancelTaskRequest &request, - const rpc::ClientCallback &callback) override { + void CancelTask(const rpc::CancelTaskRequest &request, + const rpc::ClientCallback &callback) override { kill_requests.push_front(request); - return Status::OK(); } std::list> callbacks; @@ -104,26 +100,22 @@ class MockRayletClient : public WorkerLeaseInterface { return Status::OK(); } - ray::Status RequestWorkerLease( + void RequestWorkerLease( const ray::TaskSpecification &resource_spec, const rpc::ClientCallback &callback) override { num_workers_requested += 1; callbacks.push_back(callback); - return Status::OK(); } - ray::Status ReleaseUnusedWorkers( + void ReleaseUnusedWorkers( const std::vector &workers_in_use, - const rpc::ClientCallback &callback) override { - return Status::NotImplemented("ReleaseUnusedWorkers is not supported."); - } + const rpc::ClientCallback &callback) override {} - ray::Status CancelWorkerLease( + void CancelWorkerLease( const TaskID &task_id, const rpc::ClientCallback &callback) override { num_leases_canceled += 1; cancel_callbacks.push_back(callback); - return Status::OK(); } // Trigger reply to RequestWorkerLease. diff --git a/src/ray/core_worker/test/object_recovery_manager_test.cc b/src/ray/core_worker/test/object_recovery_manager_test.cc index f0003a3f6..201fddbec 100644 --- a/src/ray/core_worker/test/object_recovery_manager_test.cc +++ b/src/ray/core_worker/test/object_recovery_manager_test.cc @@ -55,12 +55,11 @@ class MockTaskResubmitter : public TaskResubmissionInterface { class MockRayletClient : public PinObjectsInterface { public: - ray::Status PinObjectIDs( + void PinObjectIDs( const rpc::Address &caller_address, const std::vector &object_ids, const ray::rpc::ClientCallback &callback) override { RAY_LOG(INFO) << "PinObjectIDs " << object_ids.size(); callbacks.push_back(callback); - return Status::OK(); } size_t Flush() { diff --git a/src/ray/core_worker/transport/direct_actor_transport.cc b/src/ray/core_worker/transport/direct_actor_transport.cc index d1db242f9..fe3bceef2 100644 --- a/src/ray/core_worker/transport/direct_actor_transport.cc +++ b/src/ray/core_worker/transport/direct_actor_transport.cc @@ -220,7 +220,7 @@ void CoreWorkerDirectActorTaskSubmitter::SendPendingTasks(const ActorID &actor_i if (it->second.pending_force_kill) { RAY_LOG(INFO) << "Sending KillActor request to actor " << actor_id; // It's okay if this fails because this means the worker is already dead. - RAY_UNUSED(it->second.rpc_client->KillActor(*it->second.pending_force_kill, nullptr)); + it->second.rpc_client->KillActor(*it->second.pending_force_kill, nullptr); it->second.pending_force_kill.reset(); } @@ -262,7 +262,7 @@ void CoreWorkerDirectActorTaskSubmitter::PushActorTask(const ClientQueue &queue, << " actor counter " << counter << " seq no " << request->sequence_number(); rpc::Address addr(queue.rpc_client->Addr()); - RAY_UNUSED(queue.rpc_client->PushActorTask( + queue.rpc_client->PushActorTask( std::move(request), skip_queue, [this, addr, task_id, actor_id](Status status, const rpc::PushTaskReply &reply) { bool increment_completed_tasks = true; @@ -282,7 +282,7 @@ void CoreWorkerDirectActorTaskSubmitter::PushActorTask(const ClientQueue &queue, RAY_CHECK(queue != client_queues_.end()); queue->second.num_completed_tasks++; } - })); + }); } bool CoreWorkerDirectActorTaskSubmitter::IsActorAlive(const ActorID &actor_id) const { diff --git a/src/ray/core_worker/transport/direct_task_transport.cc b/src/ray/core_worker/transport/direct_task_transport.cc index 52cad6e7f..469402e24 100644 --- a/src/ray/core_worker/transport/direct_task_transport.cc +++ b/src/ray/core_worker/transport/direct_task_transport.cc @@ -168,7 +168,7 @@ void CoreWorkerDirectTaskSubmitter::CancelWorkerLeaseIfNeeded( auto &lease_client = it->second.first; auto &lease_id = it->second.second; RAY_LOG(DEBUG) << "Canceling lease request " << lease_id; - RAY_UNUSED(lease_client->CancelWorkerLease( + lease_client->CancelWorkerLease( lease_id, [this, scheduling_key](const Status &status, const rpc::CancelWorkerLeaseReply &reply) { absl::MutexLock lock(&mu_); @@ -183,7 +183,7 @@ void CoreWorkerDirectTaskSubmitter::CancelWorkerLeaseIfNeeded( // longer need to cancel. CancelWorkerLeaseIfNeeded(scheduling_key); } - })); + }); } } @@ -227,7 +227,7 @@ void CoreWorkerDirectTaskSubmitter::RequestNewWorkerIfNeeded( TaskSpecification &resource_spec = it->second.front(); TaskID task_id = resource_spec.TaskId(); RAY_LOG(DEBUG) << "Lease requested " << task_id; - RAY_UNUSED(lease_client->RequestWorkerLease( + lease_client->RequestWorkerLease( resource_spec, [this, scheduling_key](const Status &status, const rpc::RequestWorkerLeaseReply &reply) { absl::MutexLock lock(&mu_); @@ -271,7 +271,7 @@ void CoreWorkerDirectTaskSubmitter::RequestNewWorkerIfNeeded( "likely because the local raylet has crahsed."; RAY_LOG(FATAL) << status.ToString(); } - })); + }); RAY_CHECK(pending_lease_requests_ .emplace(scheduling_key, std::make_pair(lease_client, task_id)) .second); @@ -293,43 +293,42 @@ void CoreWorkerDirectTaskSubmitter::PushNormalTask( request->mutable_task_spec()->CopyFrom(task_spec.GetMessage()); request->mutable_resource_mapping()->CopyFrom(assigned_resources); request->set_intended_worker_id(addr.worker_id.Binary()); - RAY_UNUSED(client.PushNormalTask( - std::move(request), - [this, task_id, is_actor, is_actor_creation, scheduling_key, addr, - assigned_resources](Status status, const rpc::PushTaskReply &reply) { - { - absl::MutexLock lock(&mu_); - executing_tasks_.erase(task_id); + client.PushNormalTask(std::move(request), [this, task_id, is_actor, is_actor_creation, + scheduling_key, addr, assigned_resources]( + Status status, + const rpc::PushTaskReply &reply) { + { + absl::MutexLock lock(&mu_); + executing_tasks_.erase(task_id); - // Decrement the number of tasks in flight to the worker - auto &lease_entry = worker_to_lease_entry_[addr]; - RAY_CHECK(lease_entry.tasks_in_flight_ > 0); - lease_entry.tasks_in_flight_--; - } - if (reply.worker_exiting()) { - // The worker is draining and will shutdown after it is done. Don't return - // it to the Raylet since that will kill it early. - absl::MutexLock lock(&mu_); - worker_to_lease_entry_.erase(addr); - } else if (!status.ok() || !is_actor_creation) { - // Successful actor creation leases the worker indefinitely from the raylet. - absl::MutexLock lock(&mu_); - OnWorkerIdle(addr, scheduling_key, - /*error=*/!status.ok(), assigned_resources); - } - if (!status.ok()) { - // TODO: It'd be nice to differentiate here between process vs node - // failure (e.g., by contacting the raylet). If it was a process - // failure, it may have been an application-level error and it may - // not make sense to retry the task. - RAY_UNUSED(task_finisher_->PendingTaskFailed( - task_id, - is_actor ? rpc::ErrorType::ACTOR_DIED : rpc::ErrorType::WORKER_DIED, - &status)); - } else { - task_finisher_->CompletePendingTask(task_id, reply, addr.ToProto()); - } - })); + // Decrement the number of tasks in flight to the worker + auto &lease_entry = worker_to_lease_entry_[addr]; + RAY_CHECK(lease_entry.tasks_in_flight_ > 0); + lease_entry.tasks_in_flight_--; + } + if (reply.worker_exiting()) { + // The worker is draining and will shutdown after it is done. Don't return + // it to the Raylet since that will kill it early. + absl::MutexLock lock(&mu_); + worker_to_lease_entry_.erase(addr); + } else if (!status.ok() || !is_actor_creation) { + // Successful actor creation leases the worker indefinitely from the raylet. + absl::MutexLock lock(&mu_); + OnWorkerIdle(addr, scheduling_key, + /*error=*/!status.ok(), assigned_resources); + } + if (!status.ok()) { + // TODO: It'd be nice to differentiate here between process vs node + // failure (e.g., by contacting the raylet). If it was a process + // failure, it may have been an application-level error and it may + // not make sense to retry the task. + RAY_UNUSED(task_finisher_->PendingTaskFailed( + task_id, is_actor ? rpc::ErrorType::ACTOR_DIED : rpc::ErrorType::WORKER_DIED, + &status)); + } else { + task_finisher_->CompletePendingTask(task_id, reply, addr.ToProto()); + } + }); } Status CoreWorkerDirectTaskSubmitter::CancelTask(TaskSpecification task_spec, @@ -384,7 +383,7 @@ Status CoreWorkerDirectTaskSubmitter::CancelTask(TaskSpecification task_spec, auto request = rpc::CancelTaskRequest(); request.set_intended_task_id(task_spec.TaskId().Binary()); request.set_force_kill(force_kill); - RAY_UNUSED(client->CancelTask( + client->CancelTask( request, [this, task_spec, force_kill](const Status &status, const rpc::CancelTaskReply &reply) { absl::MutexLock lock(&mu_); @@ -402,7 +401,7 @@ Status CoreWorkerDirectTaskSubmitter::CancelTask(TaskSpecification task_spec, } // Retry is not attempted if !status.ok() because force-kill may kill the worker // before the reply is sent. - })); + }); return Status::OK(); } @@ -417,7 +416,8 @@ Status CoreWorkerDirectTaskSubmitter::CancelRemoteTask(const ObjectID &object_id auto request = rpc::RemoteCancelTaskRequest(); request.set_force_kill(force_kill); request.set_remote_object_id(object_id.Binary()); - return client->second->RemoteCancelTask(request, nullptr); + client->second->RemoteCancelTask(request, nullptr); + return Status::OK(); } }; // namespace ray diff --git a/src/ray/gcs/gcs_server/gcs_actor_manager.cc b/src/ray/gcs/gcs_server/gcs_actor_manager.cc index b3e231a7f..d2b15dd82 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_actor_manager.cc @@ -533,7 +533,7 @@ void GcsActorManager::PollOwnerForActorOutOfScope( rpc::WaitForActorOutOfScopeRequest wait_request; wait_request.set_intended_worker_id(owner_id.Binary()); wait_request.set_actor_id(actor_id.Binary()); - RAY_CHECK_OK(it->second.client->WaitForActorOutOfScope( + it->second.client->WaitForActorOutOfScope( wait_request, [this, owner_node_id, owner_id, actor_id]( Status status, const rpc::WaitForActorOutOfScopeReply &reply) { if (!status.ok()) { @@ -546,7 +546,7 @@ void GcsActorManager::PollOwnerForActorOutOfScope( // have already been destroyed if the owner died. DestroyActor(actor_id); } - })); + }); } void GcsActorManager::DestroyActor(const ActorID &actor_id) { diff --git a/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc b/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc index 44cceef7e..44b7e1338 100644 --- a/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc +++ b/src/ray/gcs/gcs_server/gcs_actor_scheduler.cc @@ -184,14 +184,7 @@ void GcsActorScheduler::ReleaseUnusedWorkers( // nodes do not have leased workers. In this case, GCS will send an empty list. auto workers_in_use = iter != node_to_workers.end() ? iter->second : std::vector{}; - const auto &status = lease_client->ReleaseUnusedWorkers( - workers_in_use, release_unused_workers_callback); - if (!status.ok()) { - RAY_LOG(WARNING) << "Failed to send ReleaseUnusedWorkers request to raylet because " - "raylet may be dead, node id: " - << node_id << ", status: " << status.ToString(); - nodes_of_releasing_unused_workers_.erase(node_id); - } + lease_client->ReleaseUnusedWorkers(workers_in_use, release_unused_workers_callback); } } @@ -215,7 +208,7 @@ void GcsActorScheduler::LeaseWorkerFromNode(std::shared_ptr actor, remote_address.set_ip_address(node->node_manager_address()); remote_address.set_port(node->node_manager_port()); auto lease_client = GetOrConnectLeaseClient(remote_address); - auto status = lease_client->RequestWorkerLease( + lease_client->RequestWorkerLease( actor->GetCreationTaskSpecification(), [this, node_id, actor, node](const Status &status, const rpc::RequestWorkerLeaseReply &reply) { @@ -252,10 +245,6 @@ void GcsActorScheduler::LeaseWorkerFromNode(std::shared_ptr actor, } } }); - - if (!status.ok()) { - RetryLeasingWorkerFromNode(actor, node); - } } void GcsActorScheduler::RetryLeasingWorkerFromNode( @@ -344,7 +333,7 @@ void GcsActorScheduler::CreateActorOnWorker(std::shared_ptr actor, request->mutable_resource_mapping()->CopyFrom(resources); auto client = GetOrConnectCoreWorkerClient(worker->GetAddress()); - auto status = client->PushNormalTask( + client->PushNormalTask( std::move(request), [this, actor, worker](Status status, const rpc::PushTaskReply &reply) { RAY_UNUSED(reply); @@ -377,9 +366,6 @@ void GcsActorScheduler::CreateActorOnWorker(std::shared_ptr actor, } } }); - if (!status.ok()) { - RetryCreatingActorOnWorker(actor, worker); - } } void GcsActorScheduler::RetryCreatingActorOnWorker( 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 9b9f0936e..76cf87b72 100644 --- a/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.cc +++ b/src/ray/gcs/gcs_server/gcs_placement_group_scheduler.cc @@ -164,7 +164,7 @@ void GcsPlacementGroupScheduler::ReserveResourceFromNode( auto lease_client = GetOrConnectLeaseClient(remote_address); RAY_LOG(DEBUG) << "Start leasing resource from node " << node_id << " for bundle " << bundle->BundleId().first << bundle->BundleId().second; - auto status = lease_client->RequestResourceReserve( + lease_client->RequestResourceReserve( *bundle, [this, node_id, bundle, node, callback]( const Status &status, const rpc::RequestResourceReserveReply &reply) { // TODO(AlisaWu): Add placement group cancel. @@ -192,11 +192,6 @@ void GcsPlacementGroupScheduler::ReserveResourceFromNode( } } }); - if (!status.ok()) { - rpc::RequestResourceReserveReply reply; - reply.set_success(false); - callback(status, reply); - } } void GcsPlacementGroupScheduler::CancelResourceReserve( @@ -213,7 +208,7 @@ void GcsPlacementGroupScheduler::CancelResourceReserve( remote_address.set_ip_address(node->node_manager_address()); remote_address.set_port(node->node_manager_port()); auto return_client = GetOrConnectLeaseClient(remote_address); - auto status = return_client->CancelResourceReserve( + return_client->CancelResourceReserve( *bundle_spec, [this, bundle_spec, node](const Status &status, const rpc::CancelResourceReserveReply &reply) { 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 93d2b3a76..b675303c2 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 @@ -43,18 +43,15 @@ class MockActorScheduler : public gcs::GcsActorSchedulerInterface { class MockWorkerClient : public rpc::CoreWorkerClientInterface { public: - ray::Status WaitForActorOutOfScope( + void WaitForActorOutOfScope( const rpc::WaitForActorOutOfScopeRequest &request, const rpc::ClientCallback &callback) override { callbacks.push_back(callback); - return Status::OK(); } - ray::Status KillActor( - const rpc::KillActorRequest &request, - const rpc::ClientCallback &callback) override { + void KillActor(const rpc::KillActorRequest &request, + const rpc::ClientCallback &callback) override { killed_actors.push_back(ActorID::FromBinary(request.intended_actor_id())); - return Status::OK(); } bool Reply(Status status = Status::OK()) { 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 1cf157d2b..effb7a455 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 @@ -32,11 +32,10 @@ namespace ray { struct GcsServerMocker { class MockWorkerClient : public rpc::CoreWorkerClientInterface { public: - ray::Status PushNormalTask( + void PushNormalTask( std::unique_ptr request, const rpc::ClientCallback &callback) override { callbacks.push_back(callback); - return Status::OK(); } bool ReplyPushTask(Status status = Status::OK(), bool exit = false) { @@ -68,28 +67,25 @@ struct GcsServerMocker { return Status::OK(); } - ray::Status RequestWorkerLease( + void RequestWorkerLease( const ray::TaskSpecification &resource_spec, const rpc::ClientCallback &callback) override { num_workers_requested += 1; callbacks.push_back(callback); - return Status::OK(); } - ray::Status ReleaseUnusedWorkers( + void ReleaseUnusedWorkers( const std::vector &workers_in_use, const rpc::ClientCallback &callback) override { num_release_unused_workers += 1; release_callbacks.push_back(callback); - return Status::OK(); } - ray::Status CancelWorkerLease( + void CancelWorkerLease( const TaskID &task_id, const rpc::ClientCallback &callback) override { num_leases_canceled += 1; cancel_callbacks.push_back(callback); - return Status::OK(); } bool GrantWorkerLease() { @@ -162,22 +158,20 @@ struct GcsServerMocker { class MockRayletResourceClient : public ResourceReserveInterface { public: - ray::Status RequestResourceReserve( + void RequestResourceReserve( const BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) override { num_lease_requested += 1; lease_callbacks.push_back(callback); - return Status::OK(); } - ray::Status CancelResourceReserve( + void CancelResourceReserve( BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) override { num_return_requested += 1; return_callbacks.push_back(callback); - return Status::OK(); } // Trigger reply to RequestWorkerLease. diff --git a/src/ray/raylet/node_manager.cc b/src/ray/raylet/node_manager.cc index 3ab693009..0faee2122 100644 --- a/src/ray/raylet/node_manager.cc +++ b/src/ray/raylet/node_manager.cc @@ -509,15 +509,12 @@ void NodeManager::DoLocalGC() { RAY_LOG(WARNING) << "Sending local GC request to " << all_workers.size() << " workers."; for (const auto &worker : all_workers) { rpc::LocalGCRequest request; - auto status = worker->rpc_client()->LocalGC( + worker->rpc_client()->LocalGC( request, [](const ray::Status &status, const rpc::LocalGCReply &r) { if (!status.ok()) { RAY_LOG(ERROR) << "Failed to send local GC request: " << status.ToString(); } }); - if (!status.ok()) { - RAY_LOG(ERROR) << "Failed to send local GC request: " << status.ToString(); - } } } @@ -2957,27 +2954,26 @@ void NodeManager::HandleTaskReconstruction(const TaskID &task_id, rpc::GetObjectStatusRequest request; request.set_object_id(required_object_id.Binary()); request.set_owner_worker_id(owner_addr.worker_id()); - RAY_CHECK_OK(client->GetObjectStatus( - request, [this, required_object_id](Status status, - const rpc::GetObjectStatusReply &reply) { - if (!status.ok() || - reply.status() == rpc::GetObjectStatusReply::OUT_OF_SCOPE || - reply.status() == rpc::GetObjectStatusReply::FREED) { - // The owner is gone, or the owner replied that the object has - // gone out of scope (this is an edge case in the distributed ref - // counting protocol where a borrower dies before it can notify - // the owner of another borrower), or the object value has been - // freed. Store an error in the local plasma store so that an - // exception will be thrown when the worker tries to get the - // value. - MarkObjectsAsFailed(ErrorType::OBJECT_UNRECONSTRUCTABLE, - {required_object_id}, JobID::Nil()); - } - // Do nothing if the owner replied that the object is available. The - // object manager will continue trying to fetch the object, and this - // handler will get triggered again if the object is still - // unavailable after another timeout. - })); + client->GetObjectStatus(request, [this, required_object_id]( + Status status, + const rpc::GetObjectStatusReply &reply) { + if (!status.ok() || reply.status() == rpc::GetObjectStatusReply::OUT_OF_SCOPE || + reply.status() == rpc::GetObjectStatusReply::FREED) { + // The owner is gone, or the owner replied that the object has + // gone out of scope (this is an edge case in the distributed ref + // counting protocol where a borrower dies before it can notify + // the owner of another borrower), or the object value has been + // freed. Store an error in the local plasma store so that an + // exception will be thrown when the worker tries to get the + // value. + MarkObjectsAsFailed(ErrorType::OBJECT_UNRECONSTRUCTABLE, {required_object_id}, + JobID::Nil()); + } + // Do nothing if the owner replied that the object is available. The + // object manager will continue trying to fetch the object, and this + // handler will get triggered again if the object is still + // unavailable after another timeout. + }); } } else { // We do not have the owner's address. This is either an actor creation @@ -3416,14 +3412,14 @@ ray::Status NodeManager::SetupPlasmaSubscription() { request.set_data_size(object_info.data_size); for (auto worker : waiting_workers) { - RAY_CHECK_OK(worker->rpc_client()->PlasmaObjectReady( + worker->rpc_client()->PlasmaObjectReady( request, [](Status status, const rpc::PlasmaObjectReadyReply &reply) { if (!status.ok()) { RAY_LOG(INFO) << "Problem with telling worker that plasma object is ready" << status.ToString(); } - })); + }); } }); } @@ -3562,7 +3558,7 @@ void NodeManager::HandlePinObjectIDs(const rpc::PinObjectIDsRequest &request, wait_request.set_object_id(object_id_binary); wait_request.set_intended_worker_id(request.owner_address().worker_id()); worker_rpc_clients_[worker_id].second++; - RAY_CHECK_OK(it->second.first->WaitForObjectEviction( + it->second.first->WaitForObjectEviction( wait_request, [this, worker_id, object_id]( Status status, const rpc::WaitForObjectEvictionReply &reply) { if (!status.ok()) { @@ -3586,7 +3582,7 @@ void NodeManager::HandlePinObjectIDs(const rpc::PinObjectIDsRequest &request, if (--worker_rpc_clients_[worker_id].second == 0) { worker_rpc_clients_.erase(worker_id); } - })); + }); } send_reply_callback(Status::OK(), nullptr, nullptr); } @@ -3671,7 +3667,7 @@ void NodeManager::HandleGetNodeStats(const rpc::GetNodeStatsRequest &node_stats_ rpc::GetCoreWorkerStatsRequest request; request.set_intended_worker_id(worker->WorkerId().Binary()); request.set_include_memory_info(node_stats_request.include_memory_info()); - auto status = worker->rpc_client()->GetCoreWorkerStats( + worker->rpc_client()->GetCoreWorkerStats( request, [reply, worker, all_workers, driver_ids, send_reply_callback]( const ray::Status &status, const rpc::GetCoreWorkerStatsReply &r) { auto worker_stats = reply->add_workers_stats(); @@ -3689,10 +3685,6 @@ void NodeManager::HandleGetNodeStats(const rpc::GetNodeStatsRequest &node_stats_ send_reply_callback(Status::OK(), nullptr, nullptr); } }); - if (!status.ok()) { - RAY_LOG(ERROR) << "Failed to send get core worker stats request: " - << status.ToString(); - } } } diff --git a/src/ray/raylet/worker.cc b/src/ray/raylet/worker.cc index 2f744b6d6..aea922f2c 100644 --- a/src/ray/raylet/worker.cc +++ b/src/ray/raylet/worker.cc @@ -171,8 +171,7 @@ Status Worker::AssignTask(const Task &task, const ResourceIdSet &resource_id_set task.GetTaskExecutionSpec().GetMessage()); request.set_resource_ids(resource_id_set.Serialize()); - return rpc_client_->AssignTask(request, [](Status status, - const rpc::AssignTaskReply &reply) { + rpc_client_->AssignTask(request, [](Status status, const rpc::AssignTaskReply &reply) { if (!status.ok()) { RAY_LOG(DEBUG) << "Worker failed to finish executing task: " << status.ToString(); } @@ -180,6 +179,7 @@ Status Worker::AssignTask(const Task &task, const ResourceIdSet &resource_id_set // and assigning new task will be done when raylet receives // `TaskDone` message. }); + return Status::OK(); } void Worker::DirectActorCallArgWaitComplete(int64_t tag) { @@ -187,15 +187,12 @@ void Worker::DirectActorCallArgWaitComplete(int64_t tag) { rpc::DirectActorCallArgWaitCompleteRequest request; request.set_tag(tag); request.set_intended_worker_id(worker_id_.Binary()); - auto status = rpc_client_->DirectActorCallArgWaitComplete( + rpc_client_->DirectActorCallArgWaitComplete( request, [](Status status, const rpc::DirectActorCallArgWaitCompleteReply &reply) { if (!status.ok()) { RAY_LOG(ERROR) << "Failed to send wait complete: " << status.ToString(); } }); - if (!status.ok()) { - RAY_LOG(ERROR) << "Failed to send wait complete: " << status.ToString(); - } } } // namespace raylet diff --git a/src/ray/raylet_client/raylet_client.cc b/src/ray/raylet_client/raylet_client.cc index 2a2938a06..b2b2020c3 100644 --- a/src/ray/raylet_client/raylet_client.cc +++ b/src/ray/raylet_client/raylet_client.cc @@ -296,12 +296,12 @@ Status raylet::RayletClient::SetResource(const std::string &resource_name, return conn_->WriteMessage(MessageType::SetResourceRequest, &fbb); } -Status raylet::RayletClient::RequestWorkerLease( +void raylet::RayletClient::RequestWorkerLease( const TaskSpecification &resource_spec, const rpc::ClientCallback &callback) { rpc::RequestWorkerLeaseRequest request; request.mutable_resource_spec()->CopyFrom(resource_spec.GetMessage()); - return grpc_client_->RequestWorkerLease(request, callback); + grpc_client_->RequestWorkerLease(request, callback); } Status raylet::RayletClient::ReturnWorker(int worker_port, const WorkerID &worker_id, @@ -310,22 +310,23 @@ Status raylet::RayletClient::ReturnWorker(int worker_port, const WorkerID &worke request.set_worker_port(worker_port); request.set_worker_id(worker_id.Binary()); request.set_disconnect_worker(disconnect_worker); - return grpc_client_->ReturnWorker( + grpc_client_->ReturnWorker( request, [](const Status &status, const rpc::ReturnWorkerReply &reply) { if (!status.ok()) { RAY_LOG(INFO) << "Error returning worker: " << status; } }); + return Status::OK(); } -Status raylet::RayletClient::ReleaseUnusedWorkers( +void raylet::RayletClient::ReleaseUnusedWorkers( const std::vector &workers_in_use, const rpc::ClientCallback &callback) { rpc::ReleaseUnusedWorkersRequest request; for (auto &worker_id : workers_in_use) { request.add_worker_ids_in_use(worker_id.Binary()); } - return grpc_client_->ReleaseUnusedWorkers( + grpc_client_->ReleaseUnusedWorkers( request, [callback](const Status &status, const rpc::ReleaseUnusedWorkersReply &reply) { if (!status.ok()) { @@ -337,31 +338,31 @@ Status raylet::RayletClient::ReleaseUnusedWorkers( }); } -ray::Status raylet::RayletClient::CancelWorkerLease( +void raylet::RayletClient::CancelWorkerLease( const TaskID &task_id, const rpc::ClientCallback &callback) { rpc::CancelWorkerLeaseRequest request; request.set_task_id(task_id.Binary()); - return grpc_client_->CancelWorkerLease(request, callback); + grpc_client_->CancelWorkerLease(request, callback); } -Status raylet::RayletClient::RequestResourceReserve( +void raylet::RayletClient::RequestResourceReserve( const BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) { rpc::RequestResourceReserveRequest request; request.mutable_bundle_spec()->CopyFrom(bundle_spec.GetMessage()); - return grpc_client_->RequestResourceReserve(request, callback); + grpc_client_->RequestResourceReserve(request, callback); } -Status raylet::RayletClient::CancelResourceReserve( +void raylet::RayletClient::CancelResourceReserve( BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) { rpc::CancelResourceReserveRequest request; request.mutable_bundle_spec()->CopyFrom(bundle_spec.GetMessage()); - return grpc_client_->CancelResourceReserve(request, callback); + grpc_client_->CancelResourceReserve(request, callback); } -Status raylet::RayletClient::PinObjectIDs( +void raylet::RayletClient::PinObjectIDs( const rpc::Address &caller_address, const std::vector &object_ids, const rpc::ClientCallback &callback) { rpc::PinObjectIDsRequest request; @@ -369,13 +370,13 @@ Status raylet::RayletClient::PinObjectIDs( for (const ObjectID &object_id : object_ids) { request.add_object_ids(object_id.Binary()); } - return grpc_client_->PinObjectIDs(request, callback); + grpc_client_->PinObjectIDs(request, callback); } -Status raylet::RayletClient::GlobalGC( +void raylet::RayletClient::GlobalGC( const rpc::ClientCallback &callback) { rpc::GlobalGCRequest request; - return grpc_client_->GlobalGC(request, callback); + grpc_client_->GlobalGC(request, callback); } Status raylet::RayletClient::SubscribeToPlasma(const ObjectID &object_id) { diff --git a/src/ray/raylet_client/raylet_client.h b/src/ray/raylet_client/raylet_client.h index bda70aa1b..cbb42edbe 100644 --- a/src/ray/raylet_client/raylet_client.h +++ b/src/ray/raylet_client/raylet_client.h @@ -47,7 +47,7 @@ namespace ray { class PinObjectsInterface { public: /// Request to a raylet to pin a plasma object. The callback will be sent via gRPC. - virtual ray::Status PinObjectIDs( + virtual void PinObjectIDs( const rpc::Address &caller_address, const std::vector &object_ids, const ray::rpc::ClientCallback &callback) = 0; @@ -60,7 +60,7 @@ class WorkerLeaseInterface { /// Requests a worker from the raylet. The callback will be sent via gRPC. /// \param resource_spec Resources that should be allocated for the worker. /// \return ray::Status - virtual ray::Status RequestWorkerLease( + virtual void RequestWorkerLease( const ray::TaskSpecification &resource_spec, const ray::rpc::ClientCallback &callback) = 0; @@ -76,11 +76,11 @@ class WorkerLeaseInterface { /// \param workers_in_use Workers currently in use. /// \param callback Callback that will be called after raylet completes the release of /// unused workers. \return ray::Status - virtual ray::Status ReleaseUnusedWorkers( + virtual void ReleaseUnusedWorkers( const std::vector &workers_in_use, const rpc::ClientCallback &callback) = 0; - virtual ray::Status CancelWorkerLease( + virtual void CancelWorkerLease( const TaskID &task_id, const rpc::ClientCallback &callback) = 0; @@ -93,12 +93,12 @@ class ResourceReserveInterface { /// Requests a resource from the raylet. The callback will be sent via gRPC. /// \param resource_spec Resources that should be allocated for the worker. /// \return ray::Status - virtual ray::Status RequestResourceReserve( + virtual void RequestResourceReserve( const BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) = 0; - virtual ray::Status CancelResourceReserve( + virtual void CancelResourceReserve( BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) = 0; @@ -318,7 +318,7 @@ class RayletClient : public PinObjectsInterface, const ray::ClientID &client_Id); /// Implements WorkerLeaseInterface. - ray::Status RequestWorkerLease( + void RequestWorkerLease( const ray::TaskSpecification &resource_spec, const ray::rpc::ClientCallback &callback) override; @@ -328,31 +328,31 @@ class RayletClient : public PinObjectsInterface, bool disconnect_worker) override; /// Implements WorkerLeaseInterface. - ray::Status ReleaseUnusedWorkers( + void ReleaseUnusedWorkers( const std::vector &workers_in_use, const rpc::ClientCallback &callback) override; - ray::Status CancelWorkerLease( + void CancelWorkerLease( const TaskID &task_id, const rpc::ClientCallback &callback) override; /// Implements ResourceReserveInterface. - ray::Status RequestResourceReserve( + void RequestResourceReserve( const BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) override; /// Implements ResourceReserveInterface. - ray::Status CancelResourceReserve( + void CancelResourceReserve( BundleSpecification &bundle_spec, const ray::rpc::ClientCallback &callback) override; - ray::Status PinObjectIDs( + void PinObjectIDs( const rpc::Address &caller_address, const std::vector &object_ids, const ray::rpc::ClientCallback &callback) override; - ray::Status GlobalGC(const rpc::ClientCallback &callback); + void GlobalGC(const rpc::ClientCallback &callback); // Subscribe to receive notification on plasma object ray::Status SubscribeToPlasma(const ObjectID &object_id); diff --git a/src/ray/rpc/grpc_client.h b/src/ray/rpc/grpc_client.h index 141dcae7e..c49244e48 100644 --- a/src/ray/rpc/grpc_client.h +++ b/src/ray/rpc/grpc_client.h @@ -32,17 +32,10 @@ namespace rpc { &SERVICE::Stub::PrepareAsync##METHOD, request, callback)) // Define a void RPC client method. -#define VOID_RPC_CLIENT_METHOD(SERVICE, METHOD, rpc_client, SPECS) \ - void METHOD(const METHOD##Request &request, \ - const ClientCallback &callback) SPECS { \ - RAY_UNUSED(INVOKE_RPC_CALL(SERVICE, METHOD, request, callback, rpc_client)); \ - } - -// Define a RPC client method that returns ray::Status. -#define RPC_CLIENT_METHOD(SERVICE, METHOD, rpc_client, SPECS) \ - ray::Status METHOD(const METHOD##Request &request, \ - const ClientCallback &callback) SPECS { \ - return INVOKE_RPC_CALL(SERVICE, METHOD, request, callback, rpc_client); \ +#define VOID_RPC_CLIENT_METHOD(SERVICE, METHOD, rpc_client, SPECS) \ + void METHOD(const METHOD##Request &request, \ + const ClientCallback &callback) SPECS { \ + INVOKE_RPC_CALL(SERVICE, METHOD, request, callback, rpc_client); \ } template @@ -90,12 +83,12 @@ class GrpcClient { /// /// \return Status. template - ray::Status CallMethod( + void CallMethod( const PrepareAsyncFunction prepare_async_function, const Request &request, const ClientCallback &callback) { auto call = client_call_manager_.CreateCall( *stub_, prepare_async_function, request, callback); - return call->GetStatus(); + RAY_CHECK(call != nullptr); } private: diff --git a/src/ray/rpc/node_manager/node_manager_client.h b/src/ray/rpc/node_manager/node_manager_client.h index 8f2c26b19..c908a4d79 100644 --- a/src/ray/rpc/node_manager/node_manager_client.h +++ b/src/ray/rpc/node_manager/node_manager_client.h @@ -77,28 +77,28 @@ class NodeManagerWorkerClient } /// Request a worker lease. - RPC_CLIENT_METHOD(NodeManagerService, RequestWorkerLease, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, RequestWorkerLease, grpc_client_, ) /// Return a worker lease. - RPC_CLIENT_METHOD(NodeManagerService, ReturnWorker, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, ReturnWorker, grpc_client_, ) /// Release unused workers. - RPC_CLIENT_METHOD(NodeManagerService, ReleaseUnusedWorkers, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, ReleaseUnusedWorkers, grpc_client_, ) /// Cancel a pending worker lease request. - RPC_CLIENT_METHOD(NodeManagerService, CancelWorkerLease, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, CancelWorkerLease, grpc_client_, ) /// Request resource lease. - RPC_CLIENT_METHOD(NodeManagerService, RequestResourceReserve, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, RequestResourceReserve, grpc_client_, ) /// Return resource lease. - RPC_CLIENT_METHOD(NodeManagerService, CancelResourceReserve, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, CancelResourceReserve, grpc_client_, ) /// Notify the raylet to pin the provided object IDs. - RPC_CLIENT_METHOD(NodeManagerService, PinObjectIDs, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, PinObjectIDs, grpc_client_, ) /// Trigger global GC across the cluster. - RPC_CLIENT_METHOD(NodeManagerService, GlobalGC, grpc_client_, ) + VOID_RPC_CLIENT_METHOD(NodeManagerService, GlobalGC, grpc_client_, ) private: /// Constructor. diff --git a/src/ray/rpc/worker/core_worker_client.h b/src/ray/rpc/worker/core_worker_client.h index 55e34467b..75cdcdf5a 100644 --- a/src/ray/rpc/worker/core_worker_client.h +++ b/src/ray/rpc/worker/core_worker_client.h @@ -109,10 +109,8 @@ class CoreWorkerClientInterface { /// \param[in] request The request message. /// \param[in] callback The callback function that handles reply. /// \return if the rpc call succeeds - virtual ray::Status AssignTask(const AssignTaskRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void AssignTask(const AssignTaskRequest &request, + const ClientCallback &callback) {} /// Push an actor task directly from worker to worker. /// @@ -121,89 +119,60 @@ class CoreWorkerClientInterface { /// task for execution immediately. /// \param[in] callback The callback function that handles reply. /// \return if the rpc call succeeds - virtual ray::Status PushActorTask(std::unique_ptr request, - bool skip_queue, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void PushActorTask(std::unique_ptr request, bool skip_queue, + const ClientCallback &callback) {} /// Similar to PushActorTask, but sets no ordering constraint. This is used to /// push non-actor tasks directly to a worker. - virtual ray::Status PushNormalTask(std::unique_ptr request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void PushNormalTask(std::unique_ptr request, + const ClientCallback &callback) {} /// Notify a wait has completed for direct actor call arguments. /// /// \param[in] request The request message. /// \param[in] callback The callback function that handles reply. /// \return if the rpc call succeeds - virtual ray::Status DirectActorCallArgWaitComplete( + virtual void DirectActorCallArgWaitComplete( const DirectActorCallArgWaitCompleteRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + const ClientCallback &callback) {} /// Ask the owner of an object about the object's current status. - virtual ray::Status GetObjectStatus( - const GetObjectStatusRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void GetObjectStatus(const GetObjectStatusRequest &request, + const ClientCallback &callback) {} /// Ask the actor's owner to reply when the actor has gone out of scope. - virtual ray::Status WaitForActorOutOfScope( + virtual void WaitForActorOutOfScope( const WaitForActorOutOfScopeRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + const ClientCallback &callback) {} /// Notify the owner of an object that the object has been pinned. - virtual ray::Status WaitForObjectEviction( + virtual void WaitForObjectEviction( const WaitForObjectEvictionRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + const ClientCallback &callback) {} /// Tell this actor to exit immediately. - virtual ray::Status KillActor(const KillActorRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void KillActor(const KillActorRequest &request, + const ClientCallback &callback) {} - virtual ray::Status CancelTask(const CancelTaskRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void CancelTask(const CancelTaskRequest &request, + const ClientCallback &callback) {} - virtual ray::Status RemoteCancelTask( - const RemoteCancelTaskRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } + virtual void RemoteCancelTask(const RemoteCancelTaskRequest &request, + const ClientCallback &callback) {} - virtual ray::Status GetCoreWorkerStats( + virtual void GetCoreWorkerStats( const GetCoreWorkerStatsRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); + const ClientCallback &callback) {} + + virtual void LocalGC(const LocalGCRequest &request, + const ClientCallback &callback) {} + + virtual void WaitForRefRemoved(const WaitForRefRemovedRequest &request, + const ClientCallback &callback) { } - virtual ray::Status LocalGC(const LocalGCRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } - - virtual ray::Status WaitForRefRemoved( - const WaitForRefRemovedRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); - } - - virtual ray::Status PlasmaObjectReady( - const PlasmaObjectReadyRequest &request, - const ClientCallback &callback) { - return Status::NotImplemented(""); + virtual void PlasmaObjectReady(const PlasmaObjectReadyRequest &request, + const ClientCallback &callback) { } virtual ~CoreWorkerClientInterface(){}; @@ -227,40 +196,41 @@ class CoreWorkerClient : public std::enable_shared_from_this, const rpc::Address &Addr() const override { return addr_; } - RPC_CLIENT_METHOD(CoreWorkerService, AssignTask, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, AssignTask, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, DirectActorCallArgWaitComplete, grpc_client_, - override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, DirectActorCallArgWaitComplete, grpc_client_, + override) - RPC_CLIENT_METHOD(CoreWorkerService, GetObjectStatus, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, GetObjectStatus, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, KillActor, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, KillActor, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, CancelTask, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, CancelTask, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, RemoteCancelTask, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, RemoteCancelTask, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, WaitForActorOutOfScope, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, WaitForActorOutOfScope, grpc_client_, + override) - RPC_CLIENT_METHOD(CoreWorkerService, WaitForObjectEviction, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, WaitForObjectEviction, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, GetCoreWorkerStats, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, GetCoreWorkerStats, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, LocalGC, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, LocalGC, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, WaitForRefRemoved, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, WaitForRefRemoved, grpc_client_, override) - RPC_CLIENT_METHOD(CoreWorkerService, PlasmaObjectReady, grpc_client_, override) + VOID_RPC_CLIENT_METHOD(CoreWorkerService, PlasmaObjectReady, grpc_client_, override) - ray::Status PushActorTask(std::unique_ptr request, bool skip_queue, - const ClientCallback &callback) override { + void PushActorTask(std::unique_ptr request, bool skip_queue, + const ClientCallback &callback) override { if (skip_queue) { // Set this value so that the actor does not skip any tasks when // processing this request. We could also set it to max_finished_seq_no_, // but we just set it to the default of -1 to avoid taking the lock. request->set_client_processed_up_to(-1); - return INVOKE_RPC_CALL(CoreWorkerService, PushTask, *request, callback, - grpc_client_); + INVOKE_RPC_CALL(CoreWorkerService, PushTask, *request, callback, grpc_client_); + return; } { @@ -268,14 +238,13 @@ class CoreWorkerClient : public std::enable_shared_from_this, send_queue_.push_back(std::make_pair(std::move(request), callback)); } SendRequests(); - return ray::Status::OK(); } - ray::Status PushNormalTask(std::unique_ptr request, - const ClientCallback &callback) override { + void PushNormalTask(std::unique_ptr request, + const ClientCallback &callback) override { request->set_sequence_number(-1); request->set_client_processed_up_to(-1); - return INVOKE_RPC_CALL(CoreWorkerService, PushTask, *request, callback, grpc_client_); + INVOKE_RPC_CALL(CoreWorkerService, PushTask, *request, callback, grpc_client_); } /// Send as many pending tasks as possible. This method is thread-safe.