Skip to content

Commit 8109fcd

Browse files
authored
[Controller] Add gRPC support for isStoreMigrationAllowed API (linkedin#2409)
1 parent 1225a93 commit 8109fcd

9 files changed

Lines changed: 264 additions & 3 deletions

File tree

internal/venice-common/src/main/proto/controller/ClusterAdminOpsGrpcService.proto

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,9 @@ service ClusterAdminOpsGrpcService {
1515
rpc getAdminTopicMetadata(AdminTopicMetadataGrpcRequest) returns (AdminTopicMetadataGrpcResponse) {}
1616
rpc updateAdminTopicMetadata(UpdateAdminTopicMetadataGrpcRequest) returns (AdminTopicMetadataGrpcResponse) {}
1717
rpc updateAdminOperationProtocolVersion(UpdateAdminOperationProtocolVersionGrpcRequest) returns (AdminTopicMetadataGrpcResponse) {}
18+
19+
// StoreMigration
20+
rpc isStoreMigrationAllowed(StoreMigrationCheckGrpcRequest) returns (StoreMigrationCheckGrpcResponse) {}
1821
}
1922

2023

@@ -72,4 +75,13 @@ message AdminTopicGrpcMetadata {
7275
message PubSubPositionGrpcWireFormat {
7376
int32 typeId = 1;
7477
string base64PositionBytes = 2;
75-
}
78+
}
79+
80+
message StoreMigrationCheckGrpcRequest {
81+
string clusterName = 1;
82+
}
83+
84+
message StoreMigrationCheckGrpcResponse {
85+
string clusterName = 1;
86+
bool storeMigrationAllowed = 2;
87+
}

internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/endToEnd/TestControllerGrpcEndpoints.java

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
import com.linkedin.venice.integration.utils.ServiceFactory;
1515
import com.linkedin.venice.integration.utils.VeniceClusterCreateOptions;
1616
import com.linkedin.venice.integration.utils.VeniceClusterWrapper;
17+
import com.linkedin.venice.protocols.controller.ClusterAdminOpsGrpcServiceGrpc;
1718
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
1819
import com.linkedin.venice.protocols.controller.CreateStoreGrpcRequest;
1920
import com.linkedin.venice.protocols.controller.CreateStoreGrpcResponse;
@@ -24,6 +25,8 @@
2425
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
2526
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
2627
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
28+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
29+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
2730
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
2831
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
2932
import com.linkedin.venice.protocols.controller.VeniceControllerGrpcServiceGrpc;
@@ -217,6 +220,24 @@ public void testValidateStoreDeletedGrpcEndpoint() {
217220
});
218221
}
219222

223+
@Test(timeOut = TIMEOUT_MS)
224+
public void testIsStoreMigrationAllowedGrpcEndpoint() {
225+
String controllerGrpcUrl = veniceCluster.getLeaderVeniceController().getControllerGrpcUrl();
226+
ManagedChannel channel = Grpc.newChannelBuilder(controllerGrpcUrl, InsecureChannelCredentials.create()).build();
227+
ClusterAdminOpsGrpcServiceGrpc.ClusterAdminOpsGrpcServiceBlockingStub blockingStub =
228+
ClusterAdminOpsGrpcServiceGrpc.newBlockingStub(channel);
229+
230+
StoreMigrationCheckGrpcRequest request =
231+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(veniceCluster.getClusterName()).build();
232+
233+
StoreMigrationCheckGrpcResponse response = blockingStub.isStoreMigrationAllowed(request);
234+
235+
assertNotNull(response, "Response should not be null");
236+
assertEquals(response.getClusterName(), veniceCluster.getClusterName());
237+
// By default, store migration is allowed in test clusters
238+
assertTrue(response.getStoreMigrationAllowed(), "Store migration should be allowed by default");
239+
}
240+
220241
@Test(timeOut = TIMEOUT_MS)
221242
public void testValidateStoreDeletedOverSecureGrpcChannel() {
222243
String storeName = Utils.getUniqueString("test_validate_deleted_secure");

services/venice-controller/src/main/java/com/linkedin/venice/controller/grpc/server/ClusterAdminOpsGrpcServiceImpl.java

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@
1515
import com.linkedin.venice.protocols.controller.ClusterAdminOpsGrpcServiceGrpc;
1616
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcRequest;
1717
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcResponse;
18+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
19+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
1820
import com.linkedin.venice.protocols.controller.UpdateAdminOperationProtocolVersionGrpcRequest;
1921
import com.linkedin.venice.protocols.controller.UpdateAdminTopicMetadataGrpcRequest;
2022
import io.grpc.Context;
@@ -102,4 +104,17 @@ public void updateAdminOperationProtocolVersion(
102104
request.getClusterName(),
103105
null);
104106
}
107+
108+
@Override
109+
public void isStoreMigrationAllowed(
110+
StoreMigrationCheckGrpcRequest request,
111+
StreamObserver<StoreMigrationCheckGrpcResponse> responseObserver) {
112+
LOGGER.debug("Received isStoreMigrationAllowed request: {}", request);
113+
ControllerGrpcServerUtils.handleRequest(
114+
ClusterAdminOpsGrpcServiceGrpc.getIsStoreMigrationAllowedMethod(),
115+
() -> requestHandler.isStoreMigrationAllowed(request),
116+
responseObserver,
117+
request.getClusterName(),
118+
null);
119+
}
105120
}

services/venice-controller/src/main/java/com/linkedin/venice/controller/server/AdminSparkServer.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -314,7 +314,8 @@ public boolean startInner() throws Exception {
314314
new RoutersClusterConfigRoutes(sslEnabled, accessController);
315315
MigrationRoutes migrationRoutes = new MigrationRoutes(sslEnabled, accessController);
316316
VersionRoute versionRoute = new VersionRoute(sslEnabled, accessController);
317-
ClusterRoutes clusterRoutes = new ClusterRoutes(sslEnabled, accessController);
317+
ClusterRoutes clusterRoutes =
318+
new ClusterRoutes(sslEnabled, accessController, requestHandler.getClusterAdminOpsRequestHandler());
318319
NewClusterBuildOutRoutes newClusterBuildOutRoutes = new NewClusterBuildOutRoutes(sslEnabled, accessController);
319320
DataRecoveryRoutes dataRecoveryRoutes = new DataRecoveryRoutes(sslEnabled, accessController);
320321
AdminTopicMetadataRoutes adminTopicMetadataRoutes = new AdminTopicMetadataRoutes(sslEnabled, accessController);

services/venice-controller/src/main/java/com/linkedin/venice/controller/server/ClusterAdminOpsRequestHandler.java

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@
1818
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcRequest;
1919
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcResponse;
2020
import com.linkedin.venice.protocols.controller.PubSubPositionGrpcWireFormat;
21+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
22+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
2123
import com.linkedin.venice.protocols.controller.UpdateAdminOperationProtocolVersionGrpcRequest;
2224
import com.linkedin.venice.protocols.controller.UpdateAdminTopicMetadataGrpcRequest;
2325
import com.linkedin.venice.pubsub.api.PubSubPosition;
@@ -179,4 +181,17 @@ public AdminTopicMetadataGrpcResponse updateAdminOperationProtocolVersion(
179181
.setAdminOperationProtocolVersion(adminOperationProtocolVersion);
180182
return AdminTopicMetadataGrpcResponse.newBuilder().setMetadata(adminMetadataBuilder.build()).build();
181183
}
184+
185+
public StoreMigrationCheckGrpcResponse isStoreMigrationAllowed(StoreMigrationCheckGrpcRequest request) {
186+
String clusterName = request.getClusterName();
187+
if (StringUtils.isBlank(clusterName)) {
188+
throw new IllegalArgumentException("Cluster name is required for checking if store migration is allowed");
189+
}
190+
LOGGER.info("Checking if store migration is allowed for cluster: {}", clusterName);
191+
boolean isAllowed = admin.isStoreMigrationAllowed(clusterName);
192+
return StoreMigrationCheckGrpcResponse.newBuilder()
193+
.setClusterName(clusterName)
194+
.setStoreMigrationAllowed(isAllowed)
195+
.build();
196+
}
182197
}

services/venice-controller/src/main/java/com/linkedin/venice/controller/server/ClusterRoutes.java

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@
1717
import com.linkedin.venice.controllerapi.StoreMigrationResponse;
1818
import com.linkedin.venice.controllerapi.UpdateClusterConfigQueryParams;
1919
import com.linkedin.venice.controllerapi.UpdateDarkClusterConfigQueryParams;
20+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
21+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
2022
import com.linkedin.venice.utils.Utils;
2123
import java.util.Map;
2224
import java.util.Optional;
@@ -25,8 +27,18 @@
2527

2628

2729
public class ClusterRoutes extends AbstractRoute {
30+
private final ClusterAdminOpsRequestHandler clusterAdminOpsRequestHandler;
31+
2832
public ClusterRoutes(boolean sslEnabled, Optional<DynamicAccessController> accessController) {
33+
this(sslEnabled, accessController, null);
34+
}
35+
36+
public ClusterRoutes(
37+
boolean sslEnabled,
38+
Optional<DynamicAccessController> accessController,
39+
ClusterAdminOpsRequestHandler clusterAdminOpsRequestHandler) {
2940
super(sslEnabled, accessController);
41+
this.clusterAdminOpsRequestHandler = clusterAdminOpsRequestHandler;
3042
}
3143

3244
/**
@@ -91,7 +103,14 @@ public void internalHandle(Request request, StoreMigrationResponse veniceRespons
91103
AdminSparkServer.validateParams(request, STORE_MIGRATION_ALLOWED.getParams(), admin);
92104
String clusterName = request.queryParams(CLUSTER);
93105
veniceResponse.setCluster(clusterName);
94-
veniceResponse.setStoreMigrationAllowed(admin.isStoreMigrationAllowed(clusterName));
106+
107+
StoreMigrationCheckGrpcRequest grpcRequest =
108+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(clusterName).build();
109+
110+
StoreMigrationCheckGrpcResponse grpcResponse =
111+
clusterAdminOpsRequestHandler.isStoreMigrationAllowed(grpcRequest);
112+
113+
veniceResponse.setStoreMigrationAllowed(grpcResponse.getStoreMigrationAllowed());
95114
}
96115
};
97116
}

services/venice-controller/src/test/java/com/linkedin/venice/controller/grpc/server/ClusterAdminOpsGrpcServiceImplTest.java

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@
2525
import com.linkedin.venice.protocols.controller.ClusterAdminOpsGrpcServiceGrpc.ClusterAdminOpsGrpcServiceBlockingStub;
2626
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcRequest;
2727
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcResponse;
28+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
29+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
2830
import com.linkedin.venice.protocols.controller.UpdateAdminOperationProtocolVersionGrpcRequest;
2931
import com.linkedin.venice.protocols.controller.UpdateAdminTopicMetadataGrpcRequest;
3032
import com.linkedin.venice.protocols.controller.VeniceControllerGrpcErrorInfo;
@@ -213,4 +215,51 @@ public void testUpdateAdminOperationProtocolVersionSuccess() {
213215
assertEquals(actualResponse.getMetadata().getClusterName(), TEST_CLUSTER);
214216
assertEquals(actualResponse.getMetadata().getAdminOperationProtocolVersion(), 1L);
215217
}
218+
219+
@Test
220+
public void testIsStoreMigrationAllowedSuccess() {
221+
StoreMigrationCheckGrpcResponse response = StoreMigrationCheckGrpcResponse.newBuilder()
222+
.setClusterName(TEST_CLUSTER)
223+
.setStoreMigrationAllowed(true)
224+
.build();
225+
doReturn(response).when(requestHandler).isStoreMigrationAllowed(any(StoreMigrationCheckGrpcRequest.class));
226+
227+
StoreMigrationCheckGrpcRequest request =
228+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(TEST_CLUSTER).build();
229+
230+
StoreMigrationCheckGrpcResponse actualResponse = blockingStub.isStoreMigrationAllowed(request);
231+
assertNotNull(actualResponse);
232+
assertEquals(actualResponse.getClusterName(), TEST_CLUSTER);
233+
assertTrue(actualResponse.getStoreMigrationAllowed());
234+
}
235+
236+
@Test
237+
public void testIsStoreMigrationAllowedReturnsFalse() {
238+
StoreMigrationCheckGrpcResponse response = StoreMigrationCheckGrpcResponse.newBuilder()
239+
.setClusterName(TEST_CLUSTER)
240+
.setStoreMigrationAllowed(false)
241+
.build();
242+
doReturn(response).when(requestHandler).isStoreMigrationAllowed(any(StoreMigrationCheckGrpcRequest.class));
243+
244+
StoreMigrationCheckGrpcRequest request =
245+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(TEST_CLUSTER).build();
246+
247+
StoreMigrationCheckGrpcResponse actualResponse = blockingStub.isStoreMigrationAllowed(request);
248+
assertNotNull(actualResponse);
249+
assertEquals(actualResponse.getClusterName(), TEST_CLUSTER);
250+
assertFalse(actualResponse.getStoreMigrationAllowed());
251+
}
252+
253+
@Test
254+
public void testIsStoreMigrationAllowedError() {
255+
StoreMigrationCheckGrpcRequest request =
256+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(TEST_CLUSTER).build();
257+
258+
doThrow(new VeniceException("Error")).when(requestHandler)
259+
.isStoreMigrationAllowed(any(StoreMigrationCheckGrpcRequest.class));
260+
261+
StatusRuntimeException e =
262+
expectThrows(StatusRuntimeException.class, () -> blockingStub.isStoreMigrationAllowed(request));
263+
assertEquals(e.getStatus().getCode(), Status.INTERNAL.getCode());
264+
}
216265
}

services/venice-controller/src/test/java/com/linkedin/venice/controller/server/ClusterAdminOpsRequestHandlerTest.java

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@
2525
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcRequest;
2626
import com.linkedin.venice.protocols.controller.LastSuccessfulAdminCommandExecutionGrpcResponse;
2727
import com.linkedin.venice.protocols.controller.PubSubPositionGrpcWireFormat;
28+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
29+
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
2830
import com.linkedin.venice.protocols.controller.UpdateAdminOperationProtocolVersionGrpcRequest;
2931
import com.linkedin.venice.protocols.controller.UpdateAdminTopicMetadataGrpcRequest;
3032
import com.linkedin.venice.pubsub.PubSubUtil;
@@ -316,4 +318,44 @@ public void testUpdateAdminOperationProtocolVersionInvalidInputs() {
316318
exception.getMessage()
317319
.contains("Admin operation protocol version is required and must be -1 or greater than 0"));
318320
}
321+
322+
@Test
323+
public void testIsStoreMigrationAllowedReturnsTrue() {
324+
String clusterName = "test-cluster";
325+
326+
when(mockAdmin.isStoreMigrationAllowed(clusterName)).thenReturn(true);
327+
328+
StoreMigrationCheckGrpcRequest request =
329+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(clusterName).build();
330+
331+
StoreMigrationCheckGrpcResponse response = handler.isStoreMigrationAllowed(request);
332+
333+
assertNotNull(response);
334+
assertEquals(response.getClusterName(), clusterName);
335+
assertTrue(response.getStoreMigrationAllowed());
336+
}
337+
338+
@Test
339+
public void testIsStoreMigrationAllowedReturnsFalse() {
340+
String clusterName = "test-cluster";
341+
342+
when(mockAdmin.isStoreMigrationAllowed(clusterName)).thenReturn(false);
343+
344+
StoreMigrationCheckGrpcRequest request =
345+
StoreMigrationCheckGrpcRequest.newBuilder().setClusterName(clusterName).build();
346+
347+
StoreMigrationCheckGrpcResponse response = handler.isStoreMigrationAllowed(request);
348+
349+
assertNotNull(response);
350+
assertEquals(response.getClusterName(), clusterName);
351+
assertFalse(response.getStoreMigrationAllowed());
352+
}
353+
354+
@Test
355+
public void testIsStoreMigrationAllowedInvalidCluster() {
356+
StoreMigrationCheckGrpcRequest request = StoreMigrationCheckGrpcRequest.newBuilder().setClusterName("").build();
357+
358+
Exception exception = expectThrows(IllegalArgumentException.class, () -> handler.isStoreMigrationAllowed(request));
359+
assertEquals(exception.getMessage(), "Cluster name is required for checking if store migration is allowed");
360+
}
319361
}

0 commit comments

Comments
 (0)