diff --git a/src/ray/gcs/gcs_server/gcs_node_manager.cc b/src/ray/gcs/gcs_server/gcs_node_manager.cc index dd30d3ef8..88884cb68 100644 --- a/src/ray/gcs/gcs_server/gcs_node_manager.cc +++ b/src/ray/gcs/gcs_server/gcs_node_manager.cc @@ -216,6 +216,21 @@ void GcsNodeManager::HandleReportHeartbeat(const rpc::ReportHeartbeatRequest &re ++counts_[CountType::REPORT_HEARTBEAT_REQUEST]; } +// TODO(WangTao): Implenent this to handle resource usage report. Basically move resources +// related operations in `HandleReportHeartbeat`. +void GcsNodeManager::HandleReportResourceUsage( + const rpc::ReportResourceUsageRequest &request, rpc::ReportResourceUsageReply *reply, + rpc::SendReplyCallback send_reply_callback) { + GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); +} + +// TODO(WangTao): Implement this. Basically copy from `HandleGetAllHeartbeat`. +void GcsNodeManager::HandleGetAllResourceUsage( + const rpc::GetAllResourceUsageRequest &request, rpc::GetAllResourceUsageReply *reply, + rpc::SendReplyCallback send_reply_callback) { + GCS_RPC_SEND_REPLY(send_reply_callback, reply, Status::OK()); +} + void GcsNodeManager::HandleGetResources(const rpc::GetResourcesRequest &request, rpc::GetResourcesReply *reply, rpc::SendReplyCallback send_reply_callback) { diff --git a/src/ray/gcs/gcs_server/gcs_node_manager.h b/src/ray/gcs/gcs_server/gcs_node_manager.h index 3e49db4e2..2bad2a0fb 100644 --- a/src/ray/gcs/gcs_server/gcs_node_manager.h +++ b/src/ray/gcs/gcs_server/gcs_node_manager.h @@ -65,6 +65,16 @@ class GcsNodeManager : public rpc::NodeInfoHandler { rpc::ReportHeartbeatReply *reply, rpc::SendReplyCallback send_reply_callback) override; + /// Handle report resource usage rpc come from raylet. + void HandleReportResourceUsage(const rpc::ReportResourceUsageRequest &request, + rpc::ReportResourceUsageReply *reply, + rpc::SendReplyCallback send_reply_callback) override; + + /// Handle get all resource usage rpc request. + void HandleGetAllResourceUsage(const rpc::GetAllResourceUsageRequest &request, + rpc::GetAllResourceUsageReply *reply, + rpc::SendReplyCallback send_reply_callback) override; + /// Handle get resource rpc request. void HandleGetResources(const rpc::GetResourcesRequest &request, rpc::GetResourcesReply *reply, diff --git a/src/ray/protobuf/gcs.proto b/src/ray/protobuf/gcs.proto index fd4b4cd2a..9182279e0 100644 --- a/src/ray/protobuf/gcs.proto +++ b/src/ray/protobuf/gcs.proto @@ -322,6 +322,27 @@ message HeartbeatTableData { bool should_global_gc = 8; } +message ResourcesData { + // Node manager client id + bytes client_id = 1; + // Resource capacity currently available on this node manager. + map resources_available = 2; + // Indicates whether avaialbe resources is changed. Only used when + // light heartbeat enabled. + bool resources_available_changed = 3; + // Total resource capacity configured for this node manager. + map resources_total = 4; + // Aggregate outstanding resource load on this node manager. + map resource_load = 5; + // Indicates whether resource load is changed. Only used when + // light heartbeat enabled. + bool resource_load_changed = 6; + // The resource load on this node, sorted by resource shape. + ResourceLoad resource_load_by_shape = 7; + // Whether this node manager is requesting global GC. + bool should_global_gc = 8; +} + message HeartbeatBatchTableData { repeated HeartbeatTableData batch = 1; // The total resource demand on all nodes included in the batch, sorted by diff --git a/src/ray/protobuf/gcs_service.proto b/src/ray/protobuf/gcs_service.proto index bcccc9b4e..646c1ccb9 100644 --- a/src/ray/protobuf/gcs_service.proto +++ b/src/ray/protobuf/gcs_service.proto @@ -206,6 +206,22 @@ message GetAllHeartbeatReply { HeartbeatBatchTableData heartbeat_data = 2; } +message ReportResourceUsageRequest { + ResourcesData resources = 1; +} + +message ReportResourceUsageReply { + GcsStatus status = 1; +} + +message GetAllResourceUsageRequest { +} + +message GetAllResourceUsageReply { + GcsStatus status = 1; + repeated ResourcesData resources_list = 2; +} + message GetResourcesRequest { bytes node_id = 1; } @@ -267,9 +283,13 @@ service NodeInfoGcsService { rpc GetAllNodeInfo(GetAllNodeInfoRequest) returns (GetAllNodeInfoReply); // Report heartbeat of a node to GCS Service. rpc ReportHeartbeat(ReportHeartbeatRequest) returns (ReportHeartbeatReply); - // Get node's resources from GCS Service. // Get newest heartbeat of all nodes from GCS Service. rpc GetAllHeartbeat(GetAllHeartbeatRequest) returns (GetAllHeartbeatReply); + // Report resource usage of a node to GCS Service. + rpc ReportResourceUsage(ReportResourceUsageRequest) returns (ReportResourceUsageReply); + // Get resource usage of all nodes from GCS Service. + rpc GetAllResourceUsage(GetAllResourceUsageRequest) returns (GetAllResourceUsageReply); + // Get node's resources from GCS Service. rpc GetResources(GetResourcesRequest) returns (GetResourcesReply); // Update resources of a node in GCS Service. rpc UpdateResources(UpdateResourcesRequest) returns (UpdateResourcesReply); diff --git a/src/ray/rpc/gcs_server/gcs_rpc_client.h b/src/ray/rpc/gcs_server/gcs_rpc_client.h index 4f02e0480..8539b589a 100644 --- a/src/ray/rpc/gcs_server/gcs_rpc_client.h +++ b/src/ray/rpc/gcs_server/gcs_rpc_client.h @@ -176,6 +176,14 @@ class GcsRpcClient { VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetAllHeartbeat, node_info_grpc_client_, ) + /// Report resource usage of a node to GCS Service. + VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, ReportResourceUsage, + node_info_grpc_client_, ) + + /// Get resource usage of all nodes from GCS Service. + VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetAllResourceUsage, + node_info_grpc_client_, ) + /// Get node's resources from GCS Service. VOID_GCS_RPC_CLIENT_METHOD(NodeInfoGcsService, GetResources, node_info_grpc_client_, ) diff --git a/src/ray/rpc/gcs_server/gcs_rpc_server.h b/src/ray/rpc/gcs_server/gcs_rpc_server.h index f7132f02b..d2b7bd62e 100644 --- a/src/ray/rpc/gcs_server/gcs_rpc_server.h +++ b/src/ray/rpc/gcs_server/gcs_rpc_server.h @@ -200,6 +200,14 @@ class NodeInfoGcsServiceHandler { GetAllHeartbeatReply *reply, SendReplyCallback send_reply_callback) = 0; + virtual void HandleReportResourceUsage(const ReportResourceUsageRequest &request, + ReportResourceUsageReply *reply, + SendReplyCallback send_reply_callback) = 0; + + virtual void HandleGetAllResourceUsage(const GetAllResourceUsageRequest &request, + GetAllResourceUsageReply *reply, + SendReplyCallback send_reply_callback) = 0; + virtual void HandleGetResources(const GetResourcesRequest &request, GetResourcesReply *reply, SendReplyCallback send_reply_callback) = 0; @@ -247,6 +255,8 @@ class NodeInfoGrpcService : public GrpcService { NODE_INFO_SERVICE_RPC_HANDLER(GetAllNodeInfo); NODE_INFO_SERVICE_RPC_HANDLER(ReportHeartbeat); NODE_INFO_SERVICE_RPC_HANDLER(GetAllHeartbeat); + NODE_INFO_SERVICE_RPC_HANDLER(ReportResourceUsage); + NODE_INFO_SERVICE_RPC_HANDLER(GetAllResourceUsage); NODE_INFO_SERVICE_RPC_HANDLER(GetResources); NODE_INFO_SERVICE_RPC_HANDLER(UpdateResources); NODE_INFO_SERVICE_RPC_HANDLER(DeleteResources);