diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java index acf8ea32c6..25ac8665f6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java @@ -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); + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 7479cabf46..03d63b6f11 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -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); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index 74b1349acc..2f5a1c86b5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -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 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 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 reportThreadStackTrace(Context ctx, ReportThreadStackTraceRequest request) { + String commandId = request.getCommandId(); return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index 4adde9ea49..c551135e9d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -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 queryRouteEndpointConverter; private volatile ResponseHook queryRouteHook = null; @@ -53,9 +61,10 @@ public class RouteService extends BaseService { private volatile RouteAssignmentQueueSelector assignmentQueueSelector; private volatile ResponseHook 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 queueDataList = topicRouteData.getQueueDatas(); + List brokerDataList = topicRouteData.getBrokerDatas(); List 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> brokerMap = buildBrokerMap(brokerDataList); + + for (QueueData queueData : queueDataList) { + String brokerName = queueData.getBrokerName(); + Map 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 assignments = new ArrayList<>(); List 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> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas()); + for (SelectableMessageQueue messageQueue : messageQueueList) { + Map 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> buildBrokerMap(List brokerDataList) { + Map> brokerMap = new HashMap<>(); + for (BrokerData brokerData : brokerDataList) { + Map brokerIdMap = new HashMap<>(); + String brokerName = brokerData.getBrokerName(); + for (Map.Entry 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; + } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java index 41851d9b8b..73761cd9b3 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java @@ -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 brokerDataList = new ArrayList<>(); + BrokerData brokerData = new BrokerData(); + brokerData.setCluster("cluster"); + brokerData.setBrokerName("brokerName"); + HashMap brokerAddrs = new HashMap() {{ + put(0L, brokerAddress); + }}; + brokerData.setBrokerAddrs(brokerAddrs); + brokerDataList.add(brokerData); + + List 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 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 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 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 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 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 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 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); + } + } + }