diff --git a/src/ray/gcs/gcs_server/job_info_handler_impl.cc b/src/ray/gcs/gcs_server/gcs_job_manager.cc similarity index 79% rename from src/ray/gcs/gcs_server/job_info_handler_impl.cc rename to src/ray/gcs/gcs_server/gcs_job_manager.cc index 71db7c4d9..b70b2ad91 100644 --- a/src/ray/gcs/gcs_server/job_info_handler_impl.cc +++ b/src/ray/gcs/gcs_server/gcs_job_manager.cc @@ -12,15 +12,16 @@ // See the License for the specific language governing permissions and // limitations under the License. -#include "job_info_handler_impl.h" +#include "gcs_job_manager.h" + #include "ray/gcs/pb_util.h" namespace ray { -namespace rpc { +namespace gcs { -void GcsJobInfoHandler::HandleAddJob(const rpc::AddJobRequest &request, - rpc::AddJobReply *reply, - rpc::SendReplyCallback send_reply_callback) { +void GcsJobManager::HandleAddJob(const rpc::AddJobRequest &request, + rpc::AddJobReply *reply, + rpc::SendReplyCallback send_reply_callback) { JobID job_id = JobID::FromBinary(request.data().job_id()); RAY_LOG(INFO) << "Adding job, job id = " << job_id << ", driver pid = " << request.data().driver_pid(); @@ -41,9 +42,9 @@ void GcsJobInfoHandler::HandleAddJob(const rpc::AddJobRequest &request, } } -void GcsJobInfoHandler::HandleMarkJobFinished( - const rpc::MarkJobFinishedRequest &request, rpc::MarkJobFinishedReply *reply, - rpc::SendReplyCallback send_reply_callback) { +void GcsJobManager::HandleMarkJobFinished(const rpc::MarkJobFinishedRequest &request, + rpc::MarkJobFinishedReply *reply, + rpc::SendReplyCallback send_reply_callback) { JobID job_id = JobID::FromBinary(request.job_id()); RAY_LOG(INFO) << "Marking job state, job id = " << job_id; auto job_table_data = @@ -67,7 +68,7 @@ void GcsJobInfoHandler::HandleMarkJobFinished( } } -void GcsJobInfoHandler::ClearJobInfos(const JobID &job_id) { +void GcsJobManager::ClearJobInfos(const JobID &job_id) { // Notify all listeners. for (auto &listener : job_finished_listeners_) { listener(std::make_shared(job_id)); @@ -77,15 +78,15 @@ void GcsJobInfoHandler::ClearJobInfos(const JobID &job_id) { /// Add listener to monitor the add action of nodes. /// /// \param listener The handler which process the add of nodes. -void GcsJobInfoHandler::AddJobFinishedListener( +void GcsJobManager::AddJobFinishedListener( std::function)> listener) { RAY_CHECK(listener); job_finished_listeners_.emplace_back(std::move(listener)); } -void GcsJobInfoHandler::HandleGetAllJobInfo(const rpc::GetAllJobInfoRequest &request, - rpc::GetAllJobInfoReply *reply, - rpc::SendReplyCallback send_reply_callback) { +void GcsJobManager::HandleGetAllJobInfo(const rpc::GetAllJobInfoRequest &request, + rpc::GetAllJobInfoReply *reply, + rpc::SendReplyCallback send_reply_callback) { RAY_LOG(INFO) << "Getting all job info."; auto on_done = [reply, send_reply_callback]( const std::unordered_map &result) { @@ -101,5 +102,5 @@ void GcsJobInfoHandler::HandleGetAllJobInfo(const rpc::GetAllJobInfoRequest &req } } -} // namespace rpc +} // namespace gcs } // namespace ray diff --git a/src/ray/gcs/gcs_server/job_info_handler_impl.h b/src/ray/gcs/gcs_server/gcs_job_manager.h similarity index 63% rename from src/ray/gcs/gcs_server/job_info_handler_impl.h rename to src/ray/gcs/gcs_server/gcs_job_manager.h index 0aee366c5..db9f129a5 100644 --- a/src/ray/gcs/gcs_server/job_info_handler_impl.h +++ b/src/ray/gcs/gcs_server/gcs_job_manager.h @@ -21,25 +21,26 @@ #include "ray/rpc/gcs_server/gcs_rpc_server.h" namespace ray { -namespace rpc { +namespace gcs { /// This implementation class of `JobInfoHandler`. -class GcsJobInfoHandler : public rpc::JobInfoHandler { +class GcsJobManager : public rpc::JobInfoHandler { public: - explicit GcsJobInfoHandler(std::shared_ptr gcs_table_storage, - std::shared_ptr gcs_pub_sub) + explicit GcsJobManager(std::shared_ptr gcs_table_storage, + std::shared_ptr gcs_pub_sub) : gcs_table_storage_(std::move(gcs_table_storage)), gcs_pub_sub_(std::move(gcs_pub_sub)) {} - void HandleAddJob(const AddJobRequest &request, AddJobReply *reply, - SendReplyCallback send_reply_callback) override; + void HandleAddJob(const rpc::AddJobRequest &request, rpc::AddJobReply *reply, + rpc::SendReplyCallback send_reply_callback) override; - void HandleMarkJobFinished(const MarkJobFinishedRequest &request, - MarkJobFinishedReply *reply, - SendReplyCallback send_reply_callback) override; + void HandleMarkJobFinished(const rpc::MarkJobFinishedRequest &request, + rpc::MarkJobFinishedReply *reply, + rpc::SendReplyCallback send_reply_callback) override; - void HandleGetAllJobInfo(const GetAllJobInfoRequest &request, GetAllJobInfoReply *reply, - SendReplyCallback send_reply_callback) override; + void HandleGetAllJobInfo(const rpc::GetAllJobInfoRequest &request, + rpc::GetAllJobInfoReply *reply, + rpc::SendReplyCallback send_reply_callback) override; void AddJobFinishedListener( std::function)> listener) override; @@ -54,5 +55,5 @@ class GcsJobInfoHandler : public rpc::JobInfoHandler { void ClearJobInfos(const JobID &job_id); }; -} // namespace rpc +} // namespace gcs } // namespace ray diff --git a/src/ray/gcs/gcs_server/gcs_server.cc b/src/ray/gcs/gcs_server/gcs_server.cc index 39c219244..3e5cfe2c6 100644 --- a/src/ray/gcs/gcs_server/gcs_server.cc +++ b/src/ray/gcs/gcs_server/gcs_server.cc @@ -16,9 +16,9 @@ #include "error_info_handler_impl.h" #include "gcs_actor_manager.h" +#include "gcs_job_manager.h" #include "gcs_node_manager.h" #include "gcs_object_manager.h" -#include "job_info_handler_impl.h" #include "ray/common/network_util.h" #include "ray/common/ray_config.h" #include "stats_handler_impl.h" @@ -71,8 +71,8 @@ void GcsServer::Start() { new rpc::TaskInfoGrpcService(main_service_, *task_info_handler_)); rpc_server_.RegisterService(*task_info_service_); - InitJobInfoHandler(); - job_info_service_.reset(new rpc::JobInfoGrpcService(main_service_, *job_info_handler_)); + InitGcsJobManager(); + job_info_service_.reset(new rpc::JobInfoGrpcService(main_service_, *gcs_job_manager_)); rpc_server_.RegisterService(*job_info_service_); actor_info_service_.reset( @@ -198,10 +198,10 @@ void GcsServer::InitGcsActorManager() { RAY_CHECK_OK(gcs_pub_sub_->SubscribeAll(WORKER_FAILURE_CHANNEL, on_subscribe, nullptr)); } -void GcsServer::InitJobInfoHandler() { - job_info_handler_ = std::unique_ptr( - new rpc::GcsJobInfoHandler(gcs_table_storage_, gcs_pub_sub_)); - job_info_handler_->AddJobFinishedListener([this](std::shared_ptr job_id) { +void GcsServer::InitGcsJobManager() { + gcs_job_manager_ = + std::unique_ptr(new GcsJobManager(gcs_table_storage_, gcs_pub_sub_)); + gcs_job_manager_->AddJobFinishedListener([this](std::shared_ptr job_id) { gcs_actor_manager_->OnJobFinished(*job_id); }); } diff --git a/src/ray/gcs/gcs_server/gcs_server.h b/src/ray/gcs/gcs_server/gcs_server.h index e4227f8ef..00d2ba197 100644 --- a/src/ray/gcs/gcs_server/gcs_server.h +++ b/src/ray/gcs/gcs_server/gcs_server.h @@ -38,6 +38,7 @@ struct GcsServerConfig { class GcsNodeManager; class GcsActorManager; +class GcsJobManager; /// The GcsServer will take over all requests from ServiceBasedGcsClient and transparent /// transmit the command to the backend reliable storage for the time being. @@ -80,8 +81,8 @@ class GcsServer { /// Initialize the gcs actor manager. virtual void InitGcsActorManager(); - /// The job info handler - virtual void InitJobInfoHandler(); + /// Initialize the gcs job manager. + virtual void InitGcsJobManager(); /// The object manager virtual std::unique_ptr InitObjectManager(); @@ -121,7 +122,7 @@ class GcsServer { /// The gcs actor manager std::shared_ptr gcs_actor_manager_; /// Job info handler and service - std::unique_ptr job_info_handler_; + std::unique_ptr gcs_job_manager_; std::unique_ptr job_info_service_; /// Actor info service std::unique_ptr actor_info_service_;