From f91a4c0ee283ef915a4e002f7ffa97c6f243c61d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B0=98=E5=A4=AE?= Date: Tue, 17 May 2022 17:27:05 +0800 Subject: [PATCH] Rename ProxyOutService to ProxyRelayService --- .../grpc/v2/channel/GrpcChannelManager.java | 12 ++++++------ .../proxy/grpc/v2/channel/GrpcClientChannel.java | 10 +++++----- .../proxy/grpc/v2/client/ClientActivity.java | 2 +- .../processor/DefaultMessagingProcessor.java | 5 ++--- .../proxy/processor/MessagingProcessor.java | 4 ++-- .../proxy/service/ClusterServiceManager.java | 12 ++++++------ .../proxy/service/LocalServiceManager.java | 12 ++++++------ .../rocketmq/proxy/service/ServiceManager.java | 4 ++-- .../ClusterProxyRelayService.java} | 4 ++-- .../LocalProxyRelayService.java} | 6 +++--- .../service/{out => relay}/ProxyChannel.java | 16 ++++++++-------- .../service/{out => relay}/ProxyOutResult.java | 2 +- .../ProxyRelayService.java} | 4 ++-- .../rocketmq/proxy/grpc/v2/BaseActivityTest.java | 7 +++---- .../proxy/processor/BaseProcessorTest.java | 6 +++--- 15 files changed, 52 insertions(+), 54 deletions(-) rename proxy/src/main/java/org/apache/rocketmq/proxy/service/{out/ClusterProxyOutService.java => relay/ClusterProxyRelayService.java} (93%) rename proxy/src/main/java/org/apache/rocketmq/proxy/service/{out/LocalProxyOutService.java => relay/LocalProxyRelayService.java} (94%) rename proxy/src/main/java/org/apache/rocketmq/proxy/service/{out => relay}/ProxyChannel.java (91%) rename proxy/src/main/java/org/apache/rocketmq/proxy/service/{out => relay}/ProxyOutResult.java (96%) rename proxy/src/main/java/org/apache/rocketmq/proxy/service/{out/ProxyOutService.java => relay/ProxyRelayService.java} (94%) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java index b5bb417ba6..bb7f4b203f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java @@ -33,11 +33,11 @@ import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; -import org.apache.rocketmq.proxy.service.out.ProxyOutResult; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ProxyOutResult; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; public class GrpcChannelManager implements StartAndShutdown { - private final ProxyOutService proxyOutService; + private final ProxyRelayService proxyRelayService; protected final ConcurrentMap/* clientId */> groupClientIdChannelMap = new ConcurrentHashMap<>(); protected final AtomicLong nonceIdGenerator = new AtomicLong(0); @@ -47,8 +47,8 @@ public class GrpcChannelManager implements StartAndShutdown { new ThreadFactoryImpl("GrpcChannelManager_") ); - public GrpcChannelManager(ProxyOutService proxyOutService) { - this.proxyOutService = proxyOutService; + public GrpcChannelManager(ProxyRelayService proxyRelayService) { + this.proxyRelayService = proxyRelayService; } protected void init() { @@ -63,7 +63,7 @@ public class GrpcChannelManager implements StartAndShutdown { if (clientIdMap == null) { clientIdMap = new ConcurrentHashMap<>(); } - clientIdMap.computeIfAbsent(clientId, clientIdKey -> new GrpcClientChannel(proxyOutService, this, ctx, group, clientId)); + clientIdMap.computeIfAbsent(clientId, clientIdKey -> new GrpcClientChannel(proxyRelayService, this, ctx, group, clientId)); return clientIdMap; }); return getChannel(group, clientId); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java index 1dbed8eeed..2ec136c5b0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java @@ -35,9 +35,9 @@ import org.apache.rocketmq.common.protocol.header.ConsumeMessageDirectlyResultRe import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestHeader; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter; -import org.apache.rocketmq.proxy.service.out.ProxyChannel; -import org.apache.rocketmq.proxy.service.out.ProxyOutResult; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ProxyChannel; +import org.apache.rocketmq.proxy.service.relay.ProxyOutResult; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.transaction.TransactionId; import org.apache.rocketmq.remoting.common.RemotingUtil; import org.apache.rocketmq.remoting.protocol.RemotingCommand; @@ -54,8 +54,8 @@ public class GrpcClientChannel extends ProxyChannel { private final String remoteAddress; private final String localAddress; - public GrpcClientChannel(ProxyOutService proxyOutService, GrpcChannelManager grpcChannelManager, Context ctx, String group, String clientId) { - super(proxyOutService, null, new GrpcChannelId(group, clientId)); + public GrpcClientChannel(ProxyRelayService proxyRelayService, GrpcChannelManager grpcChannelManager, Context ctx, String group, String clientId) { + super(proxyRelayService, null, new GrpcChannelId(group, clientId)); this.grpcChannelManager = grpcChannelManager; this.remoteAddress = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.REMOTE_ADDRESS); this.localAddress = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.LOCAL_ADDRESS); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java index c0fcdf016e..b707351ed6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java @@ -64,7 +64,7 @@ import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcProxyException; import org.apache.rocketmq.proxy.grpc.v2.common.ResponseBuilder; import org.apache.rocketmq.proxy.processor.MessagingProcessor; -import org.apache.rocketmq.proxy.service.out.ProxyOutResult; +import org.apache.rocketmq.proxy.service.relay.ProxyOutResult; import org.apache.rocketmq.remoting.protocol.LanguageCode; public class ClientActivity extends AbstractMessingActivity { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java index a590da8ed6..3e5d82c866 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java @@ -40,11 +40,10 @@ import org.apache.rocketmq.common.thread.ThreadPoolMonitor; import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.common.Address; import org.apache.rocketmq.proxy.common.ProxyContext; -import org.apache.rocketmq.proxy.common.utils.FutureUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.service.ServiceManager; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.ProxyTopicRouteData; import org.apache.rocketmq.proxy.service.transaction.TransactionId; import org.apache.rocketmq.remoting.RPCHook; @@ -216,7 +215,7 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen } @Override - public ProxyOutService getProxyOutService() { + public ProxyRelayService getProxyOutService() { return this.serviceManager.getProxyOutService(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java index 3f29fe2f6c..10ac52c1c0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java @@ -37,7 +37,7 @@ import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.Address; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.common.StartAndShutdown; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.ProxyTopicRouteData; import org.apache.rocketmq.proxy.service.transaction.TransactionId; import org.apache.rocketmq.remoting.exception.RemotingException; @@ -219,5 +219,5 @@ public interface MessagingProcessor extends StartAndShutdown { String topic ); - ProxyOutService getProxyOutService(); + ProxyRelayService getProxyOutService(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java index 935f25a8ae..61d8537a92 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java @@ -35,8 +35,8 @@ import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; import org.apache.rocketmq.proxy.service.message.ClusterMessageService; import org.apache.rocketmq.proxy.service.message.MessageService; import org.apache.rocketmq.proxy.service.mqclient.DoNothingClientRemotingProcessor; -import org.apache.rocketmq.proxy.service.out.ClusterProxyOutService; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ClusterProxyRelayService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.ClusterTopicRouteService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.ClusterTransactionService; @@ -51,7 +51,7 @@ public class ClusterServiceManager extends ServiceManager { private final ConsumerManager consumerManager; private final TopicRouteService topicRouteService; private final MessageService messageService; - private final ProxyOutService proxyOutService; + private final ProxyRelayService proxyRelayService; private final ScheduledExecutorService scheduledExecutorService; private final MQClientAPIFactory mqClientAPIFactory; @@ -72,7 +72,7 @@ public class ClusterServiceManager extends ServiceManager { this.topicRouteService = new ClusterTopicRouteService(rpcHook); this.messageService = new ClusterMessageService(this.topicRouteService, this.mqClientAPIFactory); this.clusterTransactionService = new ClusterTransactionService(this.topicRouteService, this.producerManager, rpcHook); - this.proxyOutService = new ClusterProxyOutService(); + this.proxyRelayService = new ClusterProxyRelayService(); this.init(); } @@ -121,8 +121,8 @@ public class ClusterServiceManager extends ServiceManager { } @Override - public ProxyOutService getProxyOutService() { - return this.proxyOutService; + public ProxyRelayService getProxyOutService() { + return this.proxyRelayService; } protected static class ConsumerIdsChangeListenerImpl implements ConsumerIdsChangeListener { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java index 97549a6eb0..1a57721feb 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java @@ -21,8 +21,8 @@ import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.proxy.service.message.LocalMessageService; import org.apache.rocketmq.proxy.service.message.MessageService; -import org.apache.rocketmq.proxy.service.out.LocalProxyOutService; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.LocalProxyRelayService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.LocalTopicRouteService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.LocalTransactionService; @@ -35,14 +35,14 @@ public class LocalServiceManager extends ServiceManager { private final TopicRouteService topicRouteService; private final MessageService messageService; private final TransactionService transactionService; - private final ProxyOutService proxyOutService; + private final ProxyRelayService proxyRelayService; public LocalServiceManager(BrokerController brokerController, RPCHook rpcHook) { this.brokerController = brokerController; this.messageService = new LocalMessageService(brokerController, rpcHook); this.topicRouteService = new LocalTopicRouteService(brokerController, rpcHook); this.transactionService = new LocalTransactionService(); - this.proxyOutService = new LocalProxyOutService(brokerController); + this.proxyRelayService = new LocalProxyRelayService(brokerController); this.init(); } @@ -77,8 +77,8 @@ public class LocalServiceManager extends ServiceManager { } @Override - public ProxyOutService getProxyOutService() { - return this.proxyOutService; + public ProxyRelayService getProxyOutService() { + return this.proxyRelayService; } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java index 820d99444b..273a515dab 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java @@ -21,7 +21,7 @@ import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.service.message.MessageService; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; import org.apache.rocketmq.remoting.RPCHook; @@ -54,5 +54,5 @@ public abstract class ServiceManager extends AbstractStartAndShutdown { public abstract TransactionService getTransactionService(); - public abstract ProxyOutService getProxyOutService(); + public abstract ProxyRelayService getProxyOutService(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ClusterProxyOutService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java similarity index 93% rename from proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ClusterProxyOutService.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java index e6f8abc192..9b1bba2c79 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ClusterProxyOutService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.service.out; +package org.apache.rocketmq.proxy.service.relay; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.protocol.body.ConsumeMessageDirectlyResult; @@ -26,7 +26,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; /** * not implement yet */ -public class ClusterProxyOutService implements ProxyOutService { +public class ClusterProxyRelayService implements ProxyRelayService { @Override public CompletableFuture> processGetConsumerRunningInfo(RemotingCommand command, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/LocalProxyOutService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/LocalProxyRelayService.java similarity index 94% rename from proxy/src/main/java/org/apache/rocketmq/proxy/service/out/LocalProxyOutService.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/LocalProxyRelayService.java index 0a243600f5..e821fb60b4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/LocalProxyOutService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/LocalProxyRelayService.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.service.out; +package org.apache.rocketmq.proxy.service.relay; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.broker.BrokerController; @@ -27,11 +27,11 @@ import org.apache.rocketmq.remoting.RemotingServer; import org.apache.rocketmq.remoting.netty.NettyRemotingAbstract; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -public class LocalProxyOutService implements ProxyOutService { +public class LocalProxyRelayService implements ProxyRelayService { private final BrokerController brokerController; - public LocalProxyOutService(BrokerController brokerController) { + public LocalProxyRelayService(BrokerController brokerController) { this.brokerController = brokerController; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java similarity index 91% rename from proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyChannel.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java index 1a85a96d65..1dcea7aa37 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.service.out; +package org.apache.rocketmq.proxy.service.relay; import io.netty.channel.AbstractChannel; import io.netty.channel.Channel; @@ -47,16 +47,16 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; public abstract class ProxyChannel extends AbstractChannel { private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); - protected final ProxyOutService proxyOutService; + protected final ProxyRelayService proxyRelayService; - protected ProxyChannel(ProxyOutService proxyOutService, Channel parent) { + protected ProxyChannel(ProxyRelayService proxyRelayService, Channel parent) { super(parent); - this.proxyOutService = proxyOutService; + this.proxyRelayService = proxyRelayService; } - protected ProxyChannel(ProxyOutService proxyOutService, Channel parent, ChannelId id) { + protected ProxyChannel(ProxyRelayService proxyRelayService, Channel parent, ChannelId id) { super(parent, id); - this.proxyOutService = proxyOutService; + this.proxyRelayService = proxyRelayService; } @Override @@ -78,14 +78,14 @@ public abstract class ProxyChannel extends AbstractChannel { case RequestCode.GET_CONSUMER_RUNNING_INFO: { GetConsumerRunningInfoRequestHeader header = (GetConsumerRunningInfoRequestHeader) command.readCustomHeader(); processFuture = this.processGetConsumerRunningInfo(command, header, - this.proxyOutService.processGetConsumerRunningInfo(command, header)); + this.proxyRelayService.processGetConsumerRunningInfo(command, header)); break; } case RequestCode.CONSUME_MESSAGE_DIRECTLY: { ConsumeMessageDirectlyResultRequestHeader header = (ConsumeMessageDirectlyResultRequestHeader) command.readCustomHeader(); MessageExt messageExt = MessageDecoder.decode(ByteBuffer.wrap(command.getBody()), true, false, false); processFuture = this.processConsumeMessageDirectly(command, header, messageExt, - this.proxyOutService.processConsumeMessageDirectly(command, header)); + this.proxyRelayService.processConsumeMessageDirectly(command, header)); break; } default: diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyOutResult.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyOutResult.java similarity index 96% rename from proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyOutResult.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyOutResult.java index efeccbafd1..985481a17f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyOutResult.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyOutResult.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.service.out; +package org.apache.rocketmq.proxy.service.relay; public class ProxyOutResult { private int code; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyOutService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyRelayService.java similarity index 94% rename from proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyOutService.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyRelayService.java index dbe0f13255..9e905d3f5b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/out/ProxyOutService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyRelayService.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.service.out; +package org.apache.rocketmq.proxy.service.relay; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.protocol.body.ConsumeMessageDirectlyResult; @@ -23,7 +23,7 @@ import org.apache.rocketmq.common.protocol.header.ConsumeMessageDirectlyResultRe import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestHeader; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -public interface ProxyOutService { +public interface ProxyRelayService { CompletableFuture> processGetConsumerRunningInfo( RemotingCommand command, diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java index 2217ce4e7a..c155b382e8 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java @@ -28,8 +28,7 @@ import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager; import org.apache.rocketmq.proxy.processor.MessagingProcessor; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; -import org.apache.rocketmq.remoting.common.RemotingUtil; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.junit.Ignore; import org.junit.runner.RunWith; import org.mockito.Mock; @@ -46,7 +45,7 @@ public class BaseActivityTest extends InitConfigAndLoggerTest { @Mock protected GrpcClientSettingsManager grpcClientSettingsManager; @Mock - protected ProxyOutService proxyOutService; + protected ProxyRelayService proxyRelayService; protected static final String REMOTE_ADDR = "192.168.0.1:8080"; protected static final String LOCAL_ADDR = "127.0.0.1:8080"; @@ -58,7 +57,7 @@ public class BaseActivityTest extends InitConfigAndLoggerTest { metadata.put(InterceptorConstants.LANGUAGE, "JAVA"); metadata.put(InterceptorConstants.REMOTE_ADDRESS, REMOTE_ADDR); metadata.put(InterceptorConstants.LOCAL_ADDRESS, LOCAL_ADDR); - when(messagingProcessor.getProxyOutService()).thenReturn(proxyOutService); + when(messagingProcessor.getProxyOutService()).thenReturn(proxyRelayService); } protected Context createContext() { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java index 5e0a256e12..29e0ab3379 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java @@ -31,7 +31,7 @@ import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.service.ServiceManager; import org.apache.rocketmq.proxy.service.message.MessageService; -import org.apache.rocketmq.proxy.service.out.ProxyOutService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; import org.junit.Ignore; @@ -61,7 +61,7 @@ public class BaseProcessorTest extends InitConfigAndLoggerTest { @Mock protected TransactionService transactionService; @Mock - protected ProxyOutService proxyOutService; + protected ProxyRelayService proxyRelayService; @Mock protected ProducerProcessor producerProcessor; @Mock @@ -78,7 +78,7 @@ public class BaseProcessorTest extends InitConfigAndLoggerTest { when(serviceManager.getProducerManager()).thenReturn(producerManager); when(serviceManager.getConsumerManager()).thenReturn(consumerManager); when(serviceManager.getTransactionService()).thenReturn(transactionService); - when(serviceManager.getProxyOutService()).thenReturn(proxyOutService); + when(serviceManager.getProxyOutService()).thenReturn(proxyRelayService); } protected static ProxyContext createContext() {