Skip to content

Commit 5daeafd

Browse files
committed
Address PR review comments from sushantmane
- Extract VersionGrpc/VersionStatusGrpc to shared ControllerGrpcRequestContext.proto - Make pushJobId optional in VersionGrpc proto to handle null values - Handler returns RepushInfo POJO directly instead of wrapping in RepushInfoResponse - Add null guard for VersionStatusGrpc.forNumber() in convertVersionToProto - Add null guard for pushJobId before passing to protobuf setter - Add null check for repushInfo in handler, throw VeniceException if null - Change LOGGER.info() to LOGGER.debug() for getRepushInfo logging - Add test for null repushInfo case in StoreRequestHandlerTest
1 parent 045ccd4 commit 5daeafd

8 files changed

Lines changed: 75 additions & 94 deletions

File tree

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

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,3 +34,24 @@ message VeniceControllerGrpcErrorInfo {
3434
optional string clusterName = 4;
3535
optional string storeName = 5;
3636
}
37+
38+
// Shared version-related types used by multiple gRPC services
39+
enum VersionStatusGrpc {
40+
NOT_CREATED = 0;
41+
STARTED = 1;
42+
PUSHED = 2;
43+
ONLINE = 3;
44+
ERROR = 4;
45+
CREATED = 5;
46+
PARTIALLY_ONLINE = 6;
47+
KILLED = 7;
48+
}
49+
50+
message VersionGrpc {
51+
int32 number = 1;
52+
int64 createdTime = 2;
53+
VersionStatusGrpc status = 3;
54+
optional string pushJobId = 4;
55+
int32 partitionCount = 5;
56+
int32 replicationFactor = 6;
57+
}

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

Lines changed: 0 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -103,26 +103,6 @@ message RepushInfoGrpc {
103103
optional string systemSchemaClusterD2ZkHost = 4;
104104
}
105105

106-
enum VersionStatusGrpc {
107-
NOT_CREATED = 0;
108-
STARTED = 1;
109-
PUSHED = 2;
110-
ONLINE = 3;
111-
ERROR = 4;
112-
CREATED = 5;
113-
PARTIALLY_ONLINE = 6;
114-
KILLED = 7;
115-
}
116-
117-
message VersionGrpc {
118-
int32 number = 1;
119-
int64 createdTime = 2;
120-
VersionStatusGrpc status = 3;
121-
string pushJobId = 4;
122-
int32 partitionCount = 5;
123-
int32 replicationFactor = 6;
124-
}
125-
126106
message GetStoreGrpcRequest {
127107
ClusterStoreGrpcInfo storeInfo = 1;
128108
}

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

Lines changed: 14 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@
77
import com.linkedin.venice.controller.server.StoreRequestHandler;
88
import com.linkedin.venice.controller.server.VeniceControllerAccessManager;
99
import com.linkedin.venice.controllerapi.RepushInfo;
10-
import com.linkedin.venice.controllerapi.RepushInfoResponse;
1110
import com.linkedin.venice.exceptions.VeniceUnauthorizedAccessException;
1211
import com.linkedin.venice.meta.StoreInfo;
1312
import com.linkedin.venice.meta.Version;
@@ -178,11 +177,10 @@ public void getRepushInfo(
178177
String storeName = storeInfo.getStoreName();
179178
Optional<String> fabric = request.hasFabric() ? Optional.of(request.getFabric()) : Optional.empty();
180179

181-
// Call handler - returns POJO
182-
RepushInfoResponse result = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
180+
// Call handler - returns RepushInfo POJO directly
181+
RepushInfo repushInfo = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
183182

184183
// Convert POJO to protobuf response
185-
RepushInfo repushInfo = result.getRepushInfo();
186184
RepushInfoGrpc.Builder repushInfoBuilder = RepushInfoGrpc.newBuilder()
187185
.setPubSubUrl(repushInfo.getKafkaBrokerUrl() != null ? repushInfo.getKafkaBrokerUrl() : "");
188186

@@ -207,14 +205,21 @@ public void getRepushInfo(
207205
* Converts a Version object to protobuf VersionGrpc.
208206
*/
209207
private VersionGrpc convertVersionToProto(Version version) {
210-
return VersionGrpc.newBuilder()
208+
VersionStatusGrpc statusGrpc = VersionStatusGrpc.forNumber(version.getStatus().getValue());
209+
if (statusGrpc == null) {
210+
throw new IllegalArgumentException(
211+
"Unknown VersionStatus value: " + version.getStatus().getValue() + " (" + version.getStatus() + ")");
212+
}
213+
VersionGrpc.Builder builder = VersionGrpc.newBuilder()
211214
.setNumber(version.getNumber())
212215
.setCreatedTime(version.getCreatedTime())
213-
.setStatus(VersionStatusGrpc.forNumber(version.getStatus().getValue()))
214-
.setPushJobId(version.getPushJobId())
216+
.setStatus(statusGrpc)
215217
.setPartitionCount(version.getPartitionCount())
216-
.setReplicationFactor(version.getReplicationFactor())
217-
.build();
218+
.setReplicationFactor(version.getReplicationFactor());
219+
if (version.getPushJobId() != null) {
220+
builder.setPushJobId(version.getPushJobId());
221+
}
222+
return builder.build();
218223
}
219224

220225
/**

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

Lines changed: 6 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,6 @@
44
import com.linkedin.venice.controller.ControllerRequestHandlerDependencies;
55
import com.linkedin.venice.controller.StoreDeletedValidation;
66
import com.linkedin.venice.controllerapi.RepushInfo;
7-
import com.linkedin.venice.controllerapi.RepushInfoResponse;
87
import com.linkedin.venice.exceptions.VeniceException;
98
import com.linkedin.venice.exceptions.VeniceNoStoreException;
109
import com.linkedin.venice.meta.Store;
@@ -267,24 +266,20 @@ public ListStoresGrpcResponse listStores(ListStoresGrpcRequest request) {
267266
* @param clusterName the cluster name
268267
* @param storeName the store name
269268
* @param fabric optional fabric for multi-region setups
270-
* @return RepushInfoResponse containing repush information including version and Kafka details
269+
* @return RepushInfo containing repush information including version and Kafka details
271270
*/
272-
public RepushInfoResponse getRepushInfo(String clusterName, String storeName, Optional<String> fabric) {
273-
LOGGER.info(
271+
public RepushInfo getRepushInfo(String clusterName, String storeName, Optional<String> fabric) {
272+
LOGGER.debug(
274273
"Getting repush info for store: {} in cluster: {} with fabric: {}",
275274
storeName,
276275
clusterName,
277276
fabric.orElse("none"));
278277

279278
RepushInfo repushInfo = admin.getRepushInfo(clusterName, storeName, fabric);
280-
281-
RepushInfoResponse response = new RepushInfoResponse();
282-
response.setCluster(clusterName);
283-
response.setName(storeName);
284-
if (repushInfo != null) {
285-
response.setRepushInfo(repushInfo);
279+
if (repushInfo == null) {
280+
throw new VeniceException("Repush info not available for store: " + storeName + " in cluster: " + clusterName);
286281
}
287-
return response;
282+
return repushInfo;
288283
}
289284

290285
/**

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

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,7 @@
8484
import com.linkedin.venice.controllerapi.OwnerResponse;
8585
import com.linkedin.venice.controllerapi.PartitionResponse;
8686
import com.linkedin.venice.controllerapi.RegionPushDetailsResponse;
87+
import com.linkedin.venice.controllerapi.RepushInfo;
8788
import com.linkedin.venice.controllerapi.RepushInfoResponse;
8889
import com.linkedin.venice.controllerapi.RepushJobResponse;
8990
import com.linkedin.venice.controllerapi.SchemaUsageResponse;
@@ -264,13 +265,12 @@ public void internalHandle(Request request, RepushInfoResponse veniceResponse) {
264265
String fabricName = request.queryParams(FABRIC);
265266
Optional<String> fabric = fabricName != null ? Optional.of(fabricName) : Optional.empty();
266267

267-
// Call handler - returns POJO directly
268-
RepushInfoResponse result = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
268+
veniceResponse.setCluster(clusterName);
269+
veniceResponse.setName(storeName);
269270

270-
// Copy result to response
271-
veniceResponse.setCluster(result.getCluster());
272-
veniceResponse.setName(result.getName());
273-
veniceResponse.setRepushInfo(result.getRepushInfo());
271+
// Call handler - returns RepushInfo POJO directly
272+
RepushInfo repushInfo = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
273+
veniceResponse.setRepushInfo(repushInfo);
274274
}
275275
};
276276
}

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

Lines changed: 2 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
import com.linkedin.venice.controller.server.StoreRequestHandler;
1818
import com.linkedin.venice.controller.server.VeniceControllerAccessManager;
1919
import com.linkedin.venice.controllerapi.RepushInfo;
20-
import com.linkedin.venice.controllerapi.RepushInfoResponse;
2120
import com.linkedin.venice.exceptions.VeniceException;
2221
import com.linkedin.venice.exceptions.VeniceNoStoreException;
2322
import com.linkedin.venice.meta.StoreInfo;
@@ -474,12 +473,7 @@ public void testGetRepushInfoReturnsSuccessfulResponse() {
474473

475474
RepushInfo repushInfo = RepushInfo.createRepushInfo(mockVersion, "kafka.broker:9092", "d2-service", "zk-host");
476475

477-
RepushInfoResponse mockResponse = new RepushInfoResponse();
478-
mockResponse.setCluster(TEST_CLUSTER);
479-
mockResponse.setName(TEST_STORE);
480-
mockResponse.setRepushInfo(repushInfo);
481-
482-
when(storeRequestHandler.getRepushInfo(anyString(), anyString(), any())).thenReturn(mockResponse);
476+
when(storeRequestHandler.getRepushInfo(anyString(), anyString(), any())).thenReturn(repushInfo);
483477

484478
GetRepushInfoGrpcResponse actualResponse = blockingStub.getRepushInfo(request);
485479

@@ -515,12 +509,7 @@ public void testGetRepushInfoWithoutFabric() {
515509

516510
RepushInfo repushInfo = RepushInfo.createRepushInfo(null, "another.kafka.broker:9092", null, null);
517511

518-
RepushInfoResponse mockResponse = new RepushInfoResponse();
519-
mockResponse.setCluster(TEST_CLUSTER);
520-
mockResponse.setName(TEST_STORE);
521-
mockResponse.setRepushInfo(repushInfo);
522-
523-
when(storeRequestHandler.getRepushInfo(anyString(), anyString(), any())).thenReturn(mockResponse);
512+
when(storeRequestHandler.getRepushInfo(anyString(), anyString(), any())).thenReturn(repushInfo);
524513

525514
GetRepushInfoGrpcResponse actualResponse = blockingStub.getRepushInfo(request);
526515

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

Lines changed: 24 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,6 @@
1212
import com.linkedin.venice.controller.Admin;
1313
import com.linkedin.venice.controller.ControllerRequestHandlerDependencies;
1414
import com.linkedin.venice.controllerapi.RepushInfo;
15-
import com.linkedin.venice.controllerapi.RepushInfoResponse;
1615
import com.linkedin.venice.exceptions.VeniceException;
1716
import com.linkedin.venice.exceptions.VeniceNoStoreException;
1817
import com.linkedin.venice.meta.DataReplicationPolicy;
@@ -384,20 +383,18 @@ public void testGetRepushInfoSuccess() {
384383

385384
when(admin.getRepushInfo(clusterName, storeName, fabric)).thenReturn(mockRepushInfo);
386385

387-
RepushInfoResponse response = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
386+
RepushInfo response = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
388387

389388
verify(admin, times(1)).getRepushInfo(clusterName, storeName, fabric);
390-
assertEquals(response.getCluster(), clusterName);
391-
assertEquals(response.getName(), storeName);
392-
assertEquals(response.getRepushInfo().getKafkaBrokerUrl(), "kafka.broker.url:9092");
393-
assertEquals(response.getRepushInfo().getVersion().getNumber(), 1);
394-
assertEquals(response.getRepushInfo().getVersion().getCreatedTime(), 123456789L);
395-
assertEquals(response.getRepushInfo().getVersion().getStatus(), VersionStatus.ONLINE);
396-
assertEquals(response.getRepushInfo().getVersion().getPushJobId(), "test-push-job-123");
397-
assertEquals(response.getRepushInfo().getVersion().getPartitionCount(), 10);
398-
assertEquals(response.getRepushInfo().getVersion().getReplicationFactor(), 3);
399-
assertEquals(response.getRepushInfo().getSystemSchemaClusterD2ServiceName(), "testD2Service");
400-
assertEquals(response.getRepushInfo().getSystemSchemaClusterD2ZkHost(), "testZkHost");
389+
assertEquals(response.getKafkaBrokerUrl(), "kafka.broker.url:9092");
390+
assertEquals(response.getVersion().getNumber(), 1);
391+
assertEquals(response.getVersion().getCreatedTime(), 123456789L);
392+
assertEquals(response.getVersion().getStatus(), VersionStatus.ONLINE);
393+
assertEquals(response.getVersion().getPushJobId(), "test-push-job-123");
394+
assertEquals(response.getVersion().getPartitionCount(), 10);
395+
assertEquals(response.getVersion().getReplicationFactor(), 3);
396+
assertEquals(response.getSystemSchemaClusterD2ServiceName(), "testD2Service");
397+
assertEquals(response.getSystemSchemaClusterD2ZkHost(), "testZkHost");
401398
}
402399

403400
@Test
@@ -419,13 +416,11 @@ public void testGetRepushInfoWithoutFabric() {
419416

420417
when(admin.getRepushInfo(clusterName, storeName, fabric)).thenReturn(mockRepushInfo);
421418

422-
RepushInfoResponse response = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
419+
RepushInfo response = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
423420

424421
verify(admin, times(1)).getRepushInfo(clusterName, storeName, fabric);
425-
assertEquals(response.getCluster(), clusterName);
426-
assertEquals(response.getName(), storeName);
427-
assertEquals(response.getRepushInfo().getKafkaBrokerUrl(), "another.kafka.broker:9092");
428-
assertEquals(response.getRepushInfo().getVersion().getNumber(), 2);
422+
assertEquals(response.getKafkaBrokerUrl(), "another.kafka.broker:9092");
423+
assertEquals(response.getVersion().getNumber(), 2);
429424
}
430425

431426
@Test
@@ -438,12 +433,18 @@ public void testGetRepushInfoWithNullVersion() {
438433

439434
when(admin.getRepushInfo(clusterName, storeName, fabric)).thenReturn(mockRepushInfo);
440435

441-
RepushInfoResponse response = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
436+
RepushInfo response = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric);
442437

443-
assertEquals(response.getRepushInfo().getKafkaBrokerUrl(), "kafka.broker:9092");
444-
assertTrue(response.getRepushInfo().getVersion() == null);
445-
assertTrue(response.getRepushInfo().getSystemSchemaClusterD2ServiceName() == null);
446-
assertTrue(response.getRepushInfo().getSystemSchemaClusterD2ZkHost() == null);
438+
assertEquals(response.getKafkaBrokerUrl(), "kafka.broker:9092");
439+
assertTrue(response.getVersion() == null);
440+
assertTrue(response.getSystemSchemaClusterD2ServiceName() == null);
441+
assertTrue(response.getSystemSchemaClusterD2ZkHost() == null);
442+
}
443+
444+
@Test(expectedExceptions = VeniceException.class, expectedExceptionsMessageRegExp = "Repush info not available.*")
445+
public void testGetRepushInfoReturnsNullThrowsException() {
446+
when(admin.getRepushInfo("testCluster", "testStore", Optional.empty())).thenReturn(null);
447+
storeRequestHandler.getRepushInfo("testCluster", "testStore", Optional.empty());
447448
}
448449

449450
@Test

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

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -714,12 +714,7 @@ public void testGetRepushInfo() throws Exception {
714714

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

717-
RepushInfoResponse mockResponse = new RepushInfoResponse();
718-
mockResponse.setCluster(TEST_CLUSTER);
719-
mockResponse.setName(TEST_STORE_NAME);
720-
mockResponse.setRepushInfo(mockRepushInfo);
721-
722-
when(mockRequestHandler.getRepushInfo(any(), any(), any())).thenReturn(mockResponse);
717+
when(mockRequestHandler.getRepushInfo(any(), any(), any())).thenReturn(mockRepushInfo);
723718

724719
RepushInfoResponse response = ObjectMapperFactory.getInstance()
725720
.readValue(route.handle(request, mock(Response.class)).toString(), RepushInfoResponse.class);
@@ -761,12 +756,7 @@ public void testGetRepushInfoWithoutFabric() throws Exception {
761756

762757
RepushInfo mockRepushInfo = RepushInfo.createRepushInfo(null, "another.kafka:9092", null, null);
763758

764-
RepushInfoResponse mockResponse = new RepushInfoResponse();
765-
mockResponse.setCluster(TEST_CLUSTER);
766-
mockResponse.setName(TEST_STORE_NAME);
767-
mockResponse.setRepushInfo(mockRepushInfo);
768-
769-
when(mockRequestHandler.getRepushInfo(any(), any(), any())).thenReturn(mockResponse);
759+
when(mockRequestHandler.getRepushInfo(any(), any(), any())).thenReturn(mockRepushInfo);
770760

771761
RepushInfoResponse response = ObjectMapperFactory.getInstance()
772762
.readValue(route.handle(request, mock(Response.class)).toString(), RepushInfoResponse.class);

0 commit comments

Comments
 (0)