diff --git a/.travis.yml b/.travis.yml index 837ae1fe78..8a2d8c7492 100644 --- a/.travis.yml +++ b/.travis.yml @@ -50,8 +50,9 @@ before_script: script: - mvn verify -DskipTests - travis_retry mvn -B clean apache-rat:check - - travis_retry mvn -B install jacoco:report coveralls:report - - travis_retry mvn -B clean install -pl test -Pit-test + - travis_retry mvn -B clean test jacoco:report coveralls:report + - travis_retry mvn -B clean test -pl test -Pit-test + - travis_retry mvn -B clean install -DskipTests after_success: - mvn sonar:sonar -Psonar-apache diff --git a/acl/src/main/java/org/apache/rocketmq/acl/common/AclClientRPCHook.java b/acl/src/main/java/org/apache/rocketmq/acl/common/AclClientRPCHook.java index 50c73ca732..9e5bf1fb5d 100644 --- a/acl/src/main/java/org/apache/rocketmq/acl/common/AclClientRPCHook.java +++ b/acl/src/main/java/org/apache/rocketmq/acl/common/AclClientRPCHook.java @@ -56,10 +56,6 @@ public class AclClientRPCHook implements RPCHook { } - @Override public void doAfterRpcFailure(String remoteAddr, RemotingCommand request, Boolean remoteTimeout) { - - } - protected SortedMap parseRequestContent(RemotingCommand request, String ak, String securityToken) { CommandCustomHeader header = request.readCustomHeader(); // Sort property diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java index ffa95c454f..df19bb4d11 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java @@ -890,10 +890,6 @@ public class BrokerController { public void doAfterResponse(String remoteAddr, RemotingCommand request, RemotingCommand response) { } - @Override - public void doAfterRpcFailure(String remoteAddr, RemotingCommand request, Boolean remoteTimeout) { - - } }); } } diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java index cf4fe19ba7..152fa07d11 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java @@ -386,7 +386,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { public void createPlainAccessConfig(final String addr, final PlainAccessConfig plainAccessConfig, final long timeoutMillis) - throws RemotingException, MQBrokerException, InterruptedException, MQClientException { + throws RemotingException, InterruptedException, MQClientException { CreateAccessConfigRequestHeader requestHeader = new CreateAccessConfigRequestHeader(); requestHeader.setAccessKey(plainAccessConfig.getAccessKey()); requestHeader.setSecretKey(plainAccessConfig.getSecretKey()); @@ -414,7 +414,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { } public void deleteAccessConfig(final String addr, final String accessKey, final long timeoutMillis) - throws RemotingException, MQBrokerException, InterruptedException, MQClientException { + throws RemotingException, InterruptedException, MQClientException { DeleteAccessConfigRequestHeader requestHeader = new DeleteAccessConfigRequestHeader(); requestHeader.setAccessKey(accessKey); @@ -435,7 +435,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { } public void updateGlobalWhiteAddrsConfig(final String addr, final String globalWhiteAddrs, final long timeoutMillis) - throws RemotingException, MQBrokerException, InterruptedException, MQClientException { + throws RemotingException, InterruptedException, MQClientException { UpdateGlobalWhiteAddrsConfigRequestHeader requestHeader = new UpdateGlobalWhiteAddrsConfigRequestHeader(); requestHeader.setGlobalWhiteAddrs(globalWhiteAddrs); @@ -490,7 +490,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { } public AclConfig getBrokerClusterConfig(final String addr, - final long timeoutMillis) throws RemotingCommandException, InterruptedException, RemotingTimeoutException, + final long timeoutMillis) throws InterruptedException, RemotingTimeoutException, RemotingSendRequestException, RemotingConnectException, MQBrokerException { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_BROKER_CLUSTER_ACL_CONFIG, null); @@ -605,12 +605,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { if (this.remotingClient instanceof NettyRemotingClient) { NettyRemotingClient remotingClient = (NettyRemotingClient) this.remotingClient; RemotingCommand response = responseFuture.getResponseCommand(); - if (response == null) { - remotingClient.doAfterRpcFailure(RemotingHelper.parseChannelRemoteAddr(responseFuture.getChannel()), responseFuture.getRequestCommand(), - responseFuture.isTimeout()); - } else { - remotingClient.doAfterRpcHooks(RemotingHelper.parseChannelRemoteAddr(responseFuture.getChannel()), responseFuture.getRequestCommand(), response); - } + remotingClient.doAfterRpcHooks(RemotingHelper.parseChannelRemoteAddr(responseFuture.getChannel()), responseFuture.getRequestCommand(), response); } } @@ -1292,7 +1287,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { (QueryConsumerOffsetResponseHeader) response.decodeCommandCustomHeader(QueryConsumerOffsetResponseHeader.class); return responseHeader.getOffset(); } - case ResponseCode.PULL_NOT_FOUND: { + case ResponseCode.QUERY_NOT_FOUND: { throw new OffsetNotFoundException(response.getCode(), response.getRemark(), addr); } default: diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStoreTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStoreTest.java index f4aead1b84..c303f2d8e5 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStoreTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStoreTest.java @@ -88,7 +88,7 @@ public class RemoteBrokerOffsetStoreTest { offsetStore.updateOffset(messageQueue, 1024, false); - doThrow(new OffsetNotFoundException(ResponseCode.PULL_NOT_FOUND, "", null)) + doThrow(new OffsetNotFoundException(ResponseCode.QUERY_NOT_FOUND, "", null)) .when(mqClientAPI).queryConsumerOffset(anyString(), any(QueryConsumerOffsetRequestHeader.class), anyLong()); assertThat(offsetStore.readOffset(messageQueue, ReadOffsetType.READ_FROM_STORE)).isEqualTo(-1); diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/RPCHook.java b/remoting/src/main/java/org/apache/rocketmq/remoting/RPCHook.java index 0f17858c38..ebaeea40af 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/RPCHook.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/RPCHook.java @@ -24,6 +24,4 @@ public interface RPCHook { void doAfterResponse(final String remoteAddr, final RemotingCommand request, final RemotingCommand response); - - void doAfterRpcFailure(final String remoteAddr, RemotingCommand request, Boolean remoteTimeout); } diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java index c6c9178e9a..ab7461658d 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java @@ -182,14 +182,6 @@ public abstract class NettyRemotingAbstract { } } - public void doAfterRpcFailure(String addr, RemotingCommand request, Boolean remoteTimeout) { - if (rpcHooks.size() > 0) { - for (RPCHook rpcHook : rpcHooks) { - rpcHook.doAfterRpcFailure(addr, request, remoteTimeout); - } - } - } - /** * Process incoming request command issued by remote peer. * @@ -495,10 +487,6 @@ public abstract class NettyRemotingAbstract { throw new RemotingSendRequestException(RemotingHelper.parseChannelRemoteAddr(channel), e); } } else { - if (this instanceof NettyRemotingClient) { - NettyRemotingClient nettyRemotingClient = (NettyRemotingClient) this; - nettyRemotingClient.doAfterRpcFailure(RemotingHelper.parseChannelRemoteAddr(channel), request, false); - } if (timeoutMillis <= 0) { throw new RemotingTooMuchRequestException("invokeAsyncImpl invoke too fast"); } else { diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java index adb387189f..61c94366e7 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java @@ -468,7 +468,6 @@ public class NettyRemotingClient extends NettyRemotingAbstract implements Remoti return response; } catch (RemotingSendRequestException e) { LOGGER.warn("invokeSync: send request exception, so close the channel[{}]", addr); - doAfterRpcFailure(addr, request, false); this.closeChannel(addr, channel); throw e; } catch (RemotingTimeoutException e) { @@ -476,7 +475,6 @@ public class NettyRemotingClient extends NettyRemotingAbstract implements Remoti this.closeChannel(addr, channel); LOGGER.warn("invokeSync: close socket because of timeout, {}ms, {}", timeoutMillis, addr); } - doAfterRpcFailure(addr, request, true); LOGGER.warn("invokeSync: wait response timeout exception, the channel[{}]", addr); throw e; } @@ -686,7 +684,6 @@ public class NettyRemotingClient extends NettyRemotingAbstract implements Remoti this.invokeOnewayImpl(channel, request, timeoutMillis); } catch (RemotingSendRequestException e) { LOGGER.warn("invokeOneway: send request exception, so close the channel[{}]", addr); - doAfterRpcFailure(addr, request, false); this.closeChannel(addr, channel); throw e; } diff --git a/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java b/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java index 1e79820a39..b39b9bb03f 100644 --- a/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest.java @@ -743,11 +743,31 @@ public class DefaultMessageStoreTest { messageStoreConfig.setTotalReplicas(2); messageStoreConfig.setInSyncReplicas(2); messageStoreConfig.setEnableAutoInSyncReplicas(false); + ((DefaultMessageStore) this.messageStore).getBrokerConfig().setEnableSlaveActingMaster(true); this.messageStore.setAliveReplicaNumInGroup(1); MessageExtBrokerInner msg = buildMessage(); PutMessageResult result = this.messageStore.putMessage(msg); assertThat(result.getPutMessageStatus()).isEqualTo(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH); + ((DefaultMessageStore) this.messageStore).getBrokerConfig().setEnableSlaveActingMaster(false); + } + + + @Test + public void testPutMsgWhenAdaptiveDegradation () { + MessageStoreConfig messageStoreConfig = ((DefaultMessageStore) this.messageStore).getMessageStoreConfig(); + messageStoreConfig.setBrokerRole(BrokerRole.SYNC_MASTER); + messageStoreConfig.setTotalReplicas(2); + messageStoreConfig.setInSyncReplicas(2); + messageStoreConfig.setEnableAutoInSyncReplicas(true); + ((DefaultMessageStore) this.messageStore).getBrokerConfig().setEnableSlaveActingMaster(true); + this.messageStore.setAliveReplicaNumInGroup(1); + + MessageExtBrokerInner msg = buildMessage(); + PutMessageResult result = this.messageStore.putMessage(msg); + assertThat(result.getPutMessageStatus()).isEqualTo(PutMessageStatus.PUT_OK); + ((DefaultMessageStore) this.messageStore).getBrokerConfig().setEnableSlaveActingMaster(false); + messageStoreConfig.setEnableAutoInSyncReplicas(false); } diff --git a/store/src/test/java/org/apache/rocketmq/store/HATest.java b/store/src/test/java/org/apache/rocketmq/store/HATest.java index bee9d64295..ea7bf725fc 100644 --- a/store/src/test/java/org/apache/rocketmq/store/HATest.java +++ b/store/src/test/java/org/apache/rocketmq/store/HATest.java @@ -146,6 +146,39 @@ public class HATest { //shutdown slave, putMessage should return FLUSH_SLAVE_TIMEOUT slaveMessageStore.shutdown(); + + //wait to let master clean the slave's connection + Thread.sleep(masterMessageStoreConfig.getHaHousekeepingInterval() + 500); + for (long i = 0; i < totalMsgs; i++) { + CompletableFuture putResultFuture = messageStore.asyncPutMessage(buildMessage()); + PutMessageResult result = putResultFuture.get(); + assertEquals(PutMessageStatus.FLUSH_SLAVE_TIMEOUT, result.getPutMessageStatus()); + } + } + + @Test + public void testSemiSyncReplicaWhenSlaveActingMaster() throws Exception { + long totalMsgs = 5; + QUEUE_TOTAL = 1; + MessageBody = StoreMessage.getBytes(); + ((DefaultMessageStore)messageStore).getBrokerConfig().setEnableSlaveActingMaster(true); + for (long i = 0; i < totalMsgs; i++) { + MessageExtBrokerInner msg = buildMessage(); + CompletableFuture putResultFuture = messageStore.asyncPutMessage(msg); + PutMessageResult result = putResultFuture.get(); + assertEquals(PutMessageStatus.PUT_OK, result.getPutMessageStatus()); + //message has been replicated to slave's commitLog, but maybe not dispatch to ConsumeQueue yet + //so direct read from commitLog by physical offset + MessageExt slaveMsg = slaveMessageStore.lookMessageByOffset(result.getAppendMessageResult().getWroteOffset()); + assertNotNull(slaveMsg); + assertTrue(Arrays.equals(msg.getBody(), slaveMsg.getBody())); + assertEquals(msg.getTopic(), slaveMsg.getTopic()); + assertEquals(msg.getTags(), slaveMsg.getTags()); + assertEquals(msg.getKeys(), slaveMsg.getKeys()); + } + + //shutdown slave, putMessage should return IN_SYNC_REPLICAS_NOT_ENOUGH + slaveMessageStore.shutdown(); messageStore.setAliveReplicaNumInGroup(1); //wait to let master clean the slave's connection @@ -155,6 +188,46 @@ public class HATest { PutMessageResult result = putResultFuture.get(); assertEquals(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, result.getPutMessageStatus()); } + + ((DefaultMessageStore)messageStore).getBrokerConfig().setEnableSlaveActingMaster(false); + } + + @Test + public void testSemiSyncReplicaWhenAdaptiveDegradation() throws Exception { + long totalMsgs = 5; + QUEUE_TOTAL = 1; + MessageBody = StoreMessage.getBytes(); + ((DefaultMessageStore)messageStore).getBrokerConfig().setEnableSlaveActingMaster(true); + messageStore.getMessageStoreConfig().setEnableAutoInSyncReplicas(true); + for (long i = 0; i < totalMsgs; i++) { + MessageExtBrokerInner msg = buildMessage(); + CompletableFuture putResultFuture = messageStore.asyncPutMessage(msg); + PutMessageResult result = putResultFuture.get(); + assertEquals(PutMessageStatus.PUT_OK, result.getPutMessageStatus()); + //message has been replicated to slave's commitLog, but maybe not dispatch to ConsumeQueue yet + //so direct read from commitLog by physical offset + MessageExt slaveMsg = slaveMessageStore.lookMessageByOffset(result.getAppendMessageResult().getWroteOffset()); + assertNotNull(slaveMsg); + assertTrue(Arrays.equals(msg.getBody(), slaveMsg.getBody())); + assertEquals(msg.getTopic(), slaveMsg.getTopic()); + assertEquals(msg.getTags(), slaveMsg.getTags()); + assertEquals(msg.getKeys(), slaveMsg.getKeys()); + } + + //shutdown slave, putMessage should return IN_SYNC_REPLICAS_NOT_ENOUGH + slaveMessageStore.shutdown(); + messageStore.setAliveReplicaNumInGroup(1); + + //wait to let master clean the slave's connection + Thread.sleep(masterMessageStoreConfig.getHaHousekeepingInterval() + 500); + for (long i = 0; i < totalMsgs; i++) { + CompletableFuture putResultFuture = messageStore.asyncPutMessage(buildMessage()); + PutMessageResult result = putResultFuture.get(); + assertEquals(PutMessageStatus.PUT_OK, result.getPutMessageStatus()); + } + + ((DefaultMessageStore)messageStore).getBrokerConfig().setEnableSlaveActingMaster(false); + messageStore.getMessageStoreConfig().setEnableAutoInSyncReplicas(false); } @After diff --git a/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java b/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java index 7f6a2b6ee4..d376cd0a12 100644 --- a/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java +++ b/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java @@ -23,12 +23,15 @@ import org.apache.commons.cli.PosixParser; import org.apache.log4j.Logger; import org.apache.rocketmq.client.exception.MQBrokerException; import org.apache.rocketmq.client.exception.MQClientException; +import org.apache.rocketmq.client.impl.factory.MQClientInstance; import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.admin.TopicStatsTable; +import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.common.protocol.body.ClusterInfo; import org.apache.rocketmq.common.protocol.route.BrokerData; import org.apache.rocketmq.common.rpc.ClientMetadata; import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; +import org.apache.rocketmq.common.statictopic.TopicQueueMappingOne; import org.apache.rocketmq.common.statictopic.TopicQueueMappingUtils; import org.apache.rocketmq.common.statictopic.TopicRemappingDetailWrapper; import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; @@ -46,6 +49,8 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ForkJoinPool; +import static org.apache.rocketmq.common.statictopic.TopicQueueMappingUtils.getMappingDetailFromConfig; + public class MQAdminTestUtils { private static Logger log = Logger.getLogger(MQAdminTestUtils.class); @@ -153,6 +158,36 @@ public class MQAdminTestUtils { return false; } + + public static boolean awaitStaticTopicMs(long timeMs, String topic, DefaultMQAdminExt defaultMQAdminExt, MQClientInstance clientInstance) throws Exception { + long start = System.currentTimeMillis(); + while (System.currentTimeMillis() - start <= timeMs) { + if (checkStaticTopic(topic, defaultMQAdminExt, clientInstance)) { + return true; + } + Thread.sleep(100); + } + return false; + } + + //Check if the client metadata is consistent with server metadata + public static boolean checkStaticTopic(String topic, DefaultMQAdminExt defaultMQAdminExt, MQClientInstance clientInstance) throws Exception { + Map brokerConfigMap = MQAdminUtils.examineTopicConfigAll(topic, defaultMQAdminExt); + assert !brokerConfigMap.isEmpty(); + TopicQueueMappingUtils.checkPhysicalQueueConsistence(brokerConfigMap); + TopicQueueMappingUtils.checkNameEpochNumConsistence(topic, brokerConfigMap); + Map globalIdMap = TopicQueueMappingUtils.checkAndBuildMappingItems(getMappingDetailFromConfig(brokerConfigMap.values()), false, true); + for (int i = 0; i < globalIdMap.size(); i++) { + TopicQueueMappingOne mappingOne = globalIdMap.get(i); + String mockBrokerName = TopicQueueMappingUtils.getMockBrokerName(mappingOne.getMappingDetail().getScope()); + String bnameFromRoute = clientInstance.getBrokerNameFromMessageQueue(new MessageQueue(topic, mockBrokerName, mappingOne.getGlobalId())); + if (!mappingOne.getBname().equals(bnameFromRoute)) { + return false; + } + } + return true; + } + //should only be test, if some middle operation failed, it dose not backup the brokerConfigMap public static Map createStaticTopic(String topic, int queueNum, Set targetBrokers, DefaultMQAdminExt defaultMQAdminExt) throws Exception { Map brokerConfigMap = MQAdminUtils.examineTopicConfigAll(topic, defaultMQAdminExt); @@ -172,6 +207,7 @@ public class MQAdminTestUtils { MQAdminUtils.remappingStaticTopic(topic, wrapper.getBrokerToMapIn(), wrapper.getBrokerToMapOut(), brokerConfigMap, TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE, false, defaultMQAdminExt); } + //for test only public static void remappingStaticTopicWithNegativeLogicOffset(String topic, Set targetBrokers, DefaultMQAdminExt defaultMQAdminExt) throws Exception { Map brokerConfigMap = MQAdminUtils.examineTopicConfigAll(topic, defaultMQAdminExt); diff --git a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java index 8420fdd923..4e29c84c6b 100644 --- a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java +++ b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java @@ -116,6 +116,25 @@ public class BaseConf { ForkJoinPool.commonPool().execute(mqAdminExt::shutdown); } + public boolean awaitDispatchMs(long timeMs) throws Exception { + long start = System.currentTimeMillis(); + while (System.currentTimeMillis() - start <= timeMs) { + boolean allOk = true; + for (BrokerController brokerController: brokerControllerList) { + if (brokerController.getMessageStore().dispatchBehindBytes() != 0) { + allOk = false; + break; + } + } + if (allOk) { + return true; + } + Thread.sleep(100); + } + return false; + } + + public static String initTopic() { String topic = MQRandomUtils.getRandomTopic(); return initTopicWithName(topic); diff --git a/test/src/test/java/org/apache/rocketmq/test/container/BrokerFailoverIT.java b/test/src/test/java/org/apache/rocketmq/test/container/BrokerFailoverIT.java index 4c0674e437..0a23601816 100644 --- a/test/src/test/java/org/apache/rocketmq/test/container/BrokerFailoverIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/container/BrokerFailoverIT.java @@ -68,11 +68,6 @@ public class BrokerFailoverIT extends ContainerIntegrationTestBase { RemotingCommand response) { } - - @Override - public void doAfterRpcFailure(String remoteAddr, RemotingCommand request, Boolean remoteTimeout) { - - } }); InnerSalveBrokerController finalTargetSlave = targetSlave; diff --git a/test/src/test/java/org/apache/rocketmq/test/offset/OffsetNotFoundIT.java b/test/src/test/java/org/apache/rocketmq/test/offset/OffsetNotFoundIT.java new file mode 100644 index 0000000000..9d45341b51 --- /dev/null +++ b/test/src/test/java/org/apache/rocketmq/test/offset/OffsetNotFoundIT.java @@ -0,0 +1,132 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.test.offset; + +import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.common.consumer.ConsumeFromWhere; +import org.apache.rocketmq.common.protocol.RequestCode; +import org.apache.rocketmq.remoting.RPCHook; +import org.apache.rocketmq.remoting.protocol.RemotingCommand; +import org.apache.rocketmq.test.base.BaseConf; +import org.apache.rocketmq.test.client.rmq.RMQNormalConsumer; +import org.apache.rocketmq.test.client.rmq.RMQNormalProducer; +import org.apache.rocketmq.test.listener.rmq.concurrent.RMQNormalListener; +import org.apache.rocketmq.test.util.VerifyUtils; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import static com.google.common.truth.Truth.assertThat; + +public class OffsetNotFoundIT extends BaseConf { + + private OffsetRpcHook offsetRpcHook = new OffsetRpcHook(); + + static class OffsetRpcHook implements RPCHook { + + private boolean throwException = false; + + private boolean addSetZeroOfNotFound = false; + + @Override public void doBeforeRequest(String remoteAddr, RemotingCommand request) { + + if (request.getCode() == RequestCode.QUERY_CONSUMER_OFFSET) { + if (throwException) { + throw new RuntimeException("Stop by rpc hook"); + } + if (addSetZeroOfNotFound) { + request.getExtFields().put("setZeroIfNotFound", "false"); + } + } + } + + @Override public void doAfterResponse(String remoteAddr, RemotingCommand request, + RemotingCommand response) { + + } + } + + @Before + public void setUp() { + for (BrokerController brokerController: brokerControllerList) { + brokerController.registerServerRPCHook(offsetRpcHook); + } + + + } + + @After + public void tearDown() { + super.shutdown(); + } + + @Test + public void testConsumeStopAndResume() { + String topic = initTopic(); + RMQNormalProducer producer = getProducer(nsAddr, topic); + int msgSize = 10; + producer.send(msgSize); + Assert.assertEquals("Not all sent succeeded", msgSize, producer.getAllUndupMsgBody().size()); + try { + offsetRpcHook.throwException = true; + RMQNormalConsumer consumer = getConsumer(nsAddr, topic, "*", new RMQNormalListener()); + consumer.getListener().waitForMessageConsume(producer.getAllMsgBody(), 15000); + Assert.assertEquals(0, consumer.getListener().getAllMsgBody().size()); + consumer.shutdown(); + } finally { + offsetRpcHook.throwException = false; + } + //test the normal + RMQNormalConsumer consumer = getConsumer(nsAddr, topic, "*", new RMQNormalListener()); + consumer.getListener().waitForMessageConsume(producer.getAllMsgBody(), 15000); + Assert.assertEquals(producer.getAllMsgBody().size(), consumer.getListener().getAllMsgBody().size()); + assertThat(VerifyUtils.getFilterdMessage(producer.getAllMsgBody(), + consumer.getListener().getAllMsgBody())) + .containsExactlyElementsIn(producer.getAllMsgBody()); + consumer.shutdown(); + } + + + @Test + public void testOffsetNotFoundException() { + String topic = initTopic(); + String group = initConsumerGroup(); + RMQNormalProducer producer = getProducer(nsAddr, topic); + int msgSize = 10; + producer.send(msgSize); + Assert.assertEquals("Not all sent succeeded", msgSize, producer.getAllUndupMsgBody().size()); + try { + offsetRpcHook.addSetZeroOfNotFound = true; + //test the normal + RMQNormalConsumer consumer = new RMQNormalConsumer(nsAddr, topic, "*", group, new RMQNormalListener()); + consumer.create(false); + consumer.getConsumer().setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); + consumer.start(); + consumer.getListener().waitForMessageConsume(producer.getAllMsgBody(), 15000); + Assert.assertEquals(producer.getAllMsgBody().size(), consumer.getListener().getAllMsgBody().size()); + assertThat(VerifyUtils.getFilterdMessage(producer.getAllMsgBody(), + consumer.getListener().getAllMsgBody())) + .containsExactlyElementsIn(producer.getAllMsgBody()); + consumer.shutdown(); + } finally { + offsetRpcHook.addSetZeroOfNotFound = false; + } + + } +} diff --git a/test/src/test/java/org/apache/rocketmq/test/statictopic/StaticTopicIT.java b/test/src/test/java/org/apache/rocketmq/test/statictopic/StaticTopicIT.java index ed93e862c3..41c9c7e4fd 100644 --- a/test/src/test/java/org/apache/rocketmq/test/statictopic/StaticTopicIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/statictopic/StaticTopicIT.java @@ -21,6 +21,9 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableSet; import org.apache.log4j.Logger; import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; +import org.apache.rocketmq.client.impl.factory.MQClientInstance; +import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.admin.ConsumeStats; import org.apache.rocketmq.common.admin.OffsetWrapper; @@ -76,6 +79,7 @@ public class StaticTopicIT extends BaseConf { defaultMQAdminExt.start(); } + @Test public void testCommandsWithCluster() throws Exception { //This case is used to mock the env to test the command manually @@ -93,12 +97,11 @@ public class StaticTopicIT extends BaseConf { } { MQAdminTestUtils.remappingStaticTopicWithCommand(topic, null, clusterName, nsAddr); - Thread.sleep(500); - sendMessagesAndCheck(producer, getBrokers(), topic, queueNum, msgEachQueue, 100); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), consumer.getConsumer(), defaultMQAdminExt); + sendMessagesAndCheck(producer, getBrokers(), topic, queueNum, msgEachQueue, msgEachQueue); } } - @Ignore @Test public void testCommandsWithBrokers() throws Exception { //This case is used to mock the env to test the command manually @@ -117,7 +120,7 @@ public class StaticTopicIT extends BaseConf { { Set brokers = ImmutableSet.of(broker2Name); MQAdminTestUtils.remappingStaticTopicWithCommand(topic, brokers, null, nsAddr); - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), consumer.getConsumer(), defaultMQAdminExt); sendMessagesAndCheck(producer, brokers, topic, queueNum, msgEachQueue, TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE); consumeMessagesAndCheck(producer, consumer, topic, queueNum, msgEachQueue, 0, 2); } @@ -172,7 +175,7 @@ public class StaticTopicIT extends BaseConf { } Assert.assertEquals(0, producer.getSendErrorMsg().size()); //leave the time to build the cq - Thread.sleep(100); + Assert.assertTrue(awaitDispatchMs(500)); for(MessageQueue messageQueue: messageQueueList) { Assert.assertEquals(0, defaultMQAdminExt.minOffset(messageQueue)); Assert.assertEquals(msgEachQueue + baseOffset, defaultMQAdminExt.maxOffset(messageQueue)); @@ -209,6 +212,7 @@ public class StaticTopicIT extends BaseConf { /*System.out.println("produce:" + producer.getAllMsgBody().size()); System.out.println("consume:" + consumer.getListener().getAllMsgBody().size());*/ + Assert.assertEquals(producer.getAllMsgBody().size(), consumer.getListener().getAllMsgBody().size()); assertThat(VerifyUtils.getFilterdMessage(producer.getAllMsgBody(), consumer.getListener().getAllMsgBody())) .containsExactlyElementsIn(producer.getAllMsgBody()); @@ -292,13 +296,48 @@ public class StaticTopicIT extends BaseConf { Assert.assertEquals(broker2Name, mappingOne.getBname()); Assert.assertEquals(TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE, mappingOne.getItems().get(mappingOne.getItems().size() - 1).getLogicOffset()); } - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), consumer.getConsumer(), defaultMQAdminExt); sendMessagesAndCheck(producer, targetBrokers, topic, queueNum, msgEachQueue, TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE); consumeMessagesAndCheck(producer, consumer, topic, queueNum, msgEachQueue, 0, 2); } } + public boolean awaitRefreshStaticTopicMetadata(long timeMs, String topic, DefaultMQProducer producer, DefaultMQPushConsumer consumer, DefaultMQAdminExt adminExt) throws Exception { + long start = System.currentTimeMillis(); + MQClientInstance currentInstance = null; + while (System.currentTimeMillis() - start <= timeMs) { + boolean allOk = true; + if (producer != null) { + currentInstance = producer.getDefaultMQProducerImpl().getmQClientFactory(); + currentInstance.updateTopicRouteInfoFromNameServer(topic); + if (!MQAdminTestUtils.checkStaticTopic(topic, adminExt, currentInstance)) { + allOk = false; + } + } + if (consumer != null) { + currentInstance = consumer.getDefaultMQPushConsumerImpl().getmQClientFactory(); + currentInstance.updateTopicRouteInfoFromNameServer(topic); + if (!MQAdminTestUtils.checkStaticTopic(topic, adminExt, currentInstance)) { + allOk = false; + } + } + if (adminExt != null) { + currentInstance = adminExt.getDefaultMQAdminExtImpl().getMqClientInstance(); + currentInstance.updateTopicRouteInfoFromNameServer(topic); + if (!MQAdminTestUtils.checkStaticTopic(topic, adminExt, currentInstance)) { + allOk = false; + } + } + if (allOk) { + return true; + } + Thread.sleep(100); + } + return false; + } + + @Test public void testDoubleReadCheckConsumerOffset() throws Exception { String topic = "static" + MQRandomUtils.getRandomTopic(); @@ -336,11 +375,11 @@ public class StaticTopicIT extends BaseConf { Set targetBrokers = ImmutableSet.of(brokers.get(i)); MQAdminTestUtils.remappingStaticTopic(topic, targetBrokers, defaultMQAdminExt); //make the metadata - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), null, defaultMQAdminExt); sendMessagesAndCheck(producer, targetBrokers, topic, queueNum, msgEachQueue, (i + 1) * TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE); } - TestUtils.waitForSeconds(20); + TestUtils.waitForSeconds(1); consumeStats = defaultMQAdminExt.examineConsumeStats(group); messageQueues = producer.getMessageQueue(); @@ -367,7 +406,7 @@ public class StaticTopicIT extends BaseConf { Set targetBrokers = ImmutableSet.of(broker1Name); MQAdminTestUtils.createStaticTopic(topic, queueNum, targetBrokers, defaultMQAdminExt); //leave the time to refresh the metadata - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), null, defaultMQAdminExt); sendMessagesAndCheck(producer, targetBrokers, topic, queueNum, msgEachQueue, 0); } @@ -376,7 +415,7 @@ public class StaticTopicIT extends BaseConf { Set targetBrokers = ImmutableSet.of(broker2Name); MQAdminTestUtils.remappingStaticTopic(topic, targetBrokers, defaultMQAdminExt); //leave the time to refresh the metadata - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), null, defaultMQAdminExt); sendMessagesAndCheck(producer, targetBrokers, topic, queueNum, msgEachQueue, 1 * TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE); } @@ -385,7 +424,7 @@ public class StaticTopicIT extends BaseConf { Set targetBrokers = ImmutableSet.of(broker3Name); MQAdminTestUtils.remappingStaticTopic(topic, targetBrokers, defaultMQAdminExt); //leave the time to refresh the metadata - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), null, defaultMQAdminExt); sendMessagesAndCheck(producer, targetBrokers, topic, queueNum, msgEachQueue, 2 * TopicQueueMappingUtils.DEFAULT_BLOCK_SEQ_SIZE); } @@ -467,7 +506,7 @@ public class StaticTopicIT extends BaseConf { Assert.assertEquals(-1, mappingOne.getItems().get(mappingOne.getItems().size() - 1).getLogicOffset()); } //leave the time to refresh the metadata - Thread.sleep(500); + awaitRefreshStaticTopicMetadata(3000, topic, producer.getProducer(), null, defaultMQAdminExt); //here the gen should be 0 sendMessagesAndCheck(producer, targetBrokers, topic, queueNum, msgEachQueue, 0); }