-
Notifications
You must be signed in to change notification settings - Fork 89
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
[controller][grpc] Add ClusterAdminOpsGrpcService with handler and tests
- Introduced `ClusterAdminOpsGrpcServiceImpl` for gRPC support of cluster admin operations. - Added `ClusterAdminOpsRequestHandler` for handling gRPC requests. - Implemented all relevant methods in the service, including admin command status, metadata handling, and execution ID retrieval. - Provided comprehensive TestNG tests for the service and request handler.
- Loading branch information
1 parent
7442860
commit 0c69ab5
Showing
15 changed files
with
1,341 additions
and
60 deletions.
There are no files selected for viewing
66 changes: 66 additions & 0 deletions
66
internal/venice-common/src/main/proto/controller/ClusterAdminOpsGrpcService.proto
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,66 @@ | ||
syntax = 'proto3'; | ||
package com.linkedin.venice.protocols.controller; | ||
|
||
|
||
import "controller/ControllerGrpcRequestContext.proto"; | ||
|
||
option java_multiple_files = true; | ||
|
||
service ClusterAdminOpsGrpcService { | ||
// AdminCommandExecution | ||
rpc getAdminCommandExecutionStatus(AdminCommandExecutionStatusGrpcRequest) returns (AdminCommandExecutionStatusGrpcResponse) {} | ||
rpc getLastSuccessfulAdminCommandExecutionId(LastSuccessfulAdminCommandExecutionGrpcRequest) returns (LastSuccessfulAdminCommandExecutionGrpcResponse) {} | ||
|
||
// AdminTopicMetadata | ||
rpc getAdminTopicMetadata(AdminTopicMetadataGrpcRequest) returns (AdminTopicMetadataGrpcResponse) {} | ||
rpc updateAdminTopicMetadata(UpdateAdminTopicMetadataGrpcRequest) returns (UpdateAdminTopicMetadataGrpcResponse) {} | ||
} | ||
|
||
|
||
message AdminCommandExecutionStatusGrpcRequest { | ||
string clusterName = 1; | ||
int64 adminCommandExecutionId = 2; | ||
} | ||
|
||
message AdminCommandExecutionStatusGrpcResponse { | ||
string clusterName = 1; | ||
int64 adminCommandExecutionId = 2; | ||
string operation = 3; | ||
string startTime = 4; | ||
map<string, string> fabricToExecutionStatusMap = 5; | ||
} | ||
|
||
message LastSuccessfulAdminCommandExecutionGrpcRequest { | ||
string clusterName = 1; | ||
} | ||
|
||
message LastSuccessfulAdminCommandExecutionGrpcResponse { | ||
string clusterName = 1; | ||
int64 lastSuccessfulAdminCommandExecutionId = 2; | ||
} | ||
|
||
message AdminTopicMetadataGrpcRequest { | ||
string clusterName = 1; | ||
optional string storeName = 2; | ||
} | ||
|
||
message AdminTopicMetadataGrpcResponse { | ||
AdminTopicGrpcMetadata metadata = 1; | ||
} | ||
|
||
message UpdateAdminTopicMetadataGrpcResponse { | ||
string clusterName = 1; | ||
optional string storeName = 2; | ||
} | ||
|
||
message UpdateAdminTopicMetadataGrpcRequest { | ||
AdminTopicGrpcMetadata metadata = 1; | ||
} | ||
|
||
message AdminTopicGrpcMetadata { | ||
string clusterName = 1; | ||
int64 executionId = 2; | ||
optional string storeName = 3; | ||
optional int64 offset = 4; | ||
optional int64 upstreamOffset = 5; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
92 changes: 92 additions & 0 deletions
92
.../main/java/com/linkedin/venice/controller/grpc/server/ClusterAdminOpsGrpcServiceImpl.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,92 @@ | ||
package com.linkedin.venice.controller.grpc.server; | ||
|
||
import static com.linkedin.venice.controller.grpc.server.ControllerGrpcServerUtils.isAllowListUser; | ||
import static com.linkedin.venice.controller.server.VeniceRouteHandler.ACL_CHECK_FAILURE_WARN_MESSAGE_PREFIX; | ||
import static com.linkedin.venice.protocols.controller.ClusterAdminOpsGrpcServiceGrpc.*; | ||
|
||
import com.linkedin.venice.controller.server.ClusterAdminOpsRequestHandler; | ||
import com.linkedin.venice.controller.server.VeniceControllerAccessManager; | ||
import com.linkedin.venice.exceptions.VeniceUnauthorizedAccessException; | ||
import com.linkedin.venice.protocols.controller.AdminCommandExecutionStatusGrpcRequest; | ||
import com.linkedin.venice.protocols.controller.AdminCommandExecutionStatusGrpcResponse; | ||
import com.linkedin.venice.protocols.controller.AdminTopicGrpcMetadata; | ||
import com.linkedin.venice.protocols.controller.AdminTopicMetadataGrpcRequest; | ||
import com.linkedin.venice.protocols.controller.AdminTopicMetadataGrpcResponse; | ||
import com.linkedin.venice.protocols.controller.ClusterAdminOpsGrpcServiceGrpc; | ||
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcRequest; | ||
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcResponse; | ||
import com.linkedin.venice.protocols.controller.UpdateAdminTopicMetadataGrpcRequest; | ||
import com.linkedin.venice.protocols.controller.UpdateAdminTopicMetadataGrpcResponse; | ||
import io.grpc.Context; | ||
import io.grpc.stub.StreamObserver; | ||
import org.apache.logging.log4j.LogManager; | ||
import org.apache.logging.log4j.Logger; | ||
|
||
|
||
public class ClusterAdminOpsGrpcServiceImpl extends ClusterAdminOpsGrpcServiceImplBase { | ||
private static final Logger LOGGER = LogManager.getLogger(ClusterAdminOpsGrpcServiceImpl.class); | ||
private final ClusterAdminOpsRequestHandler requestHandler; | ||
private final VeniceControllerAccessManager accessManager; | ||
|
||
public ClusterAdminOpsGrpcServiceImpl( | ||
ClusterAdminOpsRequestHandler requestHandler, | ||
VeniceControllerAccessManager accessManager) { | ||
this.requestHandler = requestHandler; | ||
this.accessManager = accessManager; | ||
} | ||
|
||
@Override | ||
public void getAdminCommandExecutionStatus( | ||
AdminCommandExecutionStatusGrpcRequest request, | ||
StreamObserver<AdminCommandExecutionStatusGrpcResponse> responseObserver) { | ||
LOGGER.debug("Received getAdminCommandExecutionStatus request: {}", request); | ||
ControllerGrpcServerUtils.handleRequest( | ||
ClusterAdminOpsGrpcServiceGrpc.getGetAdminCommandExecutionStatusMethod(), | ||
() -> requestHandler.getAdminCommandExecutionStatus(request), | ||
responseObserver, | ||
request.getClusterName(), | ||
null); | ||
} | ||
|
||
@Override | ||
public void getLastSuccessfulAdminCommandExecutionId( | ||
LastSuccessfulAdminCommandExecutionGrpcRequest request, | ||
StreamObserver<LastSuccessfulAdminCommandExecutionGrpcResponse> responseObserver) { | ||
LOGGER.debug("Received getLastSuccessfulAdminCommandExecutionId request: {}", request); | ||
ControllerGrpcServerUtils.handleRequest( | ||
ClusterAdminOpsGrpcServiceGrpc.getGetLastSuccessfulAdminCommandExecutionIdMethod(), | ||
() -> requestHandler.getLastSucceedExecutionId(request), | ||
responseObserver, | ||
request.getClusterName(), | ||
null); | ||
} | ||
|
||
@Override | ||
public void getAdminTopicMetadata( | ||
AdminTopicMetadataGrpcRequest request, | ||
StreamObserver<AdminTopicMetadataGrpcResponse> responseObserver) { | ||
LOGGER.debug("Received getAdminTopicMetadata request: {}", request); | ||
ControllerGrpcServerUtils.handleRequest( | ||
ClusterAdminOpsGrpcServiceGrpc.getGetAdminTopicMetadataMethod(), | ||
() -> requestHandler.getAdminTopicMetadata(request), | ||
responseObserver, | ||
request.getClusterName(), | ||
request.hasStoreName() ? request.getStoreName() : null); | ||
} | ||
|
||
@Override | ||
public void updateAdminTopicMetadata( | ||
UpdateAdminTopicMetadataGrpcRequest request, | ||
StreamObserver<UpdateAdminTopicMetadataGrpcResponse> responseObserver) { | ||
LOGGER.debug("Received updateAdminTopicMetadata request: {}", request); | ||
AdminTopicGrpcMetadata metadata = request.getMetadata(); | ||
ControllerGrpcServerUtils.handleRequest(ClusterAdminOpsGrpcServiceGrpc.getUpdateAdminTopicMetadataMethod(), () -> { | ||
if (!isAllowListUser(accessManager, request.getMetadata().getStoreName(), Context.current())) { | ||
throw new VeniceUnauthorizedAccessException( | ||
ACL_CHECK_FAILURE_WARN_MESSAGE_PREFIX | ||
+ ClusterAdminOpsGrpcServiceGrpc.getUpdateAdminTopicMetadataMethod().getFullMethodName()); | ||
} | ||
return requestHandler.updateAdminTopicMetadata(request); | ||
}, responseObserver, metadata.getClusterName(), metadata.hasStoreName() ? metadata.getStoreName() : null); | ||
} | ||
} |
Oops, something went wrong.