From 6d00642386f675b14524a54b64e1f26660b8bf17 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=AD=E5=8B=87=E5=8D=87=20pengys?= <8082209@qq.com> Date: Mon, 11 Mar 2019 14:23:57 +0800 Subject: [PATCH] Separation the remote gRPC server from the receiver gRPC server. (#2335) * #2327 Add a new module for share receiver server or to be a bridge of core server. * Rename the module name. --- docker/config/application.yml | 2 + .../library/server/grpc/GRPCServer.java | 2 + .../library/server/jetty/JettyServer.java | 2 + oap-server/server-receiver-plugin/pom.xml | 1 + .../IstioTelemetryReceiverProvider.java | 3 +- .../skywalking-jvm-receiver-plugin/pom.xml | 11 +- .../jvm/provider/JVMModuleProvider.java | 12 +- .../skywalking-mesh-receiver-plugin/pom.xml | 10 +- .../receiver/mesh/MeshReceiverProvider.java | 6 +- .../pom.xml | 8 +- .../provider/RegisterModuleProvider.java | 28 ++--- .../skywalking-sharing-server-plugin/pom.xml | 32 +++++ .../server/ReceiverGRPCHandlerRegister.java | 39 ++++++ .../server/ReceiverJettyHandlerRegister.java | 35 ++++++ .../sharing/server/SharingServerConfig.java | 37 ++++++ .../sharing/server/SharingServerModule.java | 38 ++++++ .../server/SharingServerModuleProvider.java | 116 ++++++++++++++++++ ...ing.oap.server.library.module.ModuleDefine | 19 +++ ...g.oap.server.library.module.ModuleProvider | 19 +++ .../skywalking-trace-receiver-plugin/pom.xml | 11 +- .../trace/provider/TraceModuleProvider.java | 7 +- .../src/main/assembly/application.yml | 2 + .../src/main/resources/application.yml | 2 + 23 files changed, 405 insertions(+), 37 deletions(-) create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/pom.xml create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverGRPCHandlerRegister.java create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverJettyHandlerRegister.java create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerConfig.java create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModule.java create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModuleProvider.java create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine create mode 100644 oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider diff --git a/docker/config/application.yml b/docker/config/application.yml index fa4cbb64a..4090e0f3e 100644 --- a/docker/config/application.yml +++ b/docker/config/application.yml @@ -66,6 +66,8 @@ storage: # url: ${SW_STORAGE_H2_URL:jdbc:h2:mem:skywalking-oap-db} # user: ${SW_STORAGE_H2_USER:sa} # mysql: +receiver-sharing-server: + default: receiver-register: default: receiver-trace: diff --git a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java index 86e67de9f..3d5eebab7 100644 --- a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java +++ b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java @@ -112,10 +112,12 @@ public class GRPCServer implements Server { } public void addHandler(BindableService handler) { + logger.info("Bind handler {} into gRPC server {}:{}", handler.getClass().getSimpleName(), host, port); nettyServerBuilder.addService(handler); } public void addHandler(ServerServiceDefinition definition) { + logger.info("Bind handler {} into gRPC server {}:{}", definition.getClass().getSimpleName(), host, port); nettyServerBuilder.addService(definition); } diff --git a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/jetty/JettyServer.java b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/jetty/JettyServer.java index af8c91e54..ec96327a2 100644 --- a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/jetty/JettyServer.java +++ b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/jetty/JettyServer.java @@ -71,6 +71,8 @@ public class JettyServer implements Server { } public void addHandler(JettyHandler handler) { + logger.info("Bind handler {} into jetty server {}:{}", handler.getClass().getSimpleName(), host, port); + ServletHolder servletHolder = new ServletHolder(); servletHolder.setServlet(handler); servletContextHandler.addServlet(servletHolder, handler.pathSpec()); diff --git a/oap-server/server-receiver-plugin/pom.xml b/oap-server/server-receiver-plugin/pom.xml index 20ce31523..e276a5a9d 100644 --- a/oap-server/server-receiver-plugin/pom.xml +++ b/oap-server/server-receiver-plugin/pom.xml @@ -35,6 +35,7 @@ skywalking-register-receiver-plugin skywalking-jvm-receiver-plugin envoy-metrics-receiver-plugin + skywalking-sharing-server-plugin diff --git a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/provider/IstioTelemetryReceiverProvider.java b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/provider/IstioTelemetryReceiverProvider.java index 3b5afdbc1..ee61fb680 100644 --- a/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/provider/IstioTelemetryReceiverProvider.java +++ b/oap-server/server-receiver-plugin/skywalking-istio-telemetry-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/istio/telemetry/provider/IstioTelemetryReceiverProvider.java @@ -23,6 +23,7 @@ import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister; import org.apache.skywalking.oap.server.library.module.*; import org.apache.skywalking.oap.server.receiver.istio.telemetry.module.IstioTelemetryReceiverModule; +import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; import org.apache.skywalking.oap.server.telemetry.TelemetryModule; public class IstioTelemetryReceiverProvider extends ModuleProvider { @@ -42,7 +43,7 @@ public class IstioTelemetryReceiverProvider extends ModuleProvider { } @Override public void start() throws ServiceNotProvidedException, ModuleStartException { - GRPCHandlerRegister service = getManager().find(CoreModule.NAME).provider().getService(GRPCHandlerRegister.class); + GRPCHandlerRegister service = getManager().find(SharingServerModule.NAME).provider().getService(GRPCHandlerRegister.class); service.addHandler(new IstioTelemetryGRPCHandler(getManager())); } diff --git a/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/pom.xml b/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/pom.xml index 0040b6963..7959d1976 100644 --- a/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/pom.xml +++ b/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/pom.xml @@ -17,7 +17,8 @@ ~ --> - + server-receiver-plugin org.apache.skywalking @@ -27,4 +28,12 @@ skywalking-jvm-receiver-plugin jar + + + + org.apache.skywalking + skywalking-sharing-server-plugin + ${project.version} + + \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/jvm/provider/JVMModuleProvider.java b/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/jvm/provider/JVMModuleProvider.java index d199512c3..535144b59 100644 --- a/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/jvm/provider/JVMModuleProvider.java +++ b/oap-server/server-receiver-plugin/skywalking-jvm-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/jvm/provider/JVMModuleProvider.java @@ -20,12 +20,10 @@ package org.apache.skywalking.oap.server.receiver.jvm.provider; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister; -import org.apache.skywalking.oap.server.library.module.ModuleConfig; -import org.apache.skywalking.oap.server.library.module.ModuleDefine; -import org.apache.skywalking.oap.server.library.module.ModuleProvider; +import org.apache.skywalking.oap.server.library.module.*; import org.apache.skywalking.oap.server.receiver.jvm.module.JVMModule; -import org.apache.skywalking.oap.server.receiver.jvm.provider.handler.JVMMetricReportServiceHandler; -import org.apache.skywalking.oap.server.receiver.jvm.provider.handler.JVMMetricsServiceHandler; +import org.apache.skywalking.oap.server.receiver.jvm.provider.handler.*; +import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; /** * @author peng-yongsheng @@ -48,7 +46,7 @@ public class JVMModuleProvider extends ModuleProvider { } @Override public void start() { - GRPCHandlerRegister grpcHandlerRegister = getManager().find(CoreModule.NAME).provider().getService(GRPCHandlerRegister.class); + GRPCHandlerRegister grpcHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(GRPCHandlerRegister.class); grpcHandlerRegister.addHandler(new JVMMetricsServiceHandler(getManager())); grpcHandlerRegister.addHandler(new JVMMetricReportServiceHandler(getManager())); } @@ -58,6 +56,6 @@ public class JVMModuleProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[] {CoreModule.NAME}; + return new String[] {CoreModule.NAME, SharingServerModule.NAME}; } } diff --git a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/pom.xml b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/pom.xml index cae9db7b6..501270ed9 100644 --- a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/pom.xml +++ b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/pom.xml @@ -17,7 +17,8 @@ ~ --> - + server-receiver-plugin org.apache.skywalking @@ -28,4 +29,11 @@ skywalking-mesh-receiver-plugin jar + + + org.apache.skywalking + skywalking-sharing-server-plugin + ${project.version} + + \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshReceiverProvider.java b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshReceiverProvider.java index d221fe2b1..c629db010 100644 --- a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshReceiverProvider.java +++ b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/MeshReceiverProvider.java @@ -22,7 +22,7 @@ import java.io.IOException; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister; import org.apache.skywalking.oap.server.library.module.*; -import org.apache.skywalking.oap.server.library.module.ModuleDefine; +import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; import org.apache.skywalking.oap.server.telemetry.TelemetryModule; public class MeshReceiverProvider extends ModuleProvider { @@ -56,7 +56,7 @@ public class MeshReceiverProvider extends ModuleProvider { throw new ModuleStartException(e.getMessage(), e); } CoreRegisterLinker.setModuleManager(getManager()); - GRPCHandlerRegister service = getManager().find(CoreModule.NAME).provider().getService(GRPCHandlerRegister.class); + GRPCHandlerRegister service = getManager().find(SharingServerModule.NAME).provider().getService(GRPCHandlerRegister.class); service.addHandler(new MeshGRPCHandler(getManager())); } @@ -65,6 +65,6 @@ public class MeshReceiverProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[] {TelemetryModule.NAME, CoreModule.NAME}; + return new String[] {TelemetryModule.NAME, CoreModule.NAME, SharingServerModule.NAME}; } } diff --git a/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/pom.xml b/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/pom.xml index 34e992796..9a1d1604d 100644 --- a/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/pom.xml +++ b/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/pom.xml @@ -17,7 +17,8 @@ ~ --> - + server-receiver-plugin org.apache.skywalking @@ -34,5 +35,10 @@ apm-network ${project.version} + + org.apache.skywalking + skywalking-sharing-server-plugin + ${project.version} + \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/register/provider/RegisterModuleProvider.java b/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/register/provider/RegisterModuleProvider.java index 56d6378c4..ffdba9000 100644 --- a/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/register/provider/RegisterModuleProvider.java +++ b/oap-server/server-receiver-plugin/skywalking-register-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/register/provider/RegisterModuleProvider.java @@ -19,23 +19,13 @@ package org.apache.skywalking.oap.server.receiver.register.provider; import org.apache.skywalking.oap.server.core.CoreModule; -import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister; -import org.apache.skywalking.oap.server.core.server.JettyHandlerRegister; -import org.apache.skywalking.oap.server.library.module.ModuleConfig; -import org.apache.skywalking.oap.server.library.module.ModuleDefine; -import org.apache.skywalking.oap.server.library.module.ModuleProvider; +import org.apache.skywalking.oap.server.core.server.*; +import org.apache.skywalking.oap.server.library.module.*; import org.apache.skywalking.oap.server.receiver.register.module.RegisterModule; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.grpc.ApplicationRegisterHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.grpc.InstanceDiscoveryServiceHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.grpc.NetworkAddressRegisterServiceHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.grpc.ServiceNameDiscoveryHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.rest.ApplicationRegisterServletHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.rest.InstanceDiscoveryServletHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.rest.InstanceHeartBeatServletHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.rest.NetworkAddressRegisterServletHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.rest.ServiceNameDiscoveryServiceHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v6.grpc.RegisterServiceHandler; -import org.apache.skywalking.oap.server.receiver.register.provider.handler.v6.grpc.ServiceInstancePingServiceHandler; +import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.grpc.*; +import org.apache.skywalking.oap.server.receiver.register.provider.handler.v5.rest.*; +import org.apache.skywalking.oap.server.receiver.register.provider.handler.v6.grpc.*; +import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; /** * @author peng-yongsheng @@ -58,7 +48,7 @@ public class RegisterModuleProvider extends ModuleProvider { } @Override public void start() { - GRPCHandlerRegister grpcHandlerRegister = getManager().find(CoreModule.NAME).provider().getService(GRPCHandlerRegister.class); + GRPCHandlerRegister grpcHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(GRPCHandlerRegister.class); grpcHandlerRegister.addHandler(new ApplicationRegisterHandler(getManager())); grpcHandlerRegister.addHandler(new InstanceDiscoveryServiceHandler(getManager())); grpcHandlerRegister.addHandler(new ServiceNameDiscoveryHandler(getManager())); @@ -68,7 +58,7 @@ public class RegisterModuleProvider extends ModuleProvider { grpcHandlerRegister.addHandler(new RegisterServiceHandler(getManager())); grpcHandlerRegister.addHandler(new ServiceInstancePingServiceHandler(getManager())); - JettyHandlerRegister jettyHandlerRegister = getManager().find(CoreModule.NAME).provider().getService(JettyHandlerRegister.class); + JettyHandlerRegister jettyHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(JettyHandlerRegister.class); jettyHandlerRegister.addHandler(new ApplicationRegisterServletHandler(getManager())); jettyHandlerRegister.addHandler(new InstanceDiscoveryServletHandler(getManager())); jettyHandlerRegister.addHandler(new InstanceHeartBeatServletHandler(getManager())); @@ -81,6 +71,6 @@ public class RegisterModuleProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[] {CoreModule.NAME}; + return new String[] {CoreModule.NAME, SharingServerModule.NAME}; } } diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/pom.xml b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/pom.xml new file mode 100644 index 000000000..4c9882846 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/pom.xml @@ -0,0 +1,32 @@ + + + + + + server-receiver-plugin + org.apache.skywalking + 6.1.0-SNAPSHOT + + 4.0.0 + + skywalking-sharing-server-plugin + jar + \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverGRPCHandlerRegister.java b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverGRPCHandlerRegister.java new file mode 100644 index 000000000..2638799dd --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverGRPCHandlerRegister.java @@ -0,0 +1,39 @@ +/* + * 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.skywalking.oap.server.receiver.sharing.server; + +import io.grpc.*; +import lombok.Setter; +import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister; + +/** + * @author peng-yongsheng + */ +public class ReceiverGRPCHandlerRegister implements GRPCHandlerRegister { + + @Setter private GRPCHandlerRegister grpcHandlerRegister; + + @Override public void addHandler(BindableService handler) { + grpcHandlerRegister.addHandler(handler); + } + + @Override public void addHandler(ServerServiceDefinition definition) { + grpcHandlerRegister.addHandler(definition); + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverJettyHandlerRegister.java b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverJettyHandlerRegister.java new file mode 100644 index 000000000..8d7f98124 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/ReceiverJettyHandlerRegister.java @@ -0,0 +1,35 @@ +/* + * 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.skywalking.oap.server.receiver.sharing.server; + +import lombok.Setter; +import org.apache.skywalking.oap.server.core.server.JettyHandlerRegister; +import org.apache.skywalking.oap.server.library.server.jetty.JettyHandler; + +/** + * @author peng-yongsheng + */ +public class ReceiverJettyHandlerRegister implements JettyHandlerRegister { + + @Setter private JettyHandlerRegister jettyHandlerRegister; + + @Override public void addHandler(JettyHandler serverHandler) { + jettyHandlerRegister.addHandler(serverHandler); + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerConfig.java b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerConfig.java new file mode 100644 index 000000000..7e4cb5f86 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerConfig.java @@ -0,0 +1,37 @@ +/* + * 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.skywalking.oap.server.receiver.sharing.server; + +import lombok.*; +import org.apache.skywalking.oap.server.library.module.ModuleConfig; + +/** + * @author peng-yongsheng + */ +@Getter +@Setter +public class SharingServerConfig extends ModuleConfig { + private String restHost; + private int restPort; + private String restContextPath; + private String gRPCHost; + private int gRPCPort; + private int maxConcurrentCallsPerConnection; + private int maxMessageSize; +} diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModule.java b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModule.java new file mode 100644 index 000000000..cd858dc9d --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModule.java @@ -0,0 +1,38 @@ +/* + * 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.skywalking.oap.server.receiver.sharing.server; + +import org.apache.skywalking.oap.server.core.server.*; +import org.apache.skywalking.oap.server.library.module.ModuleDefine; + +/** + * @author peng-yongsheng + */ +public class SharingServerModule extends ModuleDefine { + + public static final String NAME = "receiver-sharing-server"; + + public SharingServerModule() { + super(NAME); + } + + @Override public Class[] services() { + return new Class[] {GRPCHandlerRegister.class, JettyHandlerRegister.class}; + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModuleProvider.java b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModuleProvider.java new file mode 100644 index 000000000..adc429afb --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/sharing/server/SharingServerModuleProvider.java @@ -0,0 +1,116 @@ +/* + * 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.skywalking.oap.server.receiver.sharing.server; + +import java.util.Objects; +import org.apache.logging.log4j.util.Strings; +import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.remote.health.HealthCheckServiceHandler; +import org.apache.skywalking.oap.server.core.server.*; +import org.apache.skywalking.oap.server.library.module.*; +import org.apache.skywalking.oap.server.library.server.ServerException; +import org.apache.skywalking.oap.server.library.server.grpc.GRPCServer; +import org.apache.skywalking.oap.server.library.server.jetty.JettyServer; + +/** + * @author peng-yongsheng + */ +public class SharingServerModuleProvider extends ModuleProvider { + + private final SharingServerConfig config; + private GRPCServer grpcServer; + private JettyServer jettyServer; + private ReceiverGRPCHandlerRegister receiverGRPCHandlerRegister; + private ReceiverJettyHandlerRegister receiverJettyHandlerRegister; + + public SharingServerModuleProvider() { + super(); + this.config = new SharingServerConfig(); + } + + @Override public String name() { + return "default"; + } + + @Override public Class module() { + return SharingServerModule.class; + } + + @Override public ModuleConfig createConfigBeanIfAbsent() { + return config; + } + + @Override public void prepare() { + if (config.getRestPort() != 0) { + jettyServer = new JettyServer(Strings.isBlank(config.getRestHost()) ? "0.0.0.0" : config.getRestHost(), config.getRestPort(), config.getRestContextPath()); + jettyServer.initialize(); + + this.registerServiceImplementation(JettyHandlerRegister.class, new JettyHandlerRegisterImpl(jettyServer)); + } else { + this.receiverJettyHandlerRegister = new ReceiverJettyHandlerRegister(); + this.registerServiceImplementation(JettyHandlerRegister.class, receiverJettyHandlerRegister); + } + + if (config.getGRPCPort() != 0) { + grpcServer = new GRPCServer(Strings.isBlank(config.getGRPCHost()) ? "0.0.0.0" : config.getGRPCHost(), config.getGRPCPort()); + if (config.getMaxMessageSize() > 0) { + grpcServer.setMaxMessageSize(config.getMaxMessageSize()); + } + if (config.getMaxConcurrentCallsPerConnection() > 0) { + grpcServer.setMaxConcurrentCallsPerConnection(config.getMaxConcurrentCallsPerConnection()); + } + grpcServer.initialize(); + + this.registerServiceImplementation(GRPCHandlerRegister.class, new GRPCHandlerRegisterImpl(grpcServer)); + } else { + this.receiverGRPCHandlerRegister = new ReceiverGRPCHandlerRegister(); + this.registerServiceImplementation(GRPCHandlerRegister.class, receiverGRPCHandlerRegister); + } + } + + @Override public void start() { + if (Objects.nonNull(grpcServer)) { + grpcServer.addHandler(new HealthCheckServiceHandler()); + } + + if (Objects.nonNull(receiverGRPCHandlerRegister)) { + receiverGRPCHandlerRegister.setGrpcHandlerRegister(getManager().find(CoreModule.NAME).provider().getService(GRPCHandlerRegister.class)); + } + if (Objects.nonNull(receiverJettyHandlerRegister)) { + receiverJettyHandlerRegister.setJettyHandlerRegister(getManager().find(CoreModule.NAME).provider().getService(JettyHandlerRegister.class)); + } + } + + @Override public void notifyAfterCompleted() throws ModuleStartException { + try { + if (Objects.nonNull(grpcServer)) { + grpcServer.start(); + } + if (Objects.nonNull(jettyServer)) { + jettyServer.start(); + } + } catch (ServerException e) { + throw new ModuleStartException(e.getMessage(), e); + } + } + + @Override public String[] requiredModules() { + return new String[] {CoreModule.NAME}; + } +} diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine new file mode 100644 index 000000000..b76d82c01 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleDefine @@ -0,0 +1,19 @@ +# +# 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. +# +# + +org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider new file mode 100644 index 000000000..215a98754 --- /dev/null +++ b/oap-server/server-receiver-plugin/skywalking-sharing-server-plugin/src/main/resources/META-INF/services/org.apache.skywalking.oap.server.library.module.ModuleProvider @@ -0,0 +1,19 @@ +# +# 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. +# +# + +org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModuleProvider \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/pom.xml b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/pom.xml index 66170078f..9e0a5b7f8 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/pom.xml +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/pom.xml @@ -17,7 +17,8 @@ ~ --> - + server-receiver-plugin org.apache.skywalking @@ -27,4 +28,12 @@ skywalking-trace-receiver-plugin jar + + + + org.apache.skywalking + skywalking-sharing-server-plugin + ${project.version} + + \ No newline at end of file diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java index 367702dcc..e346503a9 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/TraceModuleProvider.java @@ -22,6 +22,7 @@ import java.io.IOException; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.server.*; import org.apache.skywalking.oap.server.library.module.*; +import org.apache.skywalking.oap.server.receiver.sharing.server.SharingServerModule; import org.apache.skywalking.oap.server.receiver.trace.module.TraceModule; import org.apache.skywalking.oap.server.receiver.trace.provider.handler.v5.grpc.TraceSegmentServiceHandler; import org.apache.skywalking.oap.server.receiver.trace.provider.handler.v5.rest.TraceSegmentServletHandler; @@ -79,8 +80,8 @@ public class TraceModuleProvider extends ModuleProvider { } @Override public void start() throws ModuleStartException { - GRPCHandlerRegister grpcHandlerRegister = getManager().find(CoreModule.NAME).provider().getService(GRPCHandlerRegister.class); - JettyHandlerRegister jettyHandlerRegister = getManager().find(CoreModule.NAME).provider().getService(JettyHandlerRegister.class); + GRPCHandlerRegister grpcHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(GRPCHandlerRegister.class); + JettyHandlerRegister jettyHandlerRegister = getManager().find(SharingServerModule.NAME).provider().getService(JettyHandlerRegister.class); try { grpcHandlerRegister.addHandler(new TraceSegmentServiceHandler(segmentProducer)); @@ -106,6 +107,6 @@ public class TraceModuleProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[] {TelemetryModule.NAME, CoreModule.NAME}; + return new String[] {TelemetryModule.NAME, CoreModule.NAME, SharingServerModule.NAME}; } } diff --git a/oap-server/server-starter/src/main/assembly/application.yml b/oap-server/server-starter/src/main/assembly/application.yml index 3f12d5aa1..b6d963916 100644 --- a/oap-server/server-starter/src/main/assembly/application.yml +++ b/oap-server/server-starter/src/main/assembly/application.yml @@ -66,6 +66,8 @@ storage: # flushInterval: ${SW_STORAGE_ES_FLUSH_INTERVAL:10} # flush the bulk every 10 seconds whatever the number of requests # concurrentRequests: ${SW_STORAGE_ES_CONCURRENT_REQUESTS:2} # the number of concurrent requests # mysql: +receiver-sharing-server: + default: receiver-register: default: receiver-trace: diff --git a/oap-server/server-starter/src/main/resources/application.yml b/oap-server/server-starter/src/main/resources/application.yml index 616cddaa1..4309e3762 100644 --- a/oap-server/server-starter/src/main/resources/application.yml +++ b/oap-server/server-starter/src/main/resources/application.yml @@ -66,6 +66,8 @@ storage: # url: ${SW_STORAGE_H2_URL:jdbc:h2:mem:skywalking-oap-db} # user: ${SW_STORAGE_H2_USER:sa} # mysql: +receiver-sharing-server: + default: receiver-register: default: receiver-trace: