diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java new file mode 100644 index 0000000000..60a971f102 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java @@ -0,0 +1,56 @@ +/* + * 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; + +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpHandler; +import com.sun.net.httpserver.HttpServer; +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import org.apache.rocketmq.proxy.configuration.ConfigurationManager; +import org.apache.rocketmq.proxy.grpc.common.StartAndShutdown; + +public class HealthCheckServer implements StartAndShutdown { + + private HttpServer healthChecker; + + @Override + public void start() throws Exception { + this.healthChecker = HttpServer.create(new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0); + this.healthChecker.createContext("/status", new HealthCheckHandler()); + this.healthChecker.setExecutor(null); + this.healthChecker.start(); + } + + @Override + public void shutdown() { + this.healthChecker.stop(0); + } + + static class HealthCheckHandler implements HttpHandler { + @Override + public void handle(HttpExchange t) throws IOException { + String response = "Hello"; + t.sendResponseHeaders(200, response.length()); + OutputStream os = t.getResponseBody(); + os.write(response.getBytes()); + os.close(); + } + } +} \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java new file mode 100644 index 0000000000..b89627fa19 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -0,0 +1,125 @@ +/* + * 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; + +import ch.qos.logback.classic.LoggerContext; +import ch.qos.logback.classic.joran.JoranConfigurator; +import ch.qos.logback.core.joran.spi.JoranException; +import java.util.Date; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.broker.BrokerStartup; +import org.apache.rocketmq.client.log.ClientLogger; +import org.apache.rocketmq.common.thread.ThreadPoolMonitor; +import org.apache.rocketmq.proxy.configuration.ConfigurationManager; +import org.apache.rocketmq.proxy.configuration.ProxyConfig; +import org.apache.rocketmq.proxy.grpc.GrpcServer; +import org.apache.rocketmq.proxy.grpc.common.ProxyMode; +import org.apache.rocketmq.proxy.grpc.service.ClusterGrpcService; +import org.apache.rocketmq.proxy.grpc.service.GrpcService; +import org.apache.rocketmq.proxy.grpc.service.LocalGrpcService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class ProxyStartup { + + private static final Logger log = LoggerFactory.getLogger(ProxyStartup.class); + + public static void main(String[] args) { + try { + ConfigurationManager.initEnv(); + initLogger(); + ConfigurationManager.intConfig(); + + // init thread pool monitor for proxy. + initThreadPoolMonitor(); + + // create and start grpcServer + GrpcServer grpcServer = createGrpcServer(); + grpcServer.start(); + + // health check server + final HealthCheckServer healthCheckServer = new HealthCheckServer(); + healthCheckServer.start(); + + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + log.info("try to shutdown server"); + + try { + healthCheckServer.shutdown(); + Thread.sleep(TimeUnit.SECONDS.toMillis(ConfigurationManager.getProxyConfig().getWaitAfterStopHealthCheckInSeconds())); + } catch (Exception e) { + log.error("err when shutdown healthCheckServer", e); + } + + try { + grpcServer.shutdown(); + } catch (Exception e) { + log.error("err when shutdown grpc server", e); + } + })); + } catch (Exception e) { + System.err.println("find a unexpect err." + e); + e.printStackTrace(); + log.error("find a unexpect err.", e); + System.exit(1); + } + + System.out.printf("%s%n", new Date() + " rmq-proxy startup successfully"); + log.info(new Date() + "rmq-proxy startup successfully"); + } + + private static GrpcServer createGrpcServer() throws RuntimeException { + GrpcService grpcService; + String proxyModeStr = ConfigurationManager.getProxyConfig().getProxyMode(); + if (ProxyMode.isClusterMode(proxyModeStr)) { + grpcService = new ClusterGrpcService(); + } else if (ProxyMode.isLocalMode(proxyModeStr)) { + BrokerController brokerController = createBrokerController(); + grpcService = new LocalGrpcService(brokerController); + } else { + throw new IllegalArgumentException("try to start grpc server with wrong mode, use 'local' or 'cluster'"); + } + + return new GrpcServer(grpcService); + + } + + private static BrokerController createBrokerController() { + String[] brokerStartupArgs = new String[] {"-c", ConfigurationManager.getProxyConfig().getBrokerConfigPath()}; + return BrokerStartup.createBrokerController(brokerStartupArgs); + } + + private static void initThreadPoolMonitor() { + ThreadPoolMonitor.init(); + ProxyConfig config = ConfigurationManager.getProxyConfig(); + ThreadPoolMonitor.config(config.isEnablePrintJstack(), config.getPrintJstackPeriodMillis()); + } + + private static void initLogger() throws JoranException { + System.setProperty(ClientLogger.CLIENT_LOG_USESLF4J, "true"); + + LoggerContext lc = (LoggerContext) LoggerFactory.getILoggerFactory(); + JoranConfigurator configurator = new JoranConfigurator(); + configurator.setContext(lc); + lc.reset(); + //https://logback.qos.ch/manual/configuration.html + lc.setPackagingDataEnabled(false); + configurator.doConfigure(ConfigurationManager.getProxyHome() + "/conf/logback.xml"); + } +} \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java index c0aa518caa..58a8a2537a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java @@ -22,24 +22,69 @@ import org.apache.rocketmq.proxy.grpc.common.ProxyMode; public class ProxyConfig { public final static String CONFIG_FILE_NAME = "rmq-proxy.json"; + /** + * Configuration for proxy + */ + private Integer healthCheckPort = 8000; + private long waitAfterStopHealthCheckInSeconds = 40; + + /** + * configuration for ThreadPoolMonitor + */ + private boolean enablePrintJstack = true; + private long printJstackPeriodMillis = 60000; + /** * gRPC */ private String proxyMode = ProxyMode.CLUSTER.name(); private Boolean startGrpcServer = true; private Integer grpcServerPort = 8081; - private String grpcTlsKeyPath = "/home/admin/rmq-gateway/conf/tls/gRPC.key.pem"; - private String grpcTlsCertPath = "/home/admin/rmq-gateway/conf/tls/gRPC.chain.cert.pem"; + private String grpcTlsKeyPath = ConfigurationManager.getProxyHome() + "/conf/tls/gRPC.key.pem"; + private String grpcTlsCertPath = ConfigurationManager.getProxyHome() + "/conf/tls/gRPC.chain.cert.pem"; private int grpcBossLoopNum = 1; private int grpcWorkerLoopNum = Runtime.getRuntime().availableProcessors() * 2; private int grpcThreadPoolNums = 16 + Runtime.getRuntime().availableProcessors() * 2; private int grpcThreadPoolQueueCapacity = 100000; + private String brokerConfigPath = ConfigurationManager.getProxyHome() + "/conf/broker.conf"; /** * gRPC max message size * 130M = 4M * 32 messages + 2M attributes */ private int grpcMaxInboundMessageSize = 130 * 1024 * 1024; + public Integer getHealthCheckPort() { + return healthCheckPort; + } + + public void setHealthCheckPort(Integer healthCheckPort) { + this.healthCheckPort = healthCheckPort; + } + + public long getWaitAfterStopHealthCheckInSeconds() { + return waitAfterStopHealthCheckInSeconds; + } + + public void setWaitAfterStopHealthCheckInSeconds(long waitAfterStopHealthCheckInSeconds) { + this.waitAfterStopHealthCheckInSeconds = waitAfterStopHealthCheckInSeconds; + } + + public boolean isEnablePrintJstack() { + return enablePrintJstack; + } + + public void setEnablePrintJstack(boolean enablePrintJstack) { + this.enablePrintJstack = enablePrintJstack; + } + + public long getPrintJstackPeriodMillis() { + return printJstackPeriodMillis; + } + + public void setPrintJstackPeriodMillis(long printJstackPeriodMillis) { + this.printJstackPeriodMillis = printJstackPeriodMillis; + } + public String getProxyMode() { return proxyMode; } @@ -112,6 +157,14 @@ public class ProxyConfig { this.grpcThreadPoolQueueCapacity = grpcThreadPoolQueueCapacity; } + public String getBrokerConfigPath() { + return brokerConfigPath; + } + + public void setBrokerConfigPath(String brokerConfigPath) { + this.brokerConfigPath = brokerConfigPath; + } + public int getGrpcMaxInboundMessageSize() { return grpcMaxInboundMessageSize; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java index a0aa501def..3a522a47f9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java @@ -42,8 +42,10 @@ public class GrpcServer { private final io.grpc.Server server; private final ThreadPoolExecutor executor; + private final GrpcService grpcService; public GrpcServer(GrpcService grpcService) { + this.grpcService = grpcService; int port = ConfigurationManager.getProxyConfig().getGrpcServerPort(); NettyServerBuilder serverBuilder = NettyServerBuilder.forPort(port); @@ -98,7 +100,11 @@ public class GrpcServer { bossLoopNum, workerLoopNum, maxInboundMessageSize); } + public void start() throws Exception { + // first to start grpc service. + this.grpcService.start(); + this.server.start(); log.info("grpc server has started"); } @@ -108,8 +114,10 @@ public class GrpcServer { this.server.shutdown().awaitTermination(30, TimeUnit.SECONDS); this.executor.shutdown(); + this.grpcService.shutdown(); + log.info("grpc server has stopped"); - } catch (InterruptedException e) { + } catch (Exception e) { e.printStackTrace(); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/StartAndShutdown.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/StartAndShutdown.java new file mode 100644 index 0000000000..ebce8dc5f6 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/StartAndShutdown.java @@ -0,0 +1,24 @@ +/* + * 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; + +public interface StartAndShutdown { + void start() throws Exception; + + void shutdown() throws Exception; +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index a49a2ecb6a..f67e0d6ec8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -64,19 +64,23 @@ public class ClusterGrpcService implements GrpcService { } - @Override public CompleteFuture queryRoute(Context ctx, QueryRouteRequest request) { + @Override + public CompleteFuture queryRoute(Context ctx, QueryRouteRequest request) { return null; } - @Override public CompleteFuture heartbeat(Context ctx, HeartbeatRequest request) { + @Override + public CompleteFuture heartbeat(Context ctx, HeartbeatRequest request) { return null; } - @Override public CompleteFuture healthCheck(Context ctx, HealthCheckRequest request) { + @Override + public CompleteFuture healthCheck(Context ctx, HealthCheckRequest request) { return null; } - @Override public CompleteFuture sendMessage(Context ctx, SendMessageRequest request) { + @Override + public CompleteFuture sendMessage(Context ctx, SendMessageRequest request) { return null; } @@ -138,4 +142,12 @@ public class ClusterGrpcService implements GrpcService { ChangeInvisibleDurationRequest request) { return null; } + + @Override + public void start() throws Exception { + } + + @Override + public void shutdown() throws Exception { + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java index 9666eeddba..8a9aac74f4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcService.java @@ -53,8 +53,9 @@ import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import io.grpc.Context; import io.netty.util.concurrent.CompleteFuture; +import org.apache.rocketmq.proxy.grpc.common.StartAndShutdown; -public interface GrpcService { +public interface GrpcService extends StartAndShutdown { CompleteFuture queryRoute(Context ctx, QueryRouteRequest request); CompleteFuture heartbeat(Context ctx, HeartbeatRequest request); 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 348d94d6f8..f4f03ade44 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 @@ -140,4 +140,12 @@ public class LocalGrpcService implements GrpcService { ChangeInvisibleDurationRequest request) { return null; } + + @Override public void start() throws Exception { + this.brokerController.start(); + } + + @Override public void shutdown() throws Exception { + this.brokerController.shutdown(); + } }