mirror of
https://github.com/wassname/ray.git
synced 2026-08-01 12:51:09 +08:00
async grpc calls should always return void (#9533)
This commit is contained in:
@@ -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<ObjectID> &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() {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<rpc::WaitForRefRemovedReply> &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() {
|
||||
|
||||
@@ -59,12 +59,10 @@ class MockWorkerClient : public rpc::CoreWorkerClientInterface {
|
||||
public:
|
||||
const rpc::Address &Addr() const override { return addr; }
|
||||
|
||||
ray::Status PushActorTask(
|
||||
std::unique_ptr<rpc::PushTaskRequest> request, bool skip_queue,
|
||||
const rpc::ClientCallback<rpc::PushTaskReply> &callback) override {
|
||||
void PushActorTask(std::unique_ptr<rpc::PushTaskRequest> request, bool skip_queue,
|
||||
const rpc::ClientCallback<rpc::PushTaskReply> &callback) override {
|
||||
received_seq_nos.push_back(request->sequence_number());
|
||||
callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
bool ReplyPushTask(Status status = Status::OK()) {
|
||||
|
||||
@@ -31,11 +31,9 @@ int64_t kLongTimeout = 1024 * 1024 * 1024;
|
||||
|
||||
class MockWorkerClient : public rpc::CoreWorkerClientInterface {
|
||||
public:
|
||||
ray::Status PushNormalTask(
|
||||
std::unique_ptr<rpc::PushTaskRequest> request,
|
||||
const rpc::ClientCallback<rpc::PushTaskReply> &callback) override {
|
||||
void PushNormalTask(std::unique_ptr<rpc::PushTaskRequest> request,
|
||||
const rpc::ClientCallback<rpc::PushTaskReply> &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<rpc::CancelTaskReply> &callback) override {
|
||||
void CancelTask(const rpc::CancelTaskRequest &request,
|
||||
const rpc::ClientCallback<rpc::CancelTaskReply> &callback) override {
|
||||
kill_requests.push_front(request);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
std::list<rpc::ClientCallback<rpc::PushTaskReply>> 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<rpc::RequestWorkerLeaseReply> &callback) override {
|
||||
num_workers_requested += 1;
|
||||
callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
ray::Status ReleaseUnusedWorkers(
|
||||
void ReleaseUnusedWorkers(
|
||||
const std::vector<WorkerID> &workers_in_use,
|
||||
const rpc::ClientCallback<rpc::ReleaseUnusedWorkersReply> &callback) override {
|
||||
return Status::NotImplemented("ReleaseUnusedWorkers is not supported.");
|
||||
}
|
||||
const rpc::ClientCallback<rpc::ReleaseUnusedWorkersReply> &callback) override {}
|
||||
|
||||
ray::Status CancelWorkerLease(
|
||||
void CancelWorkerLease(
|
||||
const TaskID &task_id,
|
||||
const rpc::ClientCallback<rpc::CancelWorkerLeaseReply> &callback) override {
|
||||
num_leases_canceled += 1;
|
||||
cancel_callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
// Trigger reply to RequestWorkerLease.
|
||||
|
||||
@@ -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<ObjectID> &object_ids,
|
||||
const ray::rpc::ClientCallback<ray::rpc::PinObjectIDsReply> &callback) override {
|
||||
RAY_LOG(INFO) << "PinObjectIDs " << object_ids.size();
|
||||
callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
size_t Flush() {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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<WorkerID>{};
|
||||
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<GcsActor> 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<GcsActor> actor,
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
if (!status.ok()) {
|
||||
RetryLeasingWorkerFromNode(actor, node);
|
||||
}
|
||||
}
|
||||
|
||||
void GcsActorScheduler::RetryLeasingWorkerFromNode(
|
||||
@@ -344,7 +333,7 @@ void GcsActorScheduler::CreateActorOnWorker(std::shared_ptr<GcsActor> 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<GcsActor> actor,
|
||||
}
|
||||
}
|
||||
});
|
||||
if (!status.ok()) {
|
||||
RetryCreatingActorOnWorker(actor, worker);
|
||||
}
|
||||
}
|
||||
|
||||
void GcsActorScheduler::RetryCreatingActorOnWorker(
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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<rpc::WaitForActorOutOfScopeReply> &callback) override {
|
||||
callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
ray::Status KillActor(
|
||||
const rpc::KillActorRequest &request,
|
||||
const rpc::ClientCallback<rpc::KillActorReply> &callback) override {
|
||||
void KillActor(const rpc::KillActorRequest &request,
|
||||
const rpc::ClientCallback<rpc::KillActorReply> &callback) override {
|
||||
killed_actors.push_back(ActorID::FromBinary(request.intended_actor_id()));
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
bool Reply(Status status = Status::OK()) {
|
||||
|
||||
@@ -32,11 +32,10 @@ namespace ray {
|
||||
struct GcsServerMocker {
|
||||
class MockWorkerClient : public rpc::CoreWorkerClientInterface {
|
||||
public:
|
||||
ray::Status PushNormalTask(
|
||||
void PushNormalTask(
|
||||
std::unique_ptr<rpc::PushTaskRequest> request,
|
||||
const rpc::ClientCallback<rpc::PushTaskReply> &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<rpc::RequestWorkerLeaseReply> &callback) override {
|
||||
num_workers_requested += 1;
|
||||
callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
ray::Status ReleaseUnusedWorkers(
|
||||
void ReleaseUnusedWorkers(
|
||||
const std::vector<WorkerID> &workers_in_use,
|
||||
const rpc::ClientCallback<rpc::ReleaseUnusedWorkersReply> &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<rpc::CancelWorkerLeaseReply> &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<ray::rpc::RequestResourceReserveReply> &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<ray::rpc::CancelResourceReserveReply> &callback)
|
||||
override {
|
||||
num_return_requested += 1;
|
||||
return_callbacks.push_back(callback);
|
||||
return Status::OK();
|
||||
}
|
||||
|
||||
// Trigger reply to RequestWorkerLease.
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<rpc::RequestWorkerLeaseReply> &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<WorkerID> &workers_in_use,
|
||||
const rpc::ClientCallback<rpc::ReleaseUnusedWorkersReply> &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<rpc::CancelWorkerLeaseReply> &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<ray::rpc::RequestResourceReserveReply> &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<ray::rpc::CancelResourceReserveReply> &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<ObjectID> &object_ids,
|
||||
const rpc::ClientCallback<rpc::PinObjectIDsReply> &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<rpc::GlobalGCReply> &callback) {
|
||||
rpc::GlobalGCRequest request;
|
||||
return grpc_client_->GlobalGC(request, callback);
|
||||
grpc_client_->GlobalGC(request, callback);
|
||||
}
|
||||
|
||||
Status raylet::RayletClient::SubscribeToPlasma(const ObjectID &object_id) {
|
||||
|
||||
@@ -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<ObjectID> &object_ids,
|
||||
const ray::rpc::ClientCallback<ray::rpc::PinObjectIDsReply> &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<ray::rpc::RequestWorkerLeaseReply> &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<WorkerID> &workers_in_use,
|
||||
const rpc::ClientCallback<rpc::ReleaseUnusedWorkersReply> &callback) = 0;
|
||||
|
||||
virtual ray::Status CancelWorkerLease(
|
||||
virtual void CancelWorkerLease(
|
||||
const TaskID &task_id,
|
||||
const rpc::ClientCallback<rpc::CancelWorkerLeaseReply> &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<ray::rpc::RequestResourceReserveReply>
|
||||
&callback) = 0;
|
||||
|
||||
virtual ray::Status CancelResourceReserve(
|
||||
virtual void CancelResourceReserve(
|
||||
BundleSpecification &bundle_spec,
|
||||
const ray::rpc::ClientCallback<ray::rpc::CancelResourceReserveReply> &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<ray::rpc::RequestWorkerLeaseReply> &callback)
|
||||
override;
|
||||
@@ -328,31 +328,31 @@ class RayletClient : public PinObjectsInterface,
|
||||
bool disconnect_worker) override;
|
||||
|
||||
/// Implements WorkerLeaseInterface.
|
||||
ray::Status ReleaseUnusedWorkers(
|
||||
void ReleaseUnusedWorkers(
|
||||
const std::vector<WorkerID> &workers_in_use,
|
||||
const rpc::ClientCallback<rpc::ReleaseUnusedWorkersReply> &callback) override;
|
||||
|
||||
ray::Status CancelWorkerLease(
|
||||
void CancelWorkerLease(
|
||||
const TaskID &task_id,
|
||||
const rpc::ClientCallback<rpc::CancelWorkerLeaseReply> &callback) override;
|
||||
|
||||
/// Implements ResourceReserveInterface.
|
||||
ray::Status RequestResourceReserve(
|
||||
void RequestResourceReserve(
|
||||
const BundleSpecification &bundle_spec,
|
||||
const ray::rpc::ClientCallback<ray::rpc::RequestResourceReserveReply> &callback)
|
||||
override;
|
||||
|
||||
/// Implements ResourceReserveInterface.
|
||||
ray::Status CancelResourceReserve(
|
||||
void CancelResourceReserve(
|
||||
BundleSpecification &bundle_spec,
|
||||
const ray::rpc::ClientCallback<ray::rpc::CancelResourceReserveReply> &callback)
|
||||
override;
|
||||
|
||||
ray::Status PinObjectIDs(
|
||||
void PinObjectIDs(
|
||||
const rpc::Address &caller_address, const std::vector<ObjectID> &object_ids,
|
||||
const ray::rpc::ClientCallback<ray::rpc::PinObjectIDsReply> &callback) override;
|
||||
|
||||
ray::Status GlobalGC(const rpc::ClientCallback<rpc::GlobalGCReply> &callback);
|
||||
void GlobalGC(const rpc::ClientCallback<rpc::GlobalGCReply> &callback);
|
||||
|
||||
// Subscribe to receive notification on plasma object
|
||||
ray::Status SubscribeToPlasma(const ObjectID &object_id);
|
||||
|
||||
@@ -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<METHOD##Reply> &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<METHOD##Reply> &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<METHOD##Reply> &callback) SPECS { \
|
||||
INVOKE_RPC_CALL(SERVICE, METHOD, request, callback, rpc_client); \
|
||||
}
|
||||
|
||||
template <class GrpcService>
|
||||
@@ -90,12 +83,12 @@ class GrpcClient {
|
||||
///
|
||||
/// \return Status.
|
||||
template <class Request, class Reply>
|
||||
ray::Status CallMethod(
|
||||
void CallMethod(
|
||||
const PrepareAsyncFunction<GrpcService, Request, Reply> prepare_async_function,
|
||||
const Request &request, const ClientCallback<Reply> &callback) {
|
||||
auto call = client_call_manager_.CreateCall<GrpcService, Request, Reply>(
|
||||
*stub_, prepare_async_function, request, callback);
|
||||
return call->GetStatus();
|
||||
RAY_CHECK(call != nullptr);
|
||||
}
|
||||
|
||||
private:
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<AssignTaskReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void AssignTask(const AssignTaskRequest &request,
|
||||
const ClientCallback<AssignTaskReply> &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<PushTaskRequest> request,
|
||||
bool skip_queue,
|
||||
const ClientCallback<PushTaskReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void PushActorTask(std::unique_ptr<PushTaskRequest> request, bool skip_queue,
|
||||
const ClientCallback<PushTaskReply> &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<PushTaskRequest> request,
|
||||
const ClientCallback<PushTaskReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void PushNormalTask(std::unique_ptr<PushTaskRequest> request,
|
||||
const ClientCallback<PushTaskReply> &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<DirectActorCallArgWaitCompleteReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
const ClientCallback<DirectActorCallArgWaitCompleteReply> &callback) {}
|
||||
|
||||
/// Ask the owner of an object about the object's current status.
|
||||
virtual ray::Status GetObjectStatus(
|
||||
const GetObjectStatusRequest &request,
|
||||
const ClientCallback<GetObjectStatusReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void GetObjectStatus(const GetObjectStatusRequest &request,
|
||||
const ClientCallback<GetObjectStatusReply> &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<WaitForActorOutOfScopeReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
const ClientCallback<WaitForActorOutOfScopeReply> &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<WaitForObjectEvictionReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
const ClientCallback<WaitForObjectEvictionReply> &callback) {}
|
||||
|
||||
/// Tell this actor to exit immediately.
|
||||
virtual ray::Status KillActor(const KillActorRequest &request,
|
||||
const ClientCallback<KillActorReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void KillActor(const KillActorRequest &request,
|
||||
const ClientCallback<KillActorReply> &callback) {}
|
||||
|
||||
virtual ray::Status CancelTask(const CancelTaskRequest &request,
|
||||
const ClientCallback<CancelTaskReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void CancelTask(const CancelTaskRequest &request,
|
||||
const ClientCallback<CancelTaskReply> &callback) {}
|
||||
|
||||
virtual ray::Status RemoteCancelTask(
|
||||
const RemoteCancelTaskRequest &request,
|
||||
const ClientCallback<RemoteCancelTaskReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
virtual void RemoteCancelTask(const RemoteCancelTaskRequest &request,
|
||||
const ClientCallback<RemoteCancelTaskReply> &callback) {}
|
||||
|
||||
virtual ray::Status GetCoreWorkerStats(
|
||||
virtual void GetCoreWorkerStats(
|
||||
const GetCoreWorkerStatsRequest &request,
|
||||
const ClientCallback<GetCoreWorkerStatsReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
const ClientCallback<GetCoreWorkerStatsReply> &callback) {}
|
||||
|
||||
virtual void LocalGC(const LocalGCRequest &request,
|
||||
const ClientCallback<LocalGCReply> &callback) {}
|
||||
|
||||
virtual void WaitForRefRemoved(const WaitForRefRemovedRequest &request,
|
||||
const ClientCallback<WaitForRefRemovedReply> &callback) {
|
||||
}
|
||||
|
||||
virtual ray::Status LocalGC(const LocalGCRequest &request,
|
||||
const ClientCallback<LocalGCReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
|
||||
virtual ray::Status WaitForRefRemoved(
|
||||
const WaitForRefRemovedRequest &request,
|
||||
const ClientCallback<WaitForRefRemovedReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
}
|
||||
|
||||
virtual ray::Status PlasmaObjectReady(
|
||||
const PlasmaObjectReadyRequest &request,
|
||||
const ClientCallback<PlasmaObjectReadyReply> &callback) {
|
||||
return Status::NotImplemented("");
|
||||
virtual void PlasmaObjectReady(const PlasmaObjectReadyRequest &request,
|
||||
const ClientCallback<PlasmaObjectReadyReply> &callback) {
|
||||
}
|
||||
|
||||
virtual ~CoreWorkerClientInterface(){};
|
||||
@@ -227,40 +196,41 @@ class CoreWorkerClient : public std::enable_shared_from_this<CoreWorkerClient>,
|
||||
|
||||
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<PushTaskRequest> request, bool skip_queue,
|
||||
const ClientCallback<PushTaskReply> &callback) override {
|
||||
void PushActorTask(std::unique_ptr<PushTaskRequest> request, bool skip_queue,
|
||||
const ClientCallback<PushTaskReply> &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<CoreWorkerClient>,
|
||||
send_queue_.push_back(std::make_pair(std::move(request), callback));
|
||||
}
|
||||
SendRequests();
|
||||
return ray::Status::OK();
|
||||
}
|
||||
|
||||
ray::Status PushNormalTask(std::unique_ptr<PushTaskRequest> request,
|
||||
const ClientCallback<PushTaskReply> &callback) override {
|
||||
void PushNormalTask(std::unique_ptr<PushTaskRequest> request,
|
||||
const ClientCallback<PushTaskReply> &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.
|
||||
|
||||
Reference in New Issue
Block a user