Skip to content

Commit d8eef14

Browse files
authored
[Controller] Add gRPC support for getClusterHealthStores API (linkedin#2421)
Add gRPC endpoint for the getClusterHealthStores API as part of the controller gRPC migration, while maintaining backward compatibility with the existing HTTP endpoint. Changes: - Add getStoreStatuses RPC with GetStoreStatusRequest/GetStoreStatusResponse and repeated StoreStatus message to StoreGrpcService.proto - Add getClusterHealthStores() handler to StoreRequestHandler that returns Map<String, String> from Admin.getAllStoreStatuses() - Add gRPC service implementation in StoreGrpcServiceImpl with null-safe status value handling - Update HTTP route in StoresRoutes to delegate to shared handler - Add unit tests for handler, gRPC service, and HTTP route - Add integration test for end-to-end gRPC endpoint validation
1 parent 9f6bb6e commit d8eef14

8 files changed

Lines changed: 292 additions & 2 deletions

File tree

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

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ service StoreGrpcService {
1414
rpc checkResourceCleanupForStoreCreation(ClusterStoreGrpcInfo) returns (ResourceCleanupCheckGrpcResponse) {}
1515
rpc validateStoreDeleted(ValidateStoreDeletedGrpcRequest) returns (ValidateStoreDeletedGrpcResponse);
1616
rpc listStores(ListStoresGrpcRequest) returns (ListStoresGrpcResponse);
17+
rpc getStoreStatuses(GetStoreStatusRequest) returns (GetStoreStatusResponse);
1718
rpc getRepushInfo(GetRepushInfoGrpcRequest) returns (GetRepushInfoGrpcResponse);
1819
rpc getStore(GetStoreGrpcRequest) returns (GetStoreGrpcResponse);
1920
}
@@ -86,6 +87,20 @@ message ListStoresGrpcResponse {
8687
repeated string storeNames = 2;
8788
}
8889

90+
message GetStoreStatusRequest {
91+
string clusterName = 1;
92+
}
93+
94+
message StoreStatus {
95+
string storeName = 1;
96+
string status = 2;
97+
}
98+
99+
message GetStoreStatusResponse {
100+
string clusterName = 1;
101+
repeated StoreStatus storeStatuses = 2;
102+
}
103+
89104
message GetRepushInfoGrpcRequest {
90105
ClusterStoreGrpcInfo storeInfo = 1;
91106
optional string fabric = 2;

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

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,8 @@
2424
import com.linkedin.venice.protocols.controller.GetKeySchemaGrpcResponse;
2525
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
2626
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
27+
import com.linkedin.venice.protocols.controller.GetStoreStatusRequest;
28+
import com.linkedin.venice.protocols.controller.GetStoreStatusResponse;
2729
import com.linkedin.venice.protocols.controller.GetValueSchemaGrpcRequest;
2830
import com.linkedin.venice.protocols.controller.GetValueSchemaGrpcResponse;
2931
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcRequest;
@@ -34,6 +36,7 @@
3436
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
3537
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcRequest;
3638
import com.linkedin.venice.protocols.controller.StoreMigrationCheckGrpcResponse;
39+
import com.linkedin.venice.protocols.controller.StoreStatus;
3740
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
3841
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
3942
import com.linkedin.venice.protocols.controller.VeniceControllerGrpcServiceGrpc;
@@ -48,6 +51,8 @@
4851
import io.grpc.ManagedChannel;
4952
import io.grpc.StatusRuntimeException;
5053
import java.security.cert.X509Certificate;
54+
import java.util.HashMap;
55+
import java.util.Map;
5156
import java.util.Properties;
5257
import java.util.Set;
5358
import java.util.concurrent.ConcurrentHashMap;
@@ -390,6 +395,46 @@ public void testGetValueSchemaGrpcEndpoint() {
390395
assertEquals(exception.getStatus().getCode(), io.grpc.Status.Code.INVALID_ARGUMENT);
391396
}
392397

398+
@Test(timeOut = TIMEOUT_MS)
399+
public void testGetStoreStatusesGrpcEndpoint() {
400+
String storeName1 = Utils.getUniqueString("test_health_stores_1");
401+
String storeName2 = Utils.getUniqueString("test_health_stores_2");
402+
String controllerGrpcUrl = veniceCluster.getLeaderVeniceController().getControllerGrpcUrl();
403+
ManagedChannel channel = Grpc.newChannelBuilder(controllerGrpcUrl, InsecureChannelCredentials.create()).build();
404+
try {
405+
StoreGrpcServiceGrpc.StoreGrpcServiceBlockingStub storeBlockingStub =
406+
StoreGrpcServiceGrpc.newBlockingStub(channel);
407+
408+
// Step 1: Create two stores
409+
createTestStore(storeBlockingStub, veniceCluster.getClusterName(), storeName1);
410+
createTestStore(storeBlockingStub, veniceCluster.getClusterName(), storeName2);
411+
412+
// Step 2: Get store statuses
413+
GetStoreStatusRequest healthRequest =
414+
GetStoreStatusRequest.newBuilder().setClusterName(veniceCluster.getClusterName()).build();
415+
416+
GetStoreStatusResponse healthResponse = storeBlockingStub.getStoreStatuses(healthRequest);
417+
assertNotNull(healthResponse, "Response should not be null");
418+
assertEquals(healthResponse.getClusterName(), veniceCluster.getClusterName());
419+
420+
// Convert repeated StoreStatus to map for easier verification
421+
Map<String, String> storeStatusMap = new HashMap<>();
422+
for (StoreStatus status: healthResponse.getStoreStatusesList()) {
423+
storeStatusMap.put(status.getStoreName(), status.getStatus());
424+
}
425+
426+
// Verify the stores we created are in the status map
427+
assertTrue(storeStatusMap.containsKey(storeName1), "Store status map should contain " + storeName1);
428+
assertTrue(storeStatusMap.containsKey(storeName2), "Store status map should contain " + storeName2);
429+
430+
// Verify the statuses are not null/empty
431+
assertNotNull(storeStatusMap.get(storeName1), "Status for " + storeName1 + " should not be null");
432+
assertNotNull(storeStatusMap.get(storeName2), "Status for " + storeName2 + " should not be null");
433+
} finally {
434+
channel.shutdownNow();
435+
}
436+
}
437+
393438
@Test(timeOut = TIMEOUT_MS)
394439
public void testGetRepushInfoGrpcEndpoint() {
395440
String storeName = Utils.getUniqueString("test_get_repush_info_store");
@@ -484,6 +529,21 @@ public void testGetRepushInfoGrpcEndpointForNonExistentStore() {
484529
assertEquals(exception.getStatus().getCode(), io.grpc.Status.Code.NOT_FOUND);
485530
}
486531

532+
private CreateStoreGrpcResponse createTestStore(
533+
StoreGrpcServiceGrpc.StoreGrpcServiceBlockingStub stub,
534+
String clusterName,
535+
String storeName) {
536+
CreateStoreGrpcRequest request = CreateStoreGrpcRequest.newBuilder()
537+
.setStoreInfo(ClusterStoreGrpcInfo.newBuilder().setClusterName(clusterName).setStoreName(storeName).build())
538+
.setOwner("owner")
539+
.setKeySchema(DEFAULT_KEY_SCHEMA)
540+
.setValueSchema("\"string\"")
541+
.build();
542+
CreateStoreGrpcResponse response = stub.createStore(request);
543+
assertNotNull(response, "Response should not be null");
544+
return response;
545+
}
546+
487547
@Test(timeOut = TIMEOUT_MS)
488548
public void testGetKeySchemaGrpcEndpoint() {
489549
String storeName = Utils.getUniqueString("test_get_key_schema_store");

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

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,12 +21,15 @@
2121
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
2222
import com.linkedin.venice.protocols.controller.GetStoreGrpcRequest;
2323
import com.linkedin.venice.protocols.controller.GetStoreGrpcResponse;
24+
import com.linkedin.venice.protocols.controller.GetStoreStatusRequest;
25+
import com.linkedin.venice.protocols.controller.GetStoreStatusResponse;
2426
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
2527
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
2628
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
2729
import com.linkedin.venice.protocols.controller.ResourceCleanupCheckGrpcResponse;
2830
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
2931
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc.StoreGrpcServiceImplBase;
32+
import com.linkedin.venice.protocols.controller.StoreStatus;
3033
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcRequest;
3134
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
3235
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
@@ -36,6 +39,7 @@
3639
import com.linkedin.venice.utils.ObjectMapperFactory;
3740
import io.grpc.Context;
3841
import io.grpc.stub.StreamObserver;
42+
import java.util.Map;
3943
import java.util.Optional;
4044
import org.apache.logging.log4j.LogManager;
4145
import org.apache.logging.log4j.Logger;
@@ -160,6 +164,33 @@ public void listStores(ListStoresGrpcRequest grpcRequest, StreamObserver<ListSto
160164
null);
161165
}
162166

167+
/**
168+
* Gets the health status of all stores in the specified cluster.
169+
* No ACL check; any user can query store health statuses.
170+
*/
171+
@Override
172+
public void getStoreStatuses(
173+
GetStoreStatusRequest grpcRequest,
174+
StreamObserver<GetStoreStatusResponse> responseObserver) {
175+
LOGGER.debug("Received getStoreStatuses with args: {}", grpcRequest);
176+
String clusterName = grpcRequest.getClusterName();
177+
handleRequest(StoreGrpcServiceGrpc.getGetStoreStatusesMethod(), () -> {
178+
Map<String, String> storeStatusMap = storeRequestHandler.getStoreStatuses(clusterName);
179+
180+
// Convert map to repeated StoreStatus entries
181+
GetStoreStatusResponse.Builder responseBuilder = GetStoreStatusResponse.newBuilder().setClusterName(clusterName);
182+
183+
if (storeStatusMap != null) {
184+
storeStatusMap.forEach((storeName, status) -> {
185+
responseBuilder.addStoreStatuses(
186+
StoreStatus.newBuilder().setStoreName(storeName).setStatus(status != null ? status : "").build());
187+
});
188+
}
189+
190+
return responseBuilder.build();
191+
}, responseObserver, clusterName, null);
192+
}
193+
163194
/**
164195
* Retrieves repush information for a store.
165196
* No ACL check is required for this operation as it only reads store metadata.

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

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import com.linkedin.venice.systemstore.schemas.StoreProperties;
2626
import java.util.ArrayList;
2727
import java.util.List;
28+
import java.util.Map;
2829
import java.util.Optional;
2930
import org.apache.avro.Schema;
3031
import org.apache.commons.lang.StringUtils;
@@ -259,6 +260,23 @@ public ListStoresGrpcResponse listStores(ListStoresGrpcRequest request) {
259260
return ListStoresGrpcResponse.newBuilder().setClusterName(clusterName).addAllStoreNames(selectedStoreNames).build();
260261
}
261262

263+
/**
264+
* Gets the health status of all stores in the specified cluster.
265+
* @param clusterName the name of the cluster
266+
* @return map of store names to their statuses
267+
*/
268+
public Map<String, String> getStoreStatuses(String clusterName) {
269+
if (StringUtils.isBlank(clusterName)) {
270+
throw new IllegalArgumentException("Cluster name is required");
271+
}
272+
273+
LOGGER.debug("Getting health status for all stores in cluster: {}", clusterName);
274+
Map<String, String> storeStatusMap = admin.getAllStoreStatuses(clusterName);
275+
LOGGER.debug("Found {} stores with health status in cluster: {}", storeStatusMap.size(), clusterName);
276+
277+
return storeStatusMap;
278+
}
279+
262280
/**
263281
* Retrieves repush information for a store.
264282
* No ACL check required - reading repush info is public.

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -228,9 +228,9 @@ public Route getAllStoresStatuses(Admin admin) {
228228
public void internalHandle(Request request, MultiStoreStatusResponse veniceResponse) {
229229
AdminSparkServer.validateParams(request, CLUSTER_HEALTH_STORES.getParams(), admin);
230230
String clusterName = request.queryParams(CLUSTER);
231+
231232
veniceResponse.setCluster(clusterName);
232-
Map<String, String> storeStatusMap = admin.getAllStoreStatuses(clusterName);
233-
veniceResponse.setStoreStatusMap(storeStatusMap);
233+
veniceResponse.setStoreStatusMap(storeRequestHandler.getStoreStatuses(clusterName));
234234
}
235235
};
236236
}

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

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,11 +34,14 @@
3434
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
3535
import com.linkedin.venice.protocols.controller.GetStoreGrpcRequest;
3636
import com.linkedin.venice.protocols.controller.GetStoreGrpcResponse;
37+
import com.linkedin.venice.protocols.controller.GetStoreStatusRequest;
38+
import com.linkedin.venice.protocols.controller.GetStoreStatusResponse;
3739
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
3840
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
3941
import com.linkedin.venice.protocols.controller.ResourceCleanupCheckGrpcResponse;
4042
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
4143
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc.StoreGrpcServiceBlockingStub;
44+
import com.linkedin.venice.protocols.controller.StoreStatus;
4245
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcRequest;
4346
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
4447
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
@@ -52,6 +55,8 @@
5255
import io.grpc.inprocess.InProcessChannelBuilder;
5356
import io.grpc.inprocess.InProcessServerBuilder;
5457
import java.util.Arrays;
58+
import java.util.HashMap;
59+
import java.util.Map;
5560
import org.testng.annotations.AfterMethod;
5661
import org.testng.annotations.BeforeMethod;
5762
import org.testng.annotations.Test;
@@ -456,6 +461,74 @@ public void testListStoresWithFilters() {
456461
assertEquals(actualResponse.getStoreNamesCount(), 1, "Should have 1 store after filtering");
457462
}
458463

464+
@Test
465+
public void testGetStoreStatusesReturnsSuccessfulResponse() {
466+
GetStoreStatusRequest request = GetStoreStatusRequest.newBuilder().setClusterName(TEST_CLUSTER).build();
467+
Map<String, String> storeStatusMap = new HashMap<>();
468+
storeStatusMap.put("store1", "ONLINE");
469+
storeStatusMap.put("store2", "DEGRADED");
470+
storeStatusMap.put("store3", "UNAVAILABLE");
471+
when(storeRequestHandler.getStoreStatuses(TEST_CLUSTER)).thenReturn(storeStatusMap);
472+
473+
GetStoreStatusResponse actualResponse = blockingStub.getStoreStatuses(request);
474+
475+
assertNotNull(actualResponse, "Response should not be null");
476+
assertEquals(actualResponse.getClusterName(), TEST_CLUSTER, "Cluster name should match");
477+
assertEquals(actualResponse.getStoreStatusesCount(), 3, "Should have 3 stores");
478+
479+
// Verify each store status
480+
Map<String, String> responseMap = new HashMap<>();
481+
for (StoreStatus status: actualResponse.getStoreStatusesList()) {
482+
responseMap.put(status.getStoreName(), status.getStatus());
483+
}
484+
assertEquals(responseMap.get("store1"), "ONLINE", "store1 should be ONLINE");
485+
assertEquals(responseMap.get("store2"), "DEGRADED", "store2 should be DEGRADED");
486+
assertEquals(responseMap.get("store3"), "UNAVAILABLE", "store3 should be UNAVAILABLE");
487+
}
488+
489+
@Test
490+
public void testGetStoreStatusesReturnsEmptyMapWhenNoStores() {
491+
GetStoreStatusRequest request = GetStoreStatusRequest.newBuilder().setClusterName(TEST_CLUSTER).build();
492+
when(storeRequestHandler.getStoreStatuses(TEST_CLUSTER)).thenReturn(new HashMap<>());
493+
494+
GetStoreStatusResponse actualResponse = blockingStub.getStoreStatuses(request);
495+
496+
assertNotNull(actualResponse, "Response should not be null");
497+
assertEquals(actualResponse.getClusterName(), TEST_CLUSTER, "Cluster name should match");
498+
assertEquals(actualResponse.getStoreStatusesCount(), 0, "Should have 0 stores");
499+
}
500+
501+
@Test
502+
public void testGetStoreStatusesReturnsErrorResponse() {
503+
GetStoreStatusRequest request = GetStoreStatusRequest.newBuilder().setClusterName(TEST_CLUSTER).build();
504+
when(storeRequestHandler.getStoreStatuses(TEST_CLUSTER))
505+
.thenThrow(new VeniceException("Failed to get store statuses"));
506+
507+
StatusRuntimeException e = expectThrows(StatusRuntimeException.class, () -> blockingStub.getStoreStatuses(request));
508+
509+
assertNotNull(e.getStatus(), "Status should not be null");
510+
assertEquals(e.getStatus().getCode(), Status.INTERNAL.getCode());
511+
VeniceControllerGrpcErrorInfo errorInfo = GrpcRequestResponseConverter.parseControllerGrpcError(e);
512+
assertNotNull(errorInfo, "Error info should not be null");
513+
assertEquals(errorInfo.getErrorType(), ControllerGrpcErrorType.GENERAL_ERROR);
514+
assertTrue(errorInfo.getErrorMessage().contains("Failed to get store statuses"));
515+
}
516+
517+
@Test
518+
public void testGetStoreStatusesReturnsBadRequestForMissingClusterName() {
519+
GetStoreStatusRequest request = GetStoreStatusRequest.newBuilder().build();
520+
when(storeRequestHandler.getStoreStatuses("")).thenThrow(new IllegalArgumentException("Cluster name is required"));
521+
522+
StatusRuntimeException e = expectThrows(StatusRuntimeException.class, () -> blockingStub.getStoreStatuses(request));
523+
524+
assertNotNull(e.getStatus(), "Status should not be null");
525+
assertEquals(e.getStatus().getCode(), Status.INVALID_ARGUMENT.getCode());
526+
VeniceControllerGrpcErrorInfo errorInfo = GrpcRequestResponseConverter.parseControllerGrpcError(e);
527+
assertNotNull(errorInfo, "Error info should not be null");
528+
assertEquals(errorInfo.getErrorType(), ControllerGrpcErrorType.BAD_REQUEST);
529+
assertTrue(errorInfo.getErrorMessage().contains("Cluster name is required"));
530+
}
531+
459532
@Test
460533
public void testGetRepushInfoReturnsSuccessfulResponse() {
461534
ClusterStoreGrpcInfo storeInfo =

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,9 +33,12 @@
3333
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
3434
import java.util.Arrays;
3535
import java.util.Collections;
36+
import java.util.HashMap;
3637
import java.util.List;
38+
import java.util.Map;
3739
import java.util.Optional;
3840
import org.testng.annotations.BeforeMethod;
41+
import org.testng.annotations.DataProvider;
3942
import org.testng.annotations.Test;
4043

4144

@@ -364,6 +367,35 @@ public void testListStoresWithDataReplicationPolicyFilterNullPolicy() {
364367
assertEquals(response.getStoreNamesCount(), 0);
365368
}
366369

370+
@DataProvider(name = "storeStatusMaps")
371+
public Object[][] storeStatusMaps() {
372+
Map<String, String> populatedMap = new HashMap<>();
373+
populatedMap.put("store1", "ONLINE");
374+
populatedMap.put("store2", "DEGRADED");
375+
populatedMap.put("store3", "UNAVAILABLE");
376+
return new Object[][] { { populatedMap }, { Collections.emptyMap() } };
377+
}
378+
379+
@Test(dataProvider = "storeStatusMaps")
380+
public void testGetStoreStatuses(Map<String, String> expectedStatusMap) {
381+
when(admin.getAllStoreStatuses("testCluster")).thenReturn(expectedStatusMap);
382+
383+
Map<String, String> response = storeRequestHandler.getStoreStatuses("testCluster");
384+
385+
verify(admin, times(1)).getAllStoreStatuses("testCluster");
386+
assertEquals(response, expectedStatusMap);
387+
}
388+
389+
@DataProvider(name = "blankClusterNames")
390+
public Object[][] blankClusterNames() {
391+
return new Object[][] { { null }, { "" } };
392+
}
393+
394+
@Test(dataProvider = "blankClusterNames", expectedExceptions = IllegalArgumentException.class, expectedExceptionsMessageRegExp = "Cluster name is required")
395+
public void testGetStoreStatusesWithBlankClusterName(String clusterName) {
396+
storeRequestHandler.getStoreStatuses(clusterName);
397+
}
398+
367399
@Test
368400
public void testGetRepushInfoSuccess() {
369401
String clusterName = "testCluster";

0 commit comments

Comments
 (0)