[ISSUE #3949] v2 support

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:18 +08:00
committed by zhouxiang
parent fc492d31de
commit 4db626bdb0
14 changed files with 309 additions and 555 deletions
+1
View File
@@ -455,6 +455,7 @@
<groupId>${project.groupId}</groupId>
<artifactId>rocketmq-proto</artifactId>
<version>2.0.0-SNAPSHOT</version>
<classifier>compatible</classifier>
</dependency>
<dependency>
<groupId>${project.groupId}</groupId>
@@ -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;
}
@@ -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<QueryOffsetResponse> responseObserver) {
CompletableFuture<QueryOffsetResponse> 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<PullMessageResponse> responseObserver) {
CompletableFuture<PullMessageResponse> 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<NotifyClientTerminationResponse> responseObserver) {
@@ -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<org.apache.rocketmq.common.protocol.heartbeat.ProducerData> 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<SubscriptionData> 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<String, String> buildUserAttributes(MessageExt messageExt) {
Map<String, String> userAttributes = new HashMap<>();
Map<String, String> properties = messageExt.getProperties();
@@ -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)
@@ -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<QueryOffsetResponse> queryOffset(Context ctx, QueryOffsetRequest request) {
return pullMessageService.queryOffset(ctx, request);
}
@Override
public CompletableFuture<PullMessageResponse> pullMessage(Context ctx, PullMessageRequest request) {
return pullMessageService.pullMessage(ctx, request);
}
@Override
public CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx,
NotifyClientTerminationRequest request) {
@@ -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<String, ClientSettings> 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<String, ActiveClientSettings> 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);
}
}
@@ -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<EndTransactionResponse> endTransaction(Context ctx, EndTransactionRequest request);
CompletableFuture<QueryOffsetResponse> queryOffset(Context ctx, QueryOffsetRequest request);
CompletableFuture<PullMessageResponse> pullMessage(Context ctx, PullMessageRequest request);
CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx, NotifyClientTerminationRequest request);
CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request);
@@ -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<HeartbeatResponse> 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<ReceiveMessageResponse> 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<ReceiveMessageResponse> future = new CompletableFuture<>();
@@ -322,8 +325,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
CompletableFuture<NackMessageResponse> 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<QueryOffsetResponse> 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<PullMessageResponse> 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<PullMessageResponse> future = new CompletableFuture<>();
InvocationContext<PullMessageRequest, PullMessageResponse> 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: {
@@ -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<TelemetryCommand> 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();
}
}
@@ -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<AckResult> 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<AckMessageResultEntry>[] 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<AckMessageResultEntry> entryList = new ArrayList<>();
for (CompletableFuture<AckMessageResultEntry> 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<AckMessageResultEntry> processAckMessage(Context ctx, AckMessageRequest request, AckMessageEntry ackMessageEntry) {
CompletableFuture<AckMessageResultEntry> 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<AckResult> 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<RemotingCommand> resultFuture = this.producer.sendMessageBack(
brokerAddr,
@@ -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<TelemetryCommand> telemetry(Context ctx, StreamObserver<TelemetryCommand> responseObserver) {
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
return new StreamObserver<TelemetryCommand>() {
@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));
}
}
@@ -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<QueryOffsetRequest, QueryOffsetResponse> queryOffsetHook;
private volatile ResponseHook<PullMessageRequest, PullMessageResponse> pullMessageHook;
public PullMessageService(ConnectorManager connectorManager) {
super(connectorManager);
this.forwardClient = connectorManager.getDefaultForwardClient();
this.readConsumer = connectorManager.getForwardReadConsumer();
}
public CompletableFuture<QueryOffsetResponse> queryOffset(Context ctx, QueryOffsetRequest request) {
CompletableFuture<QueryOffsetResponse> 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<Long> 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<PullMessageResponse> pullMessage(Context ctx, PullMessageRequest request) {
CompletableFuture<PullMessageResponse> 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<PullResult> 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<Message> 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();
}
}
}
@@ -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<PullMessageRequestHeader> 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());
}
}