diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java new file mode 100644 index 000000000..47a4803f6 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java @@ -0,0 +1,9 @@ +package com.a.eye.skywalking.collector.actor; + +import akka.actor.UntypedActor; + +/** + * @author pengys5 + */ +public abstract class AbstractWorker extends UntypedActor { +} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java index 145c2efe4..2442eb842 100644 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java @@ -16,6 +16,7 @@ public abstract class AbstractWorkerProvider { public abstract int workerNum(); public void createWorker(ActorSystem system) { + System.out.println("workerName: " + workerName()); if (StringUtil.isEmpty(workerName())) { throw new IllegalArgumentException("cannot createWorker() with anything not obtained from workerName()"); } diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorBoot.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorBootstrap.java similarity index 59% rename from skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorBoot.java rename to skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorBootstrap.java index da5df6495..500d54d4e 100644 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorBoot.java +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorBootstrap.java @@ -5,8 +5,8 @@ import akka.actor.ActorSystem; /** * @author pengys5 */ -public class CollectorBoot { +public class CollectorBootstrap { public static void main(String[] args) { - ActorSystem system = ActorSystem.create("ClusterSystem", config); +// ActorSystem system = ActorSystem.create("ClusterSystem", config); } } diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorConfig.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorConfig.java deleted file mode 100644 index 47640fb11..000000000 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorConfig.java +++ /dev/null @@ -1,19 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -/** - * Created by pengys5 on 2017/2/22 0022. - */ -public class CollectorConfig { - - public static final String appname = "CollectorSystem"; - - public static class Collector { - public static String hostname = "127.0.0.1"; - public static String port = "2551"; - public static String cluster = "127.0.0.1:2551"; - - public static class Worker { - public static int ApplicationDiscoverMetric_Num = 2; - } - } -} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorConfigInitializer.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorConfigInitializer.java deleted file mode 100644 index d1ac3af6c..000000000 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/CollectorConfigInitializer.java +++ /dev/null @@ -1,43 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -import com.a.eye.skywalking.api.logging.api.ILog; -import com.a.eye.skywalking.api.logging.api.LogManager; -import com.a.eye.skywalking.api.util.ConfigInitializer; -import com.a.eye.skywalking.api.util.StringUtil; - -import java.io.InputStream; -import java.util.Properties; - -/** - * @author pengys5 - */ -public class CollectorConfigInitializer { - - private static ILog logger = LogManager.getLogger(CollectorConfigInitializer.class); - - public static void initialize() { - InputStream configFileStream = CollectorConfigInitializer.class.getResourceAsStream("/collector.config"); - - if (configFileStream == null) { - logger.info("Not provide sky-walking certification documents, sky-walking api run in default config."); - } else { - try { - Properties properties = new Properties(); - properties.load(configFileStream); - ConfigInitializer.initialize(properties, CollectorConfig.class); - } catch (Exception e) { - logger.error("Failed to read the config file, sky-walking api run in default config.", e); - } - } - - if (!StringUtil.isEmpty(System.getProperty("collector.hostname"))) { - CollectorConfig.Collector.hostname = System.getProperty("collector.hostname"); - } - if (!StringUtil.isEmpty(System.getProperty("collector.port"))) { - CollectorConfig.Collector.port = System.getProperty("collector.port"); - } - if (!StringUtil.isEmpty(System.getProperty("collector.cluster"))) { - CollectorConfig.Collector.cluster = System.getProperty("collector.cluster"); - } - } -} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java new file mode 100644 index 000000000..40217f654 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java @@ -0,0 +1,19 @@ +package com.a.eye.skywalking.collector.actor; + +import akka.actor.ActorSystem; + +import java.util.ServiceLoader; + +/** + * @author pengys5 + */ +public enum WorkersCreator { + INSTANCE; + + public void boot(ActorSystem system) { + ServiceLoader serviceLoader = ServiceLoader.load(AbstractWorkerProvider.class); + for (AbstractWorkerProvider provider : serviceLoader) { + provider.createWorker(system); + } + } +} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/resources/services/com.a.eye.skywalking.collector.actor.AbstractWorkerProvider b/skywalking-collector/skywalking-collector-actor/src/main/resources/services/com.a.eye.skywalking.collector.actor.AbstractWorkerProvider new file mode 100644 index 000000000..8d06ff8e3 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/main/resources/services/com.a.eye.skywalking.collector.actor.AbstractWorkerProvider @@ -0,0 +1 @@ +com.a.eye.skywalking.collector.actor.SpiTestWorkerFactory \ No newline at end of file diff --git a/skywalking-collector/skywalking-collector-actor/src/main/resources/services/com.a.eye.skywalking.collector.cluster.base.IActorProvider b/skywalking-collector/skywalking-collector-actor/src/main/resources/services/com.a.eye.skywalking.collector.cluster.base.IActorProvider deleted file mode 100644 index 4a49d7b57..000000000 --- a/skywalking-collector/skywalking-collector-actor/src/main/resources/services/com.a.eye.skywalking.collector.cluster.base.IActorProvider +++ /dev/null @@ -1 +0,0 @@ -com.a.eye.skywalking.collector.cluster.manager.ActorManagerActorFactory \ No newline at end of file diff --git a/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java new file mode 100644 index 000000000..c64c1daa5 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java @@ -0,0 +1,12 @@ +package com.a.eye.skywalking.collector.actor; + +/** + * @author pengys5 + */ +public class SpiTestWorker extends AbstractWorker { + + @Override + public void onReceive(Object message) throws Throwable { + + } +} diff --git a/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java new file mode 100644 index 000000000..186cd4653 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java @@ -0,0 +1,24 @@ +package com.a.eye.skywalking.collector.actor; + +/** + * @author pengys5 + */ +public class SpiTestWorkerFactory extends AbstractWorkerProvider { + + public static final String WorkerName = "SpiTestWorker"; + + @Override + public String workerName() { + return WorkerName; + } + + @Override + public Class workerClass() { + return SpiTestWorker.class; + } + + @Override + public int workerNum() { + return 2; + } +} diff --git a/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactoryTestCase.java b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactoryTestCase.java new file mode 100644 index 000000000..5b3ed571a --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactoryTestCase.java @@ -0,0 +1,34 @@ +package com.a.eye.skywalking.collector.actor; + +import akka.actor.ActorSystem; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.mockito.Mockito; + +/** + * @author pengys5 + */ +public class SpiTestWorkerFactoryTestCase { + + ActorSystem system; + + @Before + public void createSystem() { + system = ActorSystem.create(); + } + + @After + public void terminateSystem() throws IllegalAccessException { + system.terminate(); + system.awaitTermination(); + system = null; + } + + @Test + public void testWorkerCreate() { + SpiTestWorkerFactory factory = Mockito.mock(SpiTestWorkerFactory.class); + Mockito.when(factory.workerName()).thenReturn(""); + factory.createWorker(system); + } +} diff --git a/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/WorkersCreatorTestCase.java b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/WorkersCreatorTestCase.java new file mode 100644 index 000000000..95aae37de --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/test/java/com/a/eye/skywalking/collector/actor/WorkersCreatorTestCase.java @@ -0,0 +1,31 @@ +package com.a.eye.skywalking.collector.actor; + +import akka.actor.ActorSystem; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +/** + * @author pengys5 + */ +public class WorkersCreatorTestCase { + + ActorSystem system; + + @Before + public void createSystem() { + system = ActorSystem.create(); + } + + @After + public void terminateSystem() throws IllegalAccessException { + system.terminate(); + system.awaitTermination(); + system = null; + } + + @Test + public void testBoot() { + WorkersCreator.INSTANCE.boot(system); + } +}