[ISSUE #3949] Support v2

This commit is contained in:
zhouxiang
2022-07-13 11:29:17 +08:00
parent 53fd599d80
commit e5466e2b60
9 changed files with 279 additions and 95 deletions
@@ -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) {
@@ -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;
}
@@ -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<RemotingCommand> resultFuture = this.producer.sendMessageBack(
@@ -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;
}
@@ -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<HeartbeatResponse> 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<SendMessageResponse> 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);
@@ -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;
@@ -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;
@@ -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);
}
}
@@ -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);
}