mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Add Ack in LocalGrpcTest
This commit is contained in:
@@ -25,8 +25,8 @@ import org.apache.rocketmq.common.message.MessageExt;
|
||||
|
||||
public class ReceiptHandle {
|
||||
private static final String SEPARATOR = MessageConst.KEY_SEPARATOR;
|
||||
private static final String NORMAL_TOPIC = "0";
|
||||
private static final String RETRY_TOPIC = "1";
|
||||
public static final String NORMAL_TOPIC = "0";
|
||||
public static final String RETRY_TOPIC = "1";
|
||||
private final long startOffset;
|
||||
private final long retrieveTime;
|
||||
private final long invisibleTime;
|
||||
|
||||
@@ -17,8 +17,8 @@
|
||||
|
||||
package org.apache.rocketmq.test.base;
|
||||
|
||||
import apache.rocketmq.v1.Address;
|
||||
import apache.rocketmq.v1.AddressScheme;
|
||||
import apache.rocketmq.v1.AckMessageRequest;
|
||||
import apache.rocketmq.v1.AckMessageResponse;
|
||||
import apache.rocketmq.v1.Endpoints;
|
||||
import apache.rocketmq.v1.Message;
|
||||
import apache.rocketmq.v1.Partition;
|
||||
@@ -145,6 +145,18 @@ public class GrpcBaseTest extends BaseConf {
|
||||
.build();
|
||||
}
|
||||
|
||||
public AckMessageRequest buildAckMessageRequest(String group, String topic, String receiptHandle) {
|
||||
return AckMessageRequest.newBuilder()
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName(group)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName(topic)
|
||||
.build())
|
||||
.setReceiptHandle(receiptHandle)
|
||||
.build();
|
||||
}
|
||||
|
||||
public void assertQueryRoute(QueryRouteResponse response, int brokerSize) {
|
||||
assertThat(response.getCommon().getStatus().getCode()).isEqualTo(Code.OK_VALUE);
|
||||
assertThat(response.getPartitionsList().size()).isEqualTo(brokerSize * defaultQueueNums);
|
||||
@@ -167,4 +179,10 @@ public class GrpcBaseTest extends BaseConf {
|
||||
.getSystemAttribute()
|
||||
.getMessageId()).isEqualTo(messageId);
|
||||
}
|
||||
|
||||
public void assertAck(AckMessageResponse response) {
|
||||
assertThat(response.getCommon()
|
||||
.getStatus()
|
||||
.getCode()).isEqualTo(Code.OK_VALUE);
|
||||
}
|
||||
}
|
||||
@@ -17,6 +17,7 @@
|
||||
|
||||
package org.apache.rocketmq.test.proxy;
|
||||
|
||||
import apache.rocketmq.v1.AckMessageResponse;
|
||||
import apache.rocketmq.v1.MessagingServiceGrpc;
|
||||
import apache.rocketmq.v1.QueryRouteResponse;
|
||||
import apache.rocketmq.v1.ReceiveMessageResponse;
|
||||
@@ -81,5 +82,8 @@ public class LocalGrpcTest extends GrpcBaseTest {
|
||||
ReceiveMessageResponse receiveResponse = blockingStub.withDeadlineAfter(3, TimeUnit.SECONDS)
|
||||
.receiveMessage(buildReceiveMessageRequest(group, broker1Name));
|
||||
assertReceiveMessage(receiveResponse, messageId);
|
||||
String receiptHandle = receiveResponse.getMessages(0).getSystemAttribute().getReceiptHandle();
|
||||
AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(group, broker1Name, receiptHandle));
|
||||
assertAck(ackMessageResponse);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user