From 1451bf555857030e9132bf180333cebf17d9bb45 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Tue, 17 May 2022 17:35:36 +0800 Subject: [PATCH] [ISSUE #3949] v2 support --- .../grpc/v2/consumer/ReceiveMessageActivity.java | 12 +++++++++++- .../consumer/ReceiveMessageResponseStreamWriter.java | 2 +- 2 files changed, 12 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java index bac4e3cc44..27dc411ba9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java @@ -31,6 +31,7 @@ import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.v2.AbstractMessingActivity; +import org.apache.rocketmq.proxy.grpc.v2.GrpcContextConstants; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter; import org.apache.rocketmq.proxy.processor.MessagingProcessor; @@ -65,6 +66,15 @@ public class ReceiveMessageActivity extends AbstractMessingActivity { writer.write(proxyContext, Code.MESSAGE_NOT_FOUND, "no new message"); return; } + + long invisibleTime = Durations.toMillis(request.getInvisibleDuration()); + if (request.getAutoRenew()) { + invisibleTime = Durations.toMillis( + this.grpcClientSettingsManager.getClientSettings(proxyContext.getVal(GrpcContextConstants.CLIENT_ID)) + .getSubscription().getLongPollingTimeout() + ); + } + String topic = GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic()); String group = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); FilterExpression filterExpression = request.getFilterExpression(); @@ -85,7 +95,7 @@ public class ReceiveMessageActivity extends AbstractMessingActivity { group, topic, request.getBatchSize(), - Durations.toMillis(request.getInvisibleDuration()), + invisibleTime, pollTime, ConsumeInitMode.MAX, subscriptionData, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageResponseStreamWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageResponseStreamWriter.java index 6a94f49848..52514ef0ef 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageResponseStreamWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageResponseStreamWriter.java @@ -56,7 +56,7 @@ public class ReceiveMessageResponseStreamWriter { case FOUND: if (messageFoundList.isEmpty()) { streamObserver.onNext(ReceiveMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.MESSAGE_NOT_FOUND, "no new message")) + .setStatus(ResponseBuilder.buildStatus(Code.MESSAGE_NOT_FOUND, "no match message")) .build()); } else { streamObserver.onNext(ReceiveMessageResponse.newBuilder()