Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
import org.apache.rocketmq.remoting.protocol.route.QueueData;
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import org.apache.rocketmq.remoting.protocol.statictopic.TopicQueueMappingInfo;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
Expand Down Expand Up @@ -128,6 +129,13 @@ public void init() throws Exception {
FieldUtils.writeDeclaredField(mqClientInstance, "topicRouteTable", topicRouteTable, true);
}

@After
public void tearDown() throws Exception {
brokerAddrTable.clear();
consumerTable.clear();
topicRouteTable.clear();
}

@Test
public void testFindBrokerAddressInSubscribe() {
// dledger normal case
Expand Down Expand Up @@ -229,7 +237,7 @@ public void testTopicRouteData2TopicPublishInfo() {
@Test
public void testTopicRouteData2TopicPublishInfoWithOrderTopicConf() {
TopicRouteData topicRouteData = createTopicRouteData();
when(topicRouteData.getOrderTopicConf()).thenReturn("127.0.0.1:4");
topicRouteData.setOrderTopicConf("127.0.0.1:4");
TopicPublishInfo actual = MQClientInstance.topicRouteData2TopicPublishInfo(topic, topicRouteData);
assertFalse(actual.isHaveTopicRouterInfo());
assertEquals(4, actual.getMessageQueueList().size());
Expand All @@ -238,7 +246,7 @@ public void testTopicRouteData2TopicPublishInfoWithOrderTopicConf() {
@Test
public void testTopicRouteData2TopicPublishInfoWithTopicQueueMappingByBroker() {
TopicRouteData topicRouteData = createTopicRouteData();
when(topicRouteData.getTopicQueueMappingByBroker()).thenReturn(Collections.singletonMap(topic, new TopicQueueMappingInfo()));
topicRouteData.setTopicQueueMappingByBroker(Collections.singletonMap(topic, new TopicQueueMappingInfo()));
TopicPublishInfo actual = MQClientInstance.topicRouteData2TopicPublishInfo(topic, topicRouteData);
assertFalse(actual.isHaveTopicRouterInfo());
assertEquals(0, actual.getMessageQueueList().size());
Expand All @@ -247,7 +255,7 @@ public void testTopicRouteData2TopicPublishInfoWithTopicQueueMappingByBroker() {
@Test
public void testTopicRouteData2TopicSubscribeInfo() {
TopicRouteData topicRouteData = createTopicRouteData();
when(topicRouteData.getTopicQueueMappingByBroker()).thenReturn(Collections.singletonMap(topic, new TopicQueueMappingInfo()));
topicRouteData.setTopicQueueMappingByBroker(Collections.singletonMap(topic, new TopicQueueMappingInfo()));
Set<MessageQueue> actual = MQClientInstance.topicRouteData2TopicSubscribeInfo(topic, topicRouteData);
assertNotNull(actual);
assertEquals(0, actual.size());
Expand Down Expand Up @@ -320,7 +328,8 @@ public void testUpdateTopicRouteInfoFromNameServer() throws RemotingException, I
DefaultMQProducer defaultMQProducer = mock(DefaultMQProducer.class);
TopicRouteData topicRouteData = createTopicRouteData();
when(mQClientAPIImpl.getDefaultTopicRouteInfoFromNameServer(anyLong())).thenReturn(topicRouteData);
assertFalse(mqClientInstance.updateTopicRouteInfoFromNameServer(topic, true, defaultMQProducer));
assertTrue(mqClientInstance.updateTopicRouteInfoFromNameServer(topic, true, defaultMQProducer));
assertEquals(topicRouteData, topicRouteTable.get(topic));
}

@Test
Expand Down Expand Up @@ -450,9 +459,9 @@ private MessageQueue createMessageQueue() {
}

private TopicRouteData createTopicRouteData() {
TopicRouteData result = mock(TopicRouteData.class);
when(result.getBrokerDatas()).thenReturn(createBrokerDatas());
when(result.getQueueDatas()).thenReturn(createQueueDatas());
TopicRouteData result = new TopicRouteData();
result.setBrokerDatas(createBrokerDatas());
result.setQueueDatas(createQueueDatas());
return result;
}

Expand Down
Loading