diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java index 487a526007..8c2fc13dd6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java @@ -37,69 +37,50 @@ public class ResponseBuilder { } public static Code buildCode(int remotingResponseCode) { - Code code; switch (remotingResponseCode) { case ResponseCode.SUCCESS: - case ResponseCode.NO_MESSAGE: { - code = Code.OK; - break; - } - case ResponseCode.SYSTEM_ERROR: { - code = Code.INTERNAL_SERVER_ERROR; - break; + case ResponseCode.NO_MESSAGE: + case ResponseCode.PULL_RETRY_IMMEDIATELY: { + return Code.OK; } case ResponseCode.SYSTEM_BUSY: case ResponseCode.POLLING_FULL: { - code = Code.TOO_MANY_REQUESTS; - break; + return Code.TOO_MANY_REQUESTS; } case ResponseCode.REQUEST_CODE_NOT_SUPPORTED: { - code = Code.UNRECOGNIZED; - break; + return Code.UNRECOGNIZED; } - case ResponseCode.MESSAGE_ILLEGAL: - case ResponseCode.VERSION_NOT_SUPPORTED: - case ResponseCode.SUBSCRIPTION_PARSE_FAILED: - case ResponseCode.FILTER_DATA_NOT_EXIST: { - code = Code.INVALID_ARGUMENT; - break; + case ResponseCode.MESSAGE_ILLEGAL: { + return Code.ILLEGAL_MESSAGE; } - case ResponseCode.SERVICE_NOT_AVAILABLE: - case ResponseCode.SLAVE_NOT_AVAILABLE: - case ResponseCode.PULL_RETRY_IMMEDIATELY: - case ResponseCode.PULL_OFFSET_MOVED: - case ResponseCode.SUBSCRIPTION_NOT_LATEST: - case ResponseCode.FILTER_DATA_NOT_LATEST: { - code = Code.UNAVAILABLE; - break; + case ResponseCode.VERSION_NOT_SUPPORTED: { + return Code.VERSION_UNSUPPORTED; + } + case ResponseCode.SLAVE_NOT_AVAILABLE: { + return Code.HA_NOT_AVAILABLE; + } + case ResponseCode.PULL_OFFSET_MOVED: { + return Code.ILLEGAL_MESSAGE_OFFSET; } case ResponseCode.NO_PERMISSION: { - code = Code.PERMISSION_DENIED; - break; + return Code.FORBIDDEN; } - case ResponseCode.TOPIC_NOT_EXIST: - code = Code.TOPIC_NOT_FOUND; - break; - case ResponseCode.SUBSCRIPTION_GROUP_NOT_EXIST: - case ResponseCode.SUBSCRIPTION_NOT_EXIST: - case ResponseCode.PULL_NOT_FOUND: - case ResponseCode.QUERY_NOT_FOUND: - case ResponseCode.CONSUMER_NOT_ONLINE: { - code = Code.NOT_FOUND; - break; + case ResponseCode.TOPIC_NOT_EXIST: { + return Code.TOPIC_NOT_FOUND; + } + case ResponseCode.PULL_NOT_FOUND: { + return Code.MESSAGE_NOT_FOUND; + } + case ResponseCode.FLUSH_DISK_TIMEOUT: { + return Code.MASTER_PERSISTENCE_TIMEOUT; } - case ResponseCode.POLLING_TIMEOUT: - case ResponseCode.FLUSH_DISK_TIMEOUT: case ResponseCode.FLUSH_SLAVE_TIMEOUT: { - code = Code.DEADLINE_EXCEEDED; - break; + return Code.SLAVE_PERSISTENCE_TIMEOUT; } default: { - code = Code.INTERNAL_SERVER_ERROR; + return Code.INTERNAL_SERVER_ERROR; } - } - return code; } public static String buildMessage(int responseCode, String remark) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java index 758c1ba9db..835e31f6b8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java @@ -44,11 +44,11 @@ public class BaseService { protected String getBrokerAddr(Context ctx, String brokerName) throws Exception { if (StringUtils.isBlank(brokerName)) { - throw new ProxyException(Code.INVALID_ARGUMENT, "broker name is empty"); + throw new ProxyException(Code.UNRECOGNIZED, "broker name is empty"); } String addr = this.connectorManager.getTopicRouteCache().getBrokerAddr(brokerName); if (StringUtils.isBlank(addr)) { - throw new ProxyException(Code.NOT_FOUND, brokerName + " not exist"); + throw new ProxyException(Code.UNRECOGNIZED, brokerName + " not exist"); } return addr; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index 322e670350..d696dd07bd 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -249,7 +249,7 @@ public class ConsumerService extends BaseService { ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); - Settings settings = GrpcClientManager.getClientSettings(ctx); + Settings settings = GrpcClientManager.getClientSettings(ctx).getSettings(); int maxDeliveryAttempts = settings.getSubscription().getDeadLetterPolicy().getMaxDeliveryAttempts(); if (request.getDeliveryAttempt() >= maxDeliveryAttempts) { CompletableFuture resultFuture = this.producer.sendMessageBack( 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 8da26d154e..83322e2400 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,6 +20,7 @@ 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; @@ -51,6 +52,7 @@ import org.apache.rocketmq.proxy.common.ParameterConverter; 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.ResponseHook; +import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; public class RouteService extends BaseService { private final ProxyMode mode; @@ -238,11 +240,12 @@ public class RouteService extends BaseService { } } if (ProxyMode.isClusterMode(mode)) { - Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, request.getEndpoints()); + ClientSettings clientSettings = GrpcClientManager.getClientSettings(ctx); + Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) { future.complete(QueryAssignmentResponse.newBuilder() .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + - request.getEndpoints() + " is invalidate")) + clientSettings.getAccessPoint() + " is invalidate")) .build()); return future; } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java index 0971d9a341..f3e276c91f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java @@ -29,6 +29,7 @@ import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.HeartbeatResponse; import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.MessageQueue; import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; @@ -42,6 +43,7 @@ import apache.rocketmq.v2.ReceiveMessageResponse; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.SendMessageResponse; +import apache.rocketmq.v2.SystemProperties; import com.google.protobuf.Timestamp; import com.google.protobuf.util.Durations; import io.grpc.Context; @@ -73,8 +75,8 @@ import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader; import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader; import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; -import org.apache.rocketmq.proxy.grpc.v1.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.remoting.exception.RemotingCommandException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.apache.rocketmq.store.MessageStore; @@ -146,18 +148,15 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .thenReturn(response); Mockito.when(brokerControllerMock.getClientManageProcessor()).thenReturn(clientManageProcessorMock); HeartbeatRequest request = HeartbeatRequest.newBuilder() - .setClientId("test-client") - .setConsumerData(ConsumerData.newBuilder() - .setGroup(Resource.newBuilder() - .setName("group") - .build()) + .setGroup(Resource.newBuilder() + .setName("group") .build()) .build(); CompletableFuture grpcFuture = localGrpcService.heartbeat( Context.current().withValue(InterceptorConstants.METADATA, metadata).attach(), request); HeartbeatResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()) - .isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()) + .isEqualTo(Code.OK); } @Test @@ -167,8 +166,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) .thenReturn(response); SendMessageRequest request = SendMessageRequest.newBuilder() - .setMessage(Message.newBuilder() - .setSystemAttribute(SystemAttribute.newBuilder() + .setMessages(0, Message.newBuilder() + .setSystemProperties(SystemProperties.newBuilder() .setMessageId("123") .build()) .build()) @@ -177,8 +176,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { CompletableFuture grpcFuture = localGrpcService.sendMessage( Context.current().withValue(InterceptorConstants.METADATA, metadata).attach(), request); SendMessageResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()) - .isEqualTo(Code.INTERNAL.getNumber()); + assertThat(r.getStatus().getCode()) + .isEqualTo(Code.INTERNAL_SERVER_ERROR); } @Test @@ -186,8 +185,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) .thenReturn(null); SendMessageRequest request = SendMessageRequest.newBuilder() - .setMessage(Message.newBuilder() - .setSystemAttribute(SystemAttribute.newBuilder() + .setMessages(0, Message.newBuilder() + .setSystemProperties(SystemProperties.newBuilder() .setMessageId("123") .build()) .build()) @@ -203,8 +202,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) .thenThrow(new RemotingCommandException("test")); SendMessageRequest request = SendMessageRequest.newBuilder() - .setMessage(Message.newBuilder() - .setSystemAttribute(SystemAttribute.newBuilder() + .setMessages(0, Message.newBuilder() + .setSystemProperties(SystemProperties.newBuilder() .setMessageId("123") .build()) .build()) @@ -241,7 +240,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.when(popMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) .thenReturn(remotingCommand); ReceiveMessageRequest request = ReceiveMessageRequest.newBuilder() - .setPartition(Partition.newBuilder() + .setMessageQueue(MessageQueue.newBuilder() .setTopic(Resource.newBuilder() .setName(topic) .build()) @@ -254,7 +253,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor( new ThreadFactoryImpl("test"))), request); ReceiveMessageResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); assertThat(r.getMessagesCount()).isEqualTo(1); assertThat(Durations.toMillis(r.getInvisibleDuration())).isEqualTo(invisibleTime); assertThat(GrpcConverter.wrapResourceWithNamespace(r.getMessages(0).getTopic())).isEqualTo(topic); @@ -300,7 +299,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); AckMessageResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); } @Test @@ -333,7 +332,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); NackMessageResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); } @Test @@ -361,7 +360,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); ForwardMessageToDeadLetterQueueResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); } @Test @@ -385,7 +384,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); EndTransactionResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); } @Test @@ -401,7 +400,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.when(messageStore.getOffsetInQueueByTime(Mockito.eq(topic), Mockito.eq(queueId), Mockito.anyLong())).thenReturn(timeOffset); QueryOffsetRequest request = QueryOffsetRequest.newBuilder() - .setPartition(Partition.newBuilder() + .setMessageQueue(MessageQueue.newBuilder() .setTopic(Resource.newBuilder() .setName(topic) .build()) @@ -414,11 +413,11 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); QueryOffsetResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); assertThat(r.getOffset()).isEqualTo(0); request = QueryOffsetRequest.newBuilder() - .setPartition(Partition.newBuilder() + .setMessageQueue(MessageQueue.newBuilder() .setTopic(Resource.newBuilder() .setName(topic) .build()) @@ -431,11 +430,11 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); assertThat(r.getOffset()).isEqualTo(maxOffset); request = QueryOffsetRequest.newBuilder() - .setPartition(Partition.newBuilder() + .setMessageQueue(MessageQueue.newBuilder() .setTopic(Resource.newBuilder() .setName(topic) .build()) @@ -451,7 +450,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); assertThat(r.getOffset()).isEqualTo(timeOffset); } @@ -470,7 +469,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { Mockito.any(RemotingCommand.class))) .thenReturn(response); NotifyClientTerminationRequest request = NotifyClientTerminationRequest.newBuilder() - .setProducerGroup(Resource.newBuilder() + .setGroup(Resource.newBuilder() .setName("group") .build()) .build(); @@ -515,7 +514,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withValue(InterceptorConstants.METADATA, metadata) .attach(), request); ChangeInvisibleDurationResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); ReceiptHandle handle = ReceiptHandle.decode(r.getReceiptHandle()); assertThat(handle.getInvisibleTime()).isEqualTo(invisibleTime); assertThat(handle.getQueueId()).isEqualTo(queueId); @@ -554,7 +553,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor( new ThreadFactoryImpl("test"))), request); PullMessageResponse r = grpcFuture.get(); - assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); + assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); assertThat(r.getMessagesCount()).isEqualTo(1); assertThat(GrpcConverter.wrapResourceWithNamespace(r.getMessages(0).getTopic())).isEqualTo(topic); assertThat(r.getMessages(0).getBody().toByteArray()).isEqualTo(body); diff --git a/test/src/test/java/org/apache/rocketmq/test/proxy/ClusterGrpcTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v1/ClusterGrpcTest.java similarity index 84% rename from test/src/test/java/org/apache/rocketmq/test/proxy/ClusterGrpcTest.java rename to test/src/test/java/org/apache/rocketmq/test/grpc/v1/ClusterGrpcTest.java index 9a04cfcf63..5d40e39e29 100644 --- a/test/src/test/java/org/apache/rocketmq/test/proxy/ClusterGrpcTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v1/ClusterGrpcTest.java @@ -1,4 +1,21 @@ -package org.apache.rocketmq.test.proxy; +/* + * 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.test.grpc.v1; import apache.rocketmq.v1.AckMessageResponse; import apache.rocketmq.v1.Address; @@ -16,7 +33,6 @@ import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.v1.GrpcMessagingProcessor; import org.apache.rocketmq.proxy.grpc.v1.service.ClusterGrpcService; import org.apache.rocketmq.proxy.grpc.v1.service.GrpcForwardService; -import org.apache.rocketmq.test.base.GrpcBaseTest; import org.junit.After; import org.junit.Before; import org.junit.Test; diff --git a/test/src/test/java/org/apache/rocketmq/test/base/GrpcBaseTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v1/GrpcBaseTest.java similarity index 98% rename from test/src/test/java/org/apache/rocketmq/test/base/GrpcBaseTest.java rename to test/src/test/java/org/apache/rocketmq/test/grpc/v1/GrpcBaseTest.java index 2739348ab4..5494f5405d 100644 --- a/test/src/test/java/org/apache/rocketmq/test/base/GrpcBaseTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v1/GrpcBaseTest.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.test.base; +package org.apache.rocketmq.test.grpc.v1; import apache.rocketmq.v1.AckMessageRequest; import apache.rocketmq.v1.AckMessageResponse; @@ -49,11 +49,10 @@ import io.netty.handler.ssl.util.SelfSignedCertificate; import java.io.IOException; import java.security.cert.CertificateException; import java.util.concurrent.TimeUnit; -import java.util.function.Function; -import java.util.function.Supplier; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.interceptor.ContextInterceptor; import org.apache.rocketmq.proxy.grpc.interceptor.HeaderInterceptor; +import org.apache.rocketmq.test.base.BaseConf; import org.junit.Rule; import static org.assertj.core.api.Assertions.assertThat; diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java new file mode 100644 index 0000000000..7760be6458 --- /dev/null +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java @@ -0,0 +1,187 @@ +/* + * 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.test.grpc.v2; + +import apache.rocketmq.v2.AckMessageRequest; +import apache.rocketmq.v2.AckMessageResponse; +import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.Endpoints; +import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.MessageQueue; +import apache.rocketmq.v2.MessagingServiceGrpc; +import apache.rocketmq.v2.QueryRouteRequest; +import apache.rocketmq.v2.QueryRouteResponse; +import apache.rocketmq.v2.ReceiveMessageRequest; +import apache.rocketmq.v2.ReceiveMessageResponse; +import apache.rocketmq.v2.Resource; +import apache.rocketmq.v2.SendMessageRequest; +import apache.rocketmq.v2.SendMessageResponse; +import apache.rocketmq.v2.SystemProperties; +import com.google.protobuf.ByteString; +import com.google.protobuf.Duration; +import com.google.protobuf.Timestamp; +import io.grpc.Channel; +import io.grpc.ServerInterceptors; +import io.grpc.ServerServiceDefinition; +import io.grpc.netty.shaded.io.grpc.netty.NettyChannelBuilder; +import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder; +import io.grpc.netty.shaded.io.netty.handler.ssl.ApplicationProtocolConfig; +import io.grpc.netty.shaded.io.netty.handler.ssl.SslContextBuilder; +import io.grpc.netty.shaded.io.netty.handler.ssl.SslProvider; +import io.grpc.testing.GrpcCleanupRule; +import io.netty.handler.ssl.ApplicationProtocolNames; +import io.netty.handler.ssl.util.InsecureTrustManagerFactory; +import io.netty.handler.ssl.util.SelfSignedCertificate; +import java.io.IOException; +import java.security.cert.CertificateException; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.grpc.interceptor.ContextInterceptor; +import org.apache.rocketmq.proxy.grpc.interceptor.HeaderInterceptor; +import org.apache.rocketmq.test.base.BaseConf; +import org.junit.Rule; + +import static org.assertj.core.api.Assertions.assertThat; + +public class GrpcBaseTest extends BaseConf { + /** + * This rule manages automatic graceful shutdown for the registered servers and channels at the end of test. + */ + @Rule + public final GrpcCleanupRule grpcCleanup = new GrpcCleanupRule(); + + private static final int defaultQueueNums = 8; + + protected Channel setUpServer(MessagingServiceGrpc.MessagingServiceImplBase serverImpl, + int port, boolean enableInterceptor) throws IOException, CertificateException { + SelfSignedCertificate selfSignedCertificate = new SelfSignedCertificate(); + ServerServiceDefinition serviceDefinition = ServerInterceptors.intercept(serverImpl); + if (enableInterceptor) { + serviceDefinition = ServerInterceptors.intercept(serverImpl, new ContextInterceptor(), new HeaderInterceptor()); + } + // Create a server, add service, start, and register for automatic graceful shutdown. + grpcCleanup.register(NettyServerBuilder.forPort(port) + .directExecutor() + .addService(serviceDefinition) + .useTransportSecurity(selfSignedCertificate.certificate(), selfSignedCertificate.privateKey()) + .build() + .start()); + // Create a client channel and register for automatic graceful shutdown. + return grpcCleanup.register(NettyChannelBuilder.forAddress("127.0.0.1", port) + .directExecutor() + .sslContext(SslContextBuilder + .forClient() + .sslProvider(SslProvider.OPENSSL) + .trustManager(InsecureTrustManagerFactory.INSTANCE) + .applicationProtocolConfig(new ApplicationProtocolConfig( + ApplicationProtocolConfig.Protocol.ALPN, + ApplicationProtocolConfig.SelectorFailureBehavior.NO_ADVERTISE, + ApplicationProtocolConfig.SelectedListenerFailureBehavior.ACCEPT, + ApplicationProtocolNames.HTTP_2)) + .build() + ) + .build()); + } + + public QueryRouteRequest buildQueryRouteRequest(String topic) { + return buildQueryRouteRequest(topic, Endpoints.getDefaultInstance()); + } + + public QueryRouteRequest buildQueryRouteRequest(String topic, Endpoints endpoints) { + return QueryRouteRequest.newBuilder() + .setTopic(Resource.newBuilder() + .setName(topic) + .build()) + .setEndpoints(endpoints) + .build(); + } + + public SendMessageRequest buildSendMessageRequest(String topic, String messageId) { + return SendMessageRequest.newBuilder() + .setMessages(0, Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(topic) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(messageId) + .setQueueId(0) + .build()) + .setBody(ByteString.copyFromUtf8("123")) + .build()) + .build(); + } + + public ReceiveMessageRequest buildReceiveMessageRequest(String group, String topic) { + return ReceiveMessageRequest.newBuilder() + .setGroup(Resource.newBuilder() + .setName(group) + .build()) + .setMessageQueue(MessageQueue.newBuilder() + .setTopic(Resource.newBuilder() + .setName(topic) + .build()) + .setId(0) + .build()) + .setBatchSize(16) + .setInvisibleDuration(Duration.newBuilder() + .setSeconds(3) + .build()) + .setInitializationTimestamp(Timestamp.newBuilder() + .setSeconds(TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis())) + .build()) + .build(); + } + + public AckMessageRequest buildAckMessageRequest(String group, String topic, String receiptHandle) { + return AckMessageRequest.newBuilder() + .setGroup(Resource.newBuilder() + .setName(group) + .build()) + .setTopic(Resource.newBuilder() + .setName(topic) + .build()) + .setReceiptHandle(receiptHandle) + .build(); + } + + public void assertQueryRoute(QueryRouteResponse response, int brokerSize) { + assertThat(response.getStatus().getCode()).isEqualTo(Code.OK); + assertThat(response.getMessageQueuesList().size()).isEqualTo(brokerSize * defaultQueueNums); + assertThat(response.getMessageQueues(0).getBroker().getEndpoints().getAddresses(0).getPort()).isEqualTo(ConfigurationManager.getProxyConfig().getGrpcServerPort()); + } + + public void assertSendMessage(SendMessageResponse response, String messageId) { + assertThat(response.getStatus() + .getCode()).isEqualTo(Code.OK); + assertThat(response.getReceipts(0).getMessageId()).isEqualTo(messageId); + } + + public void assertReceiveMessage(ReceiveMessageResponse response, String messageId) { + assertThat(response.getStatus() + .getCode()).isEqualTo(Code.OK); + assertThat(response.getMessagesCount()).isEqualTo(1); + assertThat(response.getMessages(0) + .getSystemProperties() + .getMessageId()).isEqualTo(messageId); + } + + public void assertAck(AckMessageResponse response) { + assertThat(response.getStatus() + .getCode()).isEqualTo(Code.OK); + } +} \ No newline at end of file diff --git a/test/src/test/java/org/apache/rocketmq/test/proxy/LocalGrpcTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java similarity index 89% rename from test/src/test/java/org/apache/rocketmq/test/proxy/LocalGrpcTest.java rename to test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java index eaf285ac87..3b2610848e 100644 --- a/test/src/test/java/org/apache/rocketmq/test/proxy/LocalGrpcTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java @@ -15,20 +15,19 @@ * limitations under the License. */ -package org.apache.rocketmq.test.proxy; +package org.apache.rocketmq.test.grpc.v2; -import apache.rocketmq.v1.AckMessageResponse; -import apache.rocketmq.v1.MessagingServiceGrpc; -import apache.rocketmq.v1.QueryRouteResponse; -import apache.rocketmq.v1.ReceiveMessageResponse; -import apache.rocketmq.v1.SendMessageResponse; +import apache.rocketmq.v2.AckMessageResponse; +import apache.rocketmq.v2.MessagingServiceGrpc; +import apache.rocketmq.v2.QueryRouteResponse; +import apache.rocketmq.v2.ReceiveMessageResponse; +import apache.rocketmq.v2.SendMessageResponse; import io.grpc.Channel; import java.net.URL; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.grpc.v1.GrpcMessagingProcessor; +import org.apache.rocketmq.proxy.grpc.v2.GrpcMessagingProcessor; import org.apache.rocketmq.proxy.grpc.v2.service.LocalGrpcService; -import org.apache.rocketmq.test.base.GrpcBaseTest; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -82,7 +81,7 @@ public class LocalGrpcTest extends GrpcBaseTest { ReceiveMessageResponse receiveResponse = blockingStub.withDeadlineAfter(3, TimeUnit.SECONDS) .receiveMessage(buildReceiveMessageRequest(group, broker1Name)); assertReceiveMessage(receiveResponse, messageId); - String receiptHandle = receiveResponse.getMessages(0).getSystemAttribute().getReceiptHandle(); + String receiptHandle = receiveResponse.getMessages(0).getSystemProperties().getReceiptHandle(); AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(group, broker1Name, receiptHandle)); assertAck(ackMessageResponse); }