diff --git a/acl/src/test/java/org/apache/rocketmq/acl/plain/AclTestHelper.java b/acl/src/test/java/org/apache/rocketmq/acl/plain/AclTestHelper.java
index 378d24bddd..250478a5b8 100644
--- a/acl/src/test/java/org/apache/rocketmq/acl/plain/AclTestHelper.java
+++ b/acl/src/test/java/org/apache/rocketmq/acl/plain/AclTestHelper.java
@@ -80,16 +80,16 @@ public final class AclTestHelper {
public static void recursiveDelete(File file) {
if (file.isFile()) {
- Assert.assertTrue(file.delete());
- return;
- }
- File[] files = file.listFiles();
- if (null != files) {
+ file.delete();
+ } else {
+ File[] files = file.listFiles();
+ if (null != files) {
for (File f : files) {
- recursiveDelete(f);
+ recursiveDelete(f);
}
+ }
+ file.delete();
}
- Assert.assertTrue(file.delete());
}
public static File copyResources(String folder) throws IOException {
diff --git a/broker/pom.xml b/broker/pom.xml
index aa4a7cc6a3..9a9abad145 100644
--- a/broker/pom.xml
+++ b/broker/pom.xml
@@ -94,6 +94,27 @@
false
+
+
+ maven-checkstyle-plugin
+ ${maven-checkstyle-plugin.version}
+
+
+ validate
+ validate
+
+ ${project.parent.basedir}/style/rmq_checkstyle.xml
+ UTF-8
+ true
+ true
+ true
+
+
+ check
+
+
+
+
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/BrokerOuterAPITest.java b/broker/src/test/java/org/apache/rocketmq/broker/BrokerOuterAPITest.java
index ffb1d95224..a353a7ad31 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/BrokerOuterAPITest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/BrokerOuterAPITest.java
@@ -197,14 +197,14 @@ public class BrokerOuterAPITest {
final ArgumentCaptor namesrvCaptor = ArgumentCaptor.forClass(String.class);
when(nettyRemotingClient.invokeSync(namesrvCaptor.capture(), any(RemotingCommand.class),
timeoutMillisCaptor.capture())).thenAnswer((Answer) invocation -> {
- final String namesrv = namesrvCaptor.getValue();
- if (nameserver1.equals(namesrv) || nameserver2.equals(namesrv)) {
+ final String namesrv = namesrvCaptor.getValue();
+ if (nameserver1.equals(namesrv) || nameserver2.equals(namesrv)) {
+ return response;
+ }
+ long delayTimeMillis = 1000;
+ TimeUnit.MILLISECONDS.sleep(timeoutMillisCaptor.getValue() + delayTimeMillis);
return response;
- }
- long delayTimeMillis = 1000;
- TimeUnit.MILLISECONDS.sleep(timeoutMillisCaptor.getValue() + delayTimeMillis);
- return response;
- });
+ });
List registerBrokerResultList = brokerOuterAPI.registerBrokerAll(clusterName, brokerAddr, brokerName, brokerId, "hasServerAddr", topicConfigSerializeWrapper, Lists.newArrayList(), false, timeOut, false, true, new BrokerIdentity());
assertEquals(2, registerBrokerResultList.size());
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java
index fd76312941..dac5468c87 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java
@@ -69,8 +69,8 @@ public class ProducerManagerTest {
assertThat(producerManager.findChannel("clientId")).isNotNull();
Field field = ProducerManager.class.getDeclaredField("CHANNEL_EXPIRED_TIMEOUT");
field.setAccessible(true);
- long CHANNEL_EXPIRED_TIMEOUT = field.getLong(producerManager);
- clientInfo.setLastUpdateTimestamp(System.currentTimeMillis() - CHANNEL_EXPIRED_TIMEOUT - 10);
+ long channelExpiredTimeout = field.getLong(producerManager);
+ clientInfo.setLastUpdateTimestamp(System.currentTimeMillis() - channelExpiredTimeout - 10);
when(channel.close()).thenReturn(mock(ChannelFuture.class));
producerManager.scanNotActiveChannel();
assertThat(producerManager.getGroupChannelTable().get(group)).isNull();
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/controller/ReplicasManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/controller/ReplicasManagerTest.java
index b7ab79edaf..85d11508d2 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/controller/ReplicasManagerTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/controller/ReplicasManagerTest.java
@@ -139,7 +139,7 @@ public class ReplicasManagerTest {
}
@Test
- public void changeBrokerRoleTest(){
+ public void changeBrokerRoleTest() {
// not equal to localAddress
Assertions.assertThatCode(() -> replicasManager.changeBrokerRole(NEW_MASTER_ADDRESS, NEW_MASTER_EPOCH, OLD_MASTER_EPOCH, SLAVE_BROKER_ID))
.doesNotThrowAnyException();
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
index 68d60092d8..a67ec7a6de 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/filter/ConsumerFilterManagerTest.java
@@ -81,9 +81,7 @@ public class ConsumerFilterManagerTest {
public void testRegister_change() {
ConsumerFilterManager filterManager = gen(10, 10);
- ConsumerFilterData filterData = filterManager.get("topic9", "CID_9");
-
- System.out.println(filterData.getCompiledExpression());
+ ConsumerFilterData filterData;
String newExpr = "a > 0 and a < 10";
@@ -92,8 +90,6 @@ public class ConsumerFilterManagerTest {
filterData = filterManager.get("topic9", "CID_9");
assertThat(newExpr).isEqualTo(filterData.getExpression());
-
- System.out.println(filterData.toString());
}
@Test
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/filter/MessageStoreWithFilterTest.java b/broker/src/test/java/org/apache/rocketmq/broker/filter/MessageStoreWithFilterTest.java
index b14005942c..8c49581242 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/filter/MessageStoreWithFilterTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/filter/MessageStoreWithFilterTest.java
@@ -56,19 +56,19 @@ import static org.awaitility.Awaitility.await;
public class MessageStoreWithFilterTest {
- private static final String msg = "Once, there was a chance for me!";
- private static final byte[] msgBody = msg.getBytes();
+ private static final String MSG = "Once, there was a chance for me!";
+ private static final byte[] MSG_BODY = MSG.getBytes();
- private static final String topic = "topic";
- private static final int queueId = 0;
- private static final String storePath = System.getProperty("java.io.tmpdir") + File.separator + "unit_test_store";
- private static final int commitLogFileSize = 1024 * 1024 * 256;
- private static final int cqFileSize = 300000 * 20;
- private static final int cqExtFileSize = 300000 * 128;
+ private static final String TOPIC = "topic";
+ private static final int QUEUE_ID = 0;
+ private static final String STORE_PATH = System.getProperty("java.io.tmpdir") + File.separator + "unit_test_store";
+ private static final int COMMIT_LOG_FILE_SIZE = 1024 * 1024 * 256;
+ private static final int CQ_FILE_SIZE = 300000 * 20;
+ private static final int CQ_EXT_FILE_SIZE = 300000 * 128;
- private static SocketAddress BornHost;
+ private static SocketAddress bornHost;
- private static SocketAddress StoreHost;
+ private static SocketAddress storeHost;
private DefaultMessageStore master;
@@ -80,11 +80,11 @@ public class MessageStoreWithFilterTest {
static {
try {
- StoreHost = new InetSocketAddress(InetAddress.getLocalHost(), 8123);
+ storeHost = new InetSocketAddress(InetAddress.getLocalHost(), 8123);
} catch (UnknownHostException e) {
}
try {
- BornHost = new InetSocketAddress(InetAddress.getByName("127.0.0.1"), 0);
+ bornHost = new InetSocketAddress(InetAddress.getByName("127.0.0.1"), 0);
} catch (UnknownHostException e) {
}
}
@@ -101,21 +101,21 @@ public class MessageStoreWithFilterTest {
master.shutdown();
master.destroy();
}
- UtilAll.deleteFile(new File(storePath));
+ UtilAll.deleteFile(new File(STORE_PATH));
}
public MessageExtBrokerInner buildMessage() {
MessageExtBrokerInner msg = new MessageExtBrokerInner();
- msg.setTopic(topic);
+ msg.setTopic(TOPIC);
msg.setTags(System.currentTimeMillis() + "TAG");
msg.setKeys("Hello");
- msg.setBody(msgBody);
+ msg.setBody(MSG_BODY);
msg.setKeys(String.valueOf(System.currentTimeMillis()));
- msg.setQueueId(queueId);
+ msg.setQueueId(QUEUE_ID);
msg.setSysFlag(0);
msg.setBornTimestamp(System.currentTimeMillis());
- msg.setStoreHost(StoreHost);
- msg.setBornHost(BornHost);
+ msg.setStoreHost(storeHost);
+ msg.setBornHost(bornHost);
for (int i = 1; i < 3; i++) {
msg.putUserProperty(String.valueOf(i), "imagoodperson" + i);
}
@@ -133,15 +133,15 @@ public class MessageStoreWithFilterTest {
messageStoreConfig.setMessageIndexEnable(false);
messageStoreConfig.setEnableConsumeQueueExt(enableCqExt);
- messageStoreConfig.setStorePathRootDir(storePath);
- messageStoreConfig.setStorePathCommitLog(storePath + File.separator + "commitlog");
+ messageStoreConfig.setStorePathRootDir(STORE_PATH);
+ messageStoreConfig.setStorePathCommitLog(STORE_PATH + File.separator + "commitlog");
return messageStoreConfig;
}
protected DefaultMessageStore gen(ConsumerFilterManager filterManager) throws Exception {
MessageStoreConfig messageStoreConfig = buildStoreConfig(
- commitLogFileSize, cqFileSize, true, cqExtFileSize
+ COMMIT_LOG_FILE_SIZE, CQ_FILE_SIZE, true, CQ_EXT_FILE_SIZE
);
BrokerConfig brokerConfig = new BrokerConfig();
@@ -182,7 +182,7 @@ public class MessageStoreWithFilterTest {
int msgCountPerTopic) throws Exception {
List msgs = new ArrayList();
for (int i = 0; i < topicCount; i++) {
- String realTopic = topic + i;
+ String realTopic = TOPIC + i;
for (int j = 0; j < msgCountPerTopic; j++) {
MessageExtBrokerInner msg = buildMessage();
msg.setTopic(realTopic);
@@ -247,7 +247,7 @@ public class MessageStoreWithFilterTest {
resetGroup, resetSubData.getSubString(), resetSubData.getExpressionType(),
System.currentTimeMillis());
- GetMessageResult resetGetResult = master.getMessage(resetGroup, topic, queueId, 0, 1000,
+ GetMessageResult resetGetResult = master.getMessage(resetGroup, topic, QUEUE_ID, 0, 1000,
new ExpressionMessageFilter(resetSubData, resetFilterData, filterManager));
try {
@@ -274,7 +274,7 @@ public class MessageStoreWithFilterTest {
List filteredMsgs = filtered(msgs, normalFilterData);
- GetMessageResult normalGetResult = master.getMessage(normalGroup, topic, queueId, 0, 1000,
+ GetMessageResult normalGetResult = master.getMessage(normalGroup, topic, QUEUE_ID, 0, 1000,
new ExpressionMessageFilter(normalSubData, normalFilterData, filterManager));
try {
@@ -293,7 +293,7 @@ public class MessageStoreWithFilterTest {
Thread.sleep(100);
for (int i = 0; i < topicCount; i++) {
- String realTopic = topic + i;
+ String realTopic = TOPIC + i;
for (int j = 0; j < msgPerTopic; j++) {
String group = "CID_" + j;
@@ -309,7 +309,7 @@ public class MessageStoreWithFilterTest {
subscriptionData.setClassFilterMode(false);
subscriptionData.setSubString(filterData.getExpression());
- GetMessageResult getMessageResult = master.getMessage(group, realTopic, queueId, 0, 10000,
+ GetMessageResult getMessageResult = master.getMessage(group, realTopic, QUEUE_ID, 0, 10000,
new ExpressionMessageFilter(subscriptionData, filterData, filterManager));
String assertMsg = group + "-" + realTopic;
try {
@@ -356,8 +356,8 @@ public class MessageStoreWithFilterTest {
@Override
public void run() throws Throwable {
for (int i = 0; i < topicCount; i++) {
- final String realTopic = topic + i;
- GetMessageResult getMessageResult = master.getMessage("test", realTopic, queueId, 0, 10000,
+ final String realTopic = TOPIC + i;
+ GetMessageResult getMessageResult = master.getMessage("test", realTopic, QUEUE_ID, 0, 10000,
new MessageFilter() {
@Override
public boolean isMatchedByConsumeQueue(Long tagsCode,
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java
index 61dd6693ef..1df1adb1f9 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java
@@ -29,27 +29,27 @@ public class ConsumerOffsetManagerTest {
private ConsumerOffsetManager consumerOffsetManager;
- private static final String key = "FooBar@FooBarGroup";
+ private static final String KEY = "FooBar@FooBarGroup";
@Before
- public void init(){
+ public void init() {
consumerOffsetManager = new ConsumerOffsetManager();
ConcurrentHashMap> offsetTable = new ConcurrentHashMap>(512);
- offsetTable.put(key,new ConcurrentHashMap(){{
- put(1,2L);
- put(2,3L);
- }});
+ offsetTable.put(KEY,new ConcurrentHashMap() {{
+ put(1,2L);
+ put(2,3L);
+ }});
consumerOffsetManager.setOffsetTable(offsetTable);
}
@Test
- public void cleanOffsetByTopic_NotExist(){
+ public void cleanOffsetByTopic_NotExist() {
consumerOffsetManager.cleanOffsetByTopic("InvalidTopic");
- assertThat(consumerOffsetManager.getOffsetTable().containsKey(key)).isTrue();
+ assertThat(consumerOffsetManager.getOffsetTable().containsKey(KEY)).isTrue();
}
@Test
- public void cleanOffsetByTopic_Exist(){
+ public void cleanOffsetByTopic_Exist() {
consumerOffsetManager.cleanOffsetByTopic("FooBar");
- assertThat(!consumerOffsetManager.getOffsetTable().containsKey(key)).isTrue();
+ assertThat(!consumerOffsetManager.getOffsetTable().containsKey(KEY)).isTrue();
}
}
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/pagecache/ManyMessageTransferTest.java b/broker/src/test/java/org/apache/rocketmq/broker/pagecache/ManyMessageTransferTest.java
index 508635c044..2617b5cee9 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/pagecache/ManyMessageTransferTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/pagecache/ManyMessageTransferTest.java
@@ -25,7 +25,7 @@ import org.junit.Test;
public class ManyMessageTransferTest {
@Test
- public void ManyMessageTransferBuilderTest(){
+ public void ManyMessageTransferBuilderTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
GetMessageResult getMessageResult = new GetMessageResult();
@@ -33,7 +33,7 @@ public class ManyMessageTransferTest {
}
@Test
- public void ManyMessageTransferPosTest(){
+ public void ManyMessageTransferPosTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
GetMessageResult getMessageResult = new GetMessageResult();
@@ -42,7 +42,7 @@ public class ManyMessageTransferTest {
}
@Test
- public void ManyMessageTransferCountTest(){
+ public void ManyMessageTransferCountTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
GetMessageResult getMessageResult = new GetMessageResult();
@@ -53,7 +53,7 @@ public class ManyMessageTransferTest {
}
@Test
- public void ManyMessageTransferCloseTest(){
+ public void ManyMessageTransferCloseTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
GetMessageResult getMessageResult = new GetMessageResult();
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/pagecache/OneMessageTransferTest.java b/broker/src/test/java/org/apache/rocketmq/broker/pagecache/OneMessageTransferTest.java
index da705843d0..1930641d7b 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/pagecache/OneMessageTransferTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/pagecache/OneMessageTransferTest.java
@@ -26,7 +26,7 @@ import org.junit.Test;
public class OneMessageTransferTest {
@Test
- public void OneMessageTransferTest(){
+ public void OneMessageTransferTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
SelectMappedBufferResult selectMappedBufferResult = new SelectMappedBufferResult(0,byteBuffer,20,new DefaultMappedFile());
@@ -34,7 +34,7 @@ public class OneMessageTransferTest {
}
@Test
- public void OneMessageTransferCountTest(){
+ public void OneMessageTransferCountTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
SelectMappedBufferResult selectMappedBufferResult = new SelectMappedBufferResult(0,byteBuffer,20,new DefaultMappedFile());
@@ -43,7 +43,7 @@ public class OneMessageTransferTest {
}
@Test
- public void OneMessageTransferPosTest(){
+ public void OneMessageTransferPosTest() {
ByteBuffer byteBuffer = ByteBuffer.allocate(20);
byteBuffer.putInt(20);
SelectMappedBufferResult selectMappedBufferResult = new SelectMappedBufferResult(0,byteBuffer,20,new DefaultMappedFile());
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
index 3d70247739..44eefbff9e 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java
@@ -103,7 +103,7 @@ public class AdminBrokerProcessorTest {
@Spy
private BrokerController
brokerController = new BrokerController(new BrokerConfig(), new NettyServerConfig(), new NettyClientConfig(),
- new MessageStoreConfig());
+ new MessageStoreConfig());
@Mock
private MessageStore messageStore;
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java
index b81fbae31a..0e7b3ced7f 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/EndTransactionProcessorTest.java
@@ -61,7 +61,7 @@ public class EndTransactionProcessorTest {
@Spy
private BrokerController
brokerController = new BrokerController(new BrokerConfig(), new NettyServerConfig(), new NettyClientConfig(),
- new MessageStoreConfig());
+ new MessageStoreConfig());
@Mock
private MessageStore messageStore;
@@ -76,7 +76,7 @@ public class EndTransactionProcessorTest {
endTransactionProcessor = new EndTransactionProcessor(brokerController);
}
- private OperationResult createResponse(int status){
+ private OperationResult createResponse(int status) {
OperationResult response = new OperationResult();
response.setPrepareMessage(createDefaultMessageExt());
response.setResponseCode(status);
@@ -87,8 +87,8 @@ public class EndTransactionProcessorTest {
@Test
public void testProcessRequest() throws RemotingCommandException {
when(transactionMsgService.commitMessage(any(EndTransactionRequestHeader.class))).thenReturn(createResponse(ResponseCode.SUCCESS));
- when(messageStore.putMessage(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(messageStore.putMessage(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
RemotingCommand request = createEndTransactionMsgCommand(MessageSysFlag.TRANSACTION_COMMIT_TYPE, false);
RemotingCommand response = endTransactionProcessor.processRequest(handlerContext, request);
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
@@ -97,8 +97,8 @@ public class EndTransactionProcessorTest {
@Test
public void testProcessRequest_CheckMessage() throws RemotingCommandException {
when(transactionMsgService.commitMessage(any(EndTransactionRequestHeader.class))).thenReturn(createResponse(ResponseCode.SUCCESS));
- when(messageStore.putMessage(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(messageStore.putMessage(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
RemotingCommand request = createEndTransactionMsgCommand(MessageSysFlag.TRANSACTION_COMMIT_TYPE, true);
RemotingCommand response = endTransactionProcessor.processRequest(handlerContext, request);
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopMessageProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopMessageProcessorTest.java
index 7ea20ceffe..bbbcf4f867 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopMessageProcessorTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopMessageProcessorTest.java
@@ -34,13 +34,9 @@ import org.apache.rocketmq.remoting.exception.RemotingCommandException;
import org.apache.rocketmq.remoting.netty.NettyClientConfig;
import org.apache.rocketmq.remoting.netty.NettyServerConfig;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
-import org.apache.rocketmq.store.AppendMessageResult;
-import org.apache.rocketmq.store.AppendMessageStatus;
import org.apache.rocketmq.store.DefaultMessageStore;
import org.apache.rocketmq.store.GetMessageResult;
import org.apache.rocketmq.store.GetMessageStatus;
-import org.apache.rocketmq.store.PutMessageResult;
-import org.apache.rocketmq.store.PutMessageStatus;
import org.apache.rocketmq.store.SelectMappedBufferResult;
import org.apache.rocketmq.store.config.MessageStoreConfig;
import org.apache.rocketmq.store.logfile.DefaultMappedFile;
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java
index e9af449cce..d839a22e85 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java
@@ -32,7 +32,6 @@ import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig;
import org.apache.rocketmq.common.utils.DataConverter;
import org.apache.rocketmq.remoting.common.RemotingUtil;
import org.apache.rocketmq.store.MessageStore;
-import org.apache.rocketmq.store.config.BrokerRole;
import org.apache.rocketmq.store.config.MessageStoreConfig;
import org.apache.rocketmq.store.pop.AckMsg;
import org.apache.rocketmq.store.pop.PopCheckPoint;
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/SendMessageProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/SendMessageProcessorTest.java
index f9dc1071b8..0822670895 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/processor/SendMessageProcessorTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/SendMessageProcessorTest.java
@@ -140,7 +140,6 @@ public class SendMessageProcessorTest {
sendMessageHookList.add(sendMessageHook);
sendMessageProcessor.registerSendMessageHook(sendMessageHookList);
assertPutResult(ResponseCode.SUCCESS);
- System.out.println(sendMessageContext[0]);
assertThat(sendMessageContext[0]).isNotNull();
assertThat(sendMessageContext[0].getTopic()).isEqualTo(topic);
assertThat(sendMessageContext[0].getProducerGroup()).isEqualTo(group);
@@ -268,7 +267,6 @@ public class SendMessageProcessorTest {
sendMessageHookList.add(sendMessageHook);
sendMessageProcessor.registerSendMessageHook(sendMessageHookList);
assertPutResult(ResponseCode.FLOW_CONTROL);
- System.out.println(sendMessageContext[0]);
assertThat(sendMessageContext[0]).isNotNull();
assertThat(sendMessageContext[0].getTopic()).isEqualTo(topic);
assertThat(sendMessageContext[0].getProducerGroup()).isEqualTo(group);
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/schedule/ScheduleMessageServiceTest.java b/broker/src/test/java/org/apache/rocketmq/broker/schedule/ScheduleMessageServiceTest.java
index d68797f8f6..1e3ee5cc1b 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/schedule/ScheduleMessageServiceTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/schedule/ScheduleMessageServiceTest.java
@@ -57,7 +57,8 @@ import static org.apache.rocketmq.common.stats.Stats.BROKER_PUT_NUMS;
import static org.apache.rocketmq.common.stats.Stats.TOPIC_PUT_NUMS;
import static org.apache.rocketmq.common.stats.Stats.TOPIC_PUT_SIZE;
import static org.assertj.core.api.Assertions.assertThat;
-import static org.junit.Assert.*;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.assertEquals;
public class ScheduleMessageServiceTest {
@@ -73,10 +74,10 @@ public class ScheduleMessageServiceTest {
*/
int delayLevel = 3;
- private static final String storePath = System.getProperty("java.io.tmpdir") + File.separator + "schedule_test#" + UUID.randomUUID();
- private static final int commitLogFileSize = 1024;
- private static final int cqFileSize = 10;
- private static final int cqExtFileSize = 10 * (ConsumeQueueExt.CqExtUnit.MIN_EXT_UNIT_SIZE + 64);
+ private static final String STORE_PATH = System.getProperty("java.io.tmpdir") + File.separator + "schedule_test#" + UUID.randomUUID();
+ private static final int COMMIT_LOG_FILE_SIZE = 1024;
+ private static final int CQ_FILE_SIZE = 10;
+ private static final int CQ_EXT_FILE_SIZE = 10 * (ConsumeQueueExt.CqExtUnit.MIN_EXT_UNIT_SIZE + 64);
private static SocketAddress bornHost;
private static SocketAddress storeHost;
@@ -106,13 +107,13 @@ public class ScheduleMessageServiceTest {
public void setUp() throws Exception {
messageStoreConfig = new MessageStoreConfig();
messageStoreConfig.setMessageDelayLevel(testMessageDelayLevel);
- messageStoreConfig.setMappedFileSizeCommitLog(commitLogFileSize);
- messageStoreConfig.setMappedFileSizeConsumeQueue(cqFileSize);
- messageStoreConfig.setMappedFileSizeConsumeQueueExt(cqExtFileSize);
+ messageStoreConfig.setMappedFileSizeCommitLog(COMMIT_LOG_FILE_SIZE);
+ messageStoreConfig.setMappedFileSizeConsumeQueue(CQ_FILE_SIZE);
+ messageStoreConfig.setMappedFileSizeConsumeQueueExt(CQ_EXT_FILE_SIZE);
messageStoreConfig.setMessageIndexEnable(false);
messageStoreConfig.setEnableConsumeQueueExt(true);
- messageStoreConfig.setStorePathRootDir(storePath);
- messageStoreConfig.setStorePathCommitLog(storePath + File.separator + "commitlog");
+ messageStoreConfig.setStorePathRootDir(STORE_PATH);
+ messageStoreConfig.setStorePathCommitLog(STORE_PATH + File.separator + "commitlog");
// Let OS pick an available port
messageStoreConfig.setHaListenPort(0);
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/topic/TopicQueueMappingManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/topic/TopicQueueMappingManagerTest.java
index 6b4faab5d5..edac5c2396 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/topic/TopicQueueMappingManagerTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/topic/TopicQueueMappingManagerTest.java
@@ -44,12 +44,12 @@ import static org.mockito.Mockito.when;
public class TopicQueueMappingManagerTest {
@Mock
private BrokerController brokerController;
- private static final String broker1Name = "broker1";
+ private static final String BROKER1_NAME = "broker1";
@Before
public void before() {
BrokerConfig brokerConfig = new BrokerConfig();
- brokerConfig.setBrokerName(broker1Name);
+ brokerConfig.setBrokerName(BROKER1_NAME);
when(brokerController.getBrokerConfig()).thenReturn(brokerConfig);
MessageStoreConfig messageStoreConfig = new MessageStoreConfig();
@@ -74,7 +74,7 @@ public class TopicQueueMappingManagerTest {
Map mappingDetailMap = new HashMap<>();
TopicQueueMappingManager topicQueueMappingManager = null;
Set brokers = new HashSet();
- brokers.add(broker1Name);
+ brokers.add(BROKER1_NAME);
{
for (int i = 0; i < 10; i++) {
String topic = UUID.randomUUID().toString();
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageBridgeTest.java b/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageBridgeTest.java
index d7dc98ed9b..6014ce966a 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageBridgeTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageBridgeTest.java
@@ -82,8 +82,8 @@ public class TransactionalMessageBridgeTest {
@Test
public void testPutHalfMessage() {
- when(messageStore.putMessage(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(messageStore.putMessage(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
PutMessageResult result = transactionBridge.putHalfMessage(createMessageBrokerInner());
assertThat(result.getPutMessageStatus()).isEqualTo(PutMessageStatus.PUT_OK);
}
@@ -133,16 +133,16 @@ public class TransactionalMessageBridgeTest {
@Test
public void testPutMessageReturnResult() {
- when(messageStore.putMessage(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(messageStore.putMessage(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
PutMessageResult result = transactionBridge.putMessageReturnResult(createMessageBrokerInner());
assertThat(result.getPutMessageStatus()).isEqualTo(PutMessageStatus.PUT_OK);
}
@Test
public void testPutMessage() {
- when(messageStore.putMessage(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(messageStore.putMessage(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
Boolean success = transactionBridge.putMessage(createMessageBrokerInner());
assertThat(success).isEqualTo(true);
}
@@ -166,11 +166,11 @@ public class TransactionalMessageBridgeTest {
MessageExt messageExt = new MessageExt();
long bornTimeStamp = messageExt.getBornTimestamp();
MessageExt messageExtRes = transactionBridge.renewHalfMessageInner(messageExt);
- assertThat( messageExtRes.getBornTimestamp()).isEqualTo(bornTimeStamp);
+ assertThat(messageExtRes.getBornTimestamp()).isEqualTo(bornTimeStamp);
}
@Test
- public void testLookMessageByOffset(){
+ public void testLookMessageByOffset() {
when(messageStore.lookMessageByOffset(anyLong())).thenReturn(new MessageExt());
MessageExt messageExt = transactionBridge.lookMessageByOffset(123);
assertThat(messageExt).isNotNull();
diff --git a/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java b/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java
index 5c32b21182..aa1c60e0db 100644
--- a/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java
+++ b/broker/src/test/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageServiceImplTest.java
@@ -86,8 +86,8 @@ public class TransactionalMessageServiceImplTest {
@Test
public void testPrepareMessage() {
MessageExtBrokerInner inner = createMessageBrokerInner();
- when(bridge.putHalfMessage(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(bridge.putHalfMessage(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
PutMessageResult result = queueTransactionMsgService.prepareMessage(inner);
assert result.isOk();
}
@@ -134,8 +134,8 @@ public class TransactionalMessageServiceImplTest {
when(bridge.getOpMessage(anyInt(), anyLong(), anyInt())).thenReturn(createPullResult(TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC, 1, "5", 0));
when(bridge.getBrokerController()).thenReturn(this.brokerController);
when(bridge.renewHalfMessageInner(any(MessageExtBrokerInner.class))).thenReturn(createMessageBrokerInner());
- when(bridge.putMessageReturnResult(any(MessageExtBrokerInner.class))).thenReturn(new PutMessageResult
- (PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
+ when(bridge.putMessageReturnResult(any(MessageExtBrokerInner.class)))
+ .thenReturn(new PutMessageResult(PutMessageStatus.PUT_OK, new AppendMessageResult(AppendMessageStatus.PUT_OK)));
long timeOut = this.brokerController.getBrokerConfig().getTransactionTimeOut();
final int checkMax = this.brokerController.getBrokerConfig().getTransactionCheckMax();
final AtomicInteger checkMessage = new AtomicInteger(0);
diff --git a/remoting/BUILD.bazel b/remoting/BUILD.bazel
index 0691bc7d1d..a7ab0df8f3 100644
--- a/remoting/BUILD.bazel
+++ b/remoting/BUILD.bazel
@@ -24,6 +24,7 @@ java_library(
"//logging",
"@maven//:com_alibaba_fastjson",
"@maven//:io_netty_netty_all",
+ "@maven//:org_apache_commons_commons_lang3",
],
)
@@ -37,6 +38,7 @@ java_library(
"@maven//:io_netty_netty_all",
"@maven//:com_google_code_gson_gson",
"@maven//:com_alibaba_fastjson",
+ "@maven//:org_apache_commons_commons_lang3",
],
resources = glob(["src/test/resources/certs/*.pem"]) + glob(["src/test/resources/certs/*.key"])
)