diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java index 55a62de108..04619890f3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java @@ -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 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() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java index 5efa590fba..5a37e6a910 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java @@ -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 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(); - } } \ No newline at end of file diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java index 32aacfb3d4..21a37ef5de 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java @@ -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() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java index 62ca554f38..44963565e6 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java @@ -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 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 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 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 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 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 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 future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() .setTopic(Resource.newBuilder()