From 80b2f95d1786f0c8f4440a10ae74f5ad9d3c666b Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Fri, 11 Mar 2022 16:43:41 +0800 Subject: [PATCH] [ISSUE #3949] Support sendMessage in Local mode --- .../rocketmq/common/message/MessageConst.java | 3 + .../proxy/grpc/adapter/InvocationContext.java | 48 ++++ .../grpc/adapter/channel/ChannelManager.java | 88 +++++++ .../adapter/channel/SendMessageChannel.java | 51 ++++ .../grpc/adapter/channel/SimpleChannel.java | 241 +++++++++++++++++ .../channel/SimpleChannelHandlerContext.java | 247 ++++++++++++++++++ .../grpc/adapter/handler/ResponseHandler.java | 25 ++ .../handler/SendMessageResponseHandler.java | 46 ++++ .../rocketmq/proxy/grpc/common/Converter.java | 131 ++++++++++ .../grpc/common/InterceptorConstants.java | 65 +++++ .../proxy/grpc/common/ResponseBuilder.java | 96 +++++++ .../proxy/grpc/common/ResponseWriter.java | 61 +++++ .../proxy/grpc/service/LocalGrpcService.java | 51 ++++ 13 files changed, 1153 insertions(+) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelManager.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannel.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannelHandlerContext.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ResponseHandler.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/SendMessageResponseHandler.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java diff --git a/common/src/main/java/org/apache/rocketmq/common/message/MessageConst.java b/common/src/main/java/org/apache/rocketmq/common/message/MessageConst.java index a823466415..d82cc07012 100644 --- a/common/src/main/java/org/apache/rocketmq/common/message/MessageConst.java +++ b/common/src/main/java/org/apache/rocketmq/common/message/MessageConst.java @@ -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. diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java new file mode 100644 index 0000000000..8b0739efd3 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java @@ -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 { + final private R request; + final private CompletableFuture response; + final private long timestamp = System.currentTimeMillis(); + + public InvocationContext(R req, CompletableFuture 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 getResponse() { + return response; + } + + public long getTimestamp() { + return timestamp; + } +} \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelManager.java new file mode 100644 index 0000000000..24dec4a302 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelManager.java @@ -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 { + private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME); + private final ConcurrentMap> clientIdChannelMap = new ConcurrentHashMap<>(); + + public SimpleChannel 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 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>> iterator = clientIdChannelMap.entrySet() + .iterator(); + while (iterator.hasNext()) { + Map.Entry> entry = iterator.next(); + if (!entry.getValue() + .isActive()) { + iterator.remove(); + } else { + entry.getValue() + .cleanExpiredRequests(); + } + } + } catch (Throwable e) { + LOGGER.error("Unexpected exception", e); + } + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java new file mode 100644 index 0000000000..646df00fc7 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java @@ -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 { + private final SendMessageResponseHandler handler; + + public static SendMessageChannel create(SimpleChannel other, SendMessageResponseHandler handler) { + return new SendMessageChannel(other, handler); + } + + private SendMessageChannel(SimpleChannel other, SendMessageResponseHandler handler) { + super(other); + this.handler = handler; + } + + @Override + public ChannelFuture writeAndFlush(Object msg) { + if (msg instanceof RemotingCommand) { + RemotingCommand responseCommand = (RemotingCommand) msg; + InvocationContext context = inFlightRequestMap.remove(responseCommand.getOpaque()); + if (null != context) { + handler.handle(responseCommand, context); + } + } + + return super.writeAndFlush(msg); + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannel.java new file mode 100644 index 0000000000..eb170ffbc1 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannel.java @@ -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 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> 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 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 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 context) { + inFlightRequestMap.put(opaque, context); + } + + public void eraseInvocationContext(int opaque) { + inFlightRequestMap.remove(opaque); + } + + public void cleanExpiredRequests() { + Iterator>> iterator = inFlightRequestMap.entrySet().iterator(); + int count = 0; + while (iterator.hasNext()) { + Map.Entry> 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 updateLastAccessTime() { + lastAccessTime = System.currentTimeMillis(); + return this; + } +} + diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannelHandlerContext.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannelHandlerContext.java new file mode 100644 index 0000000000..62d0c2a543 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SimpleChannelHandlerContext.java @@ -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 Attribute attr(AttributeKey key) { + throw new NotImplementedException("Not implemented"); + } + + + @Override + public boolean hasAttr(AttributeKey attributeKey) { + return false; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ResponseHandler.java new file mode 100644 index 0000000000..c7fbb5bf38 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ResponseHandler.java @@ -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 { + void handle(RemotingCommand responseCommand, InvocationContext context); +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/SendMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/SendMessageResponseHandler.java new file mode 100644 index 0000000000..9ff4bc33fe --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/SendMessageResponseHandler.java @@ -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 { + private final String messageId; + + public SendMessageResponseHandler(String messageId) { + this.messageId = messageId; + } + + @Override public void handle(RemotingCommand responseCommand, + InvocationContext 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); + } + } +} \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java index bb1737353f..bcdabe1663 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java @@ -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 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 buildMessageProperty(Message message) { + org.apache.rocketmq.common.message.Message messageWithHeader = new org.apache.rocketmq.common.message.Message(); + // set user properties + Map 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 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(); + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java new file mode 100644 index 0000000000..cb175d2ef9 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/InterceptorConstants.java @@ -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 = Context.key("rpc-metadata"); + + /** + * Remote address key in attributes of call + */ + public static final Metadata.Key REMOTE_ADDRESS + = Metadata.Key.of("rpc-remote-address", Metadata.ASCII_STRING_MARSHALLER); + + /** + * Local address key in attributes of call + */ + public static final Metadata.Key LOCAL_ADDRESS + = Metadata.Key.of("rpc-local-address", Metadata.ASCII_STRING_MARSHALLER); + + + public static final Metadata.Key AUTHORIZATION + = Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key NAMESPACE_ID + = Metadata.Key.of("x-mq-namespace", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key DATE_TIME + = Metadata.Key.of("x-mq-date-time", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key REQUEST_ID + = Metadata.Key.of("x-mq-request-id", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key LANGUAGE + = Metadata.Key.of("x-mq-language", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key CLIENT_VERSION + = Metadata.Key.of("x-mq-client-version", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key PROTOCOL_VERSION + = Metadata.Key.of("x-mq-protocol", Metadata.ASCII_STRING_MARSHALLER); + + public static final Metadata.Key RPC_NAME + = Metadata.Key.of("x-mq-rpc-name", Metadata.ASCII_STRING_MARSHALLER); +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java index 7265295efc..4a05227b80 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java @@ -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; + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java new file mode 100644 index 0000000000..d8744813a8 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java @@ -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 void write(StreamObserver observer, final T response) { + if (observer instanceof ServerCallStreamObserver) { + final ServerCallStreamObserver serverCallStreamObserver = (ServerCallStreamObserver) 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 void writeException(StreamObserver observer, final Exception e) { + if (observer instanceof ServerCallStreamObserver) { + final ServerCallStreamObserver serverCallStreamObserver = (ServerCallStreamObserver) 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(); + } + } +} + diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index f4f03ade44..c415fa8daf 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -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 sendChannelManager; public LocalGrpcService(BrokerController brokerController) { this.brokerController = brokerController; + this.sendChannelManager = new ChannelManager<>(); } @Override public CompleteFuture queryRoute(Context ctx, QueryRouteRequest request) { @@ -79,6 +99,31 @@ public class LocalGrpcService implements GrpcService { } @Override public CompleteFuture 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 future = new CompletableFuture<>(); + InvocationContext context + = new InvocationContext<>(request, future); + channel.registerInvocationContext(command.getOpaque(), context); + try { + CompletableFuture 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(); + } }