mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
+2
-3
@@ -20,7 +20,6 @@ import apache.rocketmq.v2.Address;
|
||||
import apache.rocketmq.v2.AddressScheme;
|
||||
import apache.rocketmq.v2.Assignment;
|
||||
import apache.rocketmq.v2.Broker;
|
||||
import apache.rocketmq.v2.ClientSettings;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.Endpoints;
|
||||
import apache.rocketmq.v2.MessageQueue;
|
||||
@@ -113,7 +112,7 @@ public class RouteService extends BaseService {
|
||||
|
||||
List<MessageQueue> messageQueueList = new ArrayList<>();
|
||||
if (ProxyMode.isClusterMode(mode.name())) {
|
||||
ClientSettings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
GrpcClientManager.ActiveClientSettings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) {
|
||||
future.complete(QueryRouteResponse.newBuilder()
|
||||
@@ -244,7 +243,7 @@ public class RouteService extends BaseService {
|
||||
}
|
||||
}
|
||||
if (ProxyMode.isClusterMode(mode)) {
|
||||
ClientSettings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
GrpcClientManager.ActiveClientSettings clientSettings = grpcClientManager.getClientSettings(ctx);
|
||||
Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) {
|
||||
future.complete(QueryAssignmentResponse.newBuilder()
|
||||
|
||||
+7
-30
@@ -1,10 +1,9 @@
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
|
||||
|
||||
import apache.rocketmq.v2.AckMessageEntry;
|
||||
import apache.rocketmq.v2.AckMessageRequest;
|
||||
import apache.rocketmq.v2.AckMessageResponse;
|
||||
import apache.rocketmq.v2.ClientSettings;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.DeadLetterPolicy;
|
||||
import apache.rocketmq.v2.FilterExpression;
|
||||
import apache.rocketmq.v2.FilterType;
|
||||
import apache.rocketmq.v2.NackMessageRequest;
|
||||
@@ -12,8 +11,6 @@ import apache.rocketmq.v2.NackMessageResponse;
|
||||
import apache.rocketmq.v2.ReceiveMessageRequest;
|
||||
import apache.rocketmq.v2.ReceiveMessageResponse;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import apache.rocketmq.v2.Subscription;
|
||||
import io.grpc.Context;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
@@ -30,6 +27,7 @@ import org.apache.rocketmq.common.message.MessageQueue;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.assertj.core.util.Lists;
|
||||
@@ -62,15 +60,6 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr");
|
||||
when(readQueueSelector.select(any(), any(), any())).thenReturn(selectableMessageQueue);
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
.setSettings(Settings.newBuilder()
|
||||
.setSubscription(Subscription.newBuilder()
|
||||
.setFifo(false)
|
||||
.build())
|
||||
.build())
|
||||
.build();
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
|
||||
List<MessageExt> messageExtList = Lists.newArrayList(
|
||||
createMessageExt("msg1", "msg1"),
|
||||
createMessageExt("msg2", "msg2")
|
||||
@@ -120,7 +109,9 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName("group")
|
||||
.build())
|
||||
.setReceiptHandle(createReceiptHandle().encode())
|
||||
.addEntries(AckMessageEntry.newBuilder()
|
||||
.setMessageId("msgId")
|
||||
.setReceiptHandle(createReceiptHandle().encode()))
|
||||
.build())
|
||||
.get();
|
||||
|
||||
@@ -137,8 +128,7 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
}).when(producerClient).sendMessageBack(anyString(), any());
|
||||
when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr");
|
||||
|
||||
ClientSettings clientSettings = createClientSettings(3);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
ConfigurationManager.getProxyConfig().setDefaultMaxDeliveryAttempts(3);
|
||||
|
||||
NackMessageResponse response = consumerService.nackMessage(Context.current(), NackMessageRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -168,8 +158,7 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
}).when(writeConsumerClient).changeInvisibleTimeAsync(anyString(), anyString(), any());
|
||||
when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr");
|
||||
|
||||
ClientSettings clientSettings = createClientSettings(3);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
ConfigurationManager.getProxyConfig().setDefaultMaxDeliveryAttempts(3);
|
||||
|
||||
NackMessageResponse response = consumerService.nackMessage(Context.current(), NackMessageRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -187,16 +176,4 @@ public class ConsumerServiceTest extends BaseServiceTest {
|
||||
assertEquals(receiptHandle.getOffset(), headerRef.get().getOffset().longValue());
|
||||
assertEquals(receiptHandle.encode(), headerRef.get().getExtraInfo());
|
||||
}
|
||||
|
||||
private ClientSettings createClientSettings(int maxDeliveryAttempts) {
|
||||
return ClientSettings.newBuilder()
|
||||
.setSettings(Settings.newBuilder()
|
||||
.setSubscription(Subscription.newBuilder()
|
||||
.setDeadLetterPolicy(DeadLetterPolicy.newBuilder()
|
||||
.setMaxDeliveryAttempts(maxDeliveryAttempts)
|
||||
.build())
|
||||
.build())
|
||||
.build())
|
||||
.build();
|
||||
}
|
||||
}
|
||||
+16
-20
@@ -1,15 +1,14 @@
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
|
||||
|
||||
import apache.rocketmq.v2.ClientSettings;
|
||||
import apache.rocketmq.v2.ActivePublishingSettings;
|
||||
import apache.rocketmq.v2.ActiveSubscriptionSettings;
|
||||
import apache.rocketmq.v2.ClientType;
|
||||
import apache.rocketmq.v2.FilterExpression;
|
||||
import apache.rocketmq.v2.FilterType;
|
||||
import apache.rocketmq.v2.HeartbeatRequest;
|
||||
import apache.rocketmq.v2.NotifyClientTerminationRequest;
|
||||
import apache.rocketmq.v2.Publishing;
|
||||
import apache.rocketmq.v2.ReportActiveSettingsCommand;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import apache.rocketmq.v2.Subscription;
|
||||
import apache.rocketmq.v2.SubscriptionEntry;
|
||||
import io.grpc.Context;
|
||||
import io.netty.channel.Channel;
|
||||
@@ -24,6 +23,7 @@ import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.junit.Test;
|
||||
|
||||
@@ -52,19 +52,17 @@ public class ForwardClientServiceTest extends BaseServiceTest {
|
||||
|
||||
@Test
|
||||
public void testProducerHeartbeat() {
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
GrpcClientManager.ActiveClientSettings clientSettings = new GrpcClientManager.ActiveClientSettings(ReportActiveSettingsCommand.newBuilder()
|
||||
.setClientType(ClientType.PRODUCER)
|
||||
.setSettings(Settings.newBuilder()
|
||||
.setPublishing(Publishing.newBuilder()
|
||||
.addTopics(Resource.newBuilder()
|
||||
.setName("topic1")
|
||||
.build())
|
||||
.addTopics(Resource.newBuilder()
|
||||
.setName("topic2")
|
||||
.build())
|
||||
.setActivePublishingSettings(ActivePublishingSettings.newBuilder()
|
||||
.addPublishingTopics(Resource.newBuilder()
|
||||
.setName("topic1")
|
||||
.build())
|
||||
.addPublishingTopics(Resource.newBuilder()
|
||||
.setName("topic2")
|
||||
.build())
|
||||
.build())
|
||||
.build();
|
||||
.build());
|
||||
when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings);
|
||||
|
||||
clientService.heartbeat(Context.current(), HeartbeatRequest.newBuilder().build());
|
||||
@@ -91,14 +89,12 @@ public class ForwardClientServiceTest extends BaseServiceTest {
|
||||
.build())
|
||||
.build());
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
GrpcClientManager.ActiveClientSettings clientSettings = new GrpcClientManager.ActiveClientSettings(ReportActiveSettingsCommand.newBuilder()
|
||||
.setClientType(ClientType.PUSH_CONSUMER)
|
||||
.setSettings(Settings.newBuilder()
|
||||
.setSubscription(Subscription.newBuilder()
|
||||
.addAllSubscriptions(subscriptionEntryList)
|
||||
.build())
|
||||
.setActiveSubscriptionSettings(ActiveSubscriptionSettings.newBuilder()
|
||||
.addAllSubscriptions(subscriptionEntryList)
|
||||
.build())
|
||||
.build();
|
||||
.build());
|
||||
when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings);
|
||||
|
||||
clientService.heartbeat(Context.current(), HeartbeatRequest.newBuilder()
|
||||
|
||||
+23
-53
@@ -20,7 +20,6 @@ package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
|
||||
import apache.rocketmq.v2.Address;
|
||||
import apache.rocketmq.v2.AddressScheme;
|
||||
import apache.rocketmq.v2.Broker;
|
||||
import apache.rocketmq.v2.ClientSettings;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.Endpoints;
|
||||
import apache.rocketmq.v2.MessageQueue;
|
||||
@@ -29,6 +28,7 @@ import apache.rocketmq.v2.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v2.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v2.QueryRouteRequest;
|
||||
import apache.rocketmq.v2.QueryRouteResponse;
|
||||
import apache.rocketmq.v2.ReportActiveSettingsCommand;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import com.google.common.net.HostAndPort;
|
||||
import io.grpc.Context;
|
||||
@@ -44,6 +44,7 @@ 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.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
@@ -62,6 +63,20 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
.setResourceNamespace(NAMESPACE)
|
||||
.build();
|
||||
|
||||
private static final GrpcClientManager.ActiveClientSettings WITH_HOST_SETTINGS = new GrpcClientManager.ActiveClientSettings(ReportActiveSettingsCommand.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.buildPartial());
|
||||
|
||||
private static final GrpcClientManager.ActiveClientSettings INVALID_HOST_SETTINGS = new GrpcClientManager.ActiveClientSettings(ReportActiveSettingsCommand.newBuilder()
|
||||
.setAccessPoint(Endpoints.getDefaultInstance())
|
||||
.buildPartial());
|
||||
|
||||
@Override
|
||||
public void beforeEach() throws Exception {
|
||||
TopicRouteData routeData = new TopicRouteData();
|
||||
@@ -150,16 +165,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testLocalModeQueryRoute() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.LOCAL, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.build();
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -177,7 +183,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testQueryRouteWithInvalidEndpoints() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(ClientSettings.getDefaultInstance());
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(INVALID_HOST_SETTINGS);
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
@@ -192,16 +198,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testQueryRoute() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.build();
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -220,16 +217,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testQueryRouteWhenTopicNotExist() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.build();
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -245,7 +233,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testQueryAssignmentInvalidEndpoints() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(ClientSettings.getDefaultInstance());
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(INVALID_HOST_SETTINGS);
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(
|
||||
Resource.newBuilder()
|
||||
@@ -262,16 +250,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testLocalModeQueryAssignment() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.LOCAL, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.build();
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
@@ -293,16 +272,7 @@ public class RouteServiceTest extends BaseServiceTest {
|
||||
public void testQueryAssignment() throws Exception {
|
||||
RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager);
|
||||
|
||||
ClientSettings clientSettings = ClientSettings.newBuilder()
|
||||
.setAccessPoint(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.build();
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings);
|
||||
when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
|
||||
Reference in New Issue
Block a user