Skip to content

Commit 32a191e

Browse files
pthirunclaude
andcommitted
Simplify getRepushInfo HTTP route to call admin directly
- Remove requestHandler parameter from getRepushInfo method - Call admin.getRepushInfo() directly instead of converting to/from gRPC format - Remove mapGrpcRepushInfoToRepushInfo helper method (no longer needed) - Update tests to mock admin.getRepushInfo() directly - Remove unused gRPC imports from StoresRoutes and StoresRoutesTest This avoids the inefficient double conversion (RepushInfo -> RepushInfoGrpc -> RepushInfo) for the HTTP code path, while keeping the gRPC path unchanged. Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
1 parent 2183569 commit 32a191e

3 files changed

Lines changed: 8 additions & 100 deletions

File tree

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

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

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

Lines changed: 3 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -108,17 +108,11 @@
108108
import com.linkedin.venice.meta.StoreDataAudit;
109109
import com.linkedin.venice.meta.StoreInfo;
110110
import com.linkedin.venice.meta.Version;
111-
import com.linkedin.venice.meta.VersionImpl;
112-
import com.linkedin.venice.meta.VersionStatus;
113111
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
114-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
115-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
116112
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
117113
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
118-
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
119114
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
120115
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
121-
import com.linkedin.venice.protocols.controller.VersionGrpc;
122116
import com.linkedin.venice.pubsub.PubSubTopicRepository;
123117
import com.linkedin.venice.pubsub.api.PubSubTopic;
124118
import com.linkedin.venice.pubsub.api.exceptions.PubSubTopicDoesNotExistException;
@@ -259,7 +253,7 @@ public void internalHandle(Request request, SchemaUsageResponse response) {
259253
/**
260254
* @see Admin#getRepushInfo(String, String, Optional)
261255
*/
262-
public Route getRepushInfo(Admin admin, StoreRequestHandler requestHandler) {
256+
public Route getRepushInfo(Admin admin) {
263257
return new VeniceRouteHandler<RepushInfoResponse>(RepushInfoResponse.class) {
264258
@Override
265259
public void internalHandle(Request request, RepushInfoResponse veniceResponse) {
@@ -268,26 +262,11 @@ public void internalHandle(Request request, RepushInfoResponse veniceResponse) {
268262
String storeName = request.queryParams(NAME);
269263
String fabricName = request.queryParams(FABRIC);
270264

271-
// Convert to gRPC request
272-
ClusterStoreGrpcInfo storeInfo =
273-
ClusterStoreGrpcInfo.newBuilder().setClusterName(clusterName).setStoreName(storeName).build();
274-
GetRepushInfoGrpcRequest.Builder grpcRequestBuilder =
275-
GetRepushInfoGrpcRequest.newBuilder().setStoreInfo(storeInfo);
276-
if (fabricName != null) {
277-
grpcRequestBuilder.setFabric(fabricName);
278-
}
279-
GetRepushInfoGrpcRequest grpcRequest = grpcRequestBuilder.build();
265+
Optional<String> fabric = fabricName != null ? Optional.of(fabricName) : Optional.empty();
266+
RepushInfo repushInfo = admin.getRepushInfo(clusterName, storeName, fabric);
280267

281-
// Call handler
282-
GetRepushInfoGrpcResponse grpcResponse = requestHandler.getRepushInfo(grpcRequest);
283-
284-
// Map response back to HTTP
285268
veniceResponse.setCluster(clusterName);
286269
veniceResponse.setName(storeName);
287-
288-
// Convert proto RepushInfo back to Java RepushInfo for HTTP response
289-
RepushInfo repushInfo = mapGrpcRepushInfoToRepushInfo(grpcResponse.getRepushInfo(), storeName);
290-
291270
veniceResponse.setRepushInfo(repushInfo);
292271
}
293272
};
@@ -1248,33 +1227,4 @@ public void internalHandle(Request request, StoreDeletedValidationResponse venic
12481227
};
12491228
}
12501229

1251-
/**
1252-
* Converts a gRPC RepushInfoGrpc message to a RepushInfo object.
1253-
* @param repushInfoProto the gRPC message
1254-
* @param storeName the store name needed for Version creation
1255-
* @return the converted RepushInfo object
1256-
*/
1257-
RepushInfo mapGrpcRepushInfoToRepushInfo(RepushInfoGrpc repushInfoProto, String storeName) {
1258-
Version version = null;
1259-
if (repushInfoProto.hasVersion()) {
1260-
VersionGrpc versionProto = repushInfoProto.getVersion();
1261-
version = new VersionImpl(
1262-
storeName,
1263-
versionProto.getNumber(),
1264-
versionProto.getCreatedTime(),
1265-
versionProto.getPushJobId(),
1266-
versionProto.getPartitionCount(),
1267-
null,
1268-
null);
1269-
version.setStatus(VersionStatus.getVersionStatusFromInt(versionProto.getStatus()));
1270-
}
1271-
1272-
return RepushInfo.createRepushInfo(
1273-
version,
1274-
repushInfoProto.getKafkaBrokerUrl(),
1275-
repushInfoProto.hasSystemSchemaClusterD2ServiceName()
1276-
? repushInfoProto.getSystemSchemaClusterD2ServiceName()
1277-
: null,
1278-
repushInfoProto.hasSystemSchemaClusterD2ZkHost() ? repushInfoProto.getSystemSchemaClusterD2ZkHost() : null);
1279-
}
12801230
}

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

Lines changed: 4 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -42,14 +42,10 @@
4242
import com.linkedin.venice.meta.VersionStatus;
4343
import com.linkedin.venice.meta.ZKStore;
4444
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
45-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
46-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
4745
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
4846
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
49-
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
5047
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
5148
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
52-
import com.linkedin.venice.protocols.controller.VersionGrpc;
5349
import com.linkedin.venice.pubsub.PubSubTopicRepository;
5450
import com.linkedin.venice.utils.ObjectMapperFactory;
5551
import java.util.Arrays;
@@ -700,7 +696,6 @@ public void testGetAllStoresWithFilters() throws Exception {
700696
@Test
701697
public void testGetRepushInfo() throws Exception {
702698
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
703-
StoreRequestHandler mockRequestHandler = mock(StoreRequestHandler.class);
704699
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
705700

706701
Request request = mock(Request.class);
@@ -717,31 +712,7 @@ public void testGetRepushInfo() throws Exception {
717712
doReturn(queryMap).when(queryParamsMap).toMap();
718713
doReturn(queryParamsMap).when(request).queryMap();
719714

720-
Route route =
721-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getRepushInfo(mockAdmin, mockRequestHandler);
722-
723-
// Test success case
724-
ClusterStoreGrpcInfo storeInfo =
725-
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE_NAME).build();
726-
727-
VersionGrpc versionGrpc = VersionGrpc.newBuilder()
728-
.setNumber(1)
729-
.setCreatedTime(123456789L)
730-
.setStatus(VersionStatus.ONLINE.getValue())
731-
.setPushJobId("test-push-job")
732-
.setPartitionCount(10)
733-
.setReplicationFactor(3)
734-
.build();
735-
736-
RepushInfoGrpc repushInfoGrpc = RepushInfoGrpc.newBuilder()
737-
.setKafkaBrokerUrl("kafka.broker:9092")
738-
.setVersion(versionGrpc)
739-
.setSystemSchemaClusterD2ServiceName("d2-service")
740-
.setSystemSchemaClusterD2ZkHost("zk-host")
741-
.build();
742-
743-
GetRepushInfoGrpcResponse grpcResponse =
744-
GetRepushInfoGrpcResponse.newBuilder().setStoreInfo(storeInfo).setRepushInfo(repushInfoGrpc).build();
715+
Route route = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getRepushInfo(mockAdmin);
745716

746717
// Create a real Version for admin.getRepushInfo() call to avoid Jackson serialization issues
747718
Version version = new VersionImpl(TEST_STORE_NAME, 1, "test-push-job", 10);
@@ -750,7 +721,6 @@ public void testGetRepushInfo() throws Exception {
750721

751722
RepushInfo mockRepushInfo = RepushInfo.createRepushInfo(version, "kafka.broker:9092", "d2-service", "zk-host");
752723

753-
when(mockRequestHandler.getRepushInfo(any(GetRepushInfoGrpcRequest.class))).thenReturn(grpcResponse);
754724
when(mockAdmin.getRepushInfo(eq(TEST_CLUSTER), eq(TEST_STORE_NAME), eq(Optional.of("test-fabric"))))
755725
.thenReturn(mockRepushInfo);
756726

@@ -764,7 +734,8 @@ public void testGetRepushInfo() throws Exception {
764734
Assert.assertEquals(response.getRepushInfo().getKafkaBrokerUrl(), "kafka.broker:9092");
765735

766736
// Test error case
767-
when(mockRequestHandler.getRepushInfo(any(GetRepushInfoGrpcRequest.class))).thenThrow(new VeniceException("Error"));
737+
when(mockAdmin.getRepushInfo(eq(TEST_CLUSTER), eq(TEST_STORE_NAME), eq(Optional.of("test-fabric"))))
738+
.thenThrow(new VeniceException("Error"));
768739
response = ObjectMapperFactory.getInstance()
769740
.readValue(route.handle(request, mock(Response.class)).toString(), RepushInfoResponse.class);
770741
Assert.assertTrue(response.isError());
@@ -773,7 +744,6 @@ public void testGetRepushInfo() throws Exception {
773744
@Test
774745
public void testGetRepushInfoWithoutFabric() throws Exception {
775746
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
776-
StoreRequestHandler mockRequestHandler = mock(StoreRequestHandler.class);
777747
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
778748

779749
Request request = mock(Request.class);
@@ -789,20 +759,10 @@ public void testGetRepushInfoWithoutFabric() throws Exception {
789759
doReturn(queryMap).when(queryParamsMap).toMap();
790760
doReturn(queryParamsMap).when(request).queryMap();
791761

792-
Route route =
793-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getRepushInfo(mockAdmin, mockRequestHandler);
794-
795-
ClusterStoreGrpcInfo storeInfo =
796-
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE_NAME).build();
797-
798-
RepushInfoGrpc repushInfoGrpc = RepushInfoGrpc.newBuilder().setKafkaBrokerUrl("another.kafka:9092").build();
799-
800-
GetRepushInfoGrpcResponse grpcResponse =
801-
GetRepushInfoGrpcResponse.newBuilder().setStoreInfo(storeInfo).setRepushInfo(repushInfoGrpc).build();
762+
Route route = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getRepushInfo(mockAdmin);
802763

803764
RepushInfo mockRepushInfo = RepushInfo.createRepushInfo(null, "another.kafka:9092", null, null);
804765

805-
when(mockRequestHandler.getRepushInfo(any(GetRepushInfoGrpcRequest.class))).thenReturn(grpcResponse);
806766
when(mockAdmin.getRepushInfo(eq(TEST_CLUSTER), eq(TEST_STORE_NAME), eq(Optional.empty())))
807767
.thenReturn(mockRepushInfo);
808768

0 commit comments

Comments
 (0)