mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Implement route in Local mode
This commit is contained in:
@@ -34,10 +34,24 @@ public enum ProxyMode {
|
||||
return CLUSTER.mode.equals(mode.toUpperCase());
|
||||
}
|
||||
|
||||
public static boolean isClusterMode(ProxyMode mode) {
|
||||
if (mode == null) {
|
||||
return false;
|
||||
}
|
||||
return CLUSTER.equals(mode);
|
||||
}
|
||||
|
||||
public static boolean isLocalMode(String mode) {
|
||||
if (mode == null) {
|
||||
return false;
|
||||
}
|
||||
return LOCAL.mode.equals(mode.toUpperCase());
|
||||
}
|
||||
|
||||
public static boolean isLocalMode(ProxyMode mode) {
|
||||
if (mode == null) {
|
||||
return false;
|
||||
}
|
||||
return LOCAL.equals(mode);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,6 +53,9 @@ import apache.rocketmq.v1.SendMessageRequest;
|
||||
import apache.rocketmq.v1.SendMessageResponse;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
@@ -61,6 +64,7 @@ import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateCheckRequest;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService;
|
||||
@@ -71,10 +75,6 @@ import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
|
||||
public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService {
|
||||
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
|
||||
@@ -95,7 +95,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker());
|
||||
this.receiveMessageService = new ReceiveMessageService(connectorManager);
|
||||
this.producerService = new ProducerService(connectorManager);
|
||||
this.routeService = new RouteService(connectorManager);
|
||||
this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager);
|
||||
this.clientService = new ClientService(connectorManager, scheduledExecutorService, channelManager);
|
||||
this.pullMessageService = new PullMessageService(connectorManager);
|
||||
this.transactionService = new TransactionService(connectorManager, channelManager);
|
||||
|
||||
@@ -77,6 +77,7 @@ import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.ReceiveMessageChannel;
|
||||
@@ -85,7 +86,9 @@ import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHand
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.slf4j.Logger;
|
||||
@@ -98,14 +101,18 @@ public class LocalGrpcService implements GrpcForwardService {
|
||||
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(
|
||||
new ThreadFactoryImpl("LocalGrpcServiceScheduledThread"));
|
||||
private final ChannelManager channelManager;
|
||||
private final RouteService routeService;
|
||||
|
||||
public LocalGrpcService(BrokerController brokerController) {
|
||||
this.brokerController = brokerController;
|
||||
this.channelManager = new ChannelManager();
|
||||
// TransactionStateChecker is not used in Local mode.
|
||||
ConnectorManager connectorManager = new ConnectorManager(null);
|
||||
this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager);
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
return null;
|
||||
return this.routeService.queryRoute(ctx, request);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -188,7 +195,7 @@ public class LocalGrpcService implements GrpcForwardService {
|
||||
|
||||
@Override
|
||||
public CompletableFuture<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request) {
|
||||
return null;
|
||||
return this.routeService.queryAssignment(ctx, request);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -376,6 +383,7 @@ public class LocalGrpcService implements GrpcForwardService {
|
||||
|
||||
@Override public CompletableFuture<ReportThreadStackTraceResponse> reportThreadStackTrace(Context ctx,
|
||||
ReportThreadStackTraceRequest request) {
|
||||
String commandId = request.getCommandId();
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
+125
-47
@@ -16,6 +16,8 @@
|
||||
*/
|
||||
package org.apache.rocketmq.proxy.grpc.service.cluster;
|
||||
|
||||
import apache.rocketmq.v1.Address;
|
||||
import apache.rocketmq.v1.AddressScheme;
|
||||
import apache.rocketmq.v1.Assignment;
|
||||
import apache.rocketmq.v1.Broker;
|
||||
import apache.rocketmq.v1.Endpoints;
|
||||
@@ -26,9 +28,17 @@ import apache.rocketmq.v1.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v1.QueryRouteRequest;
|
||||
import apache.rocketmq.v1.QueryRouteResponse;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import com.google.common.base.Preconditions;
|
||||
import com.google.common.net.HostAndPort;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.constant.PermName;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
@@ -37,14 +47,12 @@ import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ParameterConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseHook;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
|
||||
public class RouteService extends BaseService {
|
||||
private final ProxyMode mode;
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter;
|
||||
private volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook = null;
|
||||
@@ -53,9 +61,10 @@ public class RouteService extends BaseService {
|
||||
private volatile RouteAssignmentQueueSelector assignmentQueueSelector;
|
||||
private volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook = null;
|
||||
|
||||
public RouteService(ConnectorManager connectorManager) {
|
||||
public RouteService(ProxyMode mode, ConnectorManager connectorManager) {
|
||||
super(connectorManager);
|
||||
|
||||
Preconditions.checkArgument(ProxyMode.isClusterMode(mode) || ProxyMode.isLocalMode(mode));
|
||||
this.mode = mode;
|
||||
queryRouteEndpointConverter = (ctx, parameter) -> parameter;
|
||||
queryAssignmentEndpointConverter = (ctx, parameter) -> parameter;
|
||||
assignmentQueueSelector = new DefaultRouteAssignmentQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
@@ -92,30 +101,47 @@ public class RouteService extends BaseService {
|
||||
});
|
||||
|
||||
try {
|
||||
Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, request.getEndpoints());
|
||||
if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) {
|
||||
future.complete(QueryRouteResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.INVALID_ARGUMENT, "endpoint " +
|
||||
request.getEndpoints() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache()
|
||||
.getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic()));
|
||||
TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData();
|
||||
List<QueueData> queueDataList = topicRouteData.getQueueDatas();
|
||||
List<BrokerData> brokerDataList = topicRouteData.getBrokerDatas();
|
||||
|
||||
List<Partition> partitionList = new ArrayList<>();
|
||||
for (QueueData queueData : queueDataList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(queueData.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
if (ProxyMode.isClusterMode(mode.name())) {
|
||||
Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, request.getEndpoints());
|
||||
if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) {
|
||||
future.complete(QueryRouteResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.INVALID_ARGUMENT, "endpoint " +
|
||||
request.getEndpoints() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
for (QueueData queueData : queueDataList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(queueData.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
|
||||
partitionList.addAll(genPartitionFromQueueData(queueData, request.getTopic(), broker));
|
||||
partitionList.addAll(genPartitionFromQueueData(queueData, request.getTopic(), broker));
|
||||
}
|
||||
}
|
||||
if (ProxyMode.isLocalMode(mode.name())) {
|
||||
Map<String, Map<Long, Broker>> brokerMap = buildBrokerMap(brokerDataList);
|
||||
|
||||
for (QueueData queueData : queueDataList) {
|
||||
String brokerName = queueData.getBrokerName();
|
||||
Map<Long, Broker> brokerIdMap = brokerMap.get(brokerName);
|
||||
if (brokerIdMap == null) {
|
||||
break;
|
||||
}
|
||||
for (Broker broker : brokerIdMap.values()) {
|
||||
partitionList.addAll(genPartitionFromQueueData(queueData, request.getTopic(), broker));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
QueryRouteResponse response = QueryRouteResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.addAllPartitions(partitionList)
|
||||
@@ -187,36 +213,60 @@ public class RouteService extends BaseService {
|
||||
});
|
||||
|
||||
try {
|
||||
Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, request.getEndpoints());
|
||||
if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) {
|
||||
future.complete(QueryAssignmentResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.INVALID_ARGUMENT, "endpoint " +
|
||||
request.getEndpoints() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
|
||||
List<Assignment> assignments = new ArrayList<>();
|
||||
List<SelectableMessageQueue> messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request);
|
||||
if (ProxyMode.isLocalMode(mode)) {
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache()
|
||||
.getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic()));
|
||||
TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData();
|
||||
Map<String, Map<Long, Broker>> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas());
|
||||
for (SelectableMessageQueue messageQueue : messageQueueList) {
|
||||
Map<Long, Broker> brokerIdMap = brokerMap.get(messageQueue.getBrokerName());
|
||||
if (brokerIdMap != null) {
|
||||
Broker broker = brokerIdMap.get(0L);
|
||||
|
||||
for (SelectableMessageQueue messageQueue : messageQueueList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(messageQueue.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
Partition defaultPartition = Partition.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
|
||||
Partition defaultPartition = Partition.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setPartition(defaultPartition)
|
||||
.build());
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setPartition(defaultPartition)
|
||||
.build());
|
||||
}
|
||||
}
|
||||
}
|
||||
if (ProxyMode.isClusterMode(mode)) {
|
||||
Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, request.getEndpoints());
|
||||
if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) {
|
||||
future.complete(QueryAssignmentResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.INVALID_ARGUMENT, "endpoint " +
|
||||
request.getEndpoints() + " is invalidate"))
|
||||
.build());
|
||||
return future;
|
||||
}
|
||||
for (SelectableMessageQueue messageQueue : messageQueueList) {
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(messageQueue.getBrokerName())
|
||||
.setId(0)
|
||||
.setEndpoints(resEndpoints)
|
||||
.build();
|
||||
|
||||
Partition defaultPartition = Partition.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setPartition(defaultPartition)
|
||||
.build());
|
||||
}
|
||||
}
|
||||
|
||||
QueryAssignmentResponse response = QueryAssignmentResponse.newBuilder()
|
||||
.addAllAssignments(assignments)
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
@@ -227,4 +277,32 @@ public class RouteService extends BaseService {
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
private Map<String, Map<Long, Broker>> buildBrokerMap(List<BrokerData> brokerDataList) {
|
||||
Map<String, Map<Long, Broker>> brokerMap = new HashMap<>();
|
||||
for (BrokerData brokerData : brokerDataList) {
|
||||
Map<Long, Broker> brokerIdMap = new HashMap<>();
|
||||
String brokerName = brokerData.getBrokerName();
|
||||
for (Map.Entry<Long, String> entry : brokerData.getBrokerAddrs().entrySet()) {
|
||||
Long brokerId = entry.getKey();
|
||||
HostAndPort hostAndPort = HostAndPort.fromString(entry.getValue());
|
||||
Broker broker = Broker.newBuilder()
|
||||
.setName(brokerName)
|
||||
.setId(Math.toIntExact(brokerId))
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.setScheme(AddressScheme.IPv4)
|
||||
.addAddresses(
|
||||
Address.newBuilder()
|
||||
.setPort(hostAndPort.getPort())
|
||||
.setHost(hostAndPort.getHost())
|
||||
)
|
||||
.build())
|
||||
.build();
|
||||
|
||||
brokerIdMap.put(brokerId, broker);
|
||||
}
|
||||
brokerMap.put(brokerName, brokerIdMap);
|
||||
}
|
||||
return brokerMap;
|
||||
}
|
||||
}
|
||||
|
||||
+223
-15
@@ -17,39 +17,66 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.service.cluster;
|
||||
|
||||
import apache.rocketmq.v1.Address;
|
||||
import apache.rocketmq.v1.AddressScheme;
|
||||
import apache.rocketmq.v1.Broker;
|
||||
import apache.rocketmq.v1.Endpoints;
|
||||
import apache.rocketmq.v1.Partition;
|
||||
import apache.rocketmq.v1.Permission;
|
||||
import apache.rocketmq.v1.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v1.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v1.QueryRouteRequest;
|
||||
import apache.rocketmq.v1.QueryRouteResponse;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import com.google.common.net.HostAndPort;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.constant.PermName;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNull;
|
||||
|
||||
public class RouteServiceTest {
|
||||
public class RouteServiceTest extends BaseServiceTest {
|
||||
private String brokerAddress = "127.0.0.1:10911";
|
||||
public static final String BROKER_NAME = "brokerName";
|
||||
public static final String NAMESPACE = "namespace";
|
||||
public static final String TOPIC = "topic";
|
||||
public static final Broker MOCK_BROKER = Broker.newBuilder().setName(BROKER_NAME).build();
|
||||
public static final Resource MOCK_TOPIC = Resource.newBuilder()
|
||||
.setName(TOPIC)
|
||||
.setResourceNamespace(NAMESPACE)
|
||||
.build();
|
||||
.setName(TOPIC)
|
||||
.setResourceNamespace(NAMESPACE)
|
||||
.build();
|
||||
|
||||
@Before
|
||||
public void before() throws Exception {
|
||||
@Override
|
||||
public void beforeEach() {
|
||||
List<BrokerData> brokerDataList = new ArrayList<>();
|
||||
BrokerData brokerData = new BrokerData();
|
||||
brokerData.setCluster("cluster");
|
||||
brokerData.setBrokerName("brokerName");
|
||||
HashMap<Long, String> brokerAddrs = new HashMap<Long, String>() {{
|
||||
put(0L, brokerAddress);
|
||||
}};
|
||||
brokerData.setBrokerAddrs(brokerAddrs);
|
||||
brokerDataList.add(brokerData);
|
||||
|
||||
List<QueueData> queueDataList = new ArrayList<>();
|
||||
QueueData queueData = new QueueData();
|
||||
queueData.setPerm(6);
|
||||
queueData.setWriteQueueNums(8);
|
||||
queueData.setReadQueueNums(8);
|
||||
queueData.setBrokerName("brokerName");
|
||||
queueDataList.add(queueData);
|
||||
}
|
||||
|
||||
@After
|
||||
public void after() throws Exception {
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testGenPartitionFromQueueData() throws Exception {
|
||||
// test queueData with 8 read queues, 8 write queues, and rw permission, expect 8 rw queues.
|
||||
@@ -103,4 +130,185 @@ public class RouteServiceTest {
|
||||
return queueData;
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testLocalModeQueryRoute() {
|
||||
RouteService routeService = new RouteService(ProxyMode.LOCAL, this.clientManager);
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build());
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode());
|
||||
assertEquals(8, response.getPartitionsCount());
|
||||
assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getPartitions(0).getBroker()
|
||||
.getEndpoints().getAddresses(0).getHost());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRouteWithInvalidEndpoints() {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.clientManager);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.INVALID_ARGUMENT.getNumber(), response.getCommon().getStatus().getCode());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRoute() {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.clientManager);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode());
|
||||
assertEquals(8, response.getPartitionsCount());
|
||||
assertEquals("host", response.getPartitions(0).getBroker()
|
||||
.getEndpoints().getAddresses(0).getHost());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRouteWhenTopicNotExist() {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.clientManager);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("notExistTopic")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.NOT_FOUND.getNumber(), response.getCommon().getStatus().getCode());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryAssignmentInvalidEndpoints() {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.clientManager);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(
|
||||
Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build()
|
||||
)
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.INVALID_ARGUMENT.getNumber(), response.getCommon().getStatus().getCode());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testLocalModeQueryAssignment() {
|
||||
RouteService routeService = new RouteService(ProxyMode.LOCAL, this.clientManager);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName("group")
|
||||
.build())
|
||||
.setClientId("clientId")
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode());
|
||||
assertEquals(1, response.getAssignmentsCount());
|
||||
assertEquals("brokerName", response.getAssignments(0).getPartition().getBroker().getName());
|
||||
assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getAssignments(0).getPartition().getBroker().getEndpoints().getAddresses(0).getHost());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryAssignment() {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.clientManager);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName("group")
|
||||
.build())
|
||||
.setClientId("clientId")
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode());
|
||||
assertEquals(1, response.getAssignmentsCount());
|
||||
assertEquals("brokerName", response.getAssignments(0).getPartition().getBroker().getName());
|
||||
assertEquals("host", response.getAssignments(0).getPartition().getBroker().getEndpoints().getAddresses(0).getHost());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user