From 1cb05658d65605deb26e465866e817c8889fb3e0 Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Sun, 12 Nov 2017 19:47:57 +0800 Subject: [PATCH] cluster, naming, grpc manager, jetty manager, storage modules test successful. --- .../agent/grpc/AgentModuleGRPCProvider.java | 19 ++- .../naming/AgentGRPCNamingHandler.java | 53 ++++++++ .../naming/AgentGRPCNamingListener.java | 43 +++++++ .../collector-agent-jetty-provider/pom.xml | 5 + .../agent/jetty/AgentModuleJettyProvider.java | 49 +++++++- .../jetty/AgentModuleJettyRegistration.java | 41 ++++++ .../handler/TraceSegmentServletHandler.java | 73 +++++++++++ .../naming/AgentJettyNamingHandler.java | 53 ++++++++ .../naming/AgentJettyNamingListener.java | 43 +++++++ .../reader/KeyWithStringValueJsonReader.java | 54 ++++++++ .../jetty/handler/reader/LogJsonReader.java | 55 ++++++++ .../handler/reader/ReferenceJsonReader.java | 92 ++++++++++++++ .../handler/reader/SegmentJsonReader.java | 87 +++++++++++++ .../jetty/handler/reader/SpanJsonReader.java | 117 ++++++++++++++++++ .../handler/reader/StreamJsonReader.java | 29 +++++ .../jetty/handler/reader/TraceSegment.java | 47 +++++++ .../reader/TraceSegmentJsonReader.java | 65 ++++++++++ .../handler/reader/UniqueIdJsonReader.java | 40 ++++++ apm-collector/apm-collector-agent/pom.xml | 10 ++ .../src/main/resources/application.yml | 36 +++--- .../ClusterModuleStandaloneProvider.java | 31 ++++- .../ClusterStandaloneDataMonitor.java | 91 ++++++++++++++ .../StandaloneModuleListenerService.java | 39 ++++++ .../StandaloneModuleRegisterService.java | 12 +- .../collector/server/jetty/JettyServer.java | 1 + .../collector/core/module/ModuleProvider.java | 2 +- .../main/resources/application-default.yml | 26 +--- .../manager/service/GRPCManagerService.java | 2 +- .../service/GRPCManagerServiceImpl.java | 2 +- .../manager/service/JettyManagerService.java | 5 +- .../service/JettyManagerServiceImpl.java | 13 +- .../jetty/NamingModuleJettyProvider.java | 15 +-- .../NamingJettyHandlerRegisterService.java | 27 +++- .../remote/grpc/RemoteModuleGRPCProvider.java | 2 +- .../resources/META-INF/defines/es_dao.define | 50 +++++--- .../resources/META-INF/defines/h2_dao.define | 50 +++++--- .../ui/jetty/UIModuleJettyProvider.java | 6 +- 37 files changed, 1278 insertions(+), 107 deletions(-) create mode 100644 apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingHandler.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingListener.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyRegistration.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/TraceSegmentServletHandler.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingHandler.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingListener.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/KeyWithStringValueJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/LogJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/ReferenceJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SegmentJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SpanJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/StreamJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegment.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegmentJsonReader.java create mode 100644 apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/UniqueIdJsonReader.java create mode 100644 apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataMonitor.java create mode 100644 apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleListenerService.java diff --git a/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/AgentModuleGRPCProvider.java b/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/AgentModuleGRPCProvider.java index 75d5b7ff9..2956732f2 100644 --- a/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/AgentModuleGRPCProvider.java +++ b/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/AgentModuleGRPCProvider.java @@ -21,7 +21,10 @@ package org.skywalking.apm.collector.agent.grpc; import java.util.Properties; import org.skywalking.apm.collector.agent.AgentModule; import org.skywalking.apm.collector.agent.grpc.handler.JVMMetricsServiceHandler; +import org.skywalking.apm.collector.agent.grpc.handler.naming.AgentGRPCNamingHandler; +import org.skywalking.apm.collector.agent.grpc.handler.naming.AgentGRPCNamingListener; import org.skywalking.apm.collector.cluster.ClusterModule; +import org.skywalking.apm.collector.cluster.service.ModuleListenerService; import org.skywalking.apm.collector.cluster.service.ModuleRegisterService; import org.skywalking.apm.collector.core.module.Module; import org.skywalking.apm.collector.core.module.ModuleNotFoundException; @@ -29,6 +32,8 @@ import org.skywalking.apm.collector.core.module.ModuleProvider; import org.skywalking.apm.collector.core.module.ServiceNotProvidedException; import org.skywalking.apm.collector.grpc.manager.GRPCManagerModule; import org.skywalking.apm.collector.grpc.manager.service.GRPCManagerService; +import org.skywalking.apm.collector.naming.NamingModule; +import org.skywalking.apm.collector.naming.service.NamingHandlerRegisterService; import org.skywalking.apm.collector.server.Server; import org.skywalking.apm.collector.storage.StorageModule; import org.skywalking.apm.collector.storage.service.DAOService; @@ -38,11 +43,12 @@ import org.skywalking.apm.collector.storage.service.DAOService; */ public class AgentModuleGRPCProvider extends ModuleProvider { + public static final String NAME = "gRPC"; private static final String HOST = "host"; private static final String PORT = "port"; @Override public String name() { - return "gRPC"; + return NAME; } @Override public Class module() { @@ -61,10 +67,17 @@ public class AgentModuleGRPCProvider extends ModuleProvider { ModuleRegisterService moduleRegisterService = getManager().find(ClusterModule.NAME).getService(ModuleRegisterService.class); moduleRegisterService.register(AgentModule.NAME, this.name(), new AgentModuleGRPCRegistration(host, port)); + AgentGRPCNamingListener namingListener = new AgentGRPCNamingListener(); + ModuleListenerService moduleListenerService = getManager().find(ClusterModule.NAME).getService(ModuleListenerService.class); + moduleListenerService.addListener(namingListener); + + NamingHandlerRegisterService namingHandlerRegisterService = getManager().find(NamingModule.NAME).getService(NamingHandlerRegisterService.class); + namingHandlerRegisterService.register(new AgentGRPCNamingHandler(namingListener)); + DAOService daoService = getManager().find(StorageModule.NAME).getService(DAOService.class); GRPCManagerService managerService = getManager().find(GRPCManagerModule.NAME).getService(GRPCManagerService.class); - Server gRPCServer = managerService.getOrCreateIfAbsent(host, port); + Server gRPCServer = managerService.createIfAbsent(host, port); addHandlers(daoService, gRPCServer); } catch (ModuleNotFoundException e) { throw new ServiceNotProvidedException(e.getMessage()); @@ -76,7 +89,7 @@ public class AgentModuleGRPCProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[0]; + return new String[] {ClusterModule.NAME, NamingModule.NAME, StorageModule.NAME, GRPCManagerModule.NAME}; } private void addHandlers(DAOService daoService, Server gRPCServer) { diff --git a/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingHandler.java b/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingHandler.java new file mode 100644 index 000000000..e729074ee --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingHandler.java @@ -0,0 +1,53 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.grpc.handler.naming; + +import com.google.gson.JsonArray; +import com.google.gson.JsonElement; +import java.util.Set; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; + +/** + * @author peng-yongsheng + */ +public class AgentGRPCNamingHandler extends JettyHandler { + + private final AgentGRPCNamingListener namingListener; + + public AgentGRPCNamingHandler(AgentGRPCNamingListener namingListener) { + this.namingListener = namingListener; + } + + @Override public String pathSpec() { + return "/agent/gRPC"; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + Set servers = namingListener.getAddresses(); + JsonArray serverArray = new JsonArray(); + servers.forEach(serverArray::add); + return serverArray; + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingListener.java b/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingListener.java new file mode 100644 index 000000000..b75a32d33 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-grpc-provider/src/main/java/org/skywalking/apm/collector/agent/grpc/handler/naming/AgentGRPCNamingListener.java @@ -0,0 +1,43 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.grpc.handler.naming; + +import org.skywalking.apm.collector.agent.AgentModule; +import org.skywalking.apm.collector.agent.grpc.AgentModuleGRPCProvider; +import org.skywalking.apm.collector.cluster.ClusterModuleListener; + +/** + * @author peng-yongsheng + */ +public class AgentGRPCNamingListener extends ClusterModuleListener { + + public static final String PATH = "/" + AgentModule.NAME + "/" + AgentModuleGRPCProvider.NAME; + + @Override public String path() { + return PATH; + } + + @Override public void serverJoinNotify(String serverAddress) { + + } + + @Override public void serverQuitNotify(String serverAddress) { + + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/pom.xml b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/pom.xml index 76638a1c3..9406a78be 100644 --- a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/pom.xml +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/pom.xml @@ -41,5 +41,10 @@ collector-agent-stream ${project.version} + + org.skywalking + collector-jetty-manager-define + ${project.version} + diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyProvider.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyProvider.java index 023e268b7..697e576cc 100644 --- a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyProvider.java +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyProvider.java @@ -20,17 +20,36 @@ package org.skywalking.apm.collector.agent.jetty; import java.util.Properties; import org.skywalking.apm.collector.agent.AgentModule; +import org.skywalking.apm.collector.agent.jetty.handler.TraceSegmentServletHandler; +import org.skywalking.apm.collector.agent.jetty.handler.naming.AgentJettyNamingHandler; +import org.skywalking.apm.collector.agent.jetty.handler.naming.AgentJettyNamingListener; +import org.skywalking.apm.collector.cluster.ClusterModule; +import org.skywalking.apm.collector.cluster.service.ModuleListenerService; +import org.skywalking.apm.collector.cluster.service.ModuleRegisterService; import org.skywalking.apm.collector.core.module.Module; +import org.skywalking.apm.collector.core.module.ModuleNotFoundException; import org.skywalking.apm.collector.core.module.ModuleProvider; import org.skywalking.apm.collector.core.module.ServiceNotProvidedException; +import org.skywalking.apm.collector.jetty.manager.JettyManagerModule; +import org.skywalking.apm.collector.jetty.manager.service.JettyManagerService; +import org.skywalking.apm.collector.naming.NamingModule; +import org.skywalking.apm.collector.naming.service.NamingHandlerRegisterService; +import org.skywalking.apm.collector.server.Server; +import org.skywalking.apm.collector.storage.StorageModule; +import org.skywalking.apm.collector.storage.service.DAOService; /** * @author peng-yongsheng */ public class AgentModuleJettyProvider extends ModuleProvider { + public static final String NAME = "jetty"; + private static final String HOST = "host"; + private static final String PORT = "port"; + private static final String CONTEXT_PATH = "context_path"; + @Override public String name() { - return "jetty"; + return NAME; } @Override public Class module() { @@ -42,7 +61,29 @@ public class AgentModuleJettyProvider extends ModuleProvider { } @Override public void start(Properties config) throws ServiceNotProvidedException { + String host = config.getProperty(HOST); + Integer port = (Integer)config.get(PORT); + String contextPath = config.getProperty(CONTEXT_PATH); + try { + ModuleRegisterService moduleRegisterService = getManager().find(ClusterModule.NAME).getService(ModuleRegisterService.class); + moduleRegisterService.register(AgentModule.NAME, this.name(), new AgentModuleJettyRegistration(host, port, contextPath)); + + AgentJettyNamingListener namingListener = new AgentJettyNamingListener(); + ModuleListenerService moduleListenerService = getManager().find(ClusterModule.NAME).getService(ModuleListenerService.class); + moduleListenerService.addListener(namingListener); + + NamingHandlerRegisterService namingHandlerRegisterService = getManager().find(NamingModule.NAME).getService(NamingHandlerRegisterService.class); + namingHandlerRegisterService.register(new AgentJettyNamingHandler(namingListener)); + + DAOService daoService = getManager().find(StorageModule.NAME).getService(DAOService.class); + + JettyManagerService managerService = getManager().find(JettyManagerModule.NAME).getService(JettyManagerService.class); + Server jettyServer = managerService.createIfAbsent(host, port, contextPath); + addHandlers(daoService, jettyServer); + } catch (ModuleNotFoundException e) { + throw new ServiceNotProvidedException(e.getMessage()); + } } @Override public void notifyAfterCompleted() throws ServiceNotProvidedException { @@ -50,6 +91,10 @@ public class AgentModuleJettyProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[0]; + return new String[] {ClusterModule.NAME, NamingModule.NAME, StorageModule.NAME, JettyManagerModule.NAME}; + } + + private void addHandlers(DAOService daoService, Server jettyServer) { + jettyServer.addHandler(new TraceSegmentServletHandler()); } } diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyRegistration.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyRegistration.java new file mode 100644 index 000000000..85158fb1c --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/AgentModuleJettyRegistration.java @@ -0,0 +1,41 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty; + +import org.skywalking.apm.collector.cluster.ModuleRegistration; + +/** + * @author peng-yongsheng + */ +public class AgentModuleJettyRegistration extends ModuleRegistration { + + private final String host; + private final int port; + private final String contextPath; + + public AgentModuleJettyRegistration(String host, int port, String contextPath) { + this.host = host; + this.port = port; + this.contextPath = contextPath; + } + + @Override public Value buildValue() { + return new Value(host, port, contextPath); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/TraceSegmentServletHandler.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/TraceSegmentServletHandler.java new file mode 100644 index 000000000..dd3ac8ed2 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/TraceSegmentServletHandler.java @@ -0,0 +1,73 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler; + +import com.google.gson.JsonElement; +import com.google.gson.stream.JsonReader; +import java.io.BufferedReader; +import java.io.IOException; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.agent.jetty.handler.reader.TraceSegment; +import org.skywalking.apm.collector.agent.jetty.handler.reader.TraceSegmentJsonReader; +import org.skywalking.apm.collector.agent.stream.parser.SegmentParse; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class TraceSegmentServletHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(TraceSegmentServletHandler.class); + + @Override public String pathSpec() { + return "/segments"; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + logger.debug("receive stream segment"); + try { + BufferedReader bufferedReader = req.getReader(); + read(bufferedReader); + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + return null; + } + + private TraceSegmentJsonReader jsonReader = new TraceSegmentJsonReader(); + + private void read(BufferedReader bufferedReader) throws IOException { + JsonReader reader = new JsonReader(bufferedReader); + + reader.beginArray(); + while (reader.hasNext()) { + SegmentParse segmentParse = new SegmentParse(null); + TraceSegment traceSegment = jsonReader.read(reader); + segmentParse.parse(traceSegment.getUpstreamSegment(), SegmentParse.Source.Agent); + } + reader.endArray(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingHandler.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingHandler.java new file mode 100644 index 000000000..ef0895493 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingHandler.java @@ -0,0 +1,53 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.naming; + +import com.google.gson.JsonArray; +import com.google.gson.JsonElement; +import java.util.Set; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; + +/** + * @author peng-yongsheng + */ +public class AgentJettyNamingHandler extends JettyHandler { + + private final AgentJettyNamingListener namingListener; + + public AgentJettyNamingHandler(AgentJettyNamingListener namingListener) { + this.namingListener = namingListener; + } + + @Override public String pathSpec() { + return "/agent/jetty"; + } + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + Set servers = namingListener.getAddresses(); + JsonArray serverArray = new JsonArray(); + servers.forEach(serverArray::add); + return serverArray; + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingListener.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingListener.java new file mode 100644 index 000000000..1dad304a9 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/naming/AgentJettyNamingListener.java @@ -0,0 +1,43 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.naming; + +import org.skywalking.apm.collector.agent.AgentModule; +import org.skywalking.apm.collector.agent.jetty.AgentModuleJettyProvider; +import org.skywalking.apm.collector.cluster.ClusterModuleListener; + +/** + * @author peng-yongsheng + */ +public class AgentJettyNamingListener extends ClusterModuleListener { + + public static final String PATH = "/" + AgentModule.NAME + "/" + AgentModuleJettyProvider.NAME; + + @Override public String path() { + return PATH; + } + + @Override public void serverJoinNotify(String serverAddress) { + + } + + @Override public void serverQuitNotify(String serverAddress) { + + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/KeyWithStringValueJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/KeyWithStringValueJsonReader.java new file mode 100644 index 000000000..8d05026e0 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/KeyWithStringValueJsonReader.java @@ -0,0 +1,54 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.KeyWithStringValue; + +/** + * @author peng-yongsheng + */ +public class KeyWithStringValueJsonReader implements StreamJsonReader { + + private static final String KEY = "k"; + private static final String VALUE = "v"; + + @Override public KeyWithStringValue read(JsonReader reader) throws IOException { + KeyWithStringValue.Builder builder = KeyWithStringValue.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case KEY: + builder.setKey(reader.nextString()); + break; + case VALUE: + builder.setValue(reader.nextString()); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/LogJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/LogJsonReader.java new file mode 100644 index 000000000..9448192ce --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/LogJsonReader.java @@ -0,0 +1,55 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.LogMessage; + +/** + * @author peng-yongsheng + */ +public class LogJsonReader implements StreamJsonReader { + + private KeyWithStringValueJsonReader keyWithStringValueJsonReader = new KeyWithStringValueJsonReader(); + + private static final String TIME = "ti"; + private static final String LOG_DATA = "ld"; + + @Override public LogMessage read(JsonReader reader) throws IOException { + LogMessage.Builder builder = LogMessage.newBuilder(); + + while (reader.hasNext()) { + switch (reader.nextName()) { + case TIME: + builder.setTime(reader.nextLong()); + case LOG_DATA: + reader.beginArray(); + while (reader.hasNext()) { + builder.addData(keyWithStringValueJsonReader.read(reader)); + } + reader.endArray(); + default: + reader.skipValue(); + } + } + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/ReferenceJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/ReferenceJsonReader.java new file mode 100644 index 000000000..c9ed25301 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/ReferenceJsonReader.java @@ -0,0 +1,92 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.TraceSegmentReference; + +/** + * @author peng-yongsheng + */ +public class ReferenceJsonReader implements StreamJsonReader { + + private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); + + private static final String PARENT_TRACE_SEGMENT_ID = "ts"; + private static final String PARENT_APPLICATION_ID = "ai"; + private static final String PARENT_SPAN_ID = "si"; + private static final String PARENT_SERVICE_ID = "vi"; + private static final String PARENT_SERVICE_NAME = "vn"; + private static final String NETWORK_ADDRESS_ID = "ni"; + private static final String NETWORK_ADDRESS = "nn"; + private static final String ENTRY_APPLICATION_INSTANCE_ID = "ea"; + private static final String ENTRY_SERVICE_ID = "ei"; + private static final String ENTRY_SERVICE_NAME = "en"; + private static final String REF_TYPE_VALUE = "rv"; + + @Override public TraceSegmentReference read(JsonReader reader) throws IOException { + TraceSegmentReference.Builder builder = TraceSegmentReference.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case PARENT_TRACE_SEGMENT_ID: + builder.setParentTraceSegmentId(uniqueIdJsonReader.read(reader)); + break; + case PARENT_APPLICATION_ID: + builder.setParentApplicationInstanceId(reader.nextInt()); + break; + case PARENT_SPAN_ID: + builder.setParentSpanId(reader.nextInt()); + break; + case PARENT_SERVICE_ID: + builder.setParentServiceId(reader.nextInt()); + break; + case PARENT_SERVICE_NAME: + builder.setParentServiceName(reader.nextString()); + break; + case NETWORK_ADDRESS_ID: + builder.setNetworkAddressId(reader.nextInt()); + break; + case NETWORK_ADDRESS: + builder.setNetworkAddress(reader.nextString()); + break; + case ENTRY_APPLICATION_INSTANCE_ID: + builder.setEntryApplicationInstanceId(reader.nextInt()); + break; + case ENTRY_SERVICE_ID: + builder.setEntryServiceId(reader.nextInt()); + break; + case ENTRY_SERVICE_NAME: + builder.setEntryServiceName(reader.nextString()); + break; + case REF_TYPE_VALUE: + builder.setRefTypeValue(reader.nextInt()); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SegmentJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SegmentJsonReader.java new file mode 100644 index 000000000..0997202cd --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SegmentJsonReader.java @@ -0,0 +1,87 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class SegmentJsonReader implements StreamJsonReader { + + private final Logger logger = LoggerFactory.getLogger(SegmentJsonReader.class); + + private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); + private ReferenceJsonReader referenceJsonReader = new ReferenceJsonReader(); + private SpanJsonReader spanJsonReader = new SpanJsonReader(); + + private static final String TRACE_SEGMENT_ID = "ts"; + private static final String APPLICATION_ID = "ai"; + private static final String APPLICATION_INSTANCE_ID = "ii"; + private static final String TRACE_SEGMENT_REFERENCE = "rs"; + private static final String SPANS = "ss"; + + @Override public TraceSegmentObject.Builder read(JsonReader reader) throws IOException { + TraceSegmentObject.Builder builder = TraceSegmentObject.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case TRACE_SEGMENT_ID: + builder.setTraceSegmentId(uniqueIdJsonReader.read(reader)); + if (logger.isDebugEnabled()) { + StringBuilder segmentId = new StringBuilder(); + builder.getTraceSegmentId().getIdPartsList().forEach(idPart -> segmentId.append(idPart)); + logger.debug("segment id: {}", segmentId); + } + break; + case APPLICATION_ID: + builder.setApplicationId(reader.nextInt()); + break; + case APPLICATION_INSTANCE_ID: + builder.setApplicationInstanceId(reader.nextInt()); + break; + case TRACE_SEGMENT_REFERENCE: + reader.beginArray(); + while (reader.hasNext()) { + builder.addRefs(referenceJsonReader.read(reader)); + } + reader.endArray(); + break; + case SPANS: + reader.beginArray(); + while (reader.hasNext()) { + builder.addSpans(spanJsonReader.read(reader)); + } + reader.endArray(); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder; + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SpanJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SpanJsonReader.java new file mode 100644 index 000000000..121d83b6f --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/SpanJsonReader.java @@ -0,0 +1,117 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.SpanObject; + +/** + * @author peng-yongsheng + */ +public class SpanJsonReader implements StreamJsonReader { + + private KeyWithStringValueJsonReader keyWithStringValueJsonReader = new KeyWithStringValueJsonReader(); + private LogJsonReader logJsonReader = new LogJsonReader(); + + private static final String SPAN_ID = "si"; + private static final String SPAN_TYPE_VALUE = "tv"; + private static final String SPAN_LAYER_VALUE = "lv"; + private static final String PARENT_SPAN_ID = "ps"; + private static final String START_TIME = "st"; + private static final String END_TIME = "et"; + private static final String COMPONENT_ID = "ci"; + private static final String COMPONENT_NAME = "cn"; + private static final String OPERATION_NAME_ID = "oi"; + private static final String OPERATION_NAME = "on"; + private static final String PEER_ID = "pi"; + private static final String PEER = "pn"; + private static final String IS_ERROR = "ie"; + private static final String TAGS = "to"; + private static final String LOGS = "lo"; + + @Override public SpanObject read(JsonReader reader) throws IOException { + SpanObject.Builder builder = SpanObject.newBuilder(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case SPAN_ID: + builder.setSpanId(reader.nextInt()); + break; + case SPAN_TYPE_VALUE: + builder.setSpanTypeValue(reader.nextInt()); + break; + case SPAN_LAYER_VALUE: + builder.setSpanLayerValue(reader.nextInt()); + break; + case PARENT_SPAN_ID: + builder.setParentSpanId(reader.nextInt()); + break; + case START_TIME: + builder.setStartTime(reader.nextLong()); + break; + case END_TIME: + builder.setEndTime(reader.nextLong()); + break; + case COMPONENT_ID: + builder.setComponentId(reader.nextInt()); + break; + case COMPONENT_NAME: + builder.setComponent(reader.nextString()); + break; + case OPERATION_NAME_ID: + builder.setOperationNameId(reader.nextInt()); + break; + case OPERATION_NAME: + builder.setOperationName(reader.nextString()); + break; + case PEER_ID: + builder.setPeerId(reader.nextInt()); + break; + case PEER: + builder.setPeer(reader.nextString()); + break; + case IS_ERROR: + builder.setIsError(reader.nextBoolean()); + break; + case TAGS: + reader.beginArray(); + while (reader.hasNext()) { + builder.addTags(keyWithStringValueJsonReader.read(reader)); + } + reader.endArray(); + break; + case LOGS: + reader.beginArray(); + while (reader.hasNext()) { + builder.addLogs(logJsonReader.read(reader)); + } + reader.endArray(); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/StreamJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/StreamJsonReader.java new file mode 100644 index 000000000..c5a0b255f --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/StreamJsonReader.java @@ -0,0 +1,29 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; + +/** + * @author peng-yongsheng + */ +public interface StreamJsonReader { + T read(JsonReader reader) throws IOException; +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegment.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegment.java new file mode 100644 index 000000000..d7d4164dc --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegment.java @@ -0,0 +1,47 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.UniqueId; +import org.skywalking.apm.network.proto.UpstreamSegment; + +/** + * @author peng-yongsheng + */ +public class TraceSegment { + + private UpstreamSegment.Builder builder; + + public TraceSegment() { + builder = UpstreamSegment.newBuilder(); + } + + public void addGlobalTraceId(UniqueId.Builder globalTraceId) { + builder.addGlobalTraceIds(globalTraceId); + } + + public void setTraceSegmentBuilder(TraceSegmentObject.Builder traceSegmentBuilder) { + builder.setSegment(traceSegmentBuilder.build().toByteString()); + } + + public UpstreamSegment getUpstreamSegment() { + return builder.build(); + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegmentJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegmentJsonReader.java new file mode 100644 index 000000000..125cb35ba --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/TraceSegmentJsonReader.java @@ -0,0 +1,65 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class TraceSegmentJsonReader implements StreamJsonReader { + + private final Logger logger = LoggerFactory.getLogger(TraceSegmentJsonReader.class); + + private UniqueIdJsonReader uniqueIdJsonReader = new UniqueIdJsonReader(); + private SegmentJsonReader segmentJsonReader = new SegmentJsonReader(); + + private static final String GLOBAL_TRACE_IDS = "gt"; + private static final String SEGMENT = "sg"; + + @Override public TraceSegment read(JsonReader reader) throws IOException { + TraceSegment traceSegment = new TraceSegment(); + + reader.beginObject(); + while (reader.hasNext()) { + switch (reader.nextName()) { + case GLOBAL_TRACE_IDS: + reader.beginArray(); + while (reader.hasNext()) { + traceSegment.addGlobalTraceId(uniqueIdJsonReader.read(reader)); + } + reader.endArray(); + + break; + case SEGMENT: + traceSegment.setTraceSegmentBuilder(segmentJsonReader.read(reader)); + break; + default: + reader.skipValue(); + break; + } + } + reader.endObject(); + + return traceSegment; + } +} diff --git a/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/UniqueIdJsonReader.java b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/UniqueIdJsonReader.java new file mode 100644 index 000000000..a172f98e9 --- /dev/null +++ b/apm-collector/apm-collector-agent/collector-agent-jetty-provider/src/main/java/org/skywalking/apm/collector/agent/jetty/handler/reader/UniqueIdJsonReader.java @@ -0,0 +1,40 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.agent.jetty.handler.reader; + +import com.google.gson.stream.JsonReader; +import java.io.IOException; +import org.skywalking.apm.network.proto.UniqueId; + +/** + * @author peng-yongsheng + */ +public class UniqueIdJsonReader implements StreamJsonReader { + + @Override public UniqueId.Builder read(JsonReader reader) throws IOException { + UniqueId.Builder builder = UniqueId.newBuilder(); + + reader.beginArray(); + while (reader.hasNext()) { + builder.addIdParts(reader.nextLong()); + } + reader.endArray(); + return builder; + } +} diff --git a/apm-collector/apm-collector-agent/pom.xml b/apm-collector/apm-collector-agent/pom.xml index 815f54994..bca43c861 100644 --- a/apm-collector/apm-collector-agent/pom.xml +++ b/apm-collector/apm-collector-agent/pom.xml @@ -24,5 +24,15 @@ apm-collector-core ${project.version} + + org.skywalking + collector-cluster-define + ${project.version} + + + org.skywalking + collector-naming-define + ${project.version} + diff --git a/apm-collector/apm-collector-boot/src/main/resources/application.yml b/apm-collector/apm-collector-boot/src/main/resources/application.yml index ad485f949..b286f0bf6 100644 --- a/apm-collector/apm-collector-boot/src/main/resources/application.yml +++ b/apm-collector/apm-collector-boot/src/main/resources/application.yml @@ -1,36 +1,32 @@ +cluster: + zookeeper: + hostPort: localhost:2181 + sessionTimeout: 100000 naming: jetty: host: localhost port: 10800 context_path: / -#agent_stream: -# grpc: -# host: localhost -# port: 11800 -# jetty: -# host: localhost -# port: 12800 -# context_path: / -# config: -# buffer_offset_max_file_size: 10M -# buffer_segment_max_file_size: 500M +agent: + gRPC: + host: localhost + port: 11800 + jetty: + host: localhost + port: 12800 + context_path: / + config: + buffer_offset_max_file_size: 10M + buffer_segment_max_file_size: 500M ui: jetty: host: localhost port: 12800 context_path: / -#collector_inside: -# grpc: -# host: localhost -# port: 11800 storage: elasticsearch: cluster_name: CollectorDBCluster cluster_transport_sniffer: true cluster_nodes: localhost:9300 index_shards_number: 2 - index_replicas_number: 0 -#storage: -# h2: -# url: jdbc:h2:tcp://localhost/~/test -# user_name: sa \ No newline at end of file + index_replicas_number: 0 \ No newline at end of file diff --git a/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterModuleStandaloneProvider.java b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterModuleStandaloneProvider.java index 9b3f07fd9..8788acc6a 100644 --- a/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterModuleStandaloneProvider.java +++ b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterModuleStandaloneProvider.java @@ -19,18 +19,33 @@ package org.skywalking.apm.collector.cluster.standalone; import java.util.Properties; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.client.h2.H2ClientException; import org.skywalking.apm.collector.cluster.ClusterModule; +import org.skywalking.apm.collector.cluster.service.ModuleListenerService; import org.skywalking.apm.collector.cluster.service.ModuleRegisterService; +import org.skywalking.apm.collector.cluster.standalone.service.StandaloneModuleListenerService; import org.skywalking.apm.collector.cluster.standalone.service.StandaloneModuleRegisterService; import org.skywalking.apm.collector.core.module.Module; import org.skywalking.apm.collector.core.module.ModuleProvider; import org.skywalking.apm.collector.core.module.ServiceNotProvidedException; +import org.skywalking.apm.collector.core.util.Const; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class ClusterModuleStandaloneProvider extends ModuleProvider { + private final Logger logger = LoggerFactory.getLogger(ClusterModuleStandaloneProvider.class); + + private static final String URL = "url"; + private static final String USER_NAME = "user_name"; + + private H2Client h2Client; + private ClusterStandaloneDataMonitor dataMonitor; + @Override public String name() { return "standalone"; } @@ -40,11 +55,23 @@ public class ClusterModuleStandaloneProvider extends ModuleProvider { } @Override public void prepare(Properties config) throws ServiceNotProvidedException { - this.registerServiceImplementation(ModuleRegisterService.class, new StandaloneModuleRegisterService()); + this.dataMonitor = new ClusterStandaloneDataMonitor(); + + final String url = config.getProperty(URL); + final String userName = config.getProperty(USER_NAME); + h2Client = new H2Client(url, userName, Const.EMPTY_STRING); + this.dataMonitor.setClient(h2Client); + + this.registerServiceImplementation(ModuleListenerService.class, new StandaloneModuleListenerService(dataMonitor)); + this.registerServiceImplementation(ModuleRegisterService.class, new StandaloneModuleRegisterService(dataMonitor)); } @Override public void start(Properties config) throws ServiceNotProvidedException { - + try { + h2Client.initialize(); + } catch (H2ClientException e) { + logger.error(e.getMessage(), e); + } } @Override public void notifyAfterCompleted() throws ServiceNotProvidedException { diff --git a/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataMonitor.java b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataMonitor.java new file mode 100644 index 000000000..9ceeb9960 --- /dev/null +++ b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataMonitor.java @@ -0,0 +1,91 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.cluster.standalone; + +import java.util.Iterator; +import java.util.LinkedHashMap; +import java.util.Map; +import org.skywalking.apm.collector.client.Client; +import org.skywalking.apm.collector.client.ClientException; +import org.skywalking.apm.collector.client.h2.H2Client; +import org.skywalking.apm.collector.cluster.ClusterModuleListener; +import org.skywalking.apm.collector.cluster.DataMonitor; +import org.skywalking.apm.collector.cluster.ModuleRegistration; +import org.skywalking.apm.collector.core.CollectorException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author peng-yongsheng + */ +public class ClusterStandaloneDataMonitor implements DataMonitor { + + private final Logger logger = LoggerFactory.getLogger(ClusterStandaloneDataMonitor.class); + + private H2Client client; + + private Map listeners; + private Map registrations; + + public ClusterStandaloneDataMonitor() { + listeners = new LinkedHashMap<>(); + registrations = new LinkedHashMap<>(); + } + + @Override public void setClient(Client client) { + this.client = (H2Client)client; + } + + @Override + public void addListener(ClusterModuleListener listener) { + String path = BASE_CATALOG + listener.path(); + logger.info("listener path: {}", path); + listeners.put(path, listener); + } + + @Override public ClusterModuleListener getListener(String path) { + path = BASE_CATALOG + path; + return listeners.get(path); + } + + @Override public void register(String path, ModuleRegistration registration) { + registrations.put(BASE_CATALOG + path, registration); + } + + @Override public void createPath(String path) throws ClientException { + + } + + @Override public void setData(String path, String value) throws ClientException { + if (listeners.containsKey(path)) { + listeners.get(path).addAddress(value); + listeners.get(path).serverJoinNotify(value); + } + } + + public void start() throws CollectorException { + Iterator> entryIterator = registrations.entrySet().iterator(); + while (entryIterator.hasNext()) { + Map.Entry next = entryIterator.next(); + ModuleRegistration.Value value = next.getValue().buildValue(); + String contextPath = value.getContextPath() == null ? "" : value.getContextPath(); + setData(next.getKey(), value.getHostPort() + contextPath); + } + } +} diff --git a/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleListenerService.java b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleListenerService.java new file mode 100644 index 000000000..c81753ebd --- /dev/null +++ b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleListenerService.java @@ -0,0 +1,39 @@ +/* + * Copyright 2017, OpenSkywalking Organization All rights reserved. + * + * Licensed 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. + * + * Project repository: https://github.com/OpenSkywalking/skywalking + */ + +package org.skywalking.apm.collector.cluster.standalone.service; + +import org.skywalking.apm.collector.cluster.ClusterModuleListener; +import org.skywalking.apm.collector.cluster.service.ModuleListenerService; +import org.skywalking.apm.collector.cluster.standalone.ClusterStandaloneDataMonitor; + +/** + * @author peng-yongsheng + */ +public class StandaloneModuleListenerService implements ModuleListenerService { + + private final ClusterStandaloneDataMonitor dataMonitor; + + public StandaloneModuleListenerService(ClusterStandaloneDataMonitor dataMonitor) { + this.dataMonitor = dataMonitor; + } + + @Override public void addListener(ClusterModuleListener listener) { + dataMonitor.addListener(listener); + } +} diff --git a/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleRegisterService.java b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleRegisterService.java index 13d1f491d..b4bf53ce8 100644 --- a/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleRegisterService.java +++ b/apm-collector/apm-collector-cluster/collector-cluster-standalone-provider/src/main/java/org/skywalking/apm/collector/cluster/standalone/service/StandaloneModuleRegisterService.java @@ -20,13 +20,21 @@ package org.skywalking.apm.collector.cluster.standalone.service; import org.skywalking.apm.collector.cluster.ModuleRegistration; import org.skywalking.apm.collector.cluster.service.ModuleRegisterService; +import org.skywalking.apm.collector.cluster.standalone.ClusterStandaloneDataMonitor; /** * @author peng-yongsheng */ public class StandaloneModuleRegisterService implements ModuleRegisterService { - @Override public void register(String moduleName, String providerName, ModuleRegistration registration) { + private final ClusterStandaloneDataMonitor dataMonitor; + public StandaloneModuleRegisterService(ClusterStandaloneDataMonitor dataMonitor) { + this.dataMonitor = dataMonitor; } -} + + @Override public void register(String moduleName, String providerName, ModuleRegistration registration) { + String path = "/" + moduleName + "/" + providerName; + dataMonitor.register(path, registration); + } +} \ No newline at end of file diff --git a/apm-collector/apm-collector-component/server-component/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java b/apm-collector/apm-collector-component/server-component/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java index 6f8a5c2e9..a606e9726 100644 --- a/apm-collector/apm-collector-component/server-component/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java +++ b/apm-collector/apm-collector-component/server-component/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java @@ -73,6 +73,7 @@ public class JettyServer implements Server { } @Override public void start() throws ServerException { + logger.info("start server, host: {}, port: {}", host, port); try { for (ServletMapping servletMapping : servletContextHandler.getServletHandler().getServletMappings()) { logger.info("jetty servlet mappings: {} register by {}", servletMapping.getPathSpecs(), servletMapping.getServletName()); diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java index de48bbe6b..aaa7ed165 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleProvider.java @@ -96,7 +96,7 @@ public abstract class ModuleProvider { Service service) throws ServiceNotProvidedException { if (serviceType.isInstance(service)) { if (manager.isServiceInstrument()) { - service = ServiceInstrumentation.INSTANCE.buildServiceUnderMonitor(module.name(), name(), service); +// service = ServiceInstrumentation.INSTANCE.buildServiceUnderMonitor(module.name(), name(), service); } this.services.put(serviceType, service); } else { diff --git a/apm-collector/apm-collector-core/src/main/resources/application-default.yml b/apm-collector/apm-collector-core/src/main/resources/application-default.yml index 5d52df356..b7ab17dbc 100644 --- a/apm-collector/apm-collector-core/src/main/resources/application-default.yml +++ b/apm-collector/apm-collector-core/src/main/resources/application-default.yml @@ -1,23 +1,14 @@ cluster: - h2: - hostPort: localhost:2181 - sessionTimeout: 100000 + standalone: + url: jdbc:h2:~/memorydb + user_name: sa +cache: + guava: naming: jetty: host: localhost port: 10800 context_path: / -#agent_stream: -# grpc: -# host: localhost -# port: 11800 -# jetty: -# host: localhost -# port: 12800 -# context_path: / -# config: -# buffer_offset_max_file_size: 10M -# buffer_segment_max_file_size: 500M ui: jetty: host: localhost @@ -27,12 +18,7 @@ jetty_manager: jetty: gRPC_manager: gRPC: -#collector_inside: -# grpc: -# host: localhost -# port: 11800 storage: h2: url: jdbc:h2:~/memorydb - user_name: sa - password: \ No newline at end of file + user_name: sa \ No newline at end of file diff --git a/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-define/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerService.java b/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-define/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerService.java index 2b923a16a..13180e305 100644 --- a/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-define/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerService.java +++ b/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-define/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerService.java @@ -25,5 +25,5 @@ import org.skywalking.apm.collector.server.Server; * @author peng-yongsheng */ public interface GRPCManagerService extends Service { - Server getOrCreateIfAbsent(String host, int port); + Server createIfAbsent(String host, int port); } diff --git a/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-provider/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerServiceImpl.java b/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-provider/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerServiceImpl.java index a6b3c4589..b3e04df16 100644 --- a/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-provider/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerServiceImpl.java +++ b/apm-collector/apm-collector-grpc-manager/collector-grpc-manager-provider/src/main/java/org/skywalking/apm/collector/grpc/manager/service/GRPCManagerServiceImpl.java @@ -38,7 +38,7 @@ public class GRPCManagerServiceImpl implements GRPCManagerService { this.servers = servers; } - @Override public Server getOrCreateIfAbsent(String host, int port) { + @Override public Server createIfAbsent(String host, int port) { String id = host + String.valueOf(port); if (servers.containsKey(id)) { return servers.get(id); diff --git a/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-define/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerService.java b/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-define/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerService.java index a2609363f..c22242a91 100644 --- a/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-define/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerService.java +++ b/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-define/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerService.java @@ -20,10 +20,13 @@ package org.skywalking.apm.collector.jetty.manager.service; import org.skywalking.apm.collector.core.module.Service; import org.skywalking.apm.collector.server.Server; +import org.skywalking.apm.collector.server.ServerHandler; /** * @author peng-yongsheng */ public interface JettyManagerService extends Service { - Server getOrCreateIfAbsent(String host, int port, String contextPath); + Server createIfAbsent(String host, int port, String contextPath); + + void addHandler(String host, int port, ServerHandler serverHandler); } diff --git a/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-provider/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerServiceImpl.java b/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-provider/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerServiceImpl.java index e97634895..a6eb99db2 100644 --- a/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-provider/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerServiceImpl.java +++ b/apm-collector/apm-collector-jetty-manager/collector-jetty-manager-provider/src/main/java/org/skywalking/apm/collector/jetty/manager/service/JettyManagerServiceImpl.java @@ -19,8 +19,10 @@ package org.skywalking.apm.collector.jetty.manager.service; import java.util.Map; +import org.skywalking.apm.collector.core.UnexpectedException; import org.skywalking.apm.collector.server.Server; import org.skywalking.apm.collector.server.ServerException; +import org.skywalking.apm.collector.server.ServerHandler; import org.skywalking.apm.collector.server.jetty.JettyServer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -38,7 +40,7 @@ public class JettyManagerServiceImpl implements JettyManagerService { this.servers = servers; } - @Override public Server getOrCreateIfAbsent(String host, int port, String contextPath) { + @Override public Server createIfAbsent(String host, int port, String contextPath) { String id = host + String.valueOf(port); if (servers.containsKey(id)) { return servers.get(id); @@ -53,4 +55,13 @@ public class JettyManagerServiceImpl implements JettyManagerService { return server; } } + + @Override public void addHandler(String host, int port, ServerHandler serverHandler) { + String id = host + String.valueOf(port); + if (servers.containsKey(id)) { + servers.get(id).addHandler(serverHandler); + } else { + throw new UnexpectedException("Please create server before add server handler."); + } + } } diff --git a/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/NamingModuleJettyProvider.java b/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/NamingModuleJettyProvider.java index eeb8e1726..b6e63440f 100644 --- a/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/NamingModuleJettyProvider.java +++ b/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/NamingModuleJettyProvider.java @@ -18,8 +18,6 @@ package org.skywalking.apm.collector.naming.jetty; -import java.util.ArrayList; -import java.util.List; import java.util.Properties; import org.skywalking.apm.collector.cluster.ClusterModule; import org.skywalking.apm.collector.core.module.Module; @@ -31,8 +29,6 @@ import org.skywalking.apm.collector.jetty.manager.service.JettyManagerService; import org.skywalking.apm.collector.naming.NamingModule; import org.skywalking.apm.collector.naming.jetty.service.NamingJettyHandlerRegisterService; import org.skywalking.apm.collector.naming.service.NamingHandlerRegisterService; -import org.skywalking.apm.collector.server.Server; -import org.skywalking.apm.collector.server.ServerHandler; /** * @author peng-yongsheng @@ -42,7 +38,6 @@ public class NamingModuleJettyProvider extends ModuleProvider { private static final String HOST = "host"; private static final String PORT = "port"; private static final String CONTEXT_PATH = "context_path"; - private final List handlers = new ArrayList<>(); @Override public String name() { return "jetty"; @@ -53,7 +48,9 @@ public class NamingModuleJettyProvider extends ModuleProvider { } @Override public void prepare(Properties config) throws ServiceNotProvidedException { - this.registerServiceImplementation(NamingHandlerRegisterService.class, new NamingJettyHandlerRegisterService(handlers)); + final String host = config.getProperty(HOST); + final Integer port = (Integer)config.get(PORT); + this.registerServiceImplementation(NamingHandlerRegisterService.class, new NamingJettyHandlerRegisterService(host, port, getManager())); } @Override public void start(Properties config) throws ServiceNotProvidedException { @@ -63,18 +60,16 @@ public class NamingModuleJettyProvider extends ModuleProvider { try { JettyManagerService managerService = getManager().find(JettyManagerModule.NAME).getService(JettyManagerService.class); - Server jettyServer = managerService.getOrCreateIfAbsent(host, port, contextPath); - handlers.forEach(jettyServer::addHandler); + managerService.createIfAbsent(host, port, contextPath); } catch (ModuleNotFoundException e) { throw new ServiceNotProvidedException(e.getMessage()); } } @Override public void notifyAfterCompleted() throws ServiceNotProvidedException { - } @Override public String[] requiredModules() { - return new String[] {JettyManagerModule.NAME, ClusterModule.NAME}; + return new String[] {ClusterModule.NAME, JettyManagerModule.NAME}; } } diff --git a/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/service/NamingJettyHandlerRegisterService.java b/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/service/NamingJettyHandlerRegisterService.java index d3b4134b7..5a60a17a7 100644 --- a/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/service/NamingJettyHandlerRegisterService.java +++ b/apm-collector/apm-collector-naming/collector-naming-jetty-provider/src/main/java/org/skywalking/apm/collector/naming/jetty/service/NamingJettyHandlerRegisterService.java @@ -18,22 +18,39 @@ package org.skywalking.apm.collector.naming.jetty.service; -import java.util.List; +import org.skywalking.apm.collector.core.module.ModuleManager; +import org.skywalking.apm.collector.core.module.ModuleNotFoundException; +import org.skywalking.apm.collector.core.module.ServiceNotProvidedException; +import org.skywalking.apm.collector.jetty.manager.JettyManagerModule; +import org.skywalking.apm.collector.jetty.manager.service.JettyManagerService; import org.skywalking.apm.collector.naming.service.NamingHandlerRegisterService; import org.skywalking.apm.collector.server.ServerHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author peng-yongsheng */ public class NamingJettyHandlerRegisterService implements NamingHandlerRegisterService { - private final List handlers; + private final Logger logger = LoggerFactory.getLogger(NamingJettyHandlerRegisterService.class); - public NamingJettyHandlerRegisterService(List handlers) { - this.handlers = handlers; + private final ModuleManager moduleManager; + private final String host; + private final int port; + + public NamingJettyHandlerRegisterService(String host, int port, ModuleManager moduleManager) { + this.moduleManager = moduleManager; + this.host = host; + this.port = port; } @Override public void register(ServerHandler namingHandler) { - handlers.add(namingHandler); + try { + JettyManagerService managerService = moduleManager.find(JettyManagerModule.NAME).getService(JettyManagerService.class); + managerService.addHandler(this.host, this.port, namingHandler); + } catch (ModuleNotFoundException | ServiceNotProvidedException e) { + logger.error(e.getMessage(), e); + } } } diff --git a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteModuleGRPCProvider.java b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteModuleGRPCProvider.java index 6cf92764b..a6e221b96 100644 --- a/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteModuleGRPCProvider.java +++ b/apm-collector/apm-collector-remote/collector-remote-grpc-provider/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteModuleGRPCProvider.java @@ -50,7 +50,7 @@ public class RemoteModuleGRPCProvider extends ModuleProvider { try { GRPCManagerService managerService = getManager().find(GRPCManagerModule.NAME).getService(GRPCManagerService.class); - Server gRPCServer = managerService.getOrCreateIfAbsent(host, port); + Server gRPCServer = managerService.createIfAbsent(host, port); gRPCServer.addHandler(new RemoteCommonServiceHandler(listener)); ModuleRegisterService moduleRegisterService = getManager().find(ClusterModule.NAME).getService(ModuleRegisterService.class); diff --git a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/es_dao.define index e9d3626f4..8bbb83735 100644 --- a/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/es_dao.define +++ b/apm-collector/apm-collector-storage/collector-storage-es-provider/src/main/resources/META-INF/defines/es_dao.define @@ -1,19 +1,35 @@ -org.skywalking.apm.collector.storage.es.dao.CpuMetricEsDAO -org.skywalking.apm.collector.storage.es.dao.GCMetricEsDAO -org.skywalking.apm.collector.storage.es.dao.InstanceHeartBeatEsDAO -org.skywalking.apm.collector.storage.es.dao.MemoryMetricEsDAO -org.skywalking.apm.collector.storage.es.dao.MemoryPoolMetricEsDAO -org.skywalking.apm.collector.storage.es.dao.ApplicationEsDAO -org.skywalking.apm.collector.storage.es.dao.InstanceEsDAO -org.skywalking.apm.collector.storage.es.dao.ServiceNameEsDAO -org.skywalking.apm.collector.storage.es.dao.GlobalTraceEsDAO -org.skywalking.apm.collector.storage.es.dao.InstPerformanceEsDAO -org.skywalking.apm.collector.storage.es.dao.NodeComponentEsDAO -org.skywalking.apm.collector.storage.es.dao.NodeReferenceEsDAO -org.skywalking.apm.collector.storage.es.dao.SegmentCostEsDAO -org.skywalking.apm.collector.storage.es.dao.SegmentEsDAO -org.skywalking.apm.collector.storage.es.dao.ServiceEntryEsDAO -org.skywalking.apm.collector.storage.es.dao.ServiceReferenceEsDAO org.skywalking.apm.collector.storage.es.dao.ApplicationEsCacheDAO org.skywalking.apm.collector.storage.es.dao.InstanceEsCacheDAO -org.skywalking.apm.collector.storage.es.dao.ServiceNameEsCacheDAO \ No newline at end of file +org.skywalking.apm.collector.storage.es.dao.ServiceNameEsCacheDAO + +org.skywalking.apm.collector.storage.es.dao.ApplicationEsStreamDAO +org.skywalking.apm.collector.storage.es.dao.InstanceEsStreamDAO +org.skywalking.apm.collector.storage.es.dao.ServiceNameEsStreamDAO + +org.skywalking.apm.collector.storage.es.dao.CpuMetricEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.GCMetricEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.MemoryMetricEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.MemoryPoolMetricEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.GlobalTraceEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.InstanceHeartBeatEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.InstPerformanceEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.NodeComponentEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.NodeMappingEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.NodeReferenceEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.SegmentCostEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.SegmentEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.ServiceEntryEsPersistenceDAO +org.skywalking.apm.collector.storage.es.dao.ServiceReferenceEsPersistenceDAO + +org.skywalking.apm.collector.storage.es.dao.InstanceEsUIDAO +org.skywalking.apm.collector.storage.es.dao.InstPerformanceEsUIDAO +org.skywalking.apm.collector.storage.es.dao.CpuMetricEsUIDAO +org.skywalking.apm.collector.storage.es.dao.GCMetricEsUIDAO +org.skywalking.apm.collector.storage.es.dao.MemoryMetricEsUIDAO +org.skywalking.apm.collector.storage.es.dao.GlobalTraceEsUIDAO +org.skywalking.apm.collector.storage.es.dao.NodeComponentEsUIDAO +org.skywalking.apm.collector.storage.es.dao.NodeMappingEsUIDAO +org.skywalking.apm.collector.storage.es.dao.NodeReferenceEsUIDAO +org.skywalking.apm.collector.storage.es.dao.SegmentCostEsUIDAO +org.skywalking.apm.collector.storage.es.dao.SegmentEsUIDAO +org.skywalking.apm.collector.storage.es.dao.ServiceEntryEsUIDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/resources/META-INF/defines/h2_dao.define index 3541b3722..7142ddeb2 100644 --- a/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/resources/META-INF/defines/h2_dao.define +++ b/apm-collector/apm-collector-storage/collector-storage-h2-provider/src/main/resources/META-INF/defines/h2_dao.define @@ -1,19 +1,35 @@ -org.skywalking.apm.collector.storage.h2.dao.CpuMetricH2DAO -org.skywalking.apm.collector.storage.h2.dao.GCMetricH2DAO -org.skywalking.apm.collector.storage.h2.dao.InstanceHeartBeatH2DAO -org.skywalking.apm.collector.storage.h2.dao.MemoryMetricH2DAO -org.skywalking.apm.collector.storage.h2.dao.MemoryPoolMetricH2DAO -org.skywalking.apm.collector.storage.h2.dao.ApplicationH2DAO -org.skywalking.apm.collector.storage.h2.dao.InstanceH2DAO -org.skywalking.apm.collector.storage.h2.dao.ServiceNameH2DAO -org.skywalking.apm.collector.storage.h2.dao.GlobalTraceH2DAO -org.skywalking.apm.collector.storage.h2.dao.InstPerformanceH2DAO -org.skywalking.apm.collector.storage.h2.dao.NodeComponentH2DAO -org.skywalking.apm.collector.storage.h2.dao.NodeReferenceH2DAO -org.skywalking.apm.collector.storage.h2.dao.SegmentCostH2DAO -org.skywalking.apm.collector.storage.h2.dao.SegmentH2DAO -org.skywalking.apm.collector.storage.h2.dao.ServiceEntryH2DAO -org.skywalking.apm.collector.storage.h2.dao.ServiceReferenceH2DAO org.skywalking.apm.collector.storage.h2.dao.ApplicationH2CacheDAO org.skywalking.apm.collector.storage.h2.dao.InstanceH2CacheDAO -org.skywalking.apm.collector.storage.h2.dao.ServiceNameH2CacheDAO \ No newline at end of file +org.skywalking.apm.collector.storage.h2.dao.ServiceNameH2CacheDAO + +org.skywalking.apm.collector.storage.h2.dao.ApplicationH2StreamDAO +org.skywalking.apm.collector.storage.h2.dao.InstanceH2StreamDAO +org.skywalking.apm.collector.storage.h2.dao.ServiceNameH2StreamDAO + +org.skywalking.apm.collector.storage.h2.dao.CpuMetricH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.GCMetricH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.MemoryMetricH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.MemoryPoolMetricH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.GlobalTraceH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.InstanceHeartBeatH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.InstPerformanceH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.NodeComponentH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.NodeMappingH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.NodeReferenceH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.SegmentCostH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.SegmentH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.ServiceEntryH2PersistenceDAO +org.skywalking.apm.collector.storage.h2.dao.ServiceReferenceH2PersistenceDAO + +org.skywalking.apm.collector.storage.h2.dao.InstanceH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.InstPerformanceH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.CpuMetricH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.GCMetricH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.MemoryMetricH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.GlobalTraceH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.NodeComponentH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.NodeMappingH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.NodeReferenceH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.SegmentCostH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.SegmentH2UIDAO +org.skywalking.apm.collector.storage.h2.dao.ServiceEntryH2UIDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-ui/collector-ui-jetty-provider/src/main/java/org/skywalking/apm/collector/ui/jetty/UIModuleJettyProvider.java b/apm-collector/apm-collector-ui/collector-ui-jetty-provider/src/main/java/org/skywalking/apm/collector/ui/jetty/UIModuleJettyProvider.java index d83df0a7b..9b6c5fda7 100644 --- a/apm-collector/apm-collector-ui/collector-ui-jetty-provider/src/main/java/org/skywalking/apm/collector/ui/jetty/UIModuleJettyProvider.java +++ b/apm-collector/apm-collector-ui/collector-ui-jetty-provider/src/main/java/org/skywalking/apm/collector/ui/jetty/UIModuleJettyProvider.java @@ -20,6 +20,7 @@ package org.skywalking.apm.collector.ui.jetty; import java.util.Properties; import org.skywalking.apm.collector.cache.CacheModule; +import org.skywalking.apm.collector.cache.CacheServiceManager; import org.skywalking.apm.collector.cache.service.ApplicationCacheService; import org.skywalking.apm.collector.cache.service.InstanceCacheService; import org.skywalking.apm.collector.cache.service.ServiceIdCacheService; @@ -54,7 +55,6 @@ import org.skywalking.apm.collector.ui.jetty.handler.servicetree.EntryServiceGet import org.skywalking.apm.collector.ui.jetty.handler.servicetree.ServiceTreeGetByIdHandler; import org.skywalking.apm.collector.ui.jetty.handler.time.AllInstanceLastTimeGetHandler; import org.skywalking.apm.collector.ui.jetty.handler.time.OneInstanceLastTimeGetHandler; -import org.skywalking.apm.collector.cache.CacheServiceManager; /** * @author peng-yongsheng @@ -97,7 +97,7 @@ public class UIModuleJettyProvider extends ModuleProvider { DAOService daoService = getManager().find(StorageModule.NAME).getService(DAOService.class); JettyManagerService managerService = getManager().find(JettyManagerModule.NAME).getService(JettyManagerService.class); - Server jettyServer = managerService.getOrCreateIfAbsent(host, port, contextPath); + Server jettyServer = managerService.createIfAbsent(host, port, contextPath); addHandlers(daoService, jettyServer, cacheServiceManager); } catch (ModuleNotFoundException e) { throw new ServiceNotProvidedException(e.getMessage()); @@ -109,7 +109,7 @@ public class UIModuleJettyProvider extends ModuleProvider { } @Override public String[] requiredModules() { - return new String[] {ClusterModule.NAME, JettyManagerModule.NAME, NamingModule.NAME, CacheModule.NAME}; + return new String[] {ClusterModule.NAME, JettyManagerModule.NAME, NamingModule.NAME, CacheModule.NAME, StorageModule.NAME}; } private void addHandlers(DAOService daoService, Server jettyServer, CacheServiceManager cacheServiceManager) {