qianye1001 commented on code in PR #10826:
URL: https://github.com/apache/rocketmq/pull/10826#discussion_r3976535897


##########
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/admin/ProxyAdminGrpcService.java:
##########
@@ -0,0 +1,1728 @@
+/*
+ * 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.admin;
+
+import apache.rocketmq.v2.AdminGrpc;
+import apache.rocketmq.v2.AdminSendMessageRequest;
+import apache.rocketmq.v2.AdminSendMessageResponse;
+import apache.rocketmq.v2.ChangeLogLevelRequest;
+import apache.rocketmq.v2.ChangeLogLevelResponse;
+import apache.rocketmq.v2.ClientInfo;
+import apache.rocketmq.v2.Code;
+import apache.rocketmq.v2.ConsumerRunningInfo;
+import apache.rocketmq.v2.DeleteSubscriptionRequest;
+import apache.rocketmq.v2.DeleteSubscriptionResponse;
+import apache.rocketmq.v2.DescribeGroupAccumulationRequest;
+import apache.rocketmq.v2.DescribeGroupAccumulationResponse;
+import apache.rocketmq.v2.DescribeSubscriptionRequest;
+import apache.rocketmq.v2.DescribeSubscriptionResponse;
+import apache.rocketmq.v2.DescribeTopicStatusRequest;
+import apache.rocketmq.v2.DescribeTopicStatusResponse;
+import apache.rocketmq.v2.GetConsumerRunningInfoRequest;
+import apache.rocketmq.v2.GetConsumerRunningInfoResponse;
+import apache.rocketmq.v2.GetProxyRuntimeStatsRequest;
+import apache.rocketmq.v2.GetProxyRuntimeStatsResponse;
+import apache.rocketmq.v2.GetTopicRouteRequest;
+import apache.rocketmq.v2.GetTopicRouteResponse;
+import apache.rocketmq.v2.ListConsumerConnectionRequest;
+import apache.rocketmq.v2.ListConsumerConnectionResponse;
+import apache.rocketmq.v2.ListMessageRequest;
+import apache.rocketmq.v2.ListMessageResponse;
+import apache.rocketmq.v2.ListSubscriptionRequest;
+import apache.rocketmq.v2.ListSubscriptionResponse;
+import apache.rocketmq.v2.PrintThreadStackTraceRequest;
+import apache.rocketmq.v2.PrintThreadStackTraceResponse;
+import apache.rocketmq.v2.QueryTimeSpanRequest;
+import apache.rocketmq.v2.QueryTimeSpanResponse;
+import apache.rocketmq.v2.ResetGroupOffsetRequest;
+import apache.rocketmq.v2.ResetGroupOffsetResponse;
+import apache.rocketmq.v2.Resource;
+import apache.rocketmq.v2.Status;
+import apache.rocketmq.v2.SubscriptionInfo;
+import apache.rocketmq.v2.SystemProperties;
+import apache.rocketmq.v2.VerifyMessageRequest;
+import apache.rocketmq.v2.VerifyMessageResponse;
+import com.google.protobuf.Timestamp;
+import com.google.protobuf.util.Timestamps;
+import io.grpc.stub.StreamObserver;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledFuture;
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.concurrent.TimeUnit;
+import java.util.function.BiFunction;
+import org.apache.rocketmq.broker.client.ClientChannelInfo;
+import org.apache.rocketmq.broker.client.ConsumerGroupInfo;
+import org.apache.rocketmq.broker.client.ConsumerManager;
+import org.apache.rocketmq.client.producer.SendResult;
+import org.apache.rocketmq.common.KeyBuilder;
+import org.apache.rocketmq.common.MQVersion;
+import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.common.TopicConfig;
+import org.apache.rocketmq.common.attribute.TopicMessageType;
+import org.apache.rocketmq.common.message.MessageAccessor;
+import org.apache.rocketmq.common.message.MessageConst;
+import org.apache.rocketmq.common.message.MessageExt;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.logging.org.slf4j.Logger;
+import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
+import org.apache.rocketmq.proxy.common.ProxyContext;
+import org.apache.rocketmq.proxy.config.ConfigurationManager;
+import org.apache.rocketmq.proxy.grpc.v2.channel.GrpcChannelManager;
+import org.apache.rocketmq.proxy.grpc.v2.channel.GrpcClientChannel;
+import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager;
+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.ServiceManager;
+import org.apache.rocketmq.proxy.service.admin.AdminService;
+import org.apache.rocketmq.proxy.service.relay.ProxyRelayRequest;
+import org.apache.rocketmq.proxy.service.relay.ProxyRelayResult;
+import org.apache.rocketmq.proxy.service.route.AddressableMessageQueue;
+import org.apache.rocketmq.proxy.service.route.MessageQueueView;
+import org.apache.rocketmq.remoting.protocol.RequestCode;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
+import org.apache.rocketmq.remoting.protocol.admin.ConsumeStats;
+import org.apache.rocketmq.remoting.protocol.body.CMResult;
+import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
+import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection;
+import org.apache.rocketmq.remoting.protocol.body.Connection;
+import org.apache.rocketmq.remoting.protocol.body.GroupList;
+import org.apache.rocketmq.remoting.protocol.body.QueueTimeSpan;
+import org.apache.rocketmq.remoting.protocol.body.TopicList;
+import 
org.apache.rocketmq.remoting.protocol.header.ConsumeMessageDirectlyResultRequestHeader;
+import 
org.apache.rocketmq.remoting.protocol.header.GetConsumerRunningInfoRequestHeader;
+import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
+import 
org.apache.rocketmq.remoting.protocol.subscription.SimpleSubscriptionData;
+import 
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
+
+/**
+ * Proxy Admin gRPC service (control plane).
+ *
+ * <p>Design rules this class follows:
+ * <ul>
+ *   <li><b>Never block the gRPC executor.</b> Every broker hop goes through 
the asynchronous
+ *       {@link AdminService} gateway and is fanned out to all relevant 
brokers concurrently.</li>
+ *   <li><b>Cluster-wide view.</b> Subscription and connection data is read 
from the brokers (which
+ *       see consumers registered through every proxy) rather than from this 
proxy's local channel
+ *       table, and consumers owned by a peer proxy are reached by forwarding 
the whole RPC.</li>
+ *   <li><b>Honest responses.</b> A field the open-source proxy genuinely 
cannot supply is left
+ *       unset and explained in the status message, instead of being filled 
with a plausible value.</li>
+ *   <li><b>Graded errors.</b> Failures are mapped through {@link 
ResponseBuilder#buildStatus(Throwable)}
+ *       so callers can tell "topic not found" from "internal error".</li>
+ * </ul>
+ *
+ * <p>The translation between broker wire types and the v2 contract lives in
+ * {@link AdminModelConverter}; whole-message conversion reuses the data 
plane's
+ * {@link GrpcConverter} instead of duplicating it.
+ */
+public class ProxyAdminGrpcService extends AdminGrpc.AdminImplBase {
+
+    private static final Logger log = 
LoggerFactory.getLogger(ProxyAdminGrpcService.class);
+
+    /**
+     * Bounds every wait this service performs that is not already bounded by 
the callee: the answer
+     * to a relayed telemetry command, and each broker hop of a fan-out. The 
gRPC nonce sweeper and
+     * the remoting timeout cover the healthy paths, but a hop that never 
completes at all (an
+     * unwritable channel, a saturated blocking-call pool) would otherwise 
leave the RPC hanging
+     * forever with no response at all.
+     */
+    private static final ScheduledExecutorService ADMIN_TIMEOUT_SCHEDULER =

Review Comment:
   静态的线程池不好



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to