[ISSUE #3949] Support sendMessage in Local mode

This commit is contained in:
zhouxiang
2022-07-13 11:29:07 +08:00
parent ca3de6a315
commit 80b2f95d17
13 changed files with 1153 additions and 0 deletions
@@ -65,6 +65,9 @@ public class MessageConst {
public static final String PROPERTY_REDIRECT = "REDIRECT";
public static final String PROPERTY_INNER_MULTI_DISPATCH = "INNER_MULTI_DISPATCH";
public static final String PROPERTY_INNER_MULTI_QUEUE_OFFSET = "INNER_MULTI_QUEUE_OFFSET";
public static final String PROPERTY_TRACE_CONTEXT = "TRACE_CONTEXT";
public static final String PROPERTY_TIMER_DELAY_SEC = "TIMER_DELAY_SEC";
public static final String PROPERTY_TIMER_DELIVER_MS = "TIMER_DELIVER_MS";
/**
* property which name starts with "__RMQ.TRANSIENT." is called transient one that will not stored in broker disks.
@@ -0,0 +1,48 @@
/*
* 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.adapter;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public class InvocationContext<R, W> {
final private R request;
final private CompletableFuture<W> response;
final private long timestamp = System.currentTimeMillis();
public InvocationContext(R req, CompletableFuture<W> resp) {
request = req;
response = resp;
}
public boolean expired(long expiredTimeSec) {
return System.currentTimeMillis() - timestamp >= TimeUnit.SECONDS.toMillis(expiredTimeSec);
}
public R getRequest() {
return request;
}
public CompletableFuture<W> getResponse() {
return response;
}
public long getTimestamp() {
return timestamp;
}
}
@@ -0,0 +1,88 @@
/*
* 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.adapter.channel;
import com.google.common.base.Strings;
import io.grpc.Context;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.proxy.configuration.ConfigurationManager;
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class ChannelManager<R, W> {
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
private final ConcurrentMap<String, SimpleChannel<R, W>> clientIdChannelMap = new ConcurrentHashMap<>();
public SimpleChannel<R, W> createChannel() {
final String clientId = anonymousChannelId();
if (Strings.isNullOrEmpty(clientId)) {
LOGGER.warn("ClientId is unexpected null or empty");
return createChannelInner();
}
if (!clientIdChannelMap.containsKey(clientId)) {
clientIdChannelMap.putIfAbsent(clientId, createChannelInner());
}
return clientIdChannelMap.get(clientId)
.updateLastAccessTime();
}
private String anonymousChannelId() {
final String clientHost = InterceptorConstants.METADATA.get(Context.current())
.get(InterceptorConstants.REMOTE_ADDRESS);
final String localAddress = InterceptorConstants.METADATA.get(Context.current())
.get(InterceptorConstants.LOCAL_ADDRESS);
return clientHost + "@" + localAddress;
}
private SimpleChannel<R, W> createChannelInner() {
final String clientHost = InterceptorConstants.METADATA.get(Context.current())
.get(InterceptorConstants.REMOTE_ADDRESS);
final String localAddress = InterceptorConstants.METADATA.get(Context.current())
.get(InterceptorConstants.LOCAL_ADDRESS);
return new SimpleChannel<>(null, clientHost, localAddress, ConfigurationManager.getProxyConfig().getExpiredChannelTimeSec());
}
/**
* Scan and remove inactive mocking channels; Scan and clean expired requests;
*/
public void scanAndCleanChannels() {
try {
Iterator<Map.Entry<String, SimpleChannel<R, W>>> iterator = clientIdChannelMap.entrySet()
.iterator();
while (iterator.hasNext()) {
Map.Entry<String, SimpleChannel<R, W>> entry = iterator.next();
if (!entry.getValue()
.isActive()) {
iterator.remove();
} else {
entry.getValue()
.cleanExpiredRequests();
}
}
} catch (Throwable e) {
LOGGER.error("Unexpected exception", e);
}
}
}
@@ -0,0 +1,51 @@
/*
* 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.adapter.channel;
import apache.rocketmq.v1.SendMessageRequest;
import apache.rocketmq.v1.SendMessageResponse;
import io.netty.channel.ChannelFuture;
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
public class SendMessageChannel extends SimpleChannel<SendMessageRequest, SendMessageResponse> {
private final SendMessageResponseHandler handler;
public static SendMessageChannel create(SimpleChannel<SendMessageRequest, SendMessageResponse> other, SendMessageResponseHandler handler) {
return new SendMessageChannel(other, handler);
}
private SendMessageChannel(SimpleChannel<SendMessageRequest, SendMessageResponse> other, SendMessageResponseHandler handler) {
super(other);
this.handler = handler;
}
@Override
public ChannelFuture writeAndFlush(Object msg) {
if (msg instanceof RemotingCommand) {
RemotingCommand responseCommand = (RemotingCommand) msg;
InvocationContext<SendMessageRequest, SendMessageResponse> context = inFlightRequestMap.remove(responseCommand.getOpaque());
if (null != context) {
handler.handle(responseCommand, context);
}
}
return super.writeAndFlush(msg);
}
}
@@ -0,0 +1,241 @@
/*
* 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.adapter.channel;
import com.google.common.base.Strings;
import io.netty.channel.AbstractChannel;
import io.netty.channel.Channel;
import io.netty.channel.ChannelConfig;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelMetadata;
import io.netty.channel.ChannelOutboundBuffer;
import io.netty.channel.DefaultChannelPromise;
import io.netty.channel.EventLoop;
import io.netty.util.concurrent.GlobalEventExecutor;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* SimpleChannel is used to handle writeAndFlush situation in processor
* @see io.netty.channel.ChannelHandlerContext#writeAndFlush
* @see io.netty.channel.Channel#writeAndFlush
*/
public class SimpleChannel<R, W> extends AbstractChannel {
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
private final String remoteAddress;
private final String localAddress;
private final long expiredTimeSec;
private long lastAccessTime;
protected final ConcurrentMap<Integer, InvocationContext<R, W>> inFlightRequestMap;
/**
* Creates a new instance.
*
* @param parent the parent of this channel. {@code null} if there's no parent.
* @param remoteAddress Remote address
* @param localAddress Local address
* @param expiredTimeSec Expired time second for cleaning channel
*/
public SimpleChannel(Channel parent, String remoteAddress, String localAddress, long expiredTimeSec) {
super(parent);
lastAccessTime = System.currentTimeMillis();
this.remoteAddress = remoteAddress;
this.localAddress = localAddress;
this.inFlightRequestMap = new ConcurrentHashMap<>();
this.expiredTimeSec = expiredTimeSec;
}
public SimpleChannel(SimpleChannel<R, W> other) {
super(other);
lastAccessTime = System.currentTimeMillis();
this.remoteAddress = other.remoteAddress;
this.localAddress = other.localAddress;
this.inFlightRequestMap = other.inFlightRequestMap;
this.expiredTimeSec = other.expiredTimeSec;
}
@Override
protected AbstractUnsafe newUnsafe() {
return null;
}
@Override
protected boolean isCompatible(EventLoop loop) {
return false;
}
private static SocketAddress parseSocketAddress(String address) {
if (Strings.isNullOrEmpty(address)) {
return null;
}
String[] segments = address.split(":");
if (2 == segments.length) {
return new InetSocketAddress(segments[0], Integer.parseInt(segments[1]));
}
return null;
}
@Override
protected SocketAddress localAddress0() {
return parseSocketAddress(localAddress);
}
@Override
public SocketAddress localAddress() {
return localAddress0();
}
@Override
public SocketAddress remoteAddress() {
return remoteAddress0();
}
@Override
protected SocketAddress remoteAddress0() {
return parseSocketAddress(remoteAddress);
}
@Override
protected void doBind(SocketAddress localAddress) throws Exception {
}
@Override
protected void doDisconnect() throws Exception {
}
@Override
public ChannelFuture close() {
DefaultChannelPromise promise = new DefaultChannelPromise(this, GlobalEventExecutor.INSTANCE);
promise.setSuccess();
return promise;
}
@Override
protected void doClose() throws Exception {
}
@Override
protected void doBeginRead() throws Exception {
}
@Override
protected void doWrite(ChannelOutboundBuffer in) throws Exception {
}
public boolean isWritable(int opaque) {
if (!inFlightRequestMap.containsKey(opaque)) {
return false;
}
InvocationContext<R, W> invocationContext = inFlightRequestMap.get(opaque);
if (null != invocationContext) {
CompletableFuture<?> future = invocationContext.getResponse();
return null != future && !future.isCancelled() && !future.isCompletedExceptionally() && !future.isDone();
}
return false;
}
@Override
public ChannelConfig config() {
return null;
}
@Override
public boolean isOpen() {
return true;
}
@Override
public boolean isActive() {
return (System.currentTimeMillis() - lastAccessTime) <= 120L * 1000;
}
@Override
public ChannelMetadata metadata() {
return null;
}
@Override
public EventLoop eventLoop() {
return super.eventLoop();
}
@Override
public ChannelFuture writeAndFlush(Object msg) {
if (msg instanceof RemotingCommand) {
RemotingCommand responseCommand = (RemotingCommand) msg;
inFlightRequestMap.remove(responseCommand.getOpaque());
}
DefaultChannelPromise promise = new DefaultChannelPromise(this, GlobalEventExecutor.INSTANCE);
promise.setSuccess();
return promise;
}
public void registerInvocationContext(int opaque, InvocationContext<R, W> context) {
inFlightRequestMap.put(opaque, context);
}
public void eraseInvocationContext(int opaque) {
inFlightRequestMap.remove(opaque);
}
public void cleanExpiredRequests() {
Iterator<Map.Entry<Integer, InvocationContext<R, W>>> iterator = inFlightRequestMap.entrySet().iterator();
int count = 0;
while (iterator.hasNext()) {
Map.Entry<Integer, InvocationContext<R, W>> entry = iterator.next();
if (entry.getValue().expired(expiredTimeSec)) {
iterator.remove();
count++;
LOGGER.debug("An expired request is found, created time-point: {}, Request: {}",
entry.getValue().getTimestamp(), entry.getValue().getRequest());
}
}
if (count > 0) {
LOGGER.warn("[BUG] {} expired in-flight requests is cleaned.", count);
}
}
public SimpleChannel<R, W> updateLastAccessTime() {
lastAccessTime = System.currentTimeMillis();
return this;
}
}
@@ -0,0 +1,247 @@
/*
* 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.adapter.channel;
import io.netty.buffer.ByteBufAllocator;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.ChannelProgressivePromise;
import io.netty.channel.ChannelPromise;
import io.netty.util.Attribute;
import io.netty.util.AttributeKey;
import io.netty.util.concurrent.EventExecutor;
import java.net.SocketAddress;
import org.apache.commons.lang3.NotImplementedException;
public class SimpleChannelHandlerContext implements ChannelHandlerContext {
private final Channel channel;
public SimpleChannelHandlerContext(Channel channel) {
this.channel = channel;
}
@Override
public Channel channel() {
return channel;
}
@Override
public EventExecutor executor() {
throw new NotImplementedException("Not implemented");
}
@Override
public String name() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandler handler() {
throw new NotImplementedException("Not implemented");
}
@Override
public boolean isRemoved() {
return false;
}
@Override
public ChannelHandlerContext fireChannelRegistered() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireChannelUnregistered() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireChannelActive() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireChannelInactive() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireExceptionCaught(Throwable cause) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireUserEventTriggered(Object evt) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireChannelRead(Object msg) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireChannelReadComplete() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext fireChannelWritabilityChanged() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture bind(SocketAddress localAddress) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture connect(SocketAddress remoteAddress) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture connect(SocketAddress remoteAddress, SocketAddress localAddress) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture disconnect() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture close() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture deregister() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture bind(SocketAddress localAddress, ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture connect(SocketAddress remoteAddress, ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture connect(SocketAddress remoteAddress, SocketAddress localAddress, ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture disconnect(ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture close(ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture deregister(ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext read() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture write(Object msg) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture write(Object msg, ChannelPromise promise) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelHandlerContext flush() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture writeAndFlush(Object msg, ChannelPromise promise) {
return channel.writeAndFlush(msg, promise);
}
@Override
public ChannelFuture writeAndFlush(Object msg) {
return channel.writeAndFlush(msg);
}
@Override
public ChannelPipeline pipeline() {
throw new NotImplementedException("Not implemented");
}
@Override
public ByteBufAllocator alloc() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelPromise newPromise() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelProgressivePromise newProgressivePromise() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture newSucceededFuture() {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelFuture newFailedFuture(Throwable cause) {
throw new NotImplementedException("Not implemented");
}
@Override
public ChannelPromise voidPromise() {
throw new NotImplementedException("Not implemented");
}
@Override
public <T> Attribute<T> attr(AttributeKey<T> key) {
throw new NotImplementedException("Not implemented");
}
@Override
public <T> boolean hasAttr(AttributeKey<T> attributeKey) {
return false;
}
}
@@ -0,0 +1,25 @@
/*
* 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.adapter.handler;
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
public interface ResponseHandler<R, W> {
void handle(RemotingCommand responseCommand, InvocationContext<R, W> context);
}
@@ -0,0 +1,46 @@
/*
* 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.adapter.handler;
import apache.rocketmq.v1.SendMessageRequest;
import apache.rocketmq.v1.SendMessageResponse;
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
public class SendMessageResponseHandler implements ResponseHandler<SendMessageRequest, SendMessageResponse> {
private final String messageId;
public SendMessageResponseHandler(String messageId) {
this.messageId = messageId;
}
@Override public void handle(RemotingCommand responseCommand,
InvocationContext<SendMessageRequest, SendMessageResponse> context) {
// If responseCommand equals to null, then the response has been written to channel.
// org.apache.rocketmq.broker.processor.SendMessageProcessor#handlePutMessageResult
// org.apache.rocketmq.broker.processor.AbstractSendMessageProcessor#doResponse
if (null != responseCommand) {
SendMessageResponse response = ResponseBuilder.buildSendMessageResponse(responseCommand);
response = response.toBuilder()
.setMessageId(messageId)
.build();
context.getResponse().complete(response);
}
}
}
@@ -17,5 +17,136 @@
package org.apache.rocketmq.proxy.grpc.common;
import apache.rocketmq.v1.Encoding;
import apache.rocketmq.v1.Message;
import apache.rocketmq.v1.MessageType;
import apache.rocketmq.v1.Resource;
import apache.rocketmq.v1.SendMessageRequest;
import apache.rocketmq.v1.SystemAttribute;
import com.google.common.collect.Maps;
import com.google.protobuf.Duration;
import com.google.protobuf.Timestamp;
import com.google.protobuf.util.Durations;
import com.google.protobuf.util.Timestamps;
import java.util.List;
import java.util.Map;
import org.apache.rocketmq.common.message.MessageAccessor;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.protocol.NamespaceUtil;
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
public class Converter {
public static String getResourceNameWithNamespace(Resource resource) {
return NamespaceUtil.wrapNamespace(resource.getResourceNamespace(), resource.getName());
}
public static SendMessageRequestHeader buildSendMessageRequestHeader(SendMessageRequest request) {
SendMessageRequestHeader requestHeader = new SendMessageRequestHeader();
Message message = request.getMessage();
SystemAttribute systemAttribute = message.getSystemAttribute();
Map<String, String> property = buildMessageProperty(message);
requestHeader.setProducerGroup(getResourceNameWithNamespace(systemAttribute.getProducerGroup()));
requestHeader.setTopic(getResourceNameWithNamespace(message.getTopic()));
requestHeader.setDefaultTopic("");
requestHeader.setDefaultTopicQueueNums(0);
requestHeader.setQueueId(systemAttribute.getPartitionId());
// sysFlag (body encoding & message type)
int sysFlag = 0;
Encoding bodyEncoding = systemAttribute.getBodyEncoding();
if (bodyEncoding.equals(Encoding.GZIP)) {
sysFlag |= MessageSysFlag.COMPRESSED_FLAG;
}
// transaction
MessageType messageType = systemAttribute.getMessageType();
if (messageType.equals(MessageType.TRANSACTION)) {
sysFlag |= MessageSysFlag.TRANSACTION_PREPARED_TYPE;
}
requestHeader.setSysFlag(sysFlag);
requestHeader.setBornTimestamp(Timestamps.toMillis(systemAttribute.getBornTimestamp()));
requestHeader.setFlag(0);
requestHeader.setProperties(MessageDecoder.messageProperties2String(property));
requestHeader.setReconsumeTimes(systemAttribute.getDeliveryAttempt());
return requestHeader;
}
public static Map<String, String> buildMessageProperty(Message message) {
org.apache.rocketmq.common.message.Message messageWithHeader = new org.apache.rocketmq.common.message.Message();
// set user properties
Map<String, String> userProperties = message.getUserAttributeMap();
for (String key : userProperties.keySet()) {
if (MessageConst.STRING_HASH_SET.contains(key)) {
throw new IllegalArgumentException("Property is used by system: " + key);
}
}
MessageAccessor.setProperties(messageWithHeader, Maps.newHashMap(userProperties));
// set tag
String tag = message.getSystemAttribute().getTag();
if (!"".equals(tag)) {
messageWithHeader.setTags(tag);
}
// set keys
List<String> keysList = message.getSystemAttribute().getKeysList();
if (keysList.size() > 0) {
messageWithHeader.setKeys(keysList);
}
// set message id
String messageId = message.getSystemAttribute().getMessageId();
if ("".equals(messageId)) {
throw new IllegalArgumentException("message id is empty");
}
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX, messageId);
// set transaction property
MessageType messageType = message.getSystemAttribute().getMessageType();
if (messageType.equals(MessageType.TRANSACTION)) {
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_TRANSACTION_PREPARED, "true");
Duration transactionResolveDelay = message.getSystemAttribute().getOrphanedTransactionRecoveryPeriod();
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_TRANSACTION_CHECK_TIMES, String.valueOf(15));
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_CHECK_IMMUNITY_TIME_IN_SECONDS,
String.valueOf(Durations.toSeconds(transactionResolveDelay)));
}
// set delay level or deliver timestamp
switch (message.getSystemAttribute().getTimedDeliveryCase()) {
case DELAY_LEVEL:
int delayLevel = message.getSystemAttribute().getDelayLevel();
if (delayLevel > 0) {
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_DELAY_TIME_LEVEL,
String.valueOf(delayLevel));
}
break;
case DELIVERY_TIMESTAMP:
Timestamp deliveryTimestamp = message.getSystemAttribute().getDeliveryTimestamp();
String timestampString = String.valueOf(Timestamps.toMillis(deliveryTimestamp));
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_TIMER_DELIVER_MS, timestampString);
break;
case TIMEDDELIVERY_NOT_SET:
break;
default:
throw new IllegalStateException("Unexpected value: " + message.getSystemAttribute().getTimedDeliveryCase());
}
// set reconsume times
int reconsumeTimes = message.getSystemAttribute().getDeliveryAttempt();
MessageAccessor.setReconsumeTime(messageWithHeader, String.valueOf(reconsumeTimes));
// set producer group
Resource producerGroup = message.getSystemAttribute().getProducerGroup();
String producerGroupName = getResourceNameWithNamespace(producerGroup);
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_PRODUCER_GROUP, producerGroupName);
// set message group
String messageGroup = message.getSystemAttribute().getMessageGroup();
if (!messageGroup.isEmpty()) {
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_SHARDING_KEY, messageGroup);
}
// set trace context
String traceContext = message.getSystemAttribute().getTraceContext();
if (!traceContext.isEmpty()) {
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_TRACE_CONTEXT, traceContext);
}
return messageWithHeader.getProperties();
}
}
@@ -0,0 +1,65 @@
/*
* 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.common;
import io.grpc.Context;
import io.grpc.Metadata;
public class InterceptorConstants {
private InterceptorConstants() {
}
public static final Context.Key<Metadata> METADATA = Context.key("rpc-metadata");
/**
* Remote address key in attributes of call
*/
public static final Metadata.Key<String> REMOTE_ADDRESS
= Metadata.Key.of("rpc-remote-address", Metadata.ASCII_STRING_MARSHALLER);
/**
* Local address key in attributes of call
*/
public static final Metadata.Key<String> LOCAL_ADDRESS
= Metadata.Key.of("rpc-local-address", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> AUTHORIZATION
= Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> NAMESPACE_ID
= Metadata.Key.of("x-mq-namespace", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> DATE_TIME
= Metadata.Key.of("x-mq-date-time", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> REQUEST_ID
= Metadata.Key.of("x-mq-request-id", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> LANGUAGE
= Metadata.Key.of("x-mq-language", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> CLIENT_VERSION
= Metadata.Key.of("x-mq-client-version", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> PROTOCOL_VERSION
= Metadata.Key.of("x-mq-protocol", Metadata.ASCII_STRING_MARSHALLER);
public static final Metadata.Key<String> RPC_NAME
= Metadata.Key.of("x-mq-rpc-name", Metadata.ASCII_STRING_MARSHALLER);
}
@@ -18,10 +18,26 @@
package org.apache.rocketmq.proxy.grpc.common;
import apache.rocketmq.v1.ResponseCommon;
import apache.rocketmq.v1.SendMessageResponse;
import com.google.rpc.Code;
import com.google.rpc.Status;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.common.protocol.ResponseCode;
import org.apache.rocketmq.common.protocol.header.SendMessageResponseHeader;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
public class ResponseBuilder {
public static ResponseCommon buildCommon(int responseCode, String remark) {
Status status = Status.newBuilder()
.setCode(buildCode(responseCode).getNumber())
.setMessage(buildMessage(responseCode, remark))
.build();
return ResponseCommon.newBuilder()
.setStatus(status)
.build();
}
public static ResponseCommon buildCommon(Code code, String message) {
Status status = Status.newBuilder()
.setCode(code.getNumber())
@@ -32,4 +48,84 @@ public class ResponseBuilder {
.setStatus(status)
.build();
}
public static SendMessageResponse buildSendMessageResponse(RemotingCommand command) {
SendMessageResponseHeader responseHeader = (SendMessageResponseHeader) command.readCustomHeader();
return SendMessageResponse.newBuilder()
.setCommon(buildCommon(command.getCode(), command.getRemark()))
.setMessageId(StringUtils.defaultString(responseHeader.getMsgId()))
.setTransactionId(StringUtils.defaultString(responseHeader.getTransactionId()))
.build();
}
public static Code buildCode(int responseCode) {
Code code;
switch (responseCode) {
case ResponseCode.SUCCESS:
case ResponseCode.NO_MESSAGE: {
code = Code.OK;
break;
}
case ResponseCode.SYSTEM_ERROR: {
code = Code.INTERNAL;
break;
}
case ResponseCode.SYSTEM_BUSY:
case ResponseCode.POLLING_FULL: {
code = Code.RESOURCE_EXHAUSTED;
break;
}
case ResponseCode.REQUEST_CODE_NOT_SUPPORTED: {
code = Code.UNIMPLEMENTED;
break;
}
case ResponseCode.MESSAGE_ILLEGAL:
case ResponseCode.VERSION_NOT_SUPPORTED:
case ResponseCode.SUBSCRIPTION_PARSE_FAILED:
case ResponseCode.FILTER_DATA_NOT_EXIST: {
code = Code.INVALID_ARGUMENT;
break;
}
case ResponseCode.SERVICE_NOT_AVAILABLE:
case ResponseCode.SLAVE_NOT_AVAILABLE:
case ResponseCode.PULL_RETRY_IMMEDIATELY:
case ResponseCode.PULL_OFFSET_MOVED:
case ResponseCode.SUBSCRIPTION_NOT_LATEST:
case ResponseCode.FILTER_DATA_NOT_LATEST: {
code = Code.UNAVAILABLE;
break;
}
case ResponseCode.NO_PERMISSION: {
code = Code.PERMISSION_DENIED;
break;
}
case ResponseCode.TOPIC_NOT_EXIST:
case ResponseCode.SUBSCRIPTION_GROUP_NOT_EXIST:
case ResponseCode.SUBSCRIPTION_NOT_EXIST:
case ResponseCode.PULL_NOT_FOUND:
case ResponseCode.QUERY_NOT_FOUND:
case ResponseCode.CONSUMER_NOT_ONLINE: {
code = Code.NOT_FOUND;
break;
}
case ResponseCode.POLLING_TIMEOUT:
case ResponseCode.FLUSH_DISK_TIMEOUT:
case ResponseCode.FLUSH_SLAVE_TIMEOUT: {
code = Code.DEADLINE_EXCEEDED;
break;
}
default: {
code = Code.UNKNOWN;
}
}
return code;
}
public static String buildMessage(int responseCode, String remark) {
if (remark != null) {
return "ResponseCode: " + responseCode + " " + remark;
}
return "ResponseCode: " + responseCode;
}
}
@@ -0,0 +1,61 @@
/*
* 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.common;
import io.grpc.stub.ServerCallStreamObserver;
import io.grpc.stub.StreamObserver;
import org.apache.rocketmq.common.constant.LoggerName;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class ResponseWriter {
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
public static <T> void write(StreamObserver<T> observer, final T response) {
if (observer instanceof ServerCallStreamObserver) {
final ServerCallStreamObserver<T> serverCallStreamObserver = (ServerCallStreamObserver<T>) observer;
if (serverCallStreamObserver.isCancelled()) {
LOGGER.warn("client has cancelled the request. response to write: {}", response);
return;
}
LOGGER.debug("start to write response. response: {}", response);
serverCallStreamObserver.onNext(response);
serverCallStreamObserver.onCompleted();
}
}
public static <T> void writeException(StreamObserver<T> observer, final Exception e) {
if (observer instanceof ServerCallStreamObserver) {
final ServerCallStreamObserver<T> serverCallStreamObserver = (ServerCallStreamObserver<T>) observer;
if (null == e) {
return;
}
if (serverCallStreamObserver.isCancelled()) {
LOGGER.warn("Client has cancelled the request. Exception to write", e);
return;
}
LOGGER.debug("Start to write error response", e);
serverCallStreamObserver.onError(e);
serverCallStreamObserver.onCompleted();
}
}
}
@@ -29,6 +29,7 @@ import apache.rocketmq.v1.HealthCheckRequest;
import apache.rocketmq.v1.HealthCheckResponse;
import apache.rocketmq.v1.HeartbeatRequest;
import apache.rocketmq.v1.HeartbeatResponse;
import apache.rocketmq.v1.Message;
import apache.rocketmq.v1.NackMessageRequest;
import apache.rocketmq.v1.NackMessageResponse;
import apache.rocketmq.v1.NotifyClientTerminationRequest;
@@ -53,17 +54,36 @@ import apache.rocketmq.v1.SendMessageRequest;
import apache.rocketmq.v1.SendMessageResponse;
import io.grpc.Context;
import io.netty.util.concurrent.CompleteFuture;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import org.apache.rocketmq.broker.BrokerController;
import org.apache.rocketmq.common.ThreadFactoryImpl;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.common.protocol.RequestCode;
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
import org.apache.rocketmq.proxy.grpc.adapter.channel.ChannelManager;
import org.apache.rocketmq.proxy.grpc.adapter.channel.SendMessageChannel;
import org.apache.rocketmq.proxy.grpc.adapter.channel.SimpleChannelHandlerContext;
import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler;
import org.apache.rocketmq.proxy.grpc.common.Converter;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class LocalGrpcService implements GrpcService {
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
private final BrokerController brokerController;
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(
new ThreadFactoryImpl("LocalGrpcServiceScheduledThread"));
private final ChannelManager<SendMessageRequest, SendMessageResponse> sendChannelManager;
public LocalGrpcService(BrokerController brokerController) {
this.brokerController = brokerController;
this.sendChannelManager = new ChannelManager<>();
}
@Override public CompleteFuture<QueryRouteResponse> queryRoute(Context ctx, QueryRouteRequest request) {
@@ -79,6 +99,31 @@ public class LocalGrpcService implements GrpcService {
}
@Override public CompleteFuture<SendMessageResponse> sendMessage(Context ctx, SendMessageRequest request) {
SendMessageRequestHeader requestHeader = Converter.buildSendMessageRequestHeader(request);
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE, requestHeader);
Message message = request.getMessage();
command.setBody(message.getBody().toByteArray());
command.makeCustomHeaderToNet();
SendMessageResponseHandler handler = new SendMessageResponseHandler(message.getSystemAttribute().getMessageId());
SendMessageChannel channel = SendMessageChannel.create(sendChannelManager.createChannel(), handler);
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
CompletableFuture<SendMessageResponse> future = new CompletableFuture<>();
InvocationContext<SendMessageRequest, SendMessageResponse> context
= new InvocationContext<>(request, future);
channel.registerInvocationContext(command.getOpaque(), context);
try {
CompletableFuture<RemotingCommand> processorFuture = brokerController.getSendMessageProcessor()
.asyncProcessRequest(channelHandlerContext, command);
processorFuture.thenAccept(r -> {
handler.handle(r, context);
channel.eraseInvocationContext(command.getOpaque());
});
} catch (final RemotingCommandException e) {
LOGGER.error("Failed to process send message command", e);
channel.eraseInvocationContext(command.getOpaque());
future.completeExceptionally(e);
}
return null;
}
@@ -143,9 +188,15 @@ public class LocalGrpcService implements GrpcService {
@Override public void start() throws Exception {
this.brokerController.start();
this.scheduledExecutorService.scheduleWithFixedDelay(this::scanAndCleanChannels, 5, 5, TimeUnit.MINUTES);
}
@Override public void shutdown() throws Exception {
this.scheduledExecutorService.shutdown();
this.brokerController.shutdown();
}
private void scanAndCleanChannels() {
this.sendChannelManager.scanAndCleanChannels();
}
}