Skip to content

Commit eb22c62

Browse files
pthirunclaude
andcommitted
[Controller] Migrate getFutureVersion endpoint to gRPC
- Add GetFutureVersionGrpcRequest/Response proto messages - Add getFutureVersion RPC to VeniceControllerGrpcService - Add getFutureVersion handler method in VeniceControllerRequestHandler - Update VeniceControllerGrpcServiceImpl to implement the RPC - Update StoresRoutes.getFutureVersion to use the handler - Add VeniceNoStoreException handling in ControllerGrpcServerUtils - Add comprehensive unit tests Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
1 parent 17684d2 commit eb22c62

9 files changed

Lines changed: 275 additions & 44 deletions

File tree

internal/venice-common/src/main/proto/VeniceControllerGrpcService.proto

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,9 @@ service VeniceControllerGrpcService {
1414

1515
// ControllerRoutes
1616
rpc getLeaderController(LeaderControllerGrpcRequest) returns (LeaderControllerGrpcResponse);
17+
18+
// StoresRoutes
19+
rpc getFutureVersion(GetFutureVersionGrpcRequest) returns (GetFutureVersionGrpcResponse);
1720
}
1821

1922
message DiscoverClusterGrpcRequest {
@@ -40,3 +43,12 @@ message LeaderControllerGrpcResponse {
4043
string grpcUrl = 4; // gRPC URL for leader controller
4144
string secureGrpcUrl = 5; // Secure gRPC URL for leader controller
4245
}
46+
47+
message GetFutureVersionGrpcRequest {
48+
ClusterStoreGrpcInfo storeInfo = 1;
49+
}
50+
51+
message GetFutureVersionGrpcResponse {
52+
ClusterStoreGrpcInfo storeInfo = 1;
53+
map<string, string> storeVersionMap = 2; // Map of store name to future version
54+
}

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
import com.linkedin.venice.controller.grpc.GrpcRequestResponseConverter;
66
import com.linkedin.venice.controller.server.VeniceControllerAccessManager;
7+
import com.linkedin.venice.exceptions.VeniceNoStoreException;
78
import com.linkedin.venice.exceptions.VeniceUnauthorizedAccessException;
89
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
910
import com.linkedin.venice.protocols.controller.ControllerGrpcErrorType;
@@ -66,6 +67,15 @@ public static <T> void handleRequest(
6667
clusterName,
6768
storeName,
6869
responseObserver);
70+
} catch (VeniceNoStoreException e) {
71+
LOGGER.error("Store not found for method: {} on cluster: {}, store: {}", methodName, clusterName, storeName, e);
72+
GrpcRequestResponseConverter.sendErrorResponse(
73+
Status.Code.NOT_FOUND,
74+
ControllerGrpcErrorType.STORE_NOT_FOUND,
75+
e,
76+
clusterName,
77+
storeName,
78+
responseObserver);
6979
} catch (Exception e) {
7080
LOGGER.error("Error in method: {} on cluster: {}, store: {}", methodName, clusterName, storeName, e);
7181
GrpcRequestResponseConverter.sendErrorResponse(

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -336,7 +336,7 @@ public boolean startInner() throws Exception {
336336
httpService.get(STORE.getPath(), new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getStore(admin)));
337337
httpService.get(
338338
FUTURE_VERSION.getPath(),
339-
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getFutureVersion(admin)));
339+
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getFutureVersion(admin, requestHandler)));
340340
httpService.get(
341341
BACKUP_VERSION.getPath(),
342342
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getBackupVersion(admin)));

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

Lines changed: 17 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,9 @@
106106
import com.linkedin.venice.meta.VeniceUserStoreType;
107107
import com.linkedin.venice.meta.Version;
108108
import com.linkedin.venice.meta.ZKStore;
109+
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
110+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcRequest;
111+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcResponse;
109112
import com.linkedin.venice.pubsub.PubSubTopicRepository;
110113
import com.linkedin.venice.pubsub.api.PubSubTopic;
111114
import com.linkedin.venice.pubsub.api.exceptions.PubSubTopicDoesNotExistException;
@@ -376,25 +379,26 @@ public void internalHandle(Request request, StoreResponse veniceResponse) {
376379
};
377380
}
378381

379-
public Route getFutureVersion(Admin admin) {
382+
public Route getFutureVersion(Admin admin, VeniceControllerRequestHandler requestHandler) {
380383
return new VeniceRouteHandler<MultiStoreStatusResponse>(MultiStoreStatusResponse.class) {
381384
@Override
382385
public void internalHandle(Request request, MultiStoreStatusResponse veniceResponse) {
383386
AdminSparkServer.validateParams(request, FUTURE_VERSION.getParams(), admin);
384387
String clusterName = request.queryParams(CLUSTER);
385388
String storeName = request.queryParams(NAME);
386-
veniceResponse.setCluster(clusterName);
387-
Store store = admin.getStore(clusterName, storeName);
388-
if (store == null) {
389-
throw new VeniceNoStoreException(storeName);
390-
}
391-
Map<String, String> storeStatusMap = admin.getFutureVersionsForMultiColos(clusterName, storeName);
392-
if (storeStatusMap.isEmpty()) {
393-
// Non parent controllers will return an empty map, so we'll just return the children version of this api
394-
storeStatusMap =
395-
Collections.singletonMap(storeName, String.valueOf(admin.getFutureVersion(clusterName, storeName)));
396-
}
397-
veniceResponse.setStoreStatusMap(storeStatusMap);
389+
390+
// Convert to gRPC request
391+
ClusterStoreGrpcInfo storeInfo =
392+
ClusterStoreGrpcInfo.newBuilder().setClusterName(clusterName).setStoreName(storeName).build();
393+
GetFutureVersionGrpcRequest grpcRequest =
394+
GetFutureVersionGrpcRequest.newBuilder().setStoreInfo(storeInfo).build();
395+
396+
// Call handler
397+
GetFutureVersionGrpcResponse grpcResponse = requestHandler.getFutureVersion(grpcRequest);
398+
399+
// Map response back to HTTP
400+
veniceResponse.setCluster(grpcResponse.getStoreInfo().getClusterName());
401+
veniceResponse.setStoreStatusMap(grpcResponse.getStoreVersionMapMap());
398402
}
399403
};
400404
}

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

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,8 @@
44

55
import com.linkedin.venice.protocols.controller.DiscoverClusterGrpcRequest;
66
import com.linkedin.venice.protocols.controller.DiscoverClusterGrpcResponse;
7+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcRequest;
8+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcResponse;
79
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcRequest;
810
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcResponse;
911
import com.linkedin.venice.protocols.controller.VeniceControllerGrpcServiceGrpc;
@@ -52,4 +54,19 @@ public void discoverClusterForStore(
5254
null,
5355
grpcRequest.getStoreName());
5456
}
57+
58+
@Override
59+
public void getFutureVersion(
60+
GetFutureVersionGrpcRequest request,
61+
StreamObserver<GetFutureVersionGrpcResponse> responseObserver) {
62+
LOGGER.debug("Received getFutureVersion with args: {}", request);
63+
String clusterName = request.hasStoreInfo() ? request.getStoreInfo().getClusterName() : null;
64+
String storeName = request.hasStoreInfo() ? request.getStoreInfo().getStoreName() : null;
65+
handleRequest(
66+
VeniceControllerGrpcServiceGrpc.getGetFutureVersionMethod(),
67+
() -> requestHandler.getFutureVersion(request),
68+
responseObserver,
69+
clusterName,
70+
storeName);
71+
}
5572
}

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

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,19 @@
22

33
import com.linkedin.venice.controller.Admin;
44
import com.linkedin.venice.controller.ControllerRequestHandlerDependencies;
5+
import com.linkedin.venice.exceptions.VeniceNoStoreException;
56
import com.linkedin.venice.meta.Instance;
7+
import com.linkedin.venice.meta.Store;
8+
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
69
import com.linkedin.venice.protocols.controller.DiscoverClusterGrpcRequest;
710
import com.linkedin.venice.protocols.controller.DiscoverClusterGrpcResponse;
11+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcRequest;
12+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcResponse;
813
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcRequest;
914
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcResponse;
1015
import com.linkedin.venice.utils.Pair;
16+
import java.util.Collections;
17+
import java.util.Map;
1118
import org.apache.commons.lang.StringUtils;
1219
import org.apache.logging.log4j.LogManager;
1320
import org.apache.logging.log4j.Logger;
@@ -113,4 +120,37 @@ public DiscoverClusterGrpcResponse discoverCluster(DiscoverClusterGrpcRequest re
113120
public VeniceControllerAccessManager getControllerAccessManager() {
114121
return accessManager;
115122
}
123+
124+
/**
125+
* Gets the future version for a store in a cluster.
126+
* For parent controllers, returns future versions for all colos.
127+
* For child controllers, returns the local future version.
128+
* @param request the request containing cluster and store name
129+
* @return response containing a map of store/region to future version
130+
*/
131+
public GetFutureVersionGrpcResponse getFutureVersion(GetFutureVersionGrpcRequest request) {
132+
ClusterStoreGrpcInfo storeInfo = request.getStoreInfo();
133+
ControllerRequestParamValidator.validateClusterStoreInfo(storeInfo);
134+
String clusterName = storeInfo.getClusterName();
135+
String storeName = storeInfo.getStoreName();
136+
137+
LOGGER.info("Getting future version for store: {} in cluster: {}", storeName, clusterName);
138+
139+
Store store = admin.getStore(clusterName, storeName);
140+
if (store == null) {
141+
throw new VeniceNoStoreException(storeName);
142+
}
143+
144+
Map<String, String> storeVersionMap = admin.getFutureVersionsForMultiColos(clusterName, storeName);
145+
if (storeVersionMap.isEmpty()) {
146+
// Non parent controllers will return an empty map, so we'll just return the child version of this api
147+
storeVersionMap =
148+
Collections.singletonMap(storeName, String.valueOf(admin.getFutureVersion(clusterName, storeName)));
149+
}
150+
151+
return GetFutureVersionGrpcResponse.newBuilder()
152+
.setStoreInfo(storeInfo)
153+
.putAllStoreVersionMap(storeVersionMap)
154+
.build();
155+
}
116156
}

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

Lines changed: 43 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@
77
import static org.mockito.ArgumentMatchers.any;
88
import static org.mockito.ArgumentMatchers.eq;
99
import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
10-
import static org.mockito.Mockito.doCallRealMethod;
1110
import static org.mockito.Mockito.doReturn;
1211
import static org.mockito.Mockito.doThrow;
1312
import static org.mockito.Mockito.mock;
@@ -28,6 +27,9 @@
2827
import com.linkedin.venice.meta.RoutingStrategy;
2928
import com.linkedin.venice.meta.Store;
3029
import com.linkedin.venice.meta.ZKStore;
30+
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
31+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcRequest;
32+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcResponse;
3133
import com.linkedin.venice.pubsub.PubSubTopicRepository;
3234
import com.linkedin.venice.utils.ObjectMapperFactory;
3335
import java.util.Collections;
@@ -51,20 +53,30 @@ public class StoresRoutesTest {
5153
@Test
5254
public void testGetFutureVersion() throws Exception {
5355
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
56+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
5457
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
5558

56-
Store mockStore = mock(Store.class);
57-
doReturn(mockStore).when(mockAdmin).getStore(TEST_CLUSTER, TEST_STORE_NAME);
58-
5959
Map<String, String> storeStatusMap = Collections.singletonMap("dc-0", "1");
60-
doReturn(storeStatusMap).when(mockAdmin).getFutureVersionsForMultiColos(TEST_CLUSTER, TEST_STORE_NAME);
60+
61+
ClusterStoreGrpcInfo storeInfo =
62+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE_NAME).build();
63+
GetFutureVersionGrpcResponse grpcResponse =
64+
GetFutureVersionGrpcResponse.newBuilder().setStoreInfo(storeInfo).putAllStoreVersionMap(storeStatusMap).build();
65+
when(mockRequestHandler.getFutureVersion(any(GetFutureVersionGrpcRequest.class))).thenReturn(grpcResponse);
6166

6267
Request request = mock(Request.class);
6368
doReturn(TEST_CLUSTER).when(request).queryParams(eq(ControllerApiConstants.CLUSTER));
6469
doReturn(TEST_STORE_NAME).when(request).queryParams(eq(ControllerApiConstants.NAME));
6570

66-
Route getFutureVersionRoute =
67-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
71+
QueryParamsMap queryParamsMap = mock(QueryParamsMap.class);
72+
Map<String, String[]> queryMap = new HashMap<>();
73+
queryMap.put(ControllerApiConstants.CLUSTER, new String[] { TEST_CLUSTER });
74+
queryMap.put(ControllerApiConstants.NAME, new String[] { TEST_STORE_NAME });
75+
doReturn(queryMap).when(queryParamsMap).toMap();
76+
doReturn(queryParamsMap).when(request).queryMap();
77+
78+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
79+
.getFutureVersion(mockAdmin, mockRequestHandler);
6880
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
6981
.readValue(
7082
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),
@@ -199,42 +211,41 @@ public void testCleanExecutionIds() throws Exception {
199211
@Test
200212
public void testGetFutureVersionForChildController() throws Exception {
201213
Admin mockAdmin = mock(VeniceHelixAdmin.class);
214+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
202215
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
203216

204-
Store mockStore = mock(Store.class);
205-
doReturn(mockStore).when(mockAdmin).getStore(TEST_CLUSTER, TEST_STORE_NAME);
206-
207-
doCallRealMethod().when(mockAdmin).getFutureVersionsForMultiColos(TEST_CLUSTER, TEST_STORE_NAME);
208-
doReturn(1).when(mockAdmin).getFutureVersion(TEST_CLUSTER, TEST_STORE_NAME);
217+
Map<String, String> storeStatusMap = Collections.singletonMap(TEST_STORE_NAME, "1");
218+
ClusterStoreGrpcInfo storeInfo =
219+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE_NAME).build();
220+
GetFutureVersionGrpcResponse grpcResponse =
221+
GetFutureVersionGrpcResponse.newBuilder().setStoreInfo(storeInfo).putAllStoreVersionMap(storeStatusMap).build();
222+
when(mockRequestHandler.getFutureVersion(any(GetFutureVersionGrpcRequest.class))).thenReturn(grpcResponse);
209223

210224
Request request = mock(Request.class);
211225
doReturn(TEST_CLUSTER).when(request).queryParams(eq(ControllerApiConstants.CLUSTER));
212226
doReturn(TEST_STORE_NAME).when(request).queryParams(eq(ControllerApiConstants.NAME));
213227

214-
Route getFutureVersionRoute =
215-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
228+
QueryParamsMap queryParamsMap = mock(QueryParamsMap.class);
229+
Map<String, String[]> queryMap = new HashMap<>();
230+
queryMap.put(ControllerApiConstants.CLUSTER, new String[] { TEST_CLUSTER });
231+
queryMap.put(ControllerApiConstants.NAME, new String[] { TEST_STORE_NAME });
232+
doReturn(queryMap).when(queryParamsMap).toMap();
233+
doReturn(queryParamsMap).when(request).queryMap();
234+
235+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
236+
.getFutureVersion(mockAdmin, mockRequestHandler);
216237
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
217238
.readValue(
218239
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),
219240
MultiStoreStatusResponse.class);
220241
Assert.assertEquals(multiStoreStatusResponse.getCluster(), TEST_CLUSTER);
221242
Assert.assertEquals(multiStoreStatusResponse.getStoreStatusMap(), Collections.singletonMap(TEST_STORE_NAME, "1"));
222-
223-
doCallRealMethod().when(mockAdmin).getBackupVersionsForMultiColos(TEST_CLUSTER, TEST_STORE_NAME);
224-
doReturn(2).when(mockAdmin).getBackupVersion(TEST_CLUSTER, TEST_STORE_NAME);
225-
Route getBackupVersionRoute =
226-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getBackupVersion(mockAdmin);
227-
multiStoreStatusResponse = ObjectMapperFactory.getInstance()
228-
.readValue(
229-
getBackupVersionRoute.handle(request, mock(Response.class)).toString(),
230-
MultiStoreStatusResponse.class);
231-
Assert.assertEquals(multiStoreStatusResponse.getCluster(), TEST_CLUSTER);
232-
Assert.assertEquals(multiStoreStatusResponse.getStoreStatusMap(), Collections.singletonMap(TEST_STORE_NAME, "2"));
233243
}
234244

235245
@Test
236246
public void testGetFutureVersionWhenNotLeaderController() throws Exception {
237247
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
248+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
238249
doReturn(false).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
239250

240251
Store mockStore = mock(Store.class);
@@ -253,8 +264,8 @@ public void testGetFutureVersionWhenNotLeaderController() throws Exception {
253264
doReturn(queryMap).when(queryParamsMap).toMap();
254265
doReturn(queryParamsMap).when(request).queryMap();
255266

256-
Route getFutureVersionRoute =
257-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
267+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
268+
.getFutureVersion(mockAdmin, mockRequestHandler);
258269
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
259270
.readValue(
260271
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),
@@ -266,9 +277,11 @@ public void testGetFutureVersionWhenNotLeaderController() throws Exception {
266277
@Test
267278
public void testGetFutureVersionWhenStoreNotExist() throws Exception {
268279
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
280+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
269281
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
270282

271-
doReturn(null).when(mockAdmin).getStore(TEST_CLUSTER, TEST_STORE_NAME);
283+
when(mockRequestHandler.getFutureVersion(any(GetFutureVersionGrpcRequest.class)))
284+
.thenThrow(new com.linkedin.venice.exceptions.VeniceNoStoreException(TEST_STORE_NAME));
272285

273286
Request request = mock(Request.class);
274287
doReturn(TEST_CLUSTER).when(request).queryParams(eq(ControllerApiConstants.CLUSTER));
@@ -283,8 +296,8 @@ public void testGetFutureVersionWhenStoreNotExist() throws Exception {
283296
doReturn(queryMap).when(queryParamsMap).toMap();
284297
doReturn(queryParamsMap).when(request).queryMap();
285298

286-
Route getFutureVersionRoute =
287-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
299+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
300+
.getFutureVersion(mockAdmin, mockRequestHandler);
288301
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
289302
.readValue(
290303
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),

0 commit comments

Comments
 (0)