Skip to content

Commit bfc16ab

Browse files
committed
[Controller] Add gRPC support for getRepushInfo API
1 parent 8109fcd commit bfc16ab

9 files changed

Lines changed: 414 additions & 2 deletions

File tree

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

Lines changed: 27 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 getRepushInfo(GetRepushInfoGrpcRequest) returns (GetRepushInfoGrpcResponse);
1718
}
1819

1920
message CreateStoreGrpcRequest {
@@ -82,4 +83,30 @@ message ListStoresGrpcRequest {
8283
message ListStoresGrpcResponse {
8384
string clusterName = 1;
8485
repeated string storeNames = 2;
86+
}
87+
88+
message GetRepushInfoGrpcRequest {
89+
ClusterStoreGrpcInfo storeInfo = 1;
90+
optional string fabric = 2;
91+
}
92+
93+
message GetRepushInfoGrpcResponse {
94+
ClusterStoreGrpcInfo storeInfo = 1;
95+
RepushInfoGrpc repushInfo = 2;
96+
}
97+
98+
message RepushInfoGrpc {
99+
string kafkaBrokerUrl = 1;
100+
optional VersionGrpc version = 2;
101+
optional string systemSchemaClusterD2ServiceName = 3;
102+
optional string systemSchemaClusterD2ZkHost = 4;
103+
}
104+
105+
message VersionGrpc {
106+
int32 number = 1;
107+
int64 createdTime = 2;
108+
int32 status = 3; // VersionStatus enum ordinal
109+
string pushJobId = 4;
110+
int32 partitionCount = 5;
111+
int32 replicationFactor = 6;
85112
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -337,6 +337,9 @@ public void testListStoresGrpcEndpoint() {
337337
}
338338
}
339339

340+
// TODO: Integration test for getRepushInfo - requires store with specific configuration (hybrid/incremental)
341+
// Unit tests in StoreRequestHandlerTest, StoreGrpcServiceImplTest, and StoresRoutesTest provide coverage
342+
340343
private static class MockDynamicAccessController extends NoOpDynamicAccessController {
341344
private final Set<String> resourcesInAllowList = ConcurrentHashMap.newKeySet();
342345

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

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@
1414
import com.linkedin.venice.protocols.controller.DeleteAclForStoreGrpcResponse;
1515
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcRequest;
1616
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcResponse;
17+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
18+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
1719
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
1820
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
1921
import com.linkedin.venice.protocols.controller.ResourceCleanupCheckGrpcResponse;
@@ -147,4 +149,16 @@ public void listStores(ListStoresGrpcRequest grpcRequest, StreamObserver<ListSto
147149
clusterName,
148150
null);
149151
}
152+
153+
@Override
154+
public void getRepushInfo(
155+
GetRepushInfoGrpcRequest request,
156+
StreamObserver<GetRepushInfoGrpcResponse> responseObserver) {
157+
LOGGER.debug("Received getRepushInfo with args: {}", request);
158+
ControllerGrpcServerUtils.handleRequest(
159+
StoreGrpcServiceGrpc.getGetRepushInfoMethod(),
160+
() -> storeRequestHandler.getRepushInfo(request),
161+
responseObserver,
162+
request.getStoreInfo());
163+
}
150164
}

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -600,7 +600,9 @@ public boolean startInner() throws Exception {
600600
new VeniceParentControllerRegionStateHandler(admin, jobRoutes.getOngoingIncrementalPushVersions(admin)));
601601
httpService.get(
602602
GET_REPUSH_INFO.getPath(),
603-
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getRepushInfo(admin)));
603+
new VeniceParentControllerRegionStateHandler(
604+
admin,
605+
storesRoutes.getRepushInfo(admin, requestHandler.getStoreRequestHandler())));
604606
httpService.get(
605607
COMPARE_STORE.getPath(),
606608
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.compareStore(admin)));

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

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,10 @@
33
import com.linkedin.venice.controller.Admin;
44
import com.linkedin.venice.controller.ControllerRequestHandlerDependencies;
55
import com.linkedin.venice.controller.StoreDeletedValidation;
6+
import com.linkedin.venice.controllerapi.RepushInfo;
67
import com.linkedin.venice.exceptions.VeniceException;
78
import com.linkedin.venice.meta.Store;
9+
import com.linkedin.venice.meta.Version;
810
import com.linkedin.venice.meta.ZKStore;
911
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
1012
import com.linkedin.venice.protocols.controller.CreateStoreGrpcRequest;
@@ -13,12 +15,16 @@
1315
import com.linkedin.venice.protocols.controller.DeleteAclForStoreGrpcResponse;
1416
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcRequest;
1517
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcResponse;
18+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
19+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
1620
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
1721
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
22+
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
1823
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcRequest;
1924
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
2025
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
2126
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
27+
import com.linkedin.venice.protocols.controller.VersionGrpc;
2228
import com.linkedin.venice.systemstore.schemas.StoreProperties;
2329
import java.util.ArrayList;
2430
import java.util.List;
@@ -255,4 +261,54 @@ public ListStoresGrpcResponse listStores(ListStoresGrpcRequest request) {
255261
LOGGER.info("Found {} stores in cluster: {}", selectedStoreNames.size(), clusterName);
256262
return ListStoresGrpcResponse.newBuilder().setClusterName(clusterName).addAllStoreNames(selectedStoreNames).build();
257263
}
264+
265+
/**
266+
* Retrieves repush information for a store.
267+
* @param request the request containing cluster, store name, and optional fabric
268+
* @return response containing repush information including version and Kafka details
269+
*/
270+
public GetRepushInfoGrpcResponse getRepushInfo(GetRepushInfoGrpcRequest request) {
271+
ClusterStoreGrpcInfo storeInfo = request.getStoreInfo();
272+
ControllerRequestParamValidator.validateClusterStoreInfo(storeInfo);
273+
String clusterName = storeInfo.getClusterName();
274+
String storeName = storeInfo.getStoreName();
275+
Optional<String> fabric = request.hasFabric() ? Optional.of(request.getFabric()) : Optional.empty();
276+
277+
LOGGER.info(
278+
"Getting repush info for store: {} in cluster: {} with fabric: {}",
279+
storeName,
280+
clusterName,
281+
fabric.orElse("none"));
282+
283+
RepushInfo repushInfo = admin.getRepushInfo(clusterName, storeName, fabric);
284+
285+
RepushInfoGrpc.Builder repushInfoBuilder = RepushInfoGrpc.newBuilder()
286+
.setKafkaBrokerUrl(repushInfo.getKafkaBrokerUrl() != null ? repushInfo.getKafkaBrokerUrl() : "");
287+
288+
if (repushInfo.getVersion() != null) {
289+
repushInfoBuilder.setVersion(convertVersionToProto(repushInfo.getVersion()));
290+
}
291+
if (repushInfo.getSystemSchemaClusterD2ServiceName() != null) {
292+
repushInfoBuilder.setSystemSchemaClusterD2ServiceName(repushInfo.getSystemSchemaClusterD2ServiceName());
293+
}
294+
if (repushInfo.getSystemSchemaClusterD2ZkHost() != null) {
295+
repushInfoBuilder.setSystemSchemaClusterD2ZkHost(repushInfo.getSystemSchemaClusterD2ZkHost());
296+
}
297+
298+
return GetRepushInfoGrpcResponse.newBuilder()
299+
.setStoreInfo(storeInfo)
300+
.setRepushInfo(repushInfoBuilder.build())
301+
.build();
302+
}
303+
304+
private VersionGrpc convertVersionToProto(Version version) {
305+
return VersionGrpc.newBuilder()
306+
.setNumber(version.getNumber())
307+
.setCreatedTime(version.getCreatedTime())
308+
.setStatus(version.getStatus().getValue())
309+
.setPushJobId(version.getPushJobId())
310+
.setPartitionCount(version.getPartitionCount())
311+
.setReplicationFactor(version.getReplicationFactor())
312+
.build();
313+
}
258314
}

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

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,8 +109,11 @@
109109
import com.linkedin.venice.meta.StoreInfo;
110110
import com.linkedin.venice.meta.Version;
111111
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
112+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
113+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
112114
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
113115
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
116+
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
114117
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
115118
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
116119
import com.linkedin.venice.pubsub.PubSubTopicRepository;
@@ -253,7 +256,7 @@ public void internalHandle(Request request, SchemaUsageResponse response) {
253256
/**
254257
* @see Admin#getRepushInfo(String, String, Optional)
255258
*/
256-
public Route getRepushInfo(Admin admin) {
259+
public Route getRepushInfo(Admin admin, StoreRequestHandler requestHandler) {
257260
return new VeniceRouteHandler<RepushInfoResponse>(RepushInfoResponse.class) {
258261
@Override
259262
public void internalHandle(Request request, RepushInfoResponse veniceResponse) {
@@ -262,9 +265,25 @@ public void internalHandle(Request request, RepushInfoResponse veniceResponse) {
262265
String storeName = request.queryParams(NAME);
263266
String fabricName = request.queryParams(FABRIC);
264267

268+
// Convert to gRPC request
269+
ClusterStoreGrpcInfo storeInfo =
270+
ClusterStoreGrpcInfo.newBuilder().setClusterName(clusterName).setStoreName(storeName).build();
271+
GetRepushInfoGrpcRequest.Builder grpcRequestBuilder =
272+
GetRepushInfoGrpcRequest.newBuilder().setStoreInfo(storeInfo);
273+
if (fabricName != null) {
274+
grpcRequestBuilder.setFabric(fabricName);
275+
}
276+
GetRepushInfoGrpcRequest grpcRequest = grpcRequestBuilder.build();
277+
278+
// Call handler
279+
GetRepushInfoGrpcResponse grpcResponse = requestHandler.getRepushInfo(grpcRequest);
280+
281+
// Map response back to HTTP
265282
veniceResponse.setCluster(clusterName);
266283
veniceResponse.setName(storeName);
267284

285+
// Convert proto RepushInfo back to Java RepushInfo for HTTP response
286+
RepushInfoGrpc repushInfoProto = grpcResponse.getRepushInfo();
268287
RepushInfo repushInfo = admin.getRepushInfo(clusterName, storeName, Optional.ofNullable(fabricName));
269288

270289
veniceResponse.setRepushInfo(repushInfo);

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

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
import com.linkedin.venice.controller.server.StoreRequestHandler;
1818
import com.linkedin.venice.controller.server.VeniceControllerAccessManager;
1919
import com.linkedin.venice.exceptions.VeniceException;
20+
import com.linkedin.venice.meta.VersionStatus;
2021
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
2122
import com.linkedin.venice.protocols.controller.ControllerGrpcErrorType;
2223
import com.linkedin.venice.protocols.controller.CreateStoreGrpcRequest;
@@ -25,8 +26,11 @@
2526
import com.linkedin.venice.protocols.controller.DeleteAclForStoreGrpcResponse;
2627
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcRequest;
2728
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcResponse;
29+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
30+
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
2831
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
2932
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
33+
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
3034
import com.linkedin.venice.protocols.controller.ResourceCleanupCheckGrpcResponse;
3135
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
3236
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc.StoreGrpcServiceBlockingStub;
@@ -35,6 +39,7 @@
3539
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
3640
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
3741
import com.linkedin.venice.protocols.controller.VeniceControllerGrpcErrorInfo;
42+
import com.linkedin.venice.protocols.controller.VersionGrpc;
3843
import io.grpc.ManagedChannel;
3944
import io.grpc.Server;
4045
import io.grpc.Status;
@@ -445,4 +450,76 @@ public void testListStoresWithFilters() {
445450
assertEquals(actualResponse.getClusterName(), TEST_CLUSTER, "Cluster name should match");
446451
assertEquals(actualResponse.getStoreNamesCount(), 1, "Should have 1 store after filtering");
447452
}
453+
454+
@Test
455+
public void testGetRepushInfoReturnsSuccessfulResponse() {
456+
ClusterStoreGrpcInfo storeInfo =
457+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE).build();
458+
GetRepushInfoGrpcRequest request =
459+
GetRepushInfoGrpcRequest.newBuilder().setStoreInfo(storeInfo).setFabric("test-fabric").build();
460+
461+
VersionGrpc versionGrpc = VersionGrpc.newBuilder()
462+
.setNumber(1)
463+
.setCreatedTime(123456789L)
464+
.setStatus(VersionStatus.ONLINE.getValue())
465+
.setPushJobId("test-push-job")
466+
.setPartitionCount(10)
467+
.setReplicationFactor(3)
468+
.build();
469+
470+
RepushInfoGrpc repushInfoGrpc = RepushInfoGrpc.newBuilder()
471+
.setKafkaBrokerUrl("kafka.broker:9092")
472+
.setVersion(versionGrpc)
473+
.setSystemSchemaClusterD2ServiceName("d2-service")
474+
.setSystemSchemaClusterD2ZkHost("zk-host")
475+
.build();
476+
477+
GetRepushInfoGrpcResponse mockResponse =
478+
GetRepushInfoGrpcResponse.newBuilder().setStoreInfo(storeInfo).setRepushInfo(repushInfoGrpc).build();
479+
480+
when(storeRequestHandler.getRepushInfo(any(GetRepushInfoGrpcRequest.class))).thenReturn(mockResponse);
481+
482+
GetRepushInfoGrpcResponse actualResponse = blockingStub.getRepushInfo(request);
483+
484+
assertNotNull(actualResponse);
485+
assertEquals(actualResponse.getStoreInfo(), storeInfo);
486+
assertEquals(actualResponse.getRepushInfo().getKafkaBrokerUrl(), "kafka.broker:9092");
487+
assertTrue(actualResponse.getRepushInfo().hasVersion());
488+
assertEquals(actualResponse.getRepushInfo().getVersion().getNumber(), 1);
489+
assertEquals(actualResponse.getRepushInfo().getVersion().getStatus(), VersionStatus.ONLINE.getValue());
490+
}
491+
492+
@Test
493+
public void testGetRepushInfoReturnsErrorResponse() {
494+
ClusterStoreGrpcInfo storeInfo =
495+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE).build();
496+
GetRepushInfoGrpcRequest request = GetRepushInfoGrpcRequest.newBuilder().setStoreInfo(storeInfo).build();
497+
498+
when(storeRequestHandler.getRepushInfo(any(GetRepushInfoGrpcRequest.class)))
499+
.thenThrow(new VeniceException("Failed to get repush info"));
500+
501+
StatusRuntimeException e = expectThrows(StatusRuntimeException.class, () -> blockingStub.getRepushInfo(request));
502+
503+
assertEquals(e.getStatus().getCode(), Status.INTERNAL.getCode());
504+
}
505+
506+
@Test
507+
public void testGetRepushInfoWithoutFabric() {
508+
ClusterStoreGrpcInfo storeInfo =
509+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE).build();
510+
GetRepushInfoGrpcRequest request = GetRepushInfoGrpcRequest.newBuilder().setStoreInfo(storeInfo).build();
511+
512+
RepushInfoGrpc repushInfoGrpc = RepushInfoGrpc.newBuilder().setKafkaBrokerUrl("another.kafka.broker:9092").build();
513+
514+
GetRepushInfoGrpcResponse mockResponse =
515+
GetRepushInfoGrpcResponse.newBuilder().setStoreInfo(storeInfo).setRepushInfo(repushInfoGrpc).build();
516+
517+
when(storeRequestHandler.getRepushInfo(any(GetRepushInfoGrpcRequest.class))).thenReturn(mockResponse);
518+
519+
GetRepushInfoGrpcResponse actualResponse = blockingStub.getRepushInfo(request);
520+
521+
assertNotNull(actualResponse);
522+
assertEquals(actualResponse.getRepushInfo().getKafkaBrokerUrl(), "another.kafka.broker:9092");
523+
assertFalse(actualResponse.getRepushInfo().hasVersion());
524+
}
448525
}

0 commit comments

Comments
 (0)