mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
This commit is contained in:
@@ -57,6 +57,12 @@ public class RecallMessageProcessor implements NettyRequestProcessor {
|
||||
final RecallMessageRequestHeader requestHeader =
|
||||
request.decodeCommandCustomHeader(RecallMessageRequestHeader.class);
|
||||
|
||||
if (!brokerController.getBrokerConfig().isRecallMessageEnable()) {
|
||||
response.setCode(ResponseCode.NO_PERMISSION);
|
||||
response.setRemark("recall failed, operation is forbidden");
|
||||
return response;
|
||||
}
|
||||
|
||||
if (BrokerRole.SLAVE == brokerController.getMessageStoreConfig().getBrokerRole()) {
|
||||
response.setCode(ResponseCode.SLAVE_NOT_AVAILABLE);
|
||||
response.setRemark("recall failed, broker service not available");
|
||||
|
||||
+9
@@ -89,6 +89,7 @@ public class RecallMessageProcessorTest {
|
||||
when(brokerController.getMessageStore()).thenReturn(messageStore);
|
||||
when(brokerController.getBrokerConfig()).thenReturn(brokerConfig);
|
||||
when(brokerConfig.getBrokerName()).thenReturn(BROKER_NAME);
|
||||
when(brokerConfig.isRecallMessageEnable()).thenReturn(true);
|
||||
when(brokerController.getBrokerStatsManager()).thenReturn(brokerStatsManager);
|
||||
when(handlerContext.channel()).thenReturn(channel);
|
||||
recallMessageProcessor = new RecallMessageProcessor(brokerController);
|
||||
@@ -134,6 +135,14 @@ public class RecallMessageProcessorTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testProcessRequest_notEnable() throws RemotingCommandException {
|
||||
when(brokerConfig.isRecallMessageEnable()).thenReturn(false);
|
||||
RemotingCommand request = mockRequest(0, TOPIC, TOPIC, "id", BROKER_NAME);
|
||||
RemotingCommand response = recallMessageProcessor.processRequest(handlerContext, request);
|
||||
Assert.assertEquals(ResponseCode.NO_PERMISSION, response.getCode());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testProcessRequest_invalidStatus() throws RemotingCommandException {
|
||||
RemotingCommand request = mockRequest(0, TOPIC, TOPIC, "id", BROKER_NAME);
|
||||
|
||||
@@ -453,6 +453,8 @@ public class BrokerConfig extends BrokerIdentity {
|
||||
|
||||
private boolean allowRecallWhenBrokerNotWriteable = true;
|
||||
|
||||
private boolean recallMessageEnable = false;
|
||||
|
||||
public String getConfigBlackList() {
|
||||
return configBlackList;
|
||||
}
|
||||
@@ -1996,4 +1998,12 @@ public class BrokerConfig extends BrokerIdentity {
|
||||
public void setAllowRecallWhenBrokerNotWriteable(boolean allowRecallWhenBrokerNotWriteable) {
|
||||
this.allowRecallWhenBrokerNotWriteable = allowRecallWhenBrokerNotWriteable;
|
||||
}
|
||||
|
||||
public boolean isRecallMessageEnable() {
|
||||
return recallMessageEnable;
|
||||
}
|
||||
|
||||
public void setRecallMessageEnable(boolean recallMessageEnable) {
|
||||
this.recallMessageEnable = recallMessageEnable;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -138,6 +138,7 @@ public class IntegrationTestBase {
|
||||
brokerConfig.setEnableCalcFilterBitMap(true);
|
||||
brokerConfig.setAppendAckAsync(true);
|
||||
brokerConfig.setAppendCkAsync(true);
|
||||
brokerConfig.setRecallMessageEnable(true);
|
||||
storeConfig.setEnableConsumeQueueExt(true);
|
||||
brokerConfig.setLoadBalancePollNameServerInterval(500);
|
||||
storeConfig.setStorePathRootDir(baseDir);
|
||||
|
||||
Reference in New Issue
Block a user