diff --git a/pom.xml b/pom.xml
index 467138921c..3362e6c84a 100644
--- a/pom.xml
+++ b/pom.xml
@@ -455,6 +455,7 @@
${project.groupId}
rocketmq-proto
2.0.0-SNAPSHOT
+ compatible
${project.groupId}
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java
index cddce40d0b..e108f8de3a 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java
@@ -57,7 +57,12 @@ public class ProxyConfig {
*/
private int grpcMaxInboundMessageSize = 130 * 1024 * 1024;
- private int channelExpiredInSeconds = 120;
+ private int maxMessageBodyBytes = 1024 * 1024 * 4;
+ private int defaultMessageBodyCompressionBytesThreshold = 1024 * 4;
+ private int defaultTransactionRecoverySecond = 30;
+ private int defaultMaxDeliveryAttempts = 16;
+
+ private int channelExpiredInSeconds = 60;
private int forwardConsumerNum = 2;
private double forwardConsumerWorkerFactor = 0.2f;
@@ -237,6 +242,38 @@ public class ProxyConfig {
this.grpcMaxInboundMessageSize = grpcMaxInboundMessageSize;
}
+ public int getMaxMessageBodyBytes() {
+ return maxMessageBodyBytes;
+ }
+
+ public void setMaxMessageBodyBytes(int maxMessageBodyBytes) {
+ this.maxMessageBodyBytes = maxMessageBodyBytes;
+ }
+
+ public int getDefaultMessageBodyCompressionBytesThreshold() {
+ return defaultMessageBodyCompressionBytesThreshold;
+ }
+
+ public void setDefaultMessageBodyCompressionBytesThreshold(int defaultMessageBodyCompressionBytesThreshold) {
+ this.defaultMessageBodyCompressionBytesThreshold = defaultMessageBodyCompressionBytesThreshold;
+ }
+
+ public int getDefaultTransactionRecoverySecond() {
+ return defaultTransactionRecoverySecond;
+ }
+
+ public void setDefaultTransactionRecoverySecond(int defaultTransactionRecoverySecond) {
+ this.defaultTransactionRecoverySecond = defaultTransactionRecoverySecond;
+ }
+
+ public int getDefaultMaxDeliveryAttempts() {
+ return defaultMaxDeliveryAttempts;
+ }
+
+ public void setDefaultMaxDeliveryAttempts(int defaultMaxDeliveryAttempts) {
+ this.defaultMaxDeliveryAttempts = defaultMaxDeliveryAttempts;
+ }
+
public int getChannelExpiredInSeconds() {
return channelExpiredInSeconds;
}
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java
index 87cc060e09..53f9fecfde 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java
@@ -21,7 +21,6 @@ import apache.rocketmq.v2.AckMessageRequest;
import apache.rocketmq.v2.AckMessageResponse;
import apache.rocketmq.v2.ChangeInvisibleDurationRequest;
import apache.rocketmq.v2.ChangeInvisibleDurationResponse;
-import apache.rocketmq.v2.Code;
import apache.rocketmq.v2.EndTransactionRequest;
import apache.rocketmq.v2.EndTransactionResponse;
import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest;
@@ -33,12 +32,8 @@ import apache.rocketmq.v2.NackMessageRequest;
import apache.rocketmq.v2.NackMessageResponse;
import apache.rocketmq.v2.NotifyClientTerminationRequest;
import apache.rocketmq.v2.NotifyClientTerminationResponse;
-import apache.rocketmq.v2.PullMessageRequest;
-import apache.rocketmq.v2.PullMessageResponse;
import apache.rocketmq.v2.QueryAssignmentRequest;
import apache.rocketmq.v2.QueryAssignmentResponse;
-import apache.rocketmq.v2.QueryOffsetRequest;
-import apache.rocketmq.v2.QueryOffsetResponse;
import apache.rocketmq.v2.QueryRouteRequest;
import apache.rocketmq.v2.QueryRouteResponse;
import apache.rocketmq.v2.ReceiveMessageRequest;
@@ -50,8 +45,6 @@ import apache.rocketmq.v2.TelemetryCommand;
import io.grpc.Context;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.CompletionException;
-import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException;
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseWriter;
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcForwardService;
@@ -64,14 +57,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
this.grpcForwardService = grpcForwardService;
}
- public Status convertExceptionToStatus(Throwable t) {
- if (t instanceof CompletionException) {
- if (t.getCause() instanceof ProxyException) {
- ProxyException proxyException = (ProxyException) t.getCause();
- return ResponseBuilder.buildStatus(proxyException.getCode(), proxyException.getMessage());
- }
- }
- return ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "internal error");
+ protected Status convertExceptionToStatus(Throwable t) {
+ return ResponseBuilder.buildStatus(t);
}
@Override
@@ -193,32 +180,6 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
});
}
- @Override
- public void queryOffset(QueryOffsetRequest request, StreamObserver responseObserver) {
- CompletableFuture future = grpcForwardService.queryOffset(Context.current(), request);
- future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
- .exceptionally(e -> {
- ResponseWriter.write(
- responseObserver,
- QueryOffsetResponse.newBuilder().setStatus(convertExceptionToStatus(e)).build()
- );
- return null;
- });
- }
-
- @Override
- public void pullMessage(PullMessageRequest request, StreamObserver responseObserver) {
- CompletableFuture future = grpcForwardService.pullMessage(Context.current(), request);
- future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
- .exceptionally(e -> {
- ResponseWriter.write(
- responseObserver,
- PullMessageResponse.newBuilder().setStatus(convertExceptionToStatus(e)).build()
- );
- return null;
- });
- }
-
@Override
public void notifyClientTermination(NotifyClientTerminationRequest request,
StreamObserver responseObserver) {
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java
index a67844a85a..cbb5bb6dcd 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java
@@ -18,14 +18,15 @@
package org.apache.rocketmq.proxy.grpc.v2.adapter;
import apache.rocketmq.v2.AckMessageRequest;
+import apache.rocketmq.v2.ApplyPassiveSettingsCommand;
import apache.rocketmq.v2.ChangeInvisibleDurationRequest;
-import apache.rocketmq.v2.ClientSettings;
import apache.rocketmq.v2.ClientType;
import apache.rocketmq.v2.Code;
import apache.rocketmq.v2.Digest;
import apache.rocketmq.v2.DigestType;
import apache.rocketmq.v2.Encoding;
import apache.rocketmq.v2.EndTransactionRequest;
+import apache.rocketmq.v2.Endpoints;
import apache.rocketmq.v2.FilterExpression;
import apache.rocketmq.v2.FilterType;
import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest;
@@ -35,7 +36,8 @@ import apache.rocketmq.v2.MessageQueue;
import apache.rocketmq.v2.MessageType;
import apache.rocketmq.v2.NackMessageRequest;
import apache.rocketmq.v2.NotifyClientTerminationRequest;
-import apache.rocketmq.v2.PullMessageRequest;
+import apache.rocketmq.v2.PassivePublishingSettings;
+import apache.rocketmq.v2.PassiveSubscriptionSettings;
import apache.rocketmq.v2.ReceiveMessageRequest;
import apache.rocketmq.v2.Resource;
import apache.rocketmq.v2.SendMessageRequest;
@@ -78,7 +80,6 @@ import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHead
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader;
import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
-import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
import org.apache.rocketmq.common.protocol.header.UnregisterClientRequestHeader;
import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType;
@@ -86,12 +87,13 @@ import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
-import org.apache.rocketmq.common.sysflag.PullSysFlag;
import org.apache.rocketmq.common.utils.BinaryUtil;
import org.apache.rocketmq.proxy.common.DelayPolicy;
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
import org.apache.rocketmq.proxy.config.ConfigurationManager;
+import org.apache.rocketmq.proxy.config.ProxyConfig;
import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
+import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -103,13 +105,13 @@ public class GrpcConverter {
}
public static HeartbeatData buildHeartbeatData(String clientId, HeartbeatRequest request,
- ClientSettings clientSettings) {
+ GrpcClientManager.ActiveClientSettings clientSettings) {
HeartbeatData heartbeatData = new HeartbeatData();
heartbeatData.setClientID(clientId);
switch (clientSettings.getClientType()) {
case PRODUCER: {
Set producerDataSet = new HashSet<>();
- for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) {
+ for (Resource topic : clientSettings.getActivePublishingSettings().getPublishingTopicsList()) {
String topicName = wrapResourceWithNamespace(topic);
producerDataSet.add(buildProducerData(topicName));
}
@@ -137,7 +139,7 @@ public class GrpcConverter {
}
public static org.apache.rocketmq.common.protocol.heartbeat.ConsumerData buildConsumerData(String groupName,
- ClientSettings clientSettings) {
+ GrpcClientManager.ActiveClientSettings clientSettings) {
org.apache.rocketmq.common.protocol.heartbeat.ConsumerData buildConsumerData = new org.apache.rocketmq.common.protocol.heartbeat.ConsumerData();
buildConsumerData.setGroupName(groupName);
buildConsumerData.setConsumeType(buildConsumeType(clientSettings.getClientType()));
@@ -145,9 +147,7 @@ public class GrpcConverter {
buildConsumerData.setMessageModel(MessageModel.CLUSTERING);
buildConsumerData.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
Set subscriptionDataSet =
- buildSubscriptionDataSet(clientSettings.getSettings()
- .getSubscription()
- .getSubscriptionsList());
+ buildSubscriptionDataSet(clientSettings.getActiveSubscriptionSettings().getSubscriptionsList());
buildConsumerData.setSubscriptionDataSet(subscriptionDataSet);
return buildConsumerData;
}
@@ -241,10 +241,9 @@ public class GrpcConverter {
return requestHeader;
}
- public static AckMessageRequestHeader buildAckMessageRequestHeader(AckMessageRequest request) {
+ public static AckMessageRequestHeader buildAckMessageRequestHeader(AckMessageRequest request, ReceiptHandle handle) {
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
- ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle());
AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader();
ackMessageRequestHeader.setConsumerGroup(groupName);
@@ -360,32 +359,6 @@ public class GrpcConverter {
return endTransactionRequestHeader;
}
- public static PullMessageRequestHeader buildPullMessageRequestHeader(PullMessageRequest request,
- long pollTimeoutInMillis) {
- MessageQueue messageQueue = request.getMessageQueue();
- String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
- String topicName = GrpcConverter.wrapResourceWithNamespace(messageQueue.getTopic());
-
- int queueId = messageQueue.getId();
- int sysFlag = PullSysFlag.buildSysFlag(false, true, true, false, false);
- String expression = request.getFilterExpression().getExpression();
- String expressionType = GrpcConverter.buildExpressionType(request.getFilterExpression().getType());
-
- PullMessageRequestHeader requestHeader = new PullMessageRequestHeader();
- requestHeader.setConsumerGroup(groupName);
- requestHeader.setTopic(topicName);
- requestHeader.setQueueId(queueId);
- requestHeader.setQueueOffset(request.getOffset());
- requestHeader.setMaxMsgNums(request.getBatchSize());
- requestHeader.setSysFlag(sysFlag);
- requestHeader.setCommitOffset(0L);
- requestHeader.setSuspendTimeoutMillis(pollTimeoutInMillis);
- requestHeader.setSubscription(expression);
- requestHeader.setSubVersion(0L);
- requestHeader.setExpressionType(expressionType);
- return requestHeader;
- }
-
public static UnregisterClientRequestHeader buildUnregisterClientRequestHeader(String clientId,
ClientType clientType, NotifyClientTerminationRequest request) {
UnregisterClientRequestHeader header = new UnregisterClientRequestHeader();
@@ -535,6 +508,30 @@ public class GrpcConverter {
.build();
}
+ public static ApplyPassiveSettingsCommand buildDefaultPublishingSettings(String nonce, Endpoints traceEndpoint) {
+ ProxyConfig proxyConfig = ConfigurationManager.getProxyConfig();
+ return ApplyPassiveSettingsCommand.newBuilder()
+ .setNonce(nonce)
+ .setTraceAccessPoint(traceEndpoint)
+ .setPassivePublishingSettings(PassivePublishingSettings.newBuilder()
+ .setMaxMessageBodyBytes(proxyConfig.getMaxMessageBodyBytes())
+ .setMessageBodyCompressionBytesThreshold(proxyConfig.getDefaultMessageBodyCompressionBytesThreshold())
+ .setOrphanedTransactionRecoveryDuration(Durations.fromSeconds(proxyConfig.getDefaultTransactionRecoverySecond()))
+ .build())
+ .build();
+ }
+
+ public static ApplyPassiveSettingsCommand buildDefaultSubscriptionSettings(String nonce, Endpoints traceEndpoint) {
+ // TODO: read config from subscriptionGroupManager
+ return ApplyPassiveSettingsCommand.newBuilder()
+ .setNonce(nonce)
+ .setTraceAccessPoint(traceEndpoint)
+ .setPassiveSubscriptionSettings(PassiveSubscriptionSettings.newBuilder()
+ .setFifo(false)
+ .build())
+ .build();
+ }
+
protected static Map buildUserAttributes(MessageExt messageExt) {
Map userAttributes = new HashMap<>();
Map properties = messageExt.getProperties();
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java
index eecc19f53c..27be2f4333 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java
@@ -19,9 +19,22 @@ package org.apache.rocketmq.proxy.grpc.v2.adapter;
import apache.rocketmq.v2.Code;
import apache.rocketmq.v2.Status;
+import java.util.concurrent.CompletionException;
import org.apache.rocketmq.common.protocol.ResponseCode;
public class ResponseBuilder {
+
+ public static Status buildStatus(Throwable t) {
+ if (t instanceof CompletionException) {
+ t = t.getCause();
+ }
+ if (t instanceof ProxyException) {
+ ProxyException proxyException = (ProxyException) t.getCause();
+ return ResponseBuilder.buildStatus(proxyException.getCode(), proxyException.getMessage());
+ }
+ return ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "internal error");
+ }
+
public static Status buildStatus(Code code, String message) {
return Status.newBuilder()
.setCode(code)
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java
index b44d958a50..100792bf6f 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java
@@ -31,12 +31,8 @@ import apache.rocketmq.v2.NackMessageRequest;
import apache.rocketmq.v2.NackMessageResponse;
import apache.rocketmq.v2.NotifyClientTerminationRequest;
import apache.rocketmq.v2.NotifyClientTerminationResponse;
-import apache.rocketmq.v2.PullMessageRequest;
-import apache.rocketmq.v2.PullMessageResponse;
import apache.rocketmq.v2.QueryAssignmentRequest;
import apache.rocketmq.v2.QueryAssignmentResponse;
-import apache.rocketmq.v2.QueryOffsetRequest;
-import apache.rocketmq.v2.QueryOffsetResponse;
import apache.rocketmq.v2.QueryRouteRequest;
import apache.rocketmq.v2.QueryRouteResponse;
import apache.rocketmq.v2.ReceiveMessageRequest;
@@ -54,15 +50,14 @@ import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.proxy.channel.ChannelManager;
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
import org.apache.rocketmq.proxy.common.StartAndShutdown;
+import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
import org.apache.rocketmq.proxy.connector.ConnectorManager;
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateCheckRequest;
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker;
-import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ConsumerService;
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ForwardClientService;
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ProducerService;
-import org.apache.rocketmq.proxy.grpc.v2.service.cluster.PullMessageService;
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService;
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.TransactionService;
import org.slf4j.Logger;
@@ -81,7 +76,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
private final ConsumerService consumerService;
private final RouteService routeService;
private final ForwardClientService clientService;
- private final PullMessageService pullMessageService;
private final TransactionService transactionService;
private final TelemetryCommandManager pollCommandResponseManager;
private final GrpcClientManager grpcClientManager;
@@ -96,7 +90,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager, grpcClientManager);
this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService,
channelManager, grpcClientManager, pollCommandResponseManager);
- this.pullMessageService = new PullMessageService(connectorManager);
this.transactionService = new TransactionService(connectorManager, channelManager);
this.appendStartAndShutdown(new ClusterGrpcServiceStartAndShutdown());
@@ -149,16 +142,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
return transactionService.endTransaction(ctx, request);
}
- @Override
- public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) {
- return pullMessageService.queryOffset(ctx, request);
- }
-
- @Override
- public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) {
- return pullMessageService.pullMessage(ctx, request);
- }
-
@Override
public CompletableFuture notifyClientTermination(Context ctx,
NotifyClientTerminationRequest request) {
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java
index f949af8ecf..3042735a04 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java
@@ -17,7 +17,12 @@
package org.apache.rocketmq.proxy.grpc.v2.service;
-import apache.rocketmq.v2.ClientSettings;
+import apache.rocketmq.v2.ActivePublishingSettings;
+import apache.rocketmq.v2.ActiveSubscriptionSettings;
+import apache.rocketmq.v2.ClientType;
+import apache.rocketmq.v2.Endpoints;
+import apache.rocketmq.v2.ReportActiveSettingsCommand;
+import com.google.protobuf.Duration;
import io.grpc.Context;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@@ -25,22 +30,68 @@ import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
public class GrpcClientManager {
- private static final Map CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>();
+ public static class ActiveClientSettings {
+ private ClientType clientType;
+ private Endpoints accessPoint;
+ private Duration connectionTimeout;
+ private boolean traceOn = true;
+ private ActivePublishingSettings activePublishingSettings;
+ private ActiveSubscriptionSettings activeSubscriptionSettings;
- public ClientSettings getClientSettings(Context ctx) {
+ public ActiveClientSettings(ReportActiveSettingsCommand reportActiveSettingsCommand) {
+ this.clientType = reportActiveSettingsCommand.getClientType();
+ this.accessPoint = reportActiveSettingsCommand.getAccessPoint();
+ this.connectionTimeout = reportActiveSettingsCommand.getConnectionTimeout();
+ this.traceOn = reportActiveSettingsCommand.getTraceOn();
+ if (reportActiveSettingsCommand.hasActivePublishingSettings()) {
+ this.activePublishingSettings = reportActiveSettingsCommand.getActivePublishingSettings();
+ }
+ if (reportActiveSettingsCommand.hasActiveSubscriptionSettings()) {
+ this.activeSubscriptionSettings = reportActiveSettingsCommand.getActiveSubscriptionSettings();
+ }
+ }
+
+ public ClientType getClientType() {
+ return clientType;
+ }
+
+ public Endpoints getAccessPoint() {
+ return accessPoint;
+ }
+
+ public Duration getConnectionTimeout() {
+ return connectionTimeout;
+ }
+
+ public boolean isTraceOn() {
+ return traceOn;
+ }
+
+ public ActivePublishingSettings getActivePublishingSettings() {
+ return activePublishingSettings;
+ }
+
+ public ActiveSubscriptionSettings getActiveSubscriptionSettings() {
+ return activeSubscriptionSettings;
+ }
+ }
+
+ private static final Map CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>();
+
+ public ActiveClientSettings getClientSettings(Context ctx) {
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
return CLIENT_SETTINGS_MAP.get(clientId);
}
- public ClientSettings getClientSettings(String clientId) {
+ public ActiveClientSettings getClientSettings(String clientId) {
return CLIENT_SETTINGS_MAP.get(clientId);
}
- public void updateClientSettings(String clientId, ClientSettings clientSettings) {
- CLIENT_SETTINGS_MAP.put(clientId, clientSettings);
+ public void updateClientSettings(String clientId, ReportActiveSettingsCommand reportActiveSettingsCommand) {
+ CLIENT_SETTINGS_MAP.put(clientId, new ActiveClientSettings(reportActiveSettingsCommand));
}
- public ClientSettings removeClientSettings(String clientId) {
+ public ActiveClientSettings removeClientSettings(String clientId) {
return CLIENT_SETTINGS_MAP.remove(clientId);
}
}
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java
index eb11cf4ade..f7bf6f946a 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java
@@ -31,12 +31,8 @@ import apache.rocketmq.v2.NackMessageRequest;
import apache.rocketmq.v2.NackMessageResponse;
import apache.rocketmq.v2.NotifyClientTerminationRequest;
import apache.rocketmq.v2.NotifyClientTerminationResponse;
-import apache.rocketmq.v2.PullMessageRequest;
-import apache.rocketmq.v2.PullMessageResponse;
import apache.rocketmq.v2.QueryAssignmentRequest;
import apache.rocketmq.v2.QueryAssignmentResponse;
-import apache.rocketmq.v2.QueryOffsetRequest;
-import apache.rocketmq.v2.QueryOffsetResponse;
import apache.rocketmq.v2.QueryRouteRequest;
import apache.rocketmq.v2.QueryRouteResponse;
import apache.rocketmq.v2.ReceiveMessageRequest;
@@ -69,10 +65,6 @@ public interface GrpcForwardService extends StartAndShutdown {
CompletableFuture endTransaction(Context ctx, EndTransactionRequest request);
- CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request);
-
- CompletableFuture pullMessage(Context ctx, PullMessageRequest request);
-
CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request);
CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request);
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java
index 73998b4283..bca408ee58 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java
@@ -98,6 +98,7 @@ import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
import org.apache.rocketmq.proxy.common.DelayPolicy;
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
import org.apache.rocketmq.proxy.common.TelemetryCommandRecord;
+import org.apache.rocketmq.proxy.config.ConfigurationManager;
import org.apache.rocketmq.proxy.connector.ConnectorManager;
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
@@ -128,6 +129,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
private final TelemetryCommandManager telemetryCommandManager;
private final GrpcClientManager grpcClientManager;
private final RouteService routeService;
+ private final ReportActiveSettingsService reportActiveSettingsService;
private final DelayPolicy delayPolicy;
public LocalGrpcService(BrokerController brokerController) {
@@ -147,6 +149,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
this.telemetryCommandManager = telemetryCommandManager;
this.grpcClientManager = new GrpcClientManager();
this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager, grpcClientManager);
+ this.reportActiveSettingsService = new ReportActiveSettingsService(this.channelManager, this.grpcClientManager, this.telemetryCommandManager);
this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel());
this.brokerController.getConsumerManager().appendConsumerIdsChangeListener(new ConsumerIdsChangeListenerImpl());
@@ -167,7 +170,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
languageCode = LanguageCode.valueOf(language);
- ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
+ GrpcClientManager.ActiveClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
HeartbeatData heartbeatData = GrpcConverter.buildHeartbeatData(clientId, request, clientSettings);
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.HEART_BEAT, null);
command.setLanguage(languageCode);
@@ -178,7 +181,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
CompletableFuture future = new CompletableFuture<>();
switch (clientSettings.getClientType()) {
case PRODUCER: {
- for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) {
+ for (Resource topic : clientSettings.getActivePublishingSettings().getPublishingTopicsList()) {
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel);
@@ -266,14 +269,14 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) {
long pollTime = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS);
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
- ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
- PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime,
- clientSettings.getSettings().getSubscription().getFifo());
+ // TODO: get fifo config from subscriptionGroupManager
+ boolean fifo = false;
+ PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime, fifo);
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader);
command.makeCustomHeaderToNet();
ReceiveMessageResponseHandler handler = new ReceiveMessageResponseHandler(brokerController.getBrokerConfig().getBrokerName(),
- clientSettings.getSettings().getSubscription().getFifo());
+ fifo);
ReceiveMessageChannel channel = channelManager.createChannel(() -> new ReceiveMessageChannel(handler), ReceiveMessageChannel.class);
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
CompletableFuture future = new CompletableFuture<>();
@@ -322,8 +325,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
CompletableFuture future = new CompletableFuture<>();
- ClientSettings clientSettings = grpcClientManager.getClientSettings(ctx);
- int maxReconsumeTimes = clientSettings.getSettings().getSubscription().getDeadLetterPolicy().getMaxDeliveryAttempts();
+ int maxReconsumeTimes = ConfigurationManager.getProxyConfig().getDefaultMaxDeliveryAttempts();
if (request.getDeliveryAttempt() >= maxReconsumeTimes) {
ConsumerSendMsgBackRequestHeader requestHeader = GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request, maxReconsumeTimes);
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader);
@@ -414,56 +416,6 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
return future;
}
- @Override
- public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) {
- String topicName = GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic());
- int queueId = request.getMessageQueue().getId();
-
- long offset;
- if (request.getPolicy() == QueryOffsetPolicy.BEGINNING) {
- offset = 0L;
- } else if (request.getPolicy() == QueryOffsetPolicy.END) {
- offset = brokerController.getMessageStore()
- .getMaxOffsetInQueue(topicName, queueId);
- } else {
- long timestamp = Timestamps.toMillis(request.getTimePoint());
- offset = brokerController.getMessageStore()
- .getOffsetInQueueByTime(topicName, queueId, timestamp);
- }
- return CompletableFuture.completedFuture(QueryOffsetResponse.newBuilder()
- .setStatus(ResponseBuilder.buildStatus(Code.OK, "ok"))
- .setOffset(offset)
- .build());
- }
-
- @Override
- public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) {
- long pollTime = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS);
- PullMessageRequestHeader requestHeader = GrpcConverter.buildPullMessageRequestHeader(request, pollTime);
- RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.PULL_MESSAGE, requestHeader);
- command.makeCustomHeaderToNet();
-
- PullMessageResponseHandler handler = new PullMessageResponseHandler();
- PullMessageChannel channel = channelManager.createChannel(() -> new PullMessageChannel(handler), PullMessageChannel.class);
- SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
- CompletableFuture future = new CompletableFuture<>();
- InvocationContext context = new InvocationContext<>(request, future);
- channel.registerInvocationContext(command.getOpaque(), context);
- try {
- RemotingCommand response = brokerController.getPullMessageProcessor()
- .processRequest(channelHandlerContext, command);
- if (response != null) {
- handler.handle(response, context);
- channel.eraseInvocationContext(command.getOpaque());
- }
- } catch (Exception e) {
- log.error("Failed to process pull message command", e);
- channel.eraseInvocationContext(command.getOpaque());
- future.completeExceptionally(e);
- }
- return future;
- }
-
public void reportThreadStackTrace(ThreadStackTrace request) {
String nonce = request.getNonce();
String threadStack = request.getThreadStackTrace();
@@ -510,7 +462,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
Channel channel = channelManager.createChannel();
SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel);
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
- ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
+ GrpcClientManager.ActiveClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
UnregisterClientRequestHeader header = GrpcConverter.buildUnregisterClientRequestHeader(clientId, clientSettings.getClientType(), request);
RemotingCommand remotingCommand = RemotingCommand.createRequestCommand(RequestCode.UNREGISTER_CLIENT, header);
@@ -569,31 +521,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
@Override
public void onNext(TelemetryCommand request) {
switch (request.getCommandCase()) {
- case CLIENT_SETTINGS: {
- ClientSettings clientSettings = request.getClientSettings();
- grpcClientManager.updateClientSettings(clientId, clientSettings);
- Settings settings = clientSettings.getSettings();
- if (settings.hasPublishing()) {
- Publishing publishing = settings.getPublishing();
- for (Resource topic : publishing.getTopicsList()) {
- String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
- GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
- producerChannel.setClientObserver(responseObserver);
- }
- }
- if (settings.hasSubscription()) {
- Subscription subscription = settings.getSubscription();
- String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup());
- GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, telemetryCommandManager);
- consumerChannel.setClientObserver(responseObserver);
- }
- responseObserver.onNext(TelemetryCommand.newBuilder()
- .setClientOverwrittenSettings(ClientOverwrittenSettings.newBuilder()
- .setNonce(clientSettings.getNonce())
- .setDirection(Direction.RESPONSE)
- .setSettings(settings)
- .build())
- .build());
+ case REPORT_ACTIVE_SETTINGS_COMMAND: {
+ responseObserver.onNext(reportActiveSettingsService.processReportActiveSettingsCommand(ctx, request, responseObserver));
break;
}
case THREAD_STACK_TRACE: {
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ReportActiveSettingsService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ReportActiveSettingsService.java
new file mode 100644
index 0000000000..6a771fcd0c
--- /dev/null
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ReportActiveSettingsService.java
@@ -0,0 +1,77 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.rocketmq.proxy.grpc.v2.service;
+
+import apache.rocketmq.v2.ActiveSubscriptionSettings;
+import apache.rocketmq.v2.ApplyPassiveSettingsCommand;
+import apache.rocketmq.v2.ReportActiveSettingsCommand;
+import apache.rocketmq.v2.Resource;
+import apache.rocketmq.v2.TelemetryCommand;
+import io.grpc.Context;
+import io.grpc.stub.StreamObserver;
+import org.apache.rocketmq.proxy.channel.ChannelManager;
+import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
+import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
+import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
+import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
+
+public class ReportActiveSettingsService {
+
+ private final ChannelManager channelManager;
+ private final GrpcClientManager grpcClientManager;
+ private final TelemetryCommandManager telemetryCommandManager;
+
+ public ReportActiveSettingsService(ChannelManager channelManager,
+ GrpcClientManager grpcClientManager,
+ TelemetryCommandManager telemetryCommandManager) {
+ this.channelManager = channelManager;
+ this.grpcClientManager = grpcClientManager;
+ this.telemetryCommandManager = telemetryCommandManager;
+ }
+
+ public TelemetryCommand processReportActiveSettingsCommand(Context ctx, TelemetryCommand request, StreamObserver responseObserver) {
+ String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
+ ReportActiveSettingsCommand reportActiveSettings = request.getReportActiveSettingsCommand();
+ grpcClientManager.updateClientSettings(clientId, reportActiveSettings);
+ ApplyPassiveSettingsCommand applyPassiveSettingsCommand = ApplyPassiveSettingsCommand.getDefaultInstance();
+ if (reportActiveSettings.hasActivePublishingSettings()) {
+ for (Resource topic : reportActiveSettings.getActivePublishingSettings().getPublishingTopicsList()) {
+ String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
+ GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
+ producerChannel.setClientObserver(responseObserver);
+ }
+ applyPassiveSettingsCommand = GrpcConverter.buildDefaultPublishingSettings(
+ reportActiveSettings.getNonce(),
+ reportActiveSettings.getAccessPoint()
+ );
+ }
+ if (reportActiveSettings.hasActiveSubscriptionSettings()) {
+ ActiveSubscriptionSettings subscription = reportActiveSettings.getActiveSubscriptionSettings();
+ String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup());
+ GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, telemetryCommandManager);
+ consumerChannel.setClientObserver(responseObserver);
+ applyPassiveSettingsCommand = GrpcConverter.buildDefaultSubscriptionSettings(
+ reportActiveSettings.getNonce(),
+ reportActiveSettings.getAccessPoint()
+ );
+ }
+ return TelemetryCommand.newBuilder()
+ .setApplyPassiveSettingsCommand(applyPassiveSettingsCommand)
+ .build();
+ }
+}
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java
index 1c80d715c5..c9d2fceebc 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java
@@ -16,8 +16,10 @@
*/
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
+import apache.rocketmq.v2.AckMessageEntry;
import apache.rocketmq.v2.AckMessageRequest;
import apache.rocketmq.v2.AckMessageResponse;
+import apache.rocketmq.v2.AckMessageResultEntry;
import apache.rocketmq.v2.ChangeInvisibleDurationRequest;
import apache.rocketmq.v2.ChangeInvisibleDurationResponse;
import apache.rocketmq.v2.Code;
@@ -26,7 +28,6 @@ import apache.rocketmq.v2.NackMessageRequest;
import apache.rocketmq.v2.NackMessageResponse;
import apache.rocketmq.v2.ReceiveMessageRequest;
import apache.rocketmq.v2.ReceiveMessageResponse;
-import apache.rocketmq.v2.Settings;
import io.grpc.Context;
import java.util.ArrayList;
import java.util.List;
@@ -130,8 +131,8 @@ public class ConsumerService extends BaseService {
protected PopMessageRequestHeader buildPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) {
checkSubscriptionData(request.getMessageQueue().getTopic(), request.getFilterExpression());
- boolean isFifo = grpcClientManager.getClientSettings(ctx).getSettings().getSubscription().getFifo();
- return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx), isFifo);
+ // TODO: get fifo config from subscriptionGroupManager
+ return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx), false);
}
protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) {
@@ -209,40 +210,70 @@ public class ConsumerService extends BaseService {
});
try {
- ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
- String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
-
- AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request);
- CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader);
- ackResultFuture
- .thenAccept(result -> {
- try {
- future.complete(convertToAckMessageResponse(ctx, request, result));
- } catch (Throwable throwable) {
- future.completeExceptionally(throwable);
- }
- })
- .exceptionally(throwable -> {
+ CompletableFuture[] futures = new CompletableFuture[request.getEntriesCount()];
+ for (int i = 0; i < request.getEntriesCount(); i++) {
+ futures[i] = processAckMessage(ctx, request, request.getEntries(i));
+ }
+ CompletableFuture.allOf(futures).whenComplete((val, throwable) -> {
+ if (throwable != null) {
future.completeExceptionally(throwable);
- return null;
- });
+ return;
+ }
+ List entryList = new ArrayList<>();
+ for (CompletableFuture entryFuture : futures) {
+ entryFuture.thenAccept(entryList::add);
+ }
+ AckMessageResponse.Builder responseBuilder = AckMessageResponse.newBuilder()
+ .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
+ .addAllEntries(entryList);
+ future.complete(responseBuilder.build());
+ });
} catch (Throwable t) {
future.completeExceptionally(t);
}
return future;
}
- protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, AckMessageRequest request) {
- return GrpcConverter.buildAckMessageRequestHeader(request);
+ protected CompletableFuture processAckMessage(Context ctx, AckMessageRequest request, AckMessageEntry ackMessageEntry) {
+ CompletableFuture future = new CompletableFuture<>();
+ AckMessageResultEntry.Builder failResult = AckMessageResultEntry.newBuilder()
+ .setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "ack message failed"))
+ .setMessageId(ackMessageEntry.getMessageId())
+ .setReceiptHandle(ackMessageEntry.getReceiptHandle());
+
+ try {
+ ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, ackMessageEntry.getReceiptHandle());
+ String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
+
+ AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request, receiptHandle);
+ CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader);
+ ackResultFuture
+ .thenAccept(result -> future.complete(convertToAckMessageResultEntry(ctx, ackMessageEntry, result)))
+ .exceptionally(throwable -> {
+ future.complete(failResult.setStatus(ResponseBuilder.buildStatus(throwable)).build());
+ return null;
+ });
+ } catch (Throwable t) {
+ future.complete(failResult.setStatus(ResponseBuilder.buildStatus(t)).build());
+ }
+ return future;
}
- protected AckMessageResponse convertToAckMessageResponse(Context ctx, AckMessageRequest request, AckResult ackResult) {
+ protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, AckMessageRequest request, ReceiptHandle handle) {
+ return GrpcConverter.buildAckMessageRequestHeader(request, handle);
+ }
+
+ protected AckMessageResultEntry convertToAckMessageResultEntry(Context ctx, AckMessageEntry ackMessageEntry, AckResult ackResult) {
if (AckStatus.OK.equals(ackResult.getStatus())) {
- return AckMessageResponse.newBuilder()
+ return AckMessageResultEntry.newBuilder()
+ .setMessageId(ackMessageEntry.getMessageId())
+ .setReceiptHandle(ackMessageEntry.getReceiptHandle())
.setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
.build();
}
- return AckMessageResponse.newBuilder()
+ return AckMessageResultEntry.newBuilder()
+ .setMessageId(ackMessageEntry.getMessageId())
+ .setReceiptHandle(ackMessageEntry.getReceiptHandle())
.setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "ack failed: status is abnormal"))
.build();
}
@@ -258,8 +289,7 @@ public class ConsumerService extends BaseService {
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
- Settings settings = grpcClientManager.getClientSettings(ctx).getSettings();
- int maxDeliveryAttempts = settings.getSubscription().getDeadLetterPolicy().getMaxDeliveryAttempts();
+ int maxDeliveryAttempts = ConfigurationManager.getProxyConfig().getDefaultMaxDeliveryAttempts();
if (request.getDeliveryAttempt() >= maxDeliveryAttempts) {
CompletableFuture resultFuture = this.producer.sendMessageBack(
brokerAddr,
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java
index cea463d2f1..b25a3f4934 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java
@@ -16,18 +16,15 @@
*/
package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
-import apache.rocketmq.v2.ClientOverwrittenSettings;
-import apache.rocketmq.v2.ClientSettings;
+import apache.rocketmq.v2.ActiveSubscriptionSettings;
+import apache.rocketmq.v2.ApplyPassiveSettingsCommand;
import apache.rocketmq.v2.Code;
-import apache.rocketmq.v2.Direction;
import apache.rocketmq.v2.HeartbeatRequest;
import apache.rocketmq.v2.HeartbeatResponse;
import apache.rocketmq.v2.NotifyClientTerminationRequest;
import apache.rocketmq.v2.NotifyClientTerminationResponse;
-import apache.rocketmq.v2.Publishing;
+import apache.rocketmq.v2.ReportActiveSettingsCommand;
import apache.rocketmq.v2.Resource;
-import apache.rocketmq.v2.Settings;
-import apache.rocketmq.v2.Subscription;
import apache.rocketmq.v2.TelemetryCommand;
import io.grpc.Context;
import io.grpc.stub.StreamObserver;
@@ -54,6 +51,7 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException;
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager;
+import org.apache.rocketmq.proxy.grpc.v2.service.ReportActiveSettingsService;
import org.apache.rocketmq.remoting.protocol.LanguageCode;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -66,6 +64,7 @@ public class ForwardClientService extends BaseService {
private final ProducerManager producerManager;
private final GrpcClientManager grpcClientManager;
private final TelemetryCommandManager telemetryCommandManager;
+ private final ReportActiveSettingsService reportActiveSettingsService;
public ForwardClientService(
ConnectorManager connectorManager,
@@ -84,6 +83,8 @@ public class ForwardClientService extends BaseService {
this.grpcClientManager = grpcClientManager;
this.telemetryCommandManager = telemetryCommandManager;
+ this.reportActiveSettingsService = new ReportActiveSettingsService(this.channelManager, this.grpcClientManager, this.telemetryCommandManager);
+
this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListenerImpl());
this.producerManager = new ProducerManager();
this.producerManager.appendProducerChangeListener(new ProducerChangeListenerImpl());
@@ -137,15 +138,16 @@ public class ForwardClientService extends BaseService {
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
LanguageCode languageCode = LanguageCode.valueOf(language);
- ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
+ GrpcClientManager.ActiveClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
switch (clientSettings.getClientType()) {
case PRODUCER: {
- for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) {
+ for (Resource topic : clientSettings.getActivePublishingSettings().getPublishingTopicsList()) {
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal());
// use topic name as producer group
producerManager.registerProducer(topicName, clientChannelInfo);
+ connectorManager.getTransactionHeartbeatRegisterService().addProducerGroup(topicName, topicName);
}
break;
}
@@ -165,9 +167,7 @@ public class ForwardClientService extends BaseService {
GrpcConverter.buildConsumeType(clientSettings.getClientType()),
MessageModel.CLUSTERING,
ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET,
- GrpcConverter.buildSubscriptionDataSet(clientSettings.getSettings()
- .getSubscription()
- .getSubscriptionsList()),
+ GrpcConverter.buildSubscriptionDataSet(clientSettings.getActiveSubscriptionSettings().getSubscriptionsList()),
false
);
break;
@@ -191,11 +191,11 @@ public class ForwardClientService extends BaseService {
try {
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
- ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
+ GrpcClientManager.ActiveClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
switch (clientSettings.getClientType()) {
case PRODUCER:
- for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) {
+ for (Resource topic : clientSettings.getActivePublishingSettings().getPublishingTopicsList()) {
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
// user topic name as producer group
GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, topicName, clientId);
@@ -229,37 +229,11 @@ public class ForwardClientService extends BaseService {
}
public StreamObserver telemetry(Context ctx, StreamObserver responseObserver) {
- String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
return new StreamObserver() {
@Override
public void onNext(TelemetryCommand request) {
- if (request.getCommandCase() == TelemetryCommand.CommandCase.CLIENT_SETTINGS) {
- ClientSettings clientSettings = request.getClientSettings();
- grpcClientManager.updateClientSettings(clientId, clientSettings);
- Settings settings = clientSettings.getSettings();
- if (settings.hasPublishing()) {
- Publishing publishing = settings.getPublishing();
- for (Resource topic : publishing.getTopicsList()) {
- String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
- // use topic name as producer group
- connectorManager.getTransactionHeartbeatRegisterService().addProducerGroup(topicName, topicName);
- GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
- producerChannel.setClientObserver(responseObserver);
- }
- }
- if (settings.hasSubscription()) {
- Subscription subscription = settings.getSubscription();
- String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup());
- GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, telemetryCommandManager);
- consumerChannel.setClientObserver(responseObserver);
- }
- responseObserver.onNext(TelemetryCommand.newBuilder()
- .setClientOverwrittenSettings(ClientOverwrittenSettings.newBuilder()
- .setNonce(clientSettings.getNonce())
- .setDirection(Direction.RESPONSE)
- .setSettings(settings)
- .build())
- .build());
+ if (request.getCommandCase() == TelemetryCommand.CommandCase.REPORT_ACTIVE_SETTINGS_COMMAND) {
+ responseObserver.onNext(reportActiveSettingsService.processReportActiveSettingsCommand(ctx, request, responseObserver));
}
}
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/PullMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/PullMessageService.java
deleted file mode 100644
index 91194ffb66..0000000000
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/PullMessageService.java
+++ /dev/null
@@ -1,160 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
-
-import apache.rocketmq.v2.Code;
-import apache.rocketmq.v2.Message;
-import apache.rocketmq.v2.MessageQueue;
-import apache.rocketmq.v2.PullMessageRequest;
-import apache.rocketmq.v2.PullMessageResponse;
-import apache.rocketmq.v2.QueryOffsetRequest;
-import apache.rocketmq.v2.QueryOffsetResponse;
-import com.google.protobuf.util.Timestamps;
-import io.grpc.Context;
-import java.util.List;
-import java.util.concurrent.CompletableFuture;
-import java.util.stream.Collectors;
-import org.apache.rocketmq.client.consumer.PullResult;
-import org.apache.rocketmq.client.consumer.PullStatus;
-import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
-import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
-import org.apache.rocketmq.proxy.common.utils.FilterUtils;
-import org.apache.rocketmq.proxy.connector.ConnectorManager;
-import org.apache.rocketmq.proxy.connector.DefaultForwardClient;
-import org.apache.rocketmq.proxy.connector.ForwardReadConsumer;
-import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
-import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
-import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook;
-
-public class PullMessageService extends BaseService {
- private final DefaultForwardClient forwardClient;
- private final ForwardReadConsumer readConsumer;
-
- private volatile ResponseHook queryOffsetHook;
- private volatile ResponseHook pullMessageHook;
-
- public PullMessageService(ConnectorManager connectorManager) {
- super(connectorManager);
- this.forwardClient = connectorManager.getDefaultForwardClient();
- this.readConsumer = connectorManager.getForwardReadConsumer();
- }
-
- public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) {
- CompletableFuture future = new CompletableFuture<>();
- future.whenComplete((response, throwable) -> {
- if (queryOffsetHook != null) {
- queryOffsetHook.beforeResponse(ctx, request, response, throwable);
- }
- });
-
- try {
- MessageQueue messageQueue = request.getMessageQueue();
- String topic = GrpcConverter.wrapResourceWithNamespace(messageQueue.getTopic());
- String brokerName = messageQueue.getBroker().getName();
- int queueId = messageQueue.getId();
-
- CompletableFuture offsetFuture;
- switch (request.getPolicy()) {
- case BEGINNING:
- offsetFuture = CompletableFuture.completedFuture(0L);
- break;
- case END:
- offsetFuture = this.forwardClient.getMaxOffset(this.getBrokerAddr(ctx, brokerName), topic, queueId);
- break;
- default:
- long timestamp = Timestamps.toMillis(request.getTimePoint());
- offsetFuture = this.forwardClient.searchOffset(this.getBrokerAddr(ctx, brokerName), topic, queueId, timestamp);
- }
-
- offsetFuture
- .thenAccept(result -> future.complete(
- QueryOffsetResponse.newBuilder()
- .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
- .setOffset(result)
- .build()))
- .exceptionally(throwable -> {
- future.completeExceptionally(throwable);
- return null;
- });
- } catch (Throwable t) {
- future.completeExceptionally(t);
- }
- return future;
- }
-
- public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) {
- CompletableFuture future = new CompletableFuture<>();
- future.whenComplete((response, throwable) -> {
- if (pullMessageHook != null) {
- pullMessageHook.beforeResponse(ctx, request, response, throwable);
- }
- });
-
- try {
- PullMessageRequestHeader requestHeader = this.buildPullMessageRequestHeader(ctx, request);
-
- String brokerName = request.getMessageQueue().getBroker().getName();
- String brokerAddr = this.getBrokerAddr(ctx, brokerName);
-
- CompletableFuture pullResultFuture = this.readConsumer.pullMessage(brokerAddr, requestHeader);
- pullResultFuture
- .thenAccept(pullResult -> {
- try {
- future.complete(convertToPullMessageResponse(ctx, request, pullResult));
- } catch (Throwable throwable) {
- future.completeExceptionally(throwable);
- }
- })
- .exceptionally(throwable -> {
- future.completeExceptionally(throwable);
- return null;
- });
-
- } catch (Throwable t) {
- future.completeExceptionally(t);
- }
- return future;
- }
-
- protected PullMessageRequestHeader buildPullMessageRequestHeader(Context ctx, PullMessageRequest request) {
- checkSubscriptionData(request.getMessageQueue().getTopic(), request.getFilterExpression());
- return GrpcConverter.buildPullMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx));
- }
-
- protected PullMessageResponse convertToPullMessageResponse(Context ctx, PullMessageRequest request, PullResult result) {
- PullMessageResponse.Builder responseBuilder = PullMessageResponse.newBuilder()
- .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name()))
- .setMinOffset(result.getMinOffset())
- .setMaxOffset(result.getMaxOffset())
- .setNextOffset(result.getNextBeginOffset());
-
- SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(
- GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic()), request.getFilterExpression());
-
- PullStatus status = result.getPullStatus();
- if (status.equals(PullStatus.FOUND)) {
- List messageList = result.getMsgFoundList().stream()
- .filter(msg -> FilterUtils.isTagMatched(subscriptionData.getTagsSet(), msg.getTags())) // only return tag matched messages.
- .map(GrpcConverter::buildMessage)
- .collect(Collectors.toList());
-
- return responseBuilder.addAllMessages(messageList).build();
- } else {
- return responseBuilder.build();
- }
- }
-}
diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/PullMessageServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/PullMessageServiceTest.java
deleted file mode 100644
index 23060b1b86..0000000000
--- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/PullMessageServiceTest.java
+++ /dev/null
@@ -1,131 +0,0 @@
-package org.apache.rocketmq.proxy.grpc.v2.service.cluster;
-
-import apache.rocketmq.v2.Broker;
-import apache.rocketmq.v2.Code;
-import apache.rocketmq.v2.FilterExpression;
-import apache.rocketmq.v2.FilterType;
-import apache.rocketmq.v2.MessageQueue;
-import apache.rocketmq.v2.PullMessageRequest;
-import apache.rocketmq.v2.PullMessageResponse;
-import apache.rocketmq.v2.QueryOffsetPolicy;
-import apache.rocketmq.v2.QueryOffsetRequest;
-import apache.rocketmq.v2.QueryOffsetResponse;
-import apache.rocketmq.v2.Resource;
-import com.google.protobuf.util.Timestamps;
-import io.grpc.Context;
-import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.Executors;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
-import org.apache.rocketmq.client.consumer.PullResult;
-import org.apache.rocketmq.client.consumer.PullStatus;
-import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
-import org.assertj.core.util.Lists;
-import org.junit.Test;
-
-import static org.junit.Assert.assertEquals;
-import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.ArgumentMatchers.anyInt;
-import static org.mockito.ArgumentMatchers.anyLong;
-import static org.mockito.ArgumentMatchers.anyString;
-import static org.mockito.Mockito.doAnswer;
-import static org.mockito.Mockito.when;
-
-public class PullMessageServiceTest extends BaseServiceTest {
-
- private PullMessageService pullMessageService;
-
- @Override
- public void beforeEach() throws Throwable {
- pullMessageService = new PullMessageService(this.connectorManager);
- when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr");
- }
-
- @Test
- public void testQueryOffset() throws Exception {
- Context ctx = Context.current();
-
- when(defaultClient.getMaxOffset(anyString(), anyString(), anyInt())).thenReturn(CompletableFuture.completedFuture(100L));
- when(defaultClient.searchOffset(anyString(), anyString(), anyInt(), anyLong())).thenReturn(CompletableFuture.completedFuture(50L));
-
- QueryOffsetResponse response = pullMessageService.queryOffset(ctx, QueryOffsetRequest.newBuilder()
- .setMessageQueue(MessageQueue.newBuilder()
- .setTopic(Resource.newBuilder()
- .setName("topic")
- .build())
- .setBroker(Broker.newBuilder().setName("brokerName").build())
- .build())
- .setPolicy(QueryOffsetPolicy.BEGINNING)
- .build()
- ).get();
- assertEquals(Code.OK, response.getStatus().getCode());
- assertEquals(0, response.getOffset());
-
- response = pullMessageService.queryOffset(ctx, QueryOffsetRequest.newBuilder()
- .setMessageQueue(MessageQueue.newBuilder()
- .setTopic(Resource.newBuilder()
- .setName("topic")
- .build())
- .setBroker(Broker.newBuilder().setName("brokerName").build())
- .build())
- .setPolicy(QueryOffsetPolicy.END)
- .build()
- ).get();
- assertEquals(Code.OK, response.getStatus().getCode());
- assertEquals(100, response.getOffset());
-
- response = pullMessageService.queryOffset(ctx, QueryOffsetRequest.newBuilder()
- .setMessageQueue(MessageQueue.newBuilder()
- .setTopic(Resource.newBuilder()
- .setName("topic")
- .build())
- .setBroker(Broker.newBuilder().setName("brokerName").build())
- .build())
- .setTimePoint(Timestamps.fromMillis(System.currentTimeMillis()))
- .setPolicy(QueryOffsetPolicy.TIME_POINT)
- .build()
- ).get();
- assertEquals(Code.OK, response.getStatus().getCode());
- assertEquals(50, response.getOffset());
- }
-
- @Test
- public void testPullMessage() throws Exception {
- AtomicReference headerRef = new AtomicReference<>();
- PullResult pullResult = new PullResult(
- PullStatus.FOUND,
- 3,
- 0,
- 10,
- Lists.newArrayList(
- createMessageExt("msg1", "msg1"),
- createMessageExt("msg2", "msg2")
- )
- );
- doAnswer(mock -> {
- headerRef.set(mock.getArgument(1));
- return CompletableFuture.completedFuture(pullResult);
- }).when(readConsumerClient).pullMessage(anyString(), any());
-
- Context ctx = Context.current().withDeadlineAfter(3, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor());
- PullMessageResponse response = pullMessageService.pullMessage(ctx, PullMessageRequest.newBuilder()
- .setMessageQueue(MessageQueue.newBuilder()
- .setBroker(Broker.newBuilder()
- .setName("brokerName")
- .build())
- .setTopic(Resource.newBuilder()
- .setName("topic")
- .build())
- .build())
- .setFilterExpression(FilterExpression.newBuilder()
- .setExpression("msg1")
- .setType(FilterType.TAG)
- .build())
- .build())
- .get();
-
- assertEquals(Code.OK, response.getStatus().getCode());
- assertEquals(1, response.getMessagesCount());
- assertEquals("msg1", response.getMessages(0).getSystemProperties().getMessageId());
- }
-}
\ No newline at end of file