Skip to content

Commit 0ddba99

Browse files
committed
Refactor getRepushInfo to follow v2 architecture pattern
This commit updates the getRepushInfo endpoint to follow the v2 architecture where protocol-agnostic handlers take primitives + ControllerRequestContext and return POJOs instead of protobuf messages. Changes: - Added ControllerRequestContext class for transport-agnostic client identity - Updated StoreRequestHandler.getRepushInfo() to take primitives + context, return RepushInfoResponse POJO - Updated StoreGrpcServiceImpl to extract params, build context, call handler, convert to protobuf - Updated StoresRoutes to remove protobuf dependency - HTTP layer never touches protobuf - Removed unused mapGrpcRepushInfoToRepushInfo conversion method Benefits of v2 architecture: - Handlers are truly protocol-agnostic (no protobuf leakage) - ACL checks can be centralized in handlers (not needed for getRepushInfo) - HTTP and gRPC layers are thin adapters with minimal boilerplate - Better separation of concerns and testability Note: Unit tests need to be updated to reflect new v2 signatures (next commit).
1 parent 14fe1f6 commit 0ddba99

4 files changed

Lines changed: 139 additions & 98 deletions

File tree

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

Lines changed: 69 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,13 @@
44
import static com.linkedin.venice.controller.grpc.server.ControllerGrpcServerUtils.isAllowListUser;
55
import static com.linkedin.venice.controller.server.VeniceRouteHandler.ACL_CHECK_FAILURE_WARN_MESSAGE_PREFIX;
66

7+
import com.linkedin.venice.controller.server.ControllerRequestContext;
78
import com.linkedin.venice.controller.server.StoreRequestHandler;
89
import com.linkedin.venice.controller.server.VeniceControllerAccessManager;
10+
import com.linkedin.venice.controllerapi.RepushInfo;
11+
import com.linkedin.venice.controllerapi.RepushInfoResponse;
912
import com.linkedin.venice.exceptions.VeniceUnauthorizedAccessException;
13+
import com.linkedin.venice.meta.Version;
1014
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
1115
import com.linkedin.venice.protocols.controller.CreateStoreGrpcRequest;
1216
import com.linkedin.venice.protocols.controller.CreateStoreGrpcResponse;
@@ -18,15 +22,18 @@
1822
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
1923
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
2024
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
25+
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
2126
import com.linkedin.venice.protocols.controller.ResourceCleanupCheckGrpcResponse;
2227
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc;
2328
import com.linkedin.venice.protocols.controller.StoreGrpcServiceGrpc.StoreGrpcServiceImplBase;
2429
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcRequest;
2530
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
2631
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
2732
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
33+
import com.linkedin.venice.protocols.controller.VersionGrpc;
2834
import io.grpc.Context;
2935
import io.grpc.stub.StreamObserver;
36+
import java.util.Optional;
3037
import org.apache.logging.log4j.LogManager;
3138
import org.apache.logging.log4j.Logger;
3239

@@ -150,15 +157,72 @@ public void listStores(ListStoresGrpcRequest grpcRequest, StreamObserver<ListSto
150157
null);
151158
}
152159

160+
/**
161+
* Retrieves repush information for a store.
162+
* No ACL check is required for this operation as it only reads store metadata.
163+
*/
153164
@Override
154165
public void getRepushInfo(
155166
GetRepushInfoGrpcRequest request,
156167
StreamObserver<GetRepushInfoGrpcResponse> responseObserver) {
157168
LOGGER.debug("Received getRepushInfo with args: {}", request);
158-
ControllerGrpcServerUtils.handleRequest(
159-
StoreGrpcServiceGrpc.getGetRepushInfoMethod(),
160-
() -> storeRequestHandler.getRepushInfo(request),
161-
responseObserver,
162-
request.getStoreInfo());
169+
170+
ControllerGrpcServerUtils.handleRequest(StoreGrpcServiceGrpc.getGetRepushInfoMethod(), () -> {
171+
// Extract primitives from protobuf
172+
ClusterStoreGrpcInfo storeInfo = request.getStoreInfo();
173+
String clusterName = storeInfo.getClusterName();
174+
String storeName = storeInfo.getStoreName();
175+
Optional<String> fabric = request.hasFabric() ? Optional.of(request.getFabric()) : Optional.empty();
176+
177+
// Build transport-agnostic context from gRPC
178+
ControllerRequestContext context = buildRequestContext(Context.current());
179+
180+
// Call handler - returns POJO
181+
RepushInfoResponse result = storeRequestHandler.getRepushInfo(clusterName, storeName, fabric, context);
182+
183+
// Convert POJO to protobuf response
184+
RepushInfo repushInfo = result.getRepushInfo();
185+
RepushInfoGrpc.Builder repushInfoBuilder = RepushInfoGrpc.newBuilder()
186+
.setKafkaBrokerUrl(repushInfo.getKafkaBrokerUrl() != null ? repushInfo.getKafkaBrokerUrl() : "");
187+
188+
if (repushInfo.getVersion() != null) {
189+
repushInfoBuilder.setVersion(convertVersionToProto(repushInfo.getVersion()));
190+
}
191+
if (repushInfo.getSystemSchemaClusterD2ServiceName() != null) {
192+
repushInfoBuilder.setSystemSchemaClusterD2ServiceName(repushInfo.getSystemSchemaClusterD2ServiceName());
193+
}
194+
if (repushInfo.getSystemSchemaClusterD2ZkHost() != null) {
195+
repushInfoBuilder.setSystemSchemaClusterD2ZkHost(repushInfo.getSystemSchemaClusterD2ZkHost());
196+
}
197+
198+
return GetRepushInfoGrpcResponse.newBuilder()
199+
.setStoreInfo(storeInfo)
200+
.setRepushInfo(repushInfoBuilder.build())
201+
.build();
202+
}, responseObserver, request.getStoreInfo());
203+
}
204+
205+
/**
206+
* Builds a ControllerRequestContext from the gRPC context.
207+
*/
208+
private ControllerRequestContext buildRequestContext(Context context) {
209+
GrpcControllerClientDetails clientDetails = ControllerGrpcServerUtils.getClientDetails(context);
210+
return new ControllerRequestContext(
211+
clientDetails.getClientCertificate(),
212+
clientDetails.getClientAddress() != null ? clientDetails.getClientAddress() : "anonymous");
213+
}
214+
215+
/**
216+
* Converts a Version object to protobuf VersionGrpc.
217+
*/
218+
private VersionGrpc convertVersionToProto(Version version) {
219+
return VersionGrpc.newBuilder()
220+
.setNumber(version.getNumber())
221+
.setCreatedTime(version.getCreatedTime())
222+
.setStatus(version.getStatus().getValue())
223+
.setPushJobId(version.getPushJobId())
224+
.setPartitionCount(version.getPartitionCount())
225+
.setReplicationFactor(version.getReplicationFactor())
226+
.build();
163227
}
164228
}
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
package com.linkedin.venice.controller.server;
2+
3+
import java.security.cert.X509Certificate;
4+
import java.util.Optional;
5+
6+
7+
/**
8+
* Transport-agnostic request context that carries client identity information.
9+
* Both HTTP and gRPC layers populate this before calling handlers.
10+
*/
11+
public class ControllerRequestContext {
12+
private final Optional<X509Certificate> clientCertificate;
13+
private final String clientPrincipalId;
14+
15+
public ControllerRequestContext(X509Certificate clientCertificate, String clientPrincipalId) {
16+
this.clientCertificate = Optional.ofNullable(clientCertificate);
17+
this.clientPrincipalId = clientPrincipalId;
18+
}
19+
20+
public Optional<X509Certificate> getClientCertificate() {
21+
return clientCertificate;
22+
}
23+
24+
public String getClientPrincipalId() {
25+
return clientPrincipalId;
26+
}
27+
28+
/**
29+
* Creates a context for unauthenticated/internal requests
30+
*/
31+
public static ControllerRequestContext anonymous() {
32+
return new ControllerRequestContext(null, "anonymous");
33+
}
34+
}

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

Lines changed: 18 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,6 @@
66
import com.linkedin.venice.controllerapi.RepushInfo;
77
import com.linkedin.venice.exceptions.VeniceException;
88
import com.linkedin.venice.meta.Store;
9-
import com.linkedin.venice.meta.Version;
109
import com.linkedin.venice.meta.ZKStore;
1110
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
1211
import com.linkedin.venice.protocols.controller.CreateStoreGrpcRequest;
@@ -15,16 +14,12 @@
1514
import com.linkedin.venice.protocols.controller.DeleteAclForStoreGrpcResponse;
1615
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcRequest;
1716
import com.linkedin.venice.protocols.controller.GetAclForStoreGrpcResponse;
18-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
19-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
2017
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
2118
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
22-
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
2319
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcRequest;
2420
import com.linkedin.venice.protocols.controller.UpdateAclForStoreGrpcResponse;
2521
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
2622
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
27-
import com.linkedin.venice.protocols.controller.VersionGrpc;
2823
import com.linkedin.venice.systemstore.schemas.StoreProperties;
2924
import java.util.ArrayList;
3025
import java.util.List;
@@ -264,16 +259,19 @@ public ListStoresGrpcResponse listStores(ListStoresGrpcRequest request) {
264259

265260
/**
266261
* Retrieves repush information for a store.
267-
* @param request the request containing cluster, store name, and optional fabric
268-
* @return response containing repush information including version and Kafka details
262+
* No ACL check required - reading repush info is public.
263+
*
264+
* @param clusterName the cluster name
265+
* @param storeName the store name
266+
* @param fabric optional fabric for multi-region setups
267+
* @param context the request context (contains client identity)
268+
* @return RepushInfoResponse containing repush information including version and Kafka details
269269
*/
270-
public GetRepushInfoGrpcResponse getRepushInfo(GetRepushInfoGrpcRequest request) {
271-
ClusterStoreGrpcInfo storeInfo = request.getStoreInfo();
272-
ControllerRequestParamValidator.validateClusterStoreInfo(storeInfo);
273-
String clusterName = storeInfo.getClusterName();
274-
String storeName = storeInfo.getStoreName();
275-
Optional<String> fabric = request.hasFabric() ? Optional.of(request.getFabric()) : Optional.empty();
276-
270+
public com.linkedin.venice.controllerapi.RepushInfoResponse getRepushInfo(
271+
String clusterName,
272+
String storeName,
273+
Optional<String> fabric,
274+
ControllerRequestContext context) {
277275
LOGGER.info(
278276
"Getting repush info for store: {} in cluster: {} with fabric: {}",
279277
storeName,
@@ -282,33 +280,11 @@ public GetRepushInfoGrpcResponse getRepushInfo(GetRepushInfoGrpcRequest request)
282280

283281
RepushInfo repushInfo = admin.getRepushInfo(clusterName, storeName, fabric);
284282

285-
RepushInfoGrpc.Builder repushInfoBuilder = RepushInfoGrpc.newBuilder()
286-
.setKafkaBrokerUrl(repushInfo.getKafkaBrokerUrl() != null ? repushInfo.getKafkaBrokerUrl() : "");
287-
288-
if (repushInfo.getVersion() != null) {
289-
repushInfoBuilder.setVersion(convertVersionToProto(repushInfo.getVersion()));
290-
}
291-
if (repushInfo.getSystemSchemaClusterD2ServiceName() != null) {
292-
repushInfoBuilder.setSystemSchemaClusterD2ServiceName(repushInfo.getSystemSchemaClusterD2ServiceName());
293-
}
294-
if (repushInfo.getSystemSchemaClusterD2ZkHost() != null) {
295-
repushInfoBuilder.setSystemSchemaClusterD2ZkHost(repushInfo.getSystemSchemaClusterD2ZkHost());
296-
}
297-
298-
return GetRepushInfoGrpcResponse.newBuilder()
299-
.setStoreInfo(storeInfo)
300-
.setRepushInfo(repushInfoBuilder.build())
301-
.build();
302-
}
303-
304-
private VersionGrpc convertVersionToProto(Version version) {
305-
return VersionGrpc.newBuilder()
306-
.setNumber(version.getNumber())
307-
.setCreatedTime(version.getCreatedTime())
308-
.setStatus(version.getStatus().getValue())
309-
.setPushJobId(version.getPushJobId())
310-
.setPartitionCount(version.getPartitionCount())
311-
.setReplicationFactor(version.getReplicationFactor())
312-
.build();
283+
com.linkedin.venice.controllerapi.RepushInfoResponse response =
284+
new com.linkedin.venice.controllerapi.RepushInfoResponse();
285+
response.setCluster(clusterName);
286+
response.setName(storeName);
287+
response.setRepushInfo(repushInfo);
288+
return response;
313289
}
314290
}

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

Lines changed: 18 additions & 51 deletions
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,6 @@
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;
8887
import com.linkedin.venice.controllerapi.RepushInfoResponse;
8988
import com.linkedin.venice.controllerapi.RepushJobResponse;
9089
import com.linkedin.venice.controllerapi.SchemaUsageResponse;
@@ -108,17 +107,11 @@
108107
import com.linkedin.venice.meta.StoreDataAudit;
109108
import com.linkedin.venice.meta.StoreInfo;
110109
import com.linkedin.venice.meta.Version;
111-
import com.linkedin.venice.meta.VersionImpl;
112-
import com.linkedin.venice.meta.VersionStatus;
113110
import com.linkedin.venice.protocols.controller.ClusterStoreGrpcInfo;
114-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcRequest;
115-
import com.linkedin.venice.protocols.controller.GetRepushInfoGrpcResponse;
116111
import com.linkedin.venice.protocols.controller.ListStoresGrpcRequest;
117112
import com.linkedin.venice.protocols.controller.ListStoresGrpcResponse;
118-
import com.linkedin.venice.protocols.controller.RepushInfoGrpc;
119113
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcRequest;
120114
import com.linkedin.venice.protocols.controller.ValidateStoreDeletedGrpcResponse;
121-
import com.linkedin.venice.protocols.controller.VersionGrpc;
122115
import com.linkedin.venice.pubsub.PubSubTopicRepository;
123116
import com.linkedin.venice.pubsub.api.PubSubTopic;
124117
import com.linkedin.venice.pubsub.api.exceptions.PubSubTopicDoesNotExistException;
@@ -264,31 +257,23 @@ public Route getRepushInfo(Admin admin, StoreRequestHandler requestHandler) {
264257
@Override
265258
public void internalHandle(Request request, RepushInfoResponse veniceResponse) {
266259
AdminSparkServer.validateParams(request, GET_REPUSH_INFO.getParams(), admin);
260+
261+
// Extract primitives from HTTP request
267262
String clusterName = request.queryParams(CLUSTER);
268263
String storeName = request.queryParams(NAME);
269264
String fabricName = request.queryParams(FABRIC);
265+
Optional<String> fabric = fabricName != null ? Optional.of(fabricName) : Optional.empty();
270266

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();
267+
// Build transport-agnostic context
268+
ControllerRequestContext context = buildRequestContext(request);
280269

281-
// Call handler
282-
GetRepushInfoGrpcResponse grpcResponse = requestHandler.getRepushInfo(grpcRequest);
270+
// Call handler - returns POJO directly
271+
RepushInfoResponse result = requestHandler.getRepushInfo(clusterName, storeName, fabric, context);
283272

284-
// Map response back to HTTP
285-
veniceResponse.setCluster(clusterName);
286-
veniceResponse.setName(storeName);
287-
288-
// Convert proto RepushInfo back to Java RepushInfo for HTTP response
289-
RepushInfo repushInfo = mapGrpcRepushInfoToRepushInfo(grpcResponse.getRepushInfo(), storeName);
290-
291-
veniceResponse.setRepushInfo(repushInfo);
273+
// Copy result to response
274+
veniceResponse.setCluster(result.getCluster());
275+
veniceResponse.setName(result.getName());
276+
veniceResponse.setRepushInfo(result.getRepushInfo());
292277
}
293278
};
294279
}
@@ -1249,32 +1234,14 @@ public void internalHandle(Request request, StoreDeletedValidationResponse venic
12491234
}
12501235

12511236
/**
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
1237+
* Build request context from HTTP request.
12561238
*/
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()));
1239+
private ControllerRequestContext buildRequestContext(spark.Request request) {
1240+
if (!isSslEnabled()) {
1241+
return ControllerRequestContext.anonymous();
12701242
}
1271-
1272-
return RepushInfo.createRepushInfo(
1273-
version,
1274-
repushInfoProto.getKafkaBrokerUrl(),
1275-
repushInfoProto.hasSystemSchemaClusterD2ServiceName()
1276-
? repushInfoProto.getSystemSchemaClusterD2ServiceName()
1277-
: null,
1278-
repushInfoProto.hasSystemSchemaClusterD2ZkHost() ? repushInfoProto.getSystemSchemaClusterD2ZkHost() : null);
1243+
java.security.cert.X509Certificate cert = getCertificate(request);
1244+
String principalId = getPrincipalId(request);
1245+
return new ControllerRequestContext(cert, principalId);
12791246
}
12801247
}

0 commit comments

Comments
 (0)