Skip to content
Open
Show file tree
Hide file tree
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 @@ -272,6 +272,7 @@ public CompletableFuture<PopResult> popMessage(ProxyContext ctx, AddressableMess
sortMap.get(key).add(messageExt.getQueueOffset());
}
Map<String, String> map = new HashMap<>(5);
List<MessageExt> validMessageExtList = new ArrayList<>(messageExtList.size());
for (MessageExt messageExt : messageExtList) {
if (startOffsetInfo == null) {
// we should set the check point info to extraInfo field , if the command is popMsg
Expand All @@ -285,14 +286,26 @@ public CompletableFuture<PopResult> popMessage(ProxyContext ctx, AddressableMess
} else {
if (messageExt.getProperty(MessageConst.PROPERTY_POP_CK) == null) {
String key = ExtraInfoUtil.getStartOffsetInfoMapKey(messageExt.getTopic(), messageExt.getQueueId());
int index = sortMap.get(key).indexOf(messageExt.getQueueOffset());
Long msgQueueOffset = msgOffsetInfo.get(key).get(index);
List<Long> sortQueueOffsets = sortMap.get(key);
List<Long> msgQueueOffsets = msgOffsetInfo == null ? null : msgOffsetInfo.get(key);
Long startOffset = startOffsetInfo.get(key);
Comment on lines +289 to +291
if (sortQueueOffsets == null || msgQueueOffsets == null || startOffset == null) {
log.warn("Pop response offset metadata is missing, key:{}", key);
continue;
}
int index = sortQueueOffsets.indexOf(messageExt.getQueueOffset());
if (index < 0 || index >= msgQueueOffsets.size()) {
log.warn("Pop response offset metadata index is invalid, key:{}, index:{}, msgOffsetCount:{}",
key, index, msgQueueOffsets.size());
continue;
}
Comment on lines +292 to +301
Long msgQueueOffset = msgQueueOffsets.get(index);
if (msgQueueOffset != messageExt.getQueueOffset()) {
log.warn("Queue offset [{}] of msg is strange, not equal to the stored in msg, {}", msgQueueOffset, messageExt);
}

messageExt.getProperties().put(MessageConst.PROPERTY_POP_CK,
ExtraInfoUtil.buildExtraInfo(startOffsetInfo.get(key), responseHeader.getPopTime(), responseHeader.getInvisibleTime(),
ExtraInfoUtil.buildExtraInfo(startOffset, responseHeader.getPopTime(), responseHeader.getInvisibleTime(),
responseHeader.getReviveQid(), messageExt.getTopic(), messageQueue.getBrokerName(), messageExt.getQueueId(), msgQueueOffset)
);
if (requestHeader.isOrder() && orderCountInfo != null) {
Expand All @@ -306,7 +319,9 @@ public CompletableFuture<PopResult> popMessage(ProxyContext ctx, AddressableMess
messageExt.getProperties().computeIfAbsent(MessageConst.PROPERTY_FIRST_POP_TIME, k -> String.valueOf(responseHeader.getPopTime()));
messageExt.setBrokerName(messageQueue.getBrokerName());
messageExt.setTopic(messageQueue.getTopic());
validMessageExtList.add(messageExt);
}
popResult.setMsgFoundList(validMessageExtList);
Comment on lines 319 to +324
}
return popResult;
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -340,6 +340,47 @@ public void testPopMessageWriteAndFlush() throws Exception {
}
}

@Test
public void testPopMessageShouldSkipMessageWithMissingOffsetMetadata() throws Exception {
Comment on lines +343 to +344
int reviveQueueId = 1;
long popTime = System.currentTimeMillis();
long invisibleTime = 3000L;
long startOffset = 100L;
StringBuilder startOffsetStringBuilder = new StringBuilder();
ExtraInfoUtil.buildStartOffsetInfo(startOffsetStringBuilder, topic, queueId, startOffset);
MessageExt message = buildMessageExt(topic, queueId, startOffset);
byte[] body = MessageDecoder.encode(message, false);
PopMessageRequestHeader requestHeader = new PopMessageRequestHeader();
requestHeader.setInvisibleTime(invisibleTime);
Mockito.when(popMessageProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> {
boolean first = argument.getCode() == RequestCode.POP_MESSAGE;
boolean second = argument.readCustomHeader() instanceof PopMessageRequestHeader;
return first && second;
}))).thenAnswer(invocation -> {
SimpleChannelHandlerContext simpleChannelHandlerContext = invocation.getArgument(0);
RemotingCommand request = invocation.getArgument(1);
RemotingCommand response = RemotingCommand.createResponseCommand(PopMessageResponseHeader.class);
response.setOpaque(request.getOpaque());
response.setCode(ResponseCode.SUCCESS);
response.setBody(body);
PopMessageResponseHeader responseHeader = (PopMessageResponseHeader) response.readCustomHeader();
responseHeader.setStartOffsetInfo(startOffsetStringBuilder.toString());
responseHeader.setInvisibleTime(requestHeader.getInvisibleTime());
responseHeader.setPopTime(popTime);
responseHeader.setReviveQid(reviveQueueId);
simpleChannelHandlerContext.writeAndFlush(response);
return null;
});

MessageQueue messageQueue = new MessageQueue(topic, brokerName, queueId);
CompletableFuture<PopResult> future = localMessageService.popMessage(proxyContext,
new AddressableMessageQueue(messageQueue, ""), requestHeader, 1000L);
PopResult popResult = future.get();

assertThat(popResult.getPopStatus()).isEqualTo(PopStatus.FOUND);
assertThat(popResult.getMsgFoundList()).isEmpty();
}

@Test
public void testPopMessagePollingTimeout() throws Exception {
RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.POLLING_TIMEOUT, "");
Expand Down
Loading