mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Add local/RouteService
This commit is contained in:
+100
@@ -0,0 +1,100 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service;
|
||||
|
||||
import apache.rocketmq.v2.Endpoints;
|
||||
import apache.rocketmq.v2.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v2.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v2.QueryRouteRequest;
|
||||
import apache.rocketmq.v2.QueryRouteResponse;
|
||||
import io.grpc.Context;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.proxy.common.ParameterConverter;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.AssignmentQueueSelector;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.BaseService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.DefaultAssignmentQueueSelector;
|
||||
|
||||
public abstract class AbstractRouteService extends BaseService {
|
||||
protected volatile ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter;
|
||||
protected volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook;
|
||||
|
||||
protected volatile ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter;
|
||||
protected volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook;
|
||||
protected volatile AssignmentQueueSelector assignmentQueueSelector;
|
||||
|
||||
protected final GrpcClientManager grpcClientManager;
|
||||
|
||||
public AbstractRouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
super(connectorManager);
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
this.queryRouteEndpointConverter = (ctx, parameter) -> parameter;
|
||||
this.queryAssignmentEndpointConverter = (ctx, parameter) -> parameter;
|
||||
this.assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
public abstract CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request);
|
||||
|
||||
public abstract CompletableFuture<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request);
|
||||
|
||||
public ParameterConverter<Endpoints, Endpoints> getQueryRouteEndpointConverter() {
|
||||
return queryRouteEndpointConverter;
|
||||
}
|
||||
|
||||
public void setQueryRouteEndpointConverter(
|
||||
ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter) {
|
||||
this.queryRouteEndpointConverter = queryRouteEndpointConverter;
|
||||
}
|
||||
|
||||
public ResponseHook<QueryRouteRequest, QueryRouteResponse> getQueryRouteHook() {
|
||||
return queryRouteHook;
|
||||
}
|
||||
|
||||
public void setQueryRouteHook(
|
||||
ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook) {
|
||||
this.queryRouteHook = queryRouteHook;
|
||||
}
|
||||
|
||||
public ParameterConverter<Endpoints, Endpoints> getQueryAssignmentEndpointConverter() {
|
||||
return queryAssignmentEndpointConverter;
|
||||
}
|
||||
|
||||
public void setQueryAssignmentEndpointConverter(
|
||||
ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter) {
|
||||
this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter;
|
||||
}
|
||||
|
||||
public AssignmentQueueSelector getAssignmentQueueSelector() {
|
||||
return assignmentQueueSelector;
|
||||
}
|
||||
|
||||
public void setAssignmentQueueSelector(
|
||||
AssignmentQueueSelector assignmentQueueSelector) {
|
||||
this.assignmentQueueSelector = assignmentQueueSelector;
|
||||
}
|
||||
|
||||
public ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> getQueryAssignmentHook() {
|
||||
return queryAssignmentHook;
|
||||
}
|
||||
|
||||
public void setQueryAssignmentHook(
|
||||
ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook) {
|
||||
this.queryAssignmentHook = queryAssignmentHook;
|
||||
}
|
||||
}
|
||||
@@ -98,7 +98,6 @@ import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseWriter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
|
||||
@@ -106,7 +105,7 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.ReceiveMessageChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.SendMessageChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.ReceiveMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.SendMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.local.RouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.local.LocalWriteQueueSelector;
|
||||
import org.apache.rocketmq.remoting.RemotingServer;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRemotingAbstract;
|
||||
@@ -143,7 +142,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
ConnectorManager connectorManager = new ConnectorManager(null);
|
||||
this.telemetryCommandManager = telemetryCommandManager;
|
||||
this.grpcClientManager = new GrpcClientManager();
|
||||
this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager, grpcClientManager);
|
||||
this.routeService = new RouteService(connectorManager, grpcClientManager);
|
||||
this.clientSettingsService = new ClientSettingsService(this.channelManager, this.grpcClientManager, this.telemetryCommandManager);
|
||||
this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel());
|
||||
this.localWriteQueueSelector = new LocalWriteQueueSelector(brokerController.getBrokerConfig().getBrokerName(),
|
||||
|
||||
+4
-65
@@ -33,38 +33,21 @@ import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.common.ParameterConverter;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.AbstractRouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
|
||||
public class RouteService extends BaseService {
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter;
|
||||
private volatile ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook;
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter;
|
||||
private volatile AssignmentQueueSelector assignmentQueueSelector;
|
||||
private volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook;
|
||||
|
||||
protected final GrpcClientManager grpcClientManager;
|
||||
|
||||
public class RouteService extends AbstractRouteService {
|
||||
public RouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
super(connectorManager);
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
super(connectorManager, grpcClientManager);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
this.queryRouteEndpointConverter = (ctx, parameter) -> parameter;
|
||||
this.queryAssignmentEndpointConverter = (ctx, parameter) -> parameter;
|
||||
this.assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
}
|
||||
|
||||
public CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
CompletableFuture<QueryRouteResponse> future = new CompletableFuture<>();
|
||||
future.whenComplete((response, throwable) -> {
|
||||
@@ -115,6 +98,7 @@ public class RouteService extends BaseService {
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request) {
|
||||
CompletableFuture<QueryAssignmentResponse> future = new CompletableFuture<>();
|
||||
future.whenComplete((response, throwable) -> {
|
||||
@@ -164,49 +148,4 @@ public class RouteService extends BaseService {
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
public ParameterConverter<Endpoints, Endpoints> getQueryRouteEndpointConverter() {
|
||||
return queryRouteEndpointConverter;
|
||||
}
|
||||
|
||||
public void setQueryRouteEndpointConverter(
|
||||
ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter) {
|
||||
this.queryRouteEndpointConverter = queryRouteEndpointConverter;
|
||||
}
|
||||
|
||||
public ResponseHook<QueryRouteRequest, QueryRouteResponse> getQueryRouteHook() {
|
||||
return queryRouteHook;
|
||||
}
|
||||
|
||||
public void setQueryRouteHook(
|
||||
ResponseHook<QueryRouteRequest, QueryRouteResponse> queryRouteHook) {
|
||||
this.queryRouteHook = queryRouteHook;
|
||||
}
|
||||
|
||||
public ParameterConverter<Endpoints, Endpoints> getQueryAssignmentEndpointConverter() {
|
||||
return queryAssignmentEndpointConverter;
|
||||
}
|
||||
|
||||
public void setQueryAssignmentEndpointConverter(
|
||||
ParameterConverter<Endpoints, Endpoints> queryAssignmentEndpointConverter) {
|
||||
this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter;
|
||||
}
|
||||
|
||||
public AssignmentQueueSelector getAssignmentQueueSelector() {
|
||||
return assignmentQueueSelector;
|
||||
}
|
||||
|
||||
public void setAssignmentQueueSelector(
|
||||
AssignmentQueueSelector assignmentQueueSelector) {
|
||||
this.assignmentQueueSelector = assignmentQueueSelector;
|
||||
}
|
||||
|
||||
public ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> getQueryAssignmentHook() {
|
||||
return queryAssignmentHook;
|
||||
}
|
||||
|
||||
public void setQueryAssignmentHook(
|
||||
ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook) {
|
||||
this.queryAssignmentHook = queryAssignmentHook;
|
||||
}
|
||||
}
|
||||
|
||||
+175
@@ -0,0 +1,175 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.local;
|
||||
|
||||
import apache.rocketmq.v2.Address;
|
||||
import apache.rocketmq.v2.AddressScheme;
|
||||
import apache.rocketmq.v2.Assignment;
|
||||
import apache.rocketmq.v2.Broker;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.Endpoints;
|
||||
import apache.rocketmq.v2.MessageQueue;
|
||||
import apache.rocketmq.v2.Permission;
|
||||
import apache.rocketmq.v2.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v2.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v2.QueryRouteRequest;
|
||||
import apache.rocketmq.v2.QueryRouteResponse;
|
||||
import com.google.common.net.HostAndPort;
|
||||
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.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.AbstractRouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
|
||||
public class RouteService extends AbstractRouteService {
|
||||
public RouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
super(connectorManager, grpcClientManager);
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
|
||||
CompletableFuture<QueryRouteResponse> future = new CompletableFuture<>();
|
||||
future.whenComplete((response, throwable) -> {
|
||||
if (queryRouteHook != null) {
|
||||
queryRouteHook.beforeResponse(ctx, request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName);
|
||||
TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData();
|
||||
List<QueueData> queueDataList = topicRouteData.getQueueDatas();
|
||||
List<BrokerData> brokerDataList = topicRouteData.getBrokerDatas();
|
||||
|
||||
List<MessageQueue> messageQueueList = new ArrayList<>();
|
||||
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()) {
|
||||
messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker));
|
||||
}
|
||||
}
|
||||
|
||||
QueryRouteResponse response = QueryRouteResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
|
||||
.addAllMessageQueues(messageQueueList)
|
||||
.build();
|
||||
future.complete(response);
|
||||
} catch (Throwable t) {
|
||||
if (TopicRouteHelper.isTopicNotExistError(t)) {
|
||||
future.complete(QueryRouteResponse.newBuilder()
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.TOPIC_NOT_FOUND, t.getMessage()))
|
||||
.build());
|
||||
} else {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request) {
|
||||
CompletableFuture<QueryAssignmentResponse> future = new CompletableFuture<>();
|
||||
future.whenComplete((response, throwable) -> {
|
||||
if (queryAssignmentHook != null) {
|
||||
queryAssignmentHook.beforeResponse(ctx, request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
List<Assignment> assignments = new ArrayList<>();
|
||||
List<SelectableMessageQueue> messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request);
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName);
|
||||
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);
|
||||
|
||||
MessageQueue defaultMessageQueue = MessageQueue.newBuilder()
|
||||
.setTopic(request.getTopic())
|
||||
.setId(-1)
|
||||
.setPermission(Permission.READ_WRITE)
|
||||
.setBroker(broker)
|
||||
.build();
|
||||
|
||||
assignments.add(Assignment.newBuilder()
|
||||
.setMessageQueue(defaultMessageQueue)
|
||||
.build());
|
||||
}
|
||||
}
|
||||
QueryAssignmentResponse response = QueryAssignmentResponse.newBuilder()
|
||||
.addAllAssignments(assignments)
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
|
||||
.build();
|
||||
future.complete(response);
|
||||
} catch (Throwable t) {
|
||||
future.completeExceptionally(t);
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
private Map<String/*brokerName*/, Map<Long/*brokerID*/, 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(ConfigurationManager.getProxyConfig().getGrpcServerPort())
|
||||
.setHost(hostAndPort.getHost())
|
||||
)
|
||||
.build())
|
||||
.build();
|
||||
|
||||
brokerIdMap.put(brokerId, broker);
|
||||
}
|
||||
brokerMap.put(brokerName, brokerIdMap);
|
||||
}
|
||||
return brokerMap;
|
||||
}
|
||||
}
|
||||
+132
@@ -0,0 +1,132 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.local;
|
||||
|
||||
import apache.rocketmq.v2.Address;
|
||||
import apache.rocketmq.v2.AddressScheme;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.Endpoints;
|
||||
import apache.rocketmq.v2.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v2.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v2.QueryRouteRequest;
|
||||
import apache.rocketmq.v2.QueryRouteResponse;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import com.google.common.net.HostAndPort;
|
||||
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.client.exception.MQClientException;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
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.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.BaseServiceTest;
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
public class RouteServiceTest extends BaseServiceTest {
|
||||
private String brokerAddress = "127.0.0.1:10911";
|
||||
private static final Settings WITH_HOST_SETTINGS = Settings.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.build();
|
||||
|
||||
@Test
|
||||
public void testLocalModeQueryRoute() throws Exception {
|
||||
RouteService routeService = new RouteService(this.connectorManager, this.grpcClientManager);
|
||||
routeService.start();
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build());
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber());
|
||||
assertEquals(8, response.getMessageQueuesCount());
|
||||
assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getMessageQueues(0).getBroker()
|
||||
.getEndpoints().getAddresses(0).getHost());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testLocalModeQueryAssignment() throws Exception {
|
||||
RouteService routeService = new RouteService(this.connectorManager, this.grpcClientManager);
|
||||
routeService.start();
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName("group")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber());
|
||||
assertEquals(1, response.getAssignmentsCount());
|
||||
assertEquals("brokerName", response.getAssignments(0).getMessageQueue().getBroker().getName());
|
||||
assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getAssignments(0).getMessageQueue().getBroker().getEndpoints().getAddresses(0).getHost());
|
||||
}
|
||||
|
||||
@Override public void beforeEach() throws Throwable {
|
||||
TopicRouteData routeData = new TopicRouteData();
|
||||
|
||||
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);
|
||||
|
||||
routeData.setBrokerDatas(brokerDataList);
|
||||
routeData.setQueueDatas(queueDataList);
|
||||
|
||||
MessageQueueWrapper messageQueueWrapper = new MessageQueueWrapper("topic", routeData);
|
||||
when(this.topicRouteCache.getMessageQueue("topic")).thenReturn(messageQueueWrapper);
|
||||
|
||||
when(this.topicRouteCache.getMessageQueue("notExistTopic")).thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST, ""));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user