Skip to content

Commit 5240ed1

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 8109fcd commit 5240ed1

9 files changed

Lines changed: 276 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
@@ -341,7 +341,7 @@ public boolean startInner() throws Exception {
341341
httpService.get(STORE.getPath(), new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getStore(admin)));
342342
httpService.get(
343343
FUTURE_VERSION.getPath(),
344-
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getFutureVersion(admin)));
344+
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getFutureVersion(admin, requestHandler)));
345345
httpService.get(
346346
BACKUP_VERSION.getPath(),
347347
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
@@ -108,7 +108,10 @@
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.ZKStore;
111112
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
113+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcRequest;
114+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcResponse;
112115
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
113116
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
114117
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
@@ -307,25 +310,26 @@ public void internalHandle(Request request, StoreResponse veniceResponse) {
307310
};
308311
}
309312

310-
public Route getFutureVersion(Admin admin) {
313+
public Route getFutureVersion(Admin admin, VeniceControllerRequestHandler requestHandler) {
311314
return new VeniceRouteHandler<MultiStoreStatusResponse>(MultiStoreStatusResponse.class) {
312315
@Override
313316
public void internalHandle(Request request, MultiStoreStatusResponse veniceResponse) {
314317
AdminSparkServer.validateParams(request, FUTURE_VERSION.getParams(), admin);
315318
String clusterName = request.queryParams(CLUSTER);
316319
String storeName = request.queryParams(NAME);
317-
veniceResponse.setCluster(clusterName);
318-
Store store = admin.getStore(clusterName, storeName);
319-
if (store == null) {
320-
throw new VeniceNoStoreException(storeName);
321-
}
322-
Map<String, String> storeStatusMap = admin.getFutureVersionsForMultiColos(clusterName, storeName);
323-
if (storeStatusMap.isEmpty()) {
324-
// Non parent controllers will return an empty map, so we'll just return the children version of this api
325-
storeStatusMap =
326-
Collections.singletonMap(storeName, String.valueOf(admin.getFutureVersion(clusterName, storeName)));
327-
}
328-
veniceResponse.setStoreStatusMap(storeStatusMap);
320+
321+
// Convert to gRPC request
322+
ClusterStoreGrpcInfo storeInfo =
323+
ClusterStoreGrpcInfo.newBuilder().setClusterName(clusterName).setStoreName(storeName).build();
324+
GetFutureVersionGrpcRequest grpcRequest =
325+
GetFutureVersionGrpcRequest.newBuilder().setStoreInfo(storeInfo).build();
326+
327+
// Call handler
328+
GetFutureVersionGrpcResponse grpcResponse = requestHandler.getFutureVersion(grpcRequest);
329+
330+
// Map response back to HTTP
331+
veniceResponse.setCluster(grpcResponse.getStoreInfo().getClusterName());
332+
veniceResponse.setStoreStatusMap(grpcResponse.getStoreVersionMapMap());
329333
}
330334
};
331335
}

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: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,11 +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;
15+
import com.linkedin.venice.utils.Pair;
16+
import java.util.Collections;
17+
import java.util.Map;
1018
import org.apache.commons.lang.StringUtils;
1119
import org.apache.logging.log4j.LogManager;
1220
import org.apache.logging.log4j.Logger;
@@ -113,4 +121,37 @@ public DiscoverClusterGrpcResponse discoverCluster(DiscoverClusterGrpcRequest re
113121
public VeniceControllerAccessManager getControllerAccessManager() {
114122
return accessManager;
115123
}
124+
125+
/**
126+
* Gets the future version for a store in a cluster.
127+
* For parent controllers, returns future versions for all colos.
128+
* For child controllers, returns the local future version.
129+
* @param request the request containing cluster and store name
130+
* @return response containing a map of store/region to future version
131+
*/
132+
public GetFutureVersionGrpcResponse getFutureVersion(GetFutureVersionGrpcRequest request) {
133+
ClusterStoreGrpcInfo storeInfo = request.getStoreInfo();
134+
ControllerRequestParamValidator.validateClusterStoreInfo(storeInfo);
135+
String clusterName = storeInfo.getClusterName();
136+
String storeName = storeInfo.getStoreName();
137+
138+
LOGGER.info("Getting future version for store: {} in cluster: {}", storeName, clusterName);
139+
140+
Store store = admin.getStore(clusterName, storeName);
141+
if (store == null) {
142+
throw new VeniceNoStoreException(storeName);
143+
}
144+
145+
Map<String, String> storeVersionMap = admin.getFutureVersionsForMultiColos(clusterName, storeName);
146+
if (storeVersionMap.isEmpty()) {
147+
// Non parent controllers will return an empty map, so we'll just return the child version of this api
148+
storeVersionMap =
149+
Collections.singletonMap(storeName, String.valueOf(admin.getFutureVersion(clusterName, storeName)));
150+
}
151+
152+
return GetFutureVersionGrpcResponse.newBuilder()
153+
.setStoreInfo(storeInfo)
154+
.putAllStoreVersionMap(storeVersionMap)
155+
.build();
156+
}
116157
}

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

Lines changed: 42 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;
@@ -37,6 +36,8 @@
3736
import com.linkedin.venice.meta.StoreInfo;
3837
import com.linkedin.venice.meta.ZKStore;
3938
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
39+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcRequest;
40+
import com.linkedin.venice.protocols.controller.GetFutureVersionGrpcResponse;
4041
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
4142
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
4243
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
@@ -66,20 +67,30 @@ public class StoresRoutesTest {
6667
@Test
6768
public void testGetFutureVersion() throws Exception {
6869
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
70+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
6971
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
7072

71-
Store mockStore = mock(Store.class);
72-
doReturn(mockStore).when(mockAdmin).getStore(TEST_CLUSTER, TEST_STORE_NAME);
73-
7473
Map<String, String> storeStatusMap = Collections.singletonMap("dc-0", "1");
75-
doReturn(storeStatusMap).when(mockAdmin).getFutureVersionsForMultiColos(TEST_CLUSTER, TEST_STORE_NAME);
74+
75+
ClusterStoreGrpcInfo storeInfo =
76+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE_NAME).build();
77+
GetFutureVersionGrpcResponse grpcResponse =
78+
GetFutureVersionGrpcResponse.newBuilder().setStoreInfo(storeInfo).putAllStoreVersionMap(storeStatusMap).build();
79+
when(mockRequestHandler.getFutureVersion(any(GetFutureVersionGrpcRequest.class))).thenReturn(grpcResponse);
7680

7781
Request request = mock(Request.class);
7882
doReturn(TEST_CLUSTER).when(request).queryParams(eq(ControllerApiConstants.CLUSTER));
7983
doReturn(TEST_STORE_NAME).when(request).queryParams(eq(ControllerApiConstants.NAME));
8084

81-
Route getFutureVersionRoute =
82-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
85+
QueryParamsMap queryParamsMap = mock(QueryParamsMap.class);
86+
Map<String, String[]> queryMap = new HashMap<>();
87+
queryMap.put(ControllerApiConstants.CLUSTER, new String[] { TEST_CLUSTER });
88+
queryMap.put(ControllerApiConstants.NAME, new String[] { TEST_STORE_NAME });
89+
doReturn(queryMap).when(queryParamsMap).toMap();
90+
doReturn(queryParamsMap).when(request).queryMap();
91+
92+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
93+
.getFutureVersion(mockAdmin, mockRequestHandler);
8394
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
8495
.readValue(
8596
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),
@@ -214,42 +225,41 @@ public void testCleanExecutionIds() throws Exception {
214225
@Test
215226
public void testGetFutureVersionForChildController() throws Exception {
216227
Admin mockAdmin = mock(VeniceHelixAdmin.class);
228+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
217229
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
218230

219-
Store mockStore = mock(Store.class);
220-
doReturn(mockStore).when(mockAdmin).getStore(TEST_CLUSTER, TEST_STORE_NAME);
221-
222-
doCallRealMethod().when(mockAdmin).getFutureVersionsForMultiColos(TEST_CLUSTER, TEST_STORE_NAME);
223-
doReturn(1).when(mockAdmin).getFutureVersion(TEST_CLUSTER, TEST_STORE_NAME);
231+
Map<String, String> storeStatusMap = Collections.singletonMap(TEST_STORE_NAME, "1");
232+
ClusterStoreGrpcInfo storeInfo =
233+
ClusterStoreGrpcInfo.newBuilder().setClusterName(TEST_CLUSTER).setStoreName(TEST_STORE_NAME).build();
234+
GetFutureVersionGrpcResponse grpcResponse =
235+
GetFutureVersionGrpcResponse.newBuilder().setStoreInfo(storeInfo).putAllStoreVersionMap(storeStatusMap).build();
236+
when(mockRequestHandler.getFutureVersion(any(GetFutureVersionGrpcRequest.class))).thenReturn(grpcResponse);
224237

225238
Request request = mock(Request.class);
226239
doReturn(TEST_CLUSTER).when(request).queryParams(eq(ControllerApiConstants.CLUSTER));
227240
doReturn(TEST_STORE_NAME).when(request).queryParams(eq(ControllerApiConstants.NAME));
228241

229-
Route getFutureVersionRoute =
230-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
242+
QueryParamsMap queryParamsMap = mock(QueryParamsMap.class);
243+
Map<String, String[]> queryMap = new HashMap<>();
244+
queryMap.put(ControllerApiConstants.CLUSTER, new String[] { TEST_CLUSTER });
245+
queryMap.put(ControllerApiConstants.NAME, new String[] { TEST_STORE_NAME });
246+
doReturn(queryMap).when(queryParamsMap).toMap();
247+
doReturn(queryParamsMap).when(request).queryMap();
248+
249+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
250+
.getFutureVersion(mockAdmin, mockRequestHandler);
231251
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
232252
.readValue(
233253
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),
234254
MultiStoreStatusResponse.class);
235255
Assert.assertEquals(multiStoreStatusResponse.getCluster(), TEST_CLUSTER);
236256
Assert.assertEquals(multiStoreStatusResponse.getStoreStatusMap(), Collections.singletonMap(TEST_STORE_NAME, "1"));
237-
238-
doCallRealMethod().when(mockAdmin).getBackupVersionsForMultiColos(TEST_CLUSTER, TEST_STORE_NAME);
239-
doReturn(2).when(mockAdmin).getBackupVersion(TEST_CLUSTER, TEST_STORE_NAME);
240-
Route getBackupVersionRoute =
241-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getBackupVersion(mockAdmin);
242-
multiStoreStatusResponse = ObjectMapperFactory.getInstance()
243-
.readValue(
244-
getBackupVersionRoute.handle(request, mock(Response.class)).toString(),
245-
MultiStoreStatusResponse.class);
246-
Assert.assertEquals(multiStoreStatusResponse.getCluster(), TEST_CLUSTER);
247-
Assert.assertEquals(multiStoreStatusResponse.getStoreStatusMap(), Collections.singletonMap(TEST_STORE_NAME, "2"));
248257
}
249258

250259
@Test
251260
public void testGetFutureVersionWhenNotLeaderController() throws Exception {
252261
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
262+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
253263
doReturn(false).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
254264

255265
Store mockStore = mock(Store.class);
@@ -268,8 +278,8 @@ public void testGetFutureVersionWhenNotLeaderController() throws Exception {
268278
doReturn(queryMap).when(queryParamsMap).toMap();
269279
doReturn(queryParamsMap).when(request).queryMap();
270280

271-
Route getFutureVersionRoute =
272-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
281+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
282+
.getFutureVersion(mockAdmin, mockRequestHandler);
273283
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
274284
.readValue(
275285
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),
@@ -281,9 +291,11 @@ public void testGetFutureVersionWhenNotLeaderController() throws Exception {
281291
@Test
282292
public void testGetFutureVersionWhenStoreNotExist() throws Exception {
283293
Admin mockAdmin = mock(VeniceParentHelixAdmin.class);
294+
VeniceControllerRequestHandler mockRequestHandler = mock(VeniceControllerRequestHandler.class);
284295
doReturn(true).when(mockAdmin).isLeaderControllerFor(TEST_CLUSTER);
285296

286-
doReturn(null).when(mockAdmin).getStore(TEST_CLUSTER, TEST_STORE_NAME);
297+
when(mockRequestHandler.getFutureVersion(any(GetFutureVersionGrpcRequest.class)))
298+
.thenThrow(new com.linkedin.venice.exceptions.VeniceNoStoreException(TEST_STORE_NAME));
287299

288300
Request request = mock(Request.class);
289301
doReturn(TEST_CLUSTER).when(request).queryParams(eq(ControllerApiConstants.CLUSTER));
@@ -298,8 +310,8 @@ public void testGetFutureVersionWhenStoreNotExist() throws Exception {
298310
doReturn(queryMap).when(queryParamsMap).toMap();
299311
doReturn(queryParamsMap).when(request).queryMap();
300312

301-
Route getFutureVersionRoute =
302-
new StoresRoutes(false, Optional.empty(), pubSubTopicRepository).getFutureVersion(mockAdmin);
313+
Route getFutureVersionRoute = new StoresRoutes(false, Optional.empty(), pubSubTopicRepository)
314+
.getFutureVersion(mockAdmin, mockRequestHandler);
303315
MultiStoreStatusResponse multiStoreStatusResponse = ObjectMapperFactory.getInstance()
304316
.readValue(
305317
getFutureVersionRoute.handle(request, mock(Response.class)).toString(),

0 commit comments

Comments
 (0)