Skip to content

Commit 1225a93

Browse files
authored
[Controller] Add gRPC support for listStores API (linkedin#2413)
Add gRPC support for the listStores API while maintaining backward compatibility with the existing HTTP endpoint. Changes: - Add listStores RPC to StoreGrpcService.proto - Add handler method to StoreRequestHandler - Add gRPC service implementation (no ACL check - consistent with HTTP) - Update HTTP route to delegate to shared handler Tests: - Unit tests for success, error, invalid input, and filter handling - Integration tests for full lifecycle"
1 parent 4d4e34f commit 1225a93

9 files changed

Lines changed: 615 additions & 114 deletions

File tree

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ service StoreGrpcService {
1313
rpc deleteAclForStore(DeleteAclForStoreGrpcRequest) returns (DeleteAclForStoreGrpcResponse);
1414
rpc checkResourceCleanupForStoreCreation(ClusterStoreGrpcInfo) returns (ResourceCleanupCheckGrpcResponse) {}
1515
rpc validateStoreDeleted(ValidateStoreDeletedGrpcRequest) returns (ValidateStoreDeletedGrpcResponse);
16+
rpc listStores(ListStoresGrpcRequest) returns (ListStoresGrpcResponse);
1617
}
1718

1819
message CreateStoreGrpcRequest {
@@ -69,4 +70,16 @@ message ValidateStoreDeletedGrpcResponse {
6970
ClusterStoreGrpcInfo storeInfo = 1;
7071
bool storeDeleted = 2;
7172
optional string reason = 3;
73+
}
74+
75+
message ListStoresGrpcRequest {
76+
string clusterName = 1;
77+
optional bool includeSystemStores = 2;
78+
optional string storeConfigNameFilter = 3;
79+
optional string storeConfigValueFilter = 4;
80+
}
81+
82+
message ListStoresGrpcResponse {
83+
string clusterName = 1;
84+
repeated string storeNames = 2;
7285
}

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

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import static org.testng.Assert.assertEquals;
66
import static org.testng.Assert.assertFalse;
77
import static org.testng.Assert.assertNotNull;
8+
import static org.testng.Assert.assertTrue;
89

910
import com.linkedin.venice.acl.NoOpDynamicAccessController;
1011
import com.linkedin.venice.authorization.Method;
@@ -20,6 +21,8 @@
2021
import com.linkedin.venice.protocols.controller.DiscoverClusterGrpcResponse;
2122
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcRequest;
2223
import com.linkedin.venice.protocols.controller.LeaderControllerGrpcResponse;
24+
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
25+
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
2326
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
2427
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
2528
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
@@ -246,6 +249,73 @@ public void testValidateStoreDeletedOverSecureGrpcChannel() {
246249
assertEquals(response.getStoreInfo().getClusterName(), veniceCluster.getClusterName());
247250
}
248251

252+
@Test(timeOut = TIMEOUT_MS)
253+
public void testListStoresGrpcEndpoint() {
254+
String storeName1 = Utils.getUniqueString("test_list_stores_1");
255+
String storeName2 = Utils.getUniqueString("test_list_stores_2");
256+
String controllerGrpcUrl = veniceCluster.getLeaderVeniceController().getControllerGrpcUrl();
257+
ManagedChannel channel = Grpc.newChannelBuilder(controllerGrpcUrl, InsecureChannelCredentials.create()).build();
258+
StoreGrpcServiceGrpc.StoreGrpcServiceBlockingStub storeBlockingStub = StoreGrpcServiceGrpc.newBlockingStub(channel);
259+
260+
// Step 1: Create two stores
261+
CreateStoreGrpcRequest createStoreRequest1 = CreateStoreGrpcRequest.newBuilder()
262+
.setStoreInfo(
263+
ClusterStoreGrpcInfo.newBuilder()
264+
.setClusterName(veniceCluster.getClusterName())
265+
.setStoreName(storeName1)
266+
.build())
267+
.setOwner("owner")
268+
.setKeySchema(DEFAULT_KEY_SCHEMA)
269+
.setValueSchema("\"string\"")
270+
.build();
271+
CreateStoreGrpcResponse createResponse1 = storeBlockingStub.createStore(createStoreRequest1);
272+
assertNotNull(createResponse1, "Response should not be null");
273+
274+
CreateStoreGrpcRequest createStoreRequest2 = CreateStoreGrpcRequest.newBuilder()
275+
.setStoreInfo(
276+
ClusterStoreGrpcInfo.newBuilder()
277+
.setClusterName(veniceCluster.getClusterName())
278+
.setStoreName(storeName2)
279+
.build())
280+
.setOwner("owner")
281+
.setKeySchema(DEFAULT_KEY_SCHEMA)
282+
.setValueSchema("\"string\"")
283+
.build();
284+
CreateStoreGrpcResponse createResponse2 = storeBlockingStub.createStore(createStoreRequest2);
285+
assertNotNull(createResponse2, "Response should not be null");
286+
287+
// Step 2: List all stores in the cluster
288+
ListStoresGrpcRequest listStoresRequest =
289+
ListStoresGrpcRequest.newBuilder().setClusterName(veniceCluster.getClusterName()).build();
290+
291+
ListStoresGrpcResponse listStoresResponse = storeBlockingStub.listStores(listStoresRequest);
292+
assertNotNull(listStoresResponse, "Response should not be null");
293+
assertEquals(listStoresResponse.getClusterName(), veniceCluster.getClusterName());
294+
assertTrue(listStoresResponse.getStoreNamesList().contains(storeName1), "Store list should contain " + storeName1);
295+
assertTrue(listStoresResponse.getStoreNamesList().contains(storeName2), "Store list should contain " + storeName2);
296+
297+
// Step 3: List stores excluding system stores
298+
ListStoresGrpcRequest listStoresNoSystemRequest = ListStoresGrpcRequest.newBuilder()
299+
.setClusterName(veniceCluster.getClusterName())
300+
.setIncludeSystemStores(false)
301+
.build();
302+
303+
ListStoresGrpcResponse listStoresNoSystemResponse = storeBlockingStub.listStores(listStoresNoSystemRequest);
304+
assertNotNull(listStoresNoSystemResponse, "Response should not be null");
305+
assertTrue(
306+
listStoresNoSystemResponse.getStoreNamesList().contains(storeName1),
307+
"Store list should contain " + storeName1);
308+
assertTrue(
309+
listStoresNoSystemResponse.getStoreNamesList().contains(storeName2),
310+
"Store list should contain " + storeName2);
311+
// Verify no system stores are in the list
312+
for (String store: listStoresNoSystemResponse.getStoreNamesList()) {
313+
assertFalse(
314+
store.startsWith("venice_system_store") || store.contains("push_status"),
315+
"Store list should not contain system stores: " + store);
316+
}
317+
}
318+
249319
private static class MockDynamicAccessController extends NoOpDynamicAccessController {
250320
private final Set<String> resourcesInAllowList = ConcurrentHashMap.newKeySet();
251321

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

Lines changed: 18 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.ListStoresGrpcRequest;
18+
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
1719
import com.linkedin.venice.protocols.controller.ResourceCleanupCheckGrpcResponse;
1820
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
1921
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc.StoreGrpcServiceImplBase;
@@ -129,4 +131,20 @@ public void validateStoreDeleted(
129131
return storeRequestHandler.validateStoreDeleted(grpcRequest);
130132
}, responseObserver, clusterName, storeName);
131133
}
134+
135+
/**
136+
* Lists all stores in a cluster with optional filtering.
137+
* No ACL check; any user can list stores.
138+
*/
139+
@Override
140+
public void listStores(ListStoresGrpcRequest grpcRequest, StreamObserver<ListStoresGrpcResponse> responseObserver) {
141+
LOGGER.debug("Received listStores with args: {}", grpcRequest);
142+
String clusterName = grpcRequest.getClusterName();
143+
handleRequest(
144+
StoreGrpcServiceGrpc.getListStoresMethod(),
145+
() -> storeRequestHandler.listStores(grpcRequest),
146+
responseObserver,
147+
clusterName,
148+
null);
149+
}
132150
}

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

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -300,7 +300,8 @@ public boolean startInner() throws Exception {
300300
// Build all different routes
301301
ControllerRoutes controllerRoutes =
302302
new ControllerRoutes(sslEnabled, accessController, pubSubTopicRepository, requestHandler);
303-
StoresRoutes storesRoutes = new StoresRoutes(sslEnabled, accessController, pubSubTopicRepository);
303+
StoresRoutes storesRoutes =
304+
new StoresRoutes(sslEnabled, accessController, pubSubTopicRepository, requestHandler.getStoreRequestHandler());
304305
JobRoutes jobRoutes = new JobRoutes(sslEnabled, accessController);
305306
SkipAdminRoute skipAdminRoute = new SkipAdminRoute(sslEnabled, accessController);
306307
CreateVersion createVersion = new CreateVersion(sslEnabled, accessController, this.checkReadMethodForKafka);
@@ -708,9 +709,7 @@ public boolean startInner() throws Exception {
708709
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.getInUseSchemaIds(admin)));
709710
httpService.get(
710711
VALIDATE_STORE_DELETED.getPath(),
711-
new VeniceParentControllerRegionStateHandler(
712-
admin,
713-
storesRoutes.validateStoreDeleted(admin, requestHandler.getStoreRequestHandler())));
712+
new VeniceParentControllerRegionStateHandler(admin, storesRoutes.validateStoreDeleted(admin)));
714713

715714
httpService.post(
716715
CLEANUP_INSTANCE_CUSTOMIZED_STATES.getPath(),

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

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,18 +3,27 @@
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.exceptions.VeniceException;
7+
import com.linkedin.venice.meta.Store;
8+
import com.linkedin.venice.meta.ZKStore;
69
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
710
import com.linkedin.venice.protocols.controller.CreateStoreGrpcRequest;
811
import com.linkedin.venice.protocols.controller.CreateStoreGrpcResponse;
912
import com.linkedin.venice.protocols.controller.DeleteAclForStoreGrpcRequest;
1013
import com.linkedin.venice.protocols.controller.DeleteAclForStoreGrpcResponse;
1114
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcRequest;
1215
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcResponse;
16+
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
17+
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
1318
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcRequest;
1419
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
1520
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
1621
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
22+
import com.linkedin.venice.systemstore.schemas.StoreProperties;
23+
import java.util.ArrayList;
24+
import java.util.List;
1725
import java.util.Optional;
26+
import org.apache.avro.Schema;
1827
import org.apache.commons.lang.StringUtils;
1928
import org.apache.logging.log4j.LogManager;
2029
import org.apache.logging.log4j.Logger;
@@ -132,4 +141,118 @@ public ValidateStoreDeletedGrpcResponse validateStoreDeleted(ValidateStoreDelete
132141
}
133142
return responseBuilder.build();
134143
}
144+
145+
/**
146+
* Lists all stores in the specified cluster with optional filtering.
147+
* @param request the request containing cluster name and optional filters
148+
* @return response containing the list of store names
149+
*/
150+
public ListStoresGrpcResponse listStores(ListStoresGrpcRequest request) {
151+
String clusterName = request.getClusterName();
152+
if (StringUtils.isBlank(clusterName)) {
153+
throw new IllegalArgumentException("Cluster name is required");
154+
}
155+
156+
boolean includeSystemStores = !request.hasIncludeSystemStores() || request.getIncludeSystemStores();
157+
Optional<String> storeConfigNameFilter =
158+
request.hasStoreConfigNameFilter() ? Optional.of(request.getStoreConfigNameFilter()) : Optional.empty();
159+
Optional<String> storeConfigValueFilter =
160+
request.hasStoreConfigValueFilter() ? Optional.of(request.getStoreConfigValueFilter()) : Optional.empty();
161+
162+
if (storeConfigNameFilter.isPresent() ^ storeConfigValueFilter.isPresent()) {
163+
throw new VeniceException(
164+
"Missing parameter: "
165+
+ (storeConfigNameFilter.isPresent() ? "store_config_value_filter" : "store_config_name_filter"));
166+
}
167+
168+
boolean isDataReplicationPolicyConfigFilter = false;
169+
Schema.Field configFilterField = null;
170+
if (storeConfigNameFilter.isPresent()) {
171+
configFilterField = StoreProperties.getClassSchema().getField(storeConfigNameFilter.get());
172+
if (configFilterField == null) {
173+
isDataReplicationPolicyConfigFilter = storeConfigNameFilter.get().equalsIgnoreCase("dataReplicationPolicy");
174+
if (!isDataReplicationPolicyConfigFilter) {
175+
throw new VeniceException(
176+
"The config name filter " + storeConfigNameFilter.get() + " is not a valid store config.");
177+
}
178+
}
179+
}
180+
181+
LOGGER.info(
182+
"Listing stores in cluster: {} with includeSystemStores: {}, configNameFilter: {}, configValueFilter: {}",
183+
clusterName,
184+
includeSystemStores,
185+
storeConfigNameFilter.orElse("none"),
186+
storeConfigValueFilter.orElse("none"));
187+
188+
List<Store> storeList = admin.getAllStores(clusterName);
189+
List<String> selectedStoreNames = new ArrayList<>();
190+
191+
for (Store store: storeList) {
192+
if (!includeSystemStores && store.isSystemStore()) {
193+
continue;
194+
}
195+
if (storeConfigValueFilter.isPresent()) {
196+
boolean configValueMatch = false;
197+
if (isDataReplicationPolicyConfigFilter) {
198+
if (!store.isHybrid() || store.getHybridStoreConfig().getDataReplicationPolicy() == null) {
199+
continue;
200+
}
201+
configValueMatch = store.getHybridStoreConfig()
202+
.getDataReplicationPolicy()
203+
.name()
204+
.equalsIgnoreCase(storeConfigValueFilter.get());
205+
} else {
206+
ZKStore cloneStore = new ZKStore(store);
207+
Object configValue = cloneStore.dataModel().get(storeConfigNameFilter.get());
208+
if (configValue == null) {
209+
continue;
210+
}
211+
Schema fieldSchema = configFilterField.schema();
212+
switch (fieldSchema.getType()) {
213+
case BOOLEAN:
214+
configValueMatch = Boolean.valueOf(storeConfigValueFilter.get()).equals((Boolean) configValue);
215+
break;
216+
case INT:
217+
configValueMatch = Integer.valueOf(storeConfigValueFilter.get()).equals((Integer) configValue);
218+
break;
219+
case LONG:
220+
configValueMatch = Long.valueOf(storeConfigValueFilter.get()).equals((Long) configValue);
221+
break;
222+
case FLOAT:
223+
configValueMatch = Float.valueOf(storeConfigValueFilter.get()).equals(configValue);
224+
break;
225+
case DOUBLE:
226+
configValueMatch = Double.valueOf(storeConfigValueFilter.get()).equals(configValue);
227+
break;
228+
case STRING:
229+
configValueMatch = storeConfigValueFilter.get().equals(configValue);
230+
break;
231+
case ENUM:
232+
configValueMatch = storeConfigValueFilter.get().equals(configValue.toString());
233+
break;
234+
case UNION:
235+
configValueMatch = (configValue != null);
236+
break;
237+
case ARRAY:
238+
case MAP:
239+
case FIXED:
240+
case BYTES:
241+
case RECORD:
242+
case NULL:
243+
default:
244+
throw new VeniceException(
245+
"Store config filtering for Schema type " + fieldSchema.getType().toString() + " is not supported");
246+
}
247+
}
248+
if (!configValueMatch) {
249+
continue;
250+
}
251+
}
252+
selectedStoreNames.add(store.getName());
253+
}
254+
255+
LOGGER.info("Found {} stores in cluster: {}", selectedStoreNames.size(), clusterName);
256+
return ListStoresGrpcResponse.newBuilder().setClusterName(clusterName).addAllStoreNames(selectedStoreNames).build();
257+
}
135258
}

0 commit comments

Comments
 (0)