clusterServiceLoader = ServiceLoader.load(AbstractLocalWorkerProvider.class);
- for (AbstractLocalWorkerProvider provider : clusterServiceLoader) {
- logger.info("loadLocalProviders provider name: %s", provider.getClass().getName());
- provider.setClusterContext(clusterContext);
- clusterContext.putProvider(provider);
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractClusterWorker.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractClusterWorker.java
deleted file mode 100644
index 8a2289366..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractClusterWorker.java
+++ /dev/null
@@ -1,123 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import akka.actor.UntypedActor;
-import akka.cluster.Cluster;
-import akka.cluster.ClusterEvent;
-import akka.cluster.Member;
-import akka.cluster.MemberStatus;
-import org.apache.logging.log4j.Logger;
-import org.skywalking.apm.collector.cluster.WorkerListenerMessage;
-import org.skywalking.apm.collector.cluster.WorkersListener;
-import org.skywalking.apm.collector.rpc.RPCAddress;
-import org.skywalking.apm.collector.rpc.RPCAddressListener;
-import org.skywalking.apm.collector.rpc.RPCAddressListenerMessage;
-import org.skywalking.apm.collector.log.LogManager;
-
-/**
- * The AbstractClusterWorker implementations represent workers,
- * which receive remote messages.
- *
- * Usually, the implementations are doing persistent, or aggregate works.
- *
- * @author pengys5
- * @since v3.0-2017
- */
-public abstract class AbstractClusterWorker extends AbstractWorker {
-
- /**
- * Construct an AbstractClusterWorker with the worker role and context.
- *
- * @param role If multi-workers are for load balance, they should be more likely called worker instance. Meaning,
- * each worker have multi instances.
- * @param clusterContext See {@link ClusterWorkerContext}
- * @param selfContext See {@link LocalWorkerContext}
- */
- protected AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
- }
-
- /**
- * This method use for message producer to call for send message.
- *
- * @param message The persistence data or metric data.
- * @throws Exception The Exception happen in {@link #onWork(Object)}
- */
- final public void allocateJob(Object message) throws WorkerException {
- onWork(message);
- }
-
- /**
- * This method use for message receiver to analyse message.
- *
- * @param message Cast the message object to a expect subclass.
- * @throws Exception Don't handle the exception, throw it.
- */
- protected abstract void onWork(Object message) throws WorkerException;
-
- static class WorkerWithAkka extends UntypedActor {
- private Logger logger = LogManager.INSTANCE.getFormatterLogger(WorkerWithAkka.class);
-
- private Cluster cluster;
- private final AbstractClusterWorker ownerWorker;
- private final RPCAddress RPCAddress;
-
- public WorkerWithAkka(AbstractClusterWorker ownerWorker, RPCAddress RPCAddress) {
- this.ownerWorker = ownerWorker;
- cluster = Cluster.get(getContext().system());
- this.RPCAddress = RPCAddress;
- }
-
- @Override
- public void preStart() {
- cluster.subscribe(getSelf(), ClusterEvent.MemberUp.class);
- }
-
- @Override
- public void postStop() {
- cluster.unsubscribe(getSelf());
- }
-
- /**
- * Listening {@link ClusterEvent.MemberUp} and {@link ClusterEvent.CurrentClusterState}
- * cluster event, when event send from the member of {@link WorkersListener} then tell
- * the sender to register self.
- */
- @Override
- public void onReceive(Object message) throws WorkerException {
- if (message instanceof ClusterEvent.CurrentClusterState) {
- ClusterEvent.CurrentClusterState state = (ClusterEvent.CurrentClusterState)message;
- for (Member member : state.getMembers()) {
- if (member.status().equals(MemberStatus.up())) {
- register(member);
- }
- }
- } else if (message instanceof ClusterEvent.MemberUp) {
- ClusterEvent.MemberUp memberUp = (ClusterEvent.MemberUp)message;
- logger.info("receive ClusterEvent.MemberUp message, address: %s", memberUp.member().address().toString());
- register(memberUp.member());
- } else {
- logger.debug("worker class: %s, message class: %s", this.getClass().getName(), message.getClass().getName());
- ownerWorker.allocateJob(message);
- }
- }
-
- /**
- * When member role is {@link WorkersListener#WORK_NAME} then Select actor from context
- * and send register message to {@link WorkersListener}
- *
- * @param member is the new created or restart worker
- */
- void register(Member member) {
- if (member.hasRole(WorkersListener.WORK_NAME)) {
- WorkerListenerMessage.RegisterMessage registerMessage = new WorkerListenerMessage.RegisterMessage(ownerWorker.getRole());
- logger.info("member address: %s, worker path: %s", member.address().toString(), getSelf().path().toString());
- getContext().actorSelection(member.address() + "/user/" + WorkersListener.WORK_NAME).tell(registerMessage, getSelf());
- }
- if (member.hasRole(RPCAddressListener.WORK_NAME) && RPCAddress != null) {
- RPCAddressListenerMessage.ConfigMessage configMessage = new RPCAddressListenerMessage.ConfigMessage(RPCAddress);
- logger.info("member address: %s, worker path: %s", member.address().toString(), getSelf().path().toString());
- getContext().actorSelection(member.address() + "/user/" + RPCAddressListener.WORK_NAME).tell(configMessage, getSelf());
- }
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorker.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorker.java
deleted file mode 100644
index 82186c50f..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorker.java
+++ /dev/null
@@ -1,103 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import com.lmax.disruptor.EventHandler;
-import com.lmax.disruptor.RingBuffer;
-import org.skywalking.apm.collector.queue.EndOfBatchCommand;
-import org.skywalking.apm.collector.queue.MessageHolder;
-
-/**
- * The AbstractLocalAsyncWorker implementations represent workers,
- * which receive local asynchronous message.
- *
- * @author pengys5
- * @since v3.0-2017
- */
-public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
-
- /**
- * Construct an AbstractLocalAsyncWorker with the worker role and context.
- *
- * @param role The responsibility of worker in cluster, more than one workers can have same responsibility which use
- * to provide load balancing ability.
- * @param clusterContext See {@link ClusterWorkerContext}
- * @param selfContext See {@link LocalWorkerContext}
- */
- public AbstractLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
- }
-
- /**
- * The asynchronous worker always use to persistence data into db, this is the end of the streaming,
- * so usually no need to create the next worker instance at the time of this worker instance create.
- *
- * @throws ProviderNotFoundException When worker provider not found, it will be throw this exception.
- */
- @Override
- public void preStart() throws ProviderNotFoundException {
- }
-
- /**
- * Receive message
- *
- * @param message The persistence data or metric data.
- * @throws WorkerException The Exception happen in {@link #onWork(Object)}
- */
- final public void allocateJob(Object message) throws WorkerException {
- onWork(message);
- }
-
- /**
- * The data process logic in this method.
- *
- * @param message Cast the message object to a expect subclass.
- * @throws WorkerException Don't handle the exception, throw it.
- */
- protected abstract void onWork(Object message) throws WorkerException;
-
- static class WorkerWithDisruptor implements EventHandler {
-
- private RingBuffer ringBuffer;
- private AbstractLocalAsyncWorker asyncWorker;
-
- WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) {
- this.ringBuffer = ringBuffer;
- this.asyncWorker = asyncWorker;
- }
-
- /**
- * Receive the message from disruptor, when message in disruptor is empty, then send the cached data
- * to the next workers.
- *
- * @param event published to the {@link RingBuffer}
- * @param sequence of the event being processed
- * @param endOfBatch flag to indicate if this is the last event in a batch from the {@link RingBuffer}
- */
- public void onEvent(MessageHolder event, long sequence, boolean endOfBatch) {
- try {
- Object message = event.getMessage();
- event.reset();
-
- asyncWorker.allocateJob(message);
- if (endOfBatch) {
- asyncWorker.allocateJob(new EndOfBatchCommand());
- }
- } catch (Exception e) {
- asyncWorker.saveException(e);
- }
- }
-
- /**
- * Push the message into disruptor ring buffer.
- *
- * @param message of the data to process.
- */
- public void tell(Object message) {
- long sequence = ringBuffer.next();
- try {
- ringBuffer.get(sequence).setMessage(message);
- } finally {
- ringBuffer.publish(sequence);
- }
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorkerProvider.java
deleted file mode 100644
index d636d9682..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorkerProvider.java
+++ /dev/null
@@ -1,47 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import com.lmax.disruptor.RingBuffer;
-import com.lmax.disruptor.dsl.Disruptor;
-import org.skywalking.apm.collector.queue.DaemonThreadFactory;
-import org.skywalking.apm.collector.queue.MessageHolder;
-import org.skywalking.apm.collector.queue.MessageHolderFactory;
-
-/**
- * @author pengys5
- */
-public abstract class AbstractLocalAsyncWorkerProvider extends AbstractLocalWorkerProvider {
-
- public abstract int queueSize();
-
- @Override final public WorkerRef onCreate(
- LocalWorkerContext localContext) throws ProviderNotFoundException {
- T localAsyncWorker = (T)workerInstance(getClusterContext());
- localAsyncWorker.preStart();
-
- // Specify the size of the ring buffer, must be power of 2.
- int bufferSize = queueSize();
- if (!((((bufferSize - 1) & bufferSize) == 0) && bufferSize != 0)) {
- throw new IllegalArgumentException("queue size must be power of 2");
- }
-
- // Construct the Disruptor
- Disruptor disruptor = new Disruptor(MessageHolderFactory.INSTANCE, bufferSize, DaemonThreadFactory.INSTANCE);
-
- RingBuffer ringBuffer = disruptor.getRingBuffer();
- T.WorkerWithDisruptor disruptorWorker = new T.WorkerWithDisruptor(ringBuffer, localAsyncWorker);
-
- // Connect the handler
- disruptor.handleEventsWith(disruptorWorker);
-
- // Start the Disruptor, starts all threads running
- disruptor.start();
-
- LocalAsyncWorkerRef workerRef = new LocalAsyncWorkerRef(role(), disruptorWorker);
-
- if (localContext != null) {
- localContext.put(workerRef);
- }
-
- return workerRef;
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/ClusterWorkerRef.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/ClusterWorkerRef.java
deleted file mode 100644
index 2a47f19f7..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/ClusterWorkerRef.java
+++ /dev/null
@@ -1,21 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import akka.actor.ActorRef;
-
-/**
- * @author pengys5
- */
-public class ClusterWorkerRef extends WorkerRef {
-
- private ActorRef actorRef;
-
- public ClusterWorkerRef(ActorRef actorRef, Role role) {
- super(role);
- this.actorRef = actorRef;
- }
-
- @Override
- public void tell(Object message) {
- actorRef.tell(message, ActorRef.noSender());
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/LocalAsyncWorkerRef.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/LocalAsyncWorkerRef.java
deleted file mode 100644
index 0e975a241..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/LocalAsyncWorkerRef.java
+++ /dev/null
@@ -1,19 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-/**
- * @author pengys5
- */
-public class LocalAsyncWorkerRef extends WorkerRef {
-
- private AbstractLocalAsyncWorker.WorkerWithDisruptor workerWithDisruptor;
-
- public LocalAsyncWorkerRef(Role role, AbstractLocalAsyncWorker.WorkerWithDisruptor workerWithDisruptor) {
- super(role);
- this.workerWithDisruptor = workerWithDisruptor;
- }
-
- @Override
- public void tell(Object message) throws WorkerInvokeException {
- workerWithDisruptor.tell(message);
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/Role.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/Role.java
deleted file mode 100644
index 811c1a275..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/actor/Role.java
+++ /dev/null
@@ -1,13 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import org.skywalking.apm.collector.actor.selector.WorkerSelector;
-
-/**
- * @author pengys5
- */
-public interface Role {
-
- String roleName();
-
- WorkerSelector workerSelector();
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterConfig.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterConfig.java
deleted file mode 100644
index d2333319e..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterConfig.java
+++ /dev/null
@@ -1,24 +0,0 @@
-package org.skywalking.apm.collector.cluster;
-
-/**
- * A static class contains some config values of cluster.
- * {@link Cluster.Current#HOSTNAME} is a ip address of server which start this process.
- * {@link Cluster.Current#PORT} is a PORT of server use to bind
- * {@link Cluster.Current#ROLES} is a ROLES of workers that use to create workers which
- * has those role in this process.
- * {@link Cluster#SEED_NODES} is a SEED_NODES which cluster have, List of strings, e.g. SEED_NODES = "ip:PORT,ip:PORT"..
- *
- * @author pengys5
- */
-public class ClusterConfig {
-
- public static class Cluster {
- public static class Current {
- public static String HOSTNAME = "";
- public static String PORT = "";
- public static String ROLES = "";
- }
-
- public static String SEED_NODES = "";
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterConfigProvider.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterConfigProvider.java
deleted file mode 100644
index c5c3c7a85..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterConfigProvider.java
+++ /dev/null
@@ -1,31 +0,0 @@
-package org.skywalking.apm.collector.cluster;
-
-import org.skywalking.apm.util.StringUtil;
-import org.skywalking.apm.collector.config.ConfigProvider;
-
-/**
- * @author pengys5
- */
-public class ClusterConfigProvider implements ConfigProvider {
-
- @Override
- public Class configClass() {
- return ClusterConfig.class;
- }
-
- @Override
- public void cliArgs() {
- if (!StringUtil.isEmpty(System.getProperty("cluster.current.HOSTNAME"))) {
- ClusterConfig.Cluster.Current.HOSTNAME = System.getProperty("cluster.current.HOSTNAME");
- }
- if (!StringUtil.isEmpty(System.getProperty("cluster.current.PORT"))) {
- ClusterConfig.Cluster.Current.PORT = System.getProperty("cluster.current.PORT");
- }
- if (!StringUtil.isEmpty(System.getProperty("cluster.current.ROLES"))) {
- ClusterConfig.Cluster.Current.ROLES = System.getProperty("cluster.current.ROLES");
- }
- if (!StringUtil.isEmpty(System.getProperty("cluster.SEED_NODES"))) {
- ClusterConfig.Cluster.SEED_NODES = System.getProperty("cluster.SEED_NODES");
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java
new file mode 100644
index 000000000..4f5246442
--- /dev/null
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.cluster;
+
+import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
+import org.skywalking.apm.collector.core.module.ModuleInstallMode;
+
+/**
+ * @author pengys5
+ */
+public class ClusterModuleGroupDefine implements ModuleGroupDefine {
+ @Override public String name() {
+ return "cluster";
+ }
+
+ @Override public ModuleInstallMode mode() {
+ return ModuleInstallMode.Single;
+ }
+}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/Const.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/Const.java
deleted file mode 100644
index 72c501acf..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/Const.java
+++ /dev/null
@@ -1,8 +0,0 @@
-package org.skywalking.apm.collector.cluster;
-
-/**
- * @author pengys5
- */
-public class Const {
- public static final String SYSTEM_NAME = "ClusterSystem";
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/WorkerListenerMessage.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/WorkerListenerMessage.java
deleted file mode 100644
index e5fb162b1..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/WorkerListenerMessage.java
+++ /dev/null
@@ -1,27 +0,0 @@
-package org.skywalking.apm.collector.cluster;
-
-import java.io.Serializable;
-import org.skywalking.apm.collector.actor.AbstractWorker;
-import org.skywalking.apm.collector.actor.Role;
-
-/**
- * WorkerListenerMessage is a message just for the worker
- * implementation of the {@link AbstractWorker}
- * to register.
- *
- * @author pengys5
- */
-public class WorkerListenerMessage {
-
- public static class RegisterMessage implements Serializable {
- private final Role role;
-
- public RegisterMessage(Role role) {
- this.role = role;
- }
-
- public Role getRole() {
- return role;
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/WorkersListener.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/WorkersListener.java
deleted file mode 100644
index 53b1f330e..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/WorkersListener.java
+++ /dev/null
@@ -1,76 +0,0 @@
-package org.skywalking.apm.collector.cluster;
-
-import akka.actor.ActorRef;
-import akka.actor.Terminated;
-import akka.actor.UntypedActor;
-import akka.cluster.Cluster;
-import akka.cluster.ClusterEvent;
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-import org.skywalking.apm.collector.actor.AbstractWorker;
-import org.skywalking.apm.collector.actor.ClusterWorkerContext;
-import org.skywalking.apm.collector.actor.ClusterWorkerRef;
-
-import java.util.Iterator;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
-/**
- * WorkersListener listening the register message from workers
- * implementation of the {@link AbstractWorker}
- * and terminated message from akka cluster.
- *
- * when listened register message then begin to watch the state for this worker
- * and register to {@link ClusterWorkerContext} and {@link #relation}.
- *
- * when listened terminate message then unregister from {@link ClusterWorkerContext} and {@link #relation} .
- *
- * @author pengys5
- */
-public class WorkersListener extends UntypedActor {
- public static final String WORK_NAME = "WorkersListener";
-
- private final Logger logger = LogManager.getFormatterLogger(WorkersListener.class);
- private final ClusterWorkerContext clusterContext;
- private Cluster cluster = Cluster.get(getContext().system());
- private Map relation = new ConcurrentHashMap<>();
-
- public WorkersListener(ClusterWorkerContext clusterContext) {
- this.clusterContext = clusterContext;
- }
-
- @Override
- public void preStart() throws Exception {
- cluster.subscribe(getSelf(), ClusterEvent.UnreachableMember.class);
- }
-
- @Override
- public void onReceive(Object message) throws Throwable {
- if (message instanceof WorkerListenerMessage.RegisterMessage) {
- WorkerListenerMessage.RegisterMessage register = (WorkerListenerMessage.RegisterMessage) message;
- ActorRef sender = getSender();
- logger.info("register worker of role: %s, path: %s", register.getRole().roleName(), sender.toString());
- ClusterWorkerRef workerRef = new ClusterWorkerRef(sender, register.getRole());
- relation.put(sender, workerRef);
- clusterContext.put(new ClusterWorkerRef(sender, register.getRole()));
- } else if (message instanceof Terminated) {
- Terminated terminated = (Terminated) message;
- clusterContext.remove(relation.get(terminated.getActor()));
- relation.remove(terminated.getActor());
- } else if (message instanceof ClusterEvent.UnreachableMember) {
- ClusterEvent.UnreachableMember unreachableMember = (ClusterEvent.UnreachableMember) message;
-
- Iterator> iterator = relation.entrySet().iterator();
- while (iterator.hasNext()) {
- Map.Entry next = iterator.next();
-
- if (next.getKey().path().address().equals(unreachableMember.member().address())) {
- clusterContext.remove(next.getValue());
- iterator.remove();
- }
- }
- } else {
- unhandled(message);
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfig.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfig.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfig.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfig.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisDataInitializer.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisDataInitializer.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisDataInitializer.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisDataInitializer.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleRegistrationWriter.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleRegistrationWriter.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-redis/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleRegistrationWriter.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleRegistrationWriter.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneConfigParser.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneConfigParser.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneConfigParser.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneConfigParser.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataInitializer.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataInitializer.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataInitializer.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneDataInitializer.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleRegistrationWriter.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleRegistrationWriter.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-standalone/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleRegistrationWriter.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleRegistrationWriter.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfig.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfig.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfig.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfig.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java
similarity index 93%
rename from apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java
index 264099154..82e5e86ab 100644
--- a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java
@@ -14,12 +14,13 @@ public class ClusterZKConfigParser implements ModuleConfigParser {
private final String SESSION_TIMEOUT = "sessionTimeout";
@Override public void parse(Map config) throws ConfigParseException {
- if (StringUtils.isEmpty(config.get(HOST_PORT))) {
- throw new ConfigParseException("");
- }
ClusterZKConfig.HOST_PORT = (String)config.get(HOST_PORT);
ClusterZKConfig.SESSION_TIMEOUT = 1000;
+ if (StringUtils.isEmpty(ClusterZKConfig.HOST_PORT)) {
+ throw new ConfigParseException("");
+ }
+
if (!StringUtils.isEmpty(config.get(SESSION_TIMEOUT))) {
ClusterZKConfig.SESSION_TIMEOUT = (Integer)config.get(SESSION_TIMEOUT);
}
diff --git a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationReader.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationReader.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationReader.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationReader.java
diff --git a/apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationWriter.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationWriter.java
similarity index 100%
rename from apm-collector/apm-collector-cluster-new/cluster-zookeeper/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationWriter.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleRegistrationWriter.java
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/config/ConfigInitializer.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/config/ConfigInitializer.java
deleted file mode 100644
index 7746883ab..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/config/ConfigInitializer.java
+++ /dev/null
@@ -1,40 +0,0 @@
-package org.skywalking.apm.collector.config;
-
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-
-import java.io.IOException;
-import java.io.InputStream;
-import java.util.Properties;
-import java.util.ServiceLoader;
-
-/**
- * @author pengys5
- */
-public enum ConfigInitializer {
- INSTANCE;
-
- private Logger logger = LogManager.getFormatterLogger(ConfigInitializer.class);
-
- public void initialize() throws IOException, IllegalAccessException {
- InputStream configFileStream = ConfigInitializer.class.getResourceAsStream("/collector.config");
- initializeConfigFile(configFileStream);
-
- ServiceLoader configProviders = ServiceLoader.load(ConfigProvider.class);
- for (ConfigProvider provider : configProviders) {
- provider.cliArgs();
- }
- }
-
- private void initializeConfigFile(InputStream configFileStream) throws IOException, IllegalAccessException {
- ServiceLoader configProviders = ServiceLoader.load(ConfigProvider.class);
- Properties properties = new Properties();
- properties.load(configFileStream);
-
- for (ConfigProvider provider : configProviders) {
- logger.info("configProvider provider name: %s", provider.getClass().getName());
- Class configClass = provider.configClass();
- org.skywalking.apm.util.ConfigInitializer.initialize(properties, configClass);
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/config/ConfigProvider.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/config/ConfigProvider.java
deleted file mode 100644
index 8de5b2110..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/config/ConfigProvider.java
+++ /dev/null
@@ -1,10 +0,0 @@
-package org.skywalking.apm.collector.config;
-
-/**
- * @author pengys5
- */
-public interface ConfigProvider {
- Class configClass();
-
- void cliArgs();
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/log/LogManager.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/log/LogManager.java
deleted file mode 100644
index 225269550..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/log/LogManager.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package org.skywalking.apm.collector.log;
-
-import org.apache.logging.log4j.Logger;
-
-/**
- * @author pengys5
- */
-public enum LogManager {
- INSTANCE;
-
- public Logger getFormatterLogger(final Class> clazz) {
- return org.apache.logging.log4j.LogManager.getFormatterLogger(clazz);
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddress.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddress.java
deleted file mode 100644
index 6cf949b78..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddress.java
+++ /dev/null
@@ -1,22 +0,0 @@
-package org.skywalking.apm.collector.rpc;
-
-/**
- * @author pengys5
- */
-public class RPCAddress {
- private final String address;
- private final int port;
-
- public RPCAddress(String address, int port) {
- this.address = address;
- this.port = port;
- }
-
- public String getAddress() {
- return address;
- }
-
- public int getPort() {
- return port;
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressContext.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressContext.java
deleted file mode 100644
index 52dcf27e4..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressContext.java
+++ /dev/null
@@ -1,25 +0,0 @@
-package org.skywalking.apm.collector.rpc;
-
-import java.util.Collection;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
-/**
- * @author pengys5
- */
-public class RPCAddressContext {
-
- private Map rpcAddresses = new ConcurrentHashMap<>();
-
- public Collection rpcAddressCollection() {
- return rpcAddresses.values();
- }
-
- public void putAddress(String ownerAddress, RPCAddress rpcAddress) {
- rpcAddresses.put(ownerAddress, rpcAddress);
- }
-
- public void removeAddress(String ownerAddress) {
- rpcAddresses.remove(ownerAddress);
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressListener.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressListener.java
deleted file mode 100644
index d23190c95..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressListener.java
+++ /dev/null
@@ -1,51 +0,0 @@
-package org.skywalking.apm.collector.rpc;
-
-import akka.actor.ActorRef;
-import akka.actor.Terminated;
-import akka.actor.UntypedActor;
-import akka.cluster.Cluster;
-import akka.cluster.ClusterEvent;
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-import org.skywalking.apm.collector.actor.ClusterWorkerContext;
-
-/**
- * @author pengys5
- */
-public class RPCAddressListener extends UntypedActor {
-
- private final Logger logger = LogManager.getFormatterLogger(RPCAddressListener.class);
-
- public static final String WORK_NAME = "RPCAddressListener";
-
- private final ClusterWorkerContext clusterContext;
- private Cluster cluster = Cluster.get(getContext().system());
-
- public RPCAddressListener(ClusterWorkerContext clusterContext) {
- this.clusterContext = clusterContext;
- }
-
- @Override
- public void preStart() throws Exception {
- cluster.subscribe(getSelf(), ClusterEvent.UnreachableMember.class);
- }
-
- @Override
- public void onReceive(Object message) throws Throwable {
- if (message instanceof RPCAddressListenerMessage.ConfigMessage) {
- RPCAddressListenerMessage.ConfigMessage configMessage = (RPCAddressListenerMessage.ConfigMessage)message;
- ActorRef sender = getSender();
- logger.info("address: %s, port: %s", configMessage.getConfig().getAddress(), configMessage.getConfig().getPort());
- String ownerAddress = sender.path().address().hostPort();
- clusterContext.getRpcContext().putAddress(ownerAddress, configMessage.getConfig());
- } else if (message instanceof Terminated) {
- Terminated terminated = (Terminated)message;
- clusterContext.getRpcContext().removeAddress(terminated.getActor().path().address().hostPort());
- } else if (message instanceof ClusterEvent.UnreachableMember) {
- ClusterEvent.UnreachableMember unreachableMember = (ClusterEvent.UnreachableMember)message;
- clusterContext.getRpcContext().removeAddress(unreachableMember.member().address().hostPort());
- } else {
- unhandled(message);
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressListenerMessage.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressListenerMessage.java
deleted file mode 100644
index bcf443024..000000000
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/rpc/RPCAddressListenerMessage.java
+++ /dev/null
@@ -1,21 +0,0 @@
-package org.skywalking.apm.collector.rpc;
-
-import java.io.Serializable;
-
-/**
- * @author pengys5
- */
-public class RPCAddressListenerMessage {
-
- public static class ConfigMessage implements Serializable {
- private final RPCAddress config;
-
- public ConfigMessage(RPCAddress config) {
- this.config = config;
- }
-
- public RPCAddress getConfig() {
- return config;
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/main/resources/META-INF/defines/group.define b/apm-collector/apm-collector-cluster/src/main/resources/META-INF/defines/group.define
new file mode 100644
index 000000000..333959bcb
--- /dev/null
+++ b/apm-collector/apm-collector-cluster/src/main/resources/META-INF/defines/group.define
@@ -0,0 +1 @@
+org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-cluster/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-cluster/src/main/resources/META-INF/defines/module.define
new file mode 100644
index 000000000..64e4e0bf7
--- /dev/null
+++ b/apm-collector/apm-collector-cluster/src/main/resources/META-INF/defines/module.define
@@ -0,0 +1,3 @@
+org.skywalking.apm.collector.cluster.zookeeper.ClusterZKModuleDefine
+org.skywalking.apm.collector.cluster.standalone.ClusterStandaloneModuleDefine
+org.skywalking.apm.collector.cluster.redis.ClusterRedisModuleDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractClusterWorkerProviderTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractClusterWorkerProviderTestCase.java
deleted file mode 100644
index b1b3f1cab..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractClusterWorkerProviderTestCase.java
+++ /dev/null
@@ -1,61 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import akka.actor.ActorSystem;
-import org.apache.logging.log4j.Logger;
-import org.mockito.Mockito;
-import org.powermock.reflect.Whitebox;
-import org.skywalking.apm.collector.actor.selector.RollingSelector;
-import org.skywalking.apm.collector.actor.selector.WorkerSelector;
-import org.skywalking.apm.collector.log.LogManager;
-
-/**
- * @author pengys5
- */
-//@RunWith(PowerMockRunner.class)
-//@PrepareForTest({LogManager.class})
-public class AbstractClusterWorkerProviderTestCase {
-
- // @Test
- public void testOnCreate() throws ProviderNotFoundException {
- LogManager logManager = Mockito.mock(LogManager.class);
- Whitebox.setInternalState(LogManager.class, "INSTANCE", logManager);
- Logger logger = Mockito.mock(Logger.class);
- Mockito.when(logManager.getFormatterLogger(Mockito.any())).thenReturn(logger);
-
- ActorSystem actorSystem = Mockito.mock(ActorSystem.class);
- ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(actorSystem);
- Impl impl = new Impl();
- impl.onCreate(null);
- }
-
- class Impl extends AbstractClusterWorkerProvider {
- @Override
- public Role role() {
- return Role.INSTANCE;
- }
-
- @Override
- public AbstractClusterWorkerTestCase.Impl workerInstance(ClusterWorkerContext clusterContext) {
- return new AbstractClusterWorkerTestCase.Impl(role(), clusterContext, new LocalWorkerContext());
- }
-
- @Override
- public int workerNum() {
- return 0;
- }
- }
-
- enum Role implements org.skywalking.apm.collector.actor.Role {
- INSTANCE;
-
- @Override
- public String roleName() {
- return AbstractClusterWorkerTestCase.Impl.class.getSimpleName();
- }
-
- @Override
- public WorkerSelector workerSelector() {
- return new RollingSelector();
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractClusterWorkerTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractClusterWorkerTestCase.java
deleted file mode 100644
index ff62a9517..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractClusterWorkerTestCase.java
+++ /dev/null
@@ -1,100 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import akka.actor.Address;
-import akka.cluster.ClusterEvent;
-import akka.cluster.Member;
-import org.apache.logging.log4j.Logger;
-import org.junit.Before;
-import org.junit.Test;
-import org.junit.runner.RunWith;
-import org.powermock.api.mockito.PowerMockito;
-import org.powermock.core.classloader.annotations.PowerMockIgnore;
-import org.powermock.core.classloader.annotations.PrepareForTest;
-import org.powermock.modules.junit4.PowerMockRunner;
-import org.powermock.reflect.Whitebox;
-import org.skywalking.apm.collector.actor.selector.RollingSelector;
-import org.skywalking.apm.collector.actor.selector.WorkerSelector;
-
-import static org.mockito.Mockito.*;
-
-/**
- * @author pengys5
- */
-@RunWith(PowerMockRunner.class)
-@PrepareForTest( {ClusterEvent.MemberUp.class, Address.class})
-@PowerMockIgnore( {"javax.management.*"})
-public class AbstractClusterWorkerTestCase {
-
- private AbstractClusterWorker.WorkerWithAkka workerWithAkka = mock(AbstractClusterWorker.WorkerWithAkka.class, CALLS_REAL_METHODS);
- private AbstractClusterWorker worker = PowerMockito.spy(new Impl(WorkerRole.INSTANCE, null, null));
-
- @Before
- public void init() {
- Logger logger = mock(Logger.class);
- Whitebox.setInternalState(workerWithAkka, "logger", logger);
- Whitebox.setInternalState(workerWithAkka, "ownerWorker", worker);
- }
-
- @Test
- public void testAllocateJob() throws Exception {
-
- String jobStr = "TestJob";
- worker.allocateJob(jobStr);
-
- verify(worker).onWork(jobStr);
- }
-
- @Test
- public void testMemberUp() throws Throwable {
- ClusterEvent.MemberUp memberUp = mock(ClusterEvent.MemberUp.class);
-
- Address address = mock(Address.class);
- when(address.toString()).thenReturn("address");
-
- Member member = mock(Member.class);
- when(member.address()).thenReturn(address);
-
- when(memberUp.member()).thenReturn(member);
-
- workerWithAkka.onReceive(memberUp);
-
- verify(workerWithAkka).register(member);
- }
-
- @Test
- public void testMessage() throws Throwable {
- String message = "test";
- workerWithAkka.onReceive(message);
-
- verify(worker).allocateJob(message);
- }
-
- static class Impl extends AbstractClusterWorker {
- @Override
- public void preStart() throws ProviderNotFoundException {
- }
-
- public Impl(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
- }
-
- @Override
- protected void onWork(Object message) throws IllegalArgumentException {
-
- }
- }
-
- public enum WorkerRole implements Role {
- INSTANCE;
-
- @Override
- public String roleName() {
- return Impl.class.getSimpleName();
- }
-
- @Override
- public WorkerSelector workerSelector() {
- return new RollingSelector();
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorkerTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorkerTestCase.java
deleted file mode 100644
index f41feddf1..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractLocalAsyncWorkerTestCase.java
+++ /dev/null
@@ -1,59 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import org.junit.Test;
-import org.mockito.ArgumentMatcher;
-import org.skywalking.apm.collector.queue.EndOfBatchCommand;
-import org.skywalking.apm.collector.queue.MessageHolder;
-
-import static org.mockito.Mockito.*;
-
-/**
- * @author pengys5
- */
-public class AbstractLocalAsyncWorkerTestCase {
-
- @Test
- public void testAllocateJob() throws Exception {
- AbstractLocalAsyncWorker worker = mock(AbstractLocalAsyncWorker.class);
-
- String message = "Test";
- worker.allocateJob(message);
- verify(worker).onWork(message);
- }
-
- @Test
- public void testOnEventWhenNotEnd() throws Exception {
- AbstractLocalAsyncWorker worker = mock(AbstractLocalAsyncWorker.class);
-
- AbstractLocalAsyncWorker.WorkerWithDisruptor disruptor = new AbstractLocalAsyncWorker.WorkerWithDisruptor(null, worker);
-
- MessageHolder holder = new MessageHolder();
- String message = "Test";
- holder.setMessage(message);
- disruptor.onEvent(holder, 0, false);
-
- verify(worker).onWork(message);
- }
-
- @Test
- public void testOnEventWhenEnd() throws Exception {
- AbstractLocalAsyncWorker worker = mock(AbstractLocalAsyncWorker.class);
-
- AbstractLocalAsyncWorker.WorkerWithDisruptor disruptor = new AbstractLocalAsyncWorker.WorkerWithDisruptor(null, worker);
-
- MessageHolder holder = new MessageHolder();
- String message = "Test";
- holder.setMessage(message);
- disruptor.onEvent(holder, 0, true);
-
- verify(worker, times(1)).onWork(message);
- verify(worker, times(1)).onWork(argThat(new IsEndOfBatchCommandClass()));
- }
-
- class IsEndOfBatchCommandClass extends ArgumentMatcher {
- public boolean matches(Object para) {
- return para.getClass() == EndOfBatchCommand.class;
- }
- }
-
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractWorkerProviderTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractWorkerProviderTestCase.java
deleted file mode 100644
index d6fbfff66..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/AbstractWorkerProviderTestCase.java
+++ /dev/null
@@ -1,77 +0,0 @@
-package org.skywalking.apm.collector.actor;
-
-import org.junit.Test;
-import org.junit.runner.RunWith;
-import org.mockito.Mockito;
-import org.powermock.core.classloader.annotations.PowerMockIgnore;
-import org.powermock.core.classloader.annotations.PrepareForTest;
-import org.powermock.modules.junit4.PowerMockRunner;
-
-import static org.powermock.api.mockito.PowerMockito.mock;
-import static org.powermock.api.mockito.PowerMockito.when;
-
-/**
- * @author pengys5
- */
-@RunWith(PowerMockRunner.class)
-@PrepareForTest( {AbstractWorker.class})
-@PowerMockIgnore( {"javax.management.*"})
-public class AbstractWorkerProviderTestCase {
-
- @Test(expected = IllegalArgumentException.class)
- public void testNullWorkerInstanceCreate() throws ProviderNotFoundException {
- AbstractWorkerProvider provider = mock(AbstractWorkerProvider.class);
- when(provider.workerInstance(null)).thenReturn(null);
-
- AbstractWorker worker = mock(AbstractWorker.class);
- provider.create(worker);
- }
-
- @Test
- public void testNoneWorkerOwner() throws ProviderNotFoundException {
- AbstractWorkerProvider provider = mock(AbstractWorkerProvider.class);
-
- ClusterWorkerContext context = mock(ClusterWorkerContext.class);
- provider.setClusterContext(context);
-
- AbstractWorker worker = mock(AbstractWorker.class);
- when(provider.workerInstance(context)).thenReturn(worker);
-
- provider.create(null);
- Mockito.verify(provider).onCreate(null);
- }
-
- @Test
- public void testHasWorkerOwner() throws ProviderNotFoundException {
- AbstractWorkerProvider provider = mock(AbstractWorkerProvider.class);
-
- ClusterWorkerContext context = mock(ClusterWorkerContext.class);
- provider.setClusterContext(context);
-
- AbstractWorker worker = mock(AbstractWorker.class);
- when(provider.workerInstance(context)).thenReturn(worker);
-
- AbstractWorker workerOwner = mock(AbstractWorker.class);
- LocalWorkerContext localWorkerContext = mock(LocalWorkerContext.class);
- when(workerOwner.getSelfContext()).thenReturn(localWorkerContext);
-
- provider.create(workerOwner);
- Mockito.verify(provider).onCreate(localWorkerContext);
- }
-
- @Test(expected = IllegalArgumentException.class)
- public void testHasWorkerOwnerButNoneContext() throws ProviderNotFoundException {
- AbstractWorkerProvider provider = mock(AbstractWorkerProvider.class);
-
- ClusterWorkerContext context = mock(ClusterWorkerContext.class);
- provider.setClusterContext(context);
-
- AbstractWorker worker = mock(AbstractWorker.class);
- when(provider.workerInstance(context)).thenReturn(worker);
-
- AbstractWorker workerOwner = mock(AbstractWorker.class);
- when(workerOwner.getSelfContext()).thenReturn(null);
-
- provider.create(workerOwner);
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/AbstractHashMessageTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/AbstractHashMessageTestCase.java
deleted file mode 100644
index ea3557cc9..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/AbstractHashMessageTestCase.java
+++ /dev/null
@@ -1,24 +0,0 @@
-package org.skywalking.apm.collector.actor.selector;
-
-import org.junit.Assert;
-import org.junit.Test;
-
-/**
- * @author pengys5
- */
-public class AbstractHashMessageTestCase {
-
- @Test
- public void testGetHashCode() {
- String key = "key";
-
- Impl impl = new Impl(key);
- Assert.assertEquals(key.hashCode(), impl.getHashCode());
- }
-
- class Impl extends AbstractHashMessage {
- public Impl(String key) {
- super(key);
- }
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/HashCodeSelectorTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/HashCodeSelectorTestCase.java
deleted file mode 100644
index 8f669e791..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/HashCodeSelectorTestCase.java
+++ /dev/null
@@ -1,55 +0,0 @@
-package org.skywalking.apm.collector.actor.selector;
-
-import org.junit.Assert;
-import org.junit.Test;
-import org.skywalking.apm.collector.actor.WorkerRef;
-
-import java.util.ArrayList;
-import java.util.List;
-
-import static org.powermock.api.mockito.PowerMockito.mock;
-import static org.powermock.api.mockito.PowerMockito.when;
-
-/**
- * @author pengys5
- */
-public class HashCodeSelectorTestCase {
-
- @Test
- public void testSelect() {
- List members = new ArrayList<>();
- WorkerRef workerRef_1 = mock(WorkerRef.class);
- WorkerRef workerRef_2 = mock(WorkerRef.class);
- WorkerRef workerRef_3 = mock(WorkerRef.class);
-
- members.add(workerRef_1);
- members.add(workerRef_2);
- members.add(workerRef_3);
-
- AbstractHashMessage message_1 = mock(AbstractHashMessage.class);
- when(message_1.getHashCode()).thenReturn(9);
-
- AbstractHashMessage message_2 = mock(AbstractHashMessage.class);
- when(message_2.getHashCode()).thenReturn(10);
-
- AbstractHashMessage message_3 = mock(AbstractHashMessage.class);
- when(message_3.getHashCode()).thenReturn(11);
-
- HashCodeSelector selector = new HashCodeSelector();
-
- WorkerRef select_1 = selector.select(members, message_1);
- Assert.assertEquals(workerRef_1.hashCode(), select_1.hashCode());
-
- WorkerRef select_2 = selector.select(members, message_2);
- Assert.assertEquals(workerRef_2.hashCode(), select_2.hashCode());
-
- WorkerRef select_3 = selector.select(members, message_3);
- Assert.assertEquals(workerRef_3.hashCode(), select_3.hashCode());
- }
-
- @Test(expected = IllegalArgumentException.class)
- public void testSelectError() {
- HashCodeSelector selector = new HashCodeSelector();
- selector.select(null, new Object());
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/RollingSelectorTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/RollingSelectorTestCase.java
deleted file mode 100644
index 442f46730..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/actor/selector/RollingSelectorTestCase.java
+++ /dev/null
@@ -1,41 +0,0 @@
-package org.skywalking.apm.collector.actor.selector;
-
-import org.junit.Assert;
-import org.junit.Test;
-import org.skywalking.apm.collector.actor.WorkerRef;
-
-import java.util.ArrayList;
-import java.util.List;
-
-import static org.powermock.api.mockito.PowerMockito.mock;
-
-/**
- * @author pengys5
- */
-public class RollingSelectorTestCase {
-
- @Test
- public void testSelect() {
- List members = new ArrayList<>();
- WorkerRef workerRef_1 = mock(WorkerRef.class);
- WorkerRef workerRef_2 = mock(WorkerRef.class);
- WorkerRef workerRef_3 = mock(WorkerRef.class);
-
- members.add(workerRef_1);
- members.add(workerRef_2);
- members.add(workerRef_3);
-
- Object message = new Object();
-
- RollingSelector selector = new RollingSelector();
-
- WorkerRef selected_1 = selector.select(members, message);
- Assert.assertEquals(workerRef_2.hashCode(), selected_1.hashCode());
-
- WorkerRef selected_2 = selector.select(members, message);
- Assert.assertEquals(workerRef_3.hashCode(), selected_2.hashCode());
-
- WorkerRef selected_3 = selector.select(members, message);
- Assert.assertEquals(workerRef_1.hashCode(), selected_3.hashCode());
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/config/ConfigInitializerTestCase.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/config/ConfigInitializerTestCase.java
deleted file mode 100644
index 88f99a623..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/config/ConfigInitializerTestCase.java
+++ /dev/null
@@ -1,45 +0,0 @@
-package org.skywalking.apm.collector.config;
-
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
-import org.skywalking.apm.collector.cluster.ClusterConfig;
-
-/**
- * @author pengys5
- */
-public class ConfigInitializerTestCase {
-
- @Before
- public void clear() {
- System.clearProperty("cluster.current.HOSTNAME");
- System.clearProperty("cluster.current.PORT");
- System.clearProperty("cluster.current.ROLES");
- System.clearProperty("cluster.SEED_NODES");
- }
-
- @Test
- public void testInitialize() throws Exception {
- ConfigInitializer.INSTANCE.initialize();
-
- Assert.assertEquals("127.0.0.1", ClusterConfig.Cluster.Current.HOSTNAME);
- Assert.assertEquals("1000", ClusterConfig.Cluster.Current.PORT);
- Assert.assertEquals("WorkersListener", ClusterConfig.Cluster.Current.ROLES);
- Assert.assertEquals("127.0.0.1:1000", ClusterConfig.Cluster.SEED_NODES);
- }
-
- @Test
- public void testInitializeWithCli() throws Exception {
- System.setProperty("cluster.current.HOSTNAME", "127.0.0.2");
- System.setProperty("cluster.current.PORT", "1001");
- System.setProperty("cluster.current.ROLES", "Test1, Test2");
- System.setProperty("cluster.SEED_NODES", "127.0.0.1:1000, 127.0.0.1:1001");
-
- ConfigInitializer.INSTANCE.initialize();
-
- Assert.assertEquals("127.0.0.2", ClusterConfig.Cluster.Current.HOSTNAME);
- Assert.assertEquals("1001", ClusterConfig.Cluster.Current.PORT);
- Assert.assertEquals("Test1, Test2", ClusterConfig.Cluster.Current.ROLES);
- Assert.assertEquals("127.0.0.1:1000, 127.0.0.1:1001", ClusterConfig.Cluster.SEED_NODES);
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/log/MockLog.java b/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/log/MockLog.java
deleted file mode 100644
index ba71fd0ca..000000000
--- a/apm-collector/apm-collector-cluster/src/test/java/org/skywalking/apm/collector/log/MockLog.java
+++ /dev/null
@@ -1,18 +0,0 @@
-package org.skywalking.apm.collector.log;
-
-import org.apache.logging.log4j.Logger;
-import org.mockito.Mockito;
-import org.powermock.api.mockito.PowerMockito;
-
-/**
- * @author pengys5
- */
-public class MockLog {
-
- public Logger mockito() {
- LogManager logManager = PowerMockito.mock(LogManager.class);
- Logger logger = Mockito.mock(Logger.class);
- Mockito.when(logManager.getFormatterLogger(Mockito.any())).thenReturn(logger);
- return logger;
- }
-}
diff --git a/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractClusterWorkerProvider b/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractClusterWorkerProvider
deleted file mode 100644
index e69de29bb..000000000
diff --git a/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractLocalWorkerProvider b/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractLocalWorkerProvider
deleted file mode 100644
index e69de29bb..000000000
diff --git a/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider b/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider
deleted file mode 100644
index fd524035e..000000000
--- a/apm-collector/apm-collector-cluster/src/test/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider
+++ /dev/null
@@ -1 +0,0 @@
-org.skywalking.apm.collector.cluster.ClusterConfigProvider
\ No newline at end of file
diff --git a/apm-collector/apm-collector-cluster/src/test/resources/collector.config b/apm-collector/apm-collector-cluster/src/test/resources/collector.config
deleted file mode 100644
index 9f0ea3fbd..000000000
--- a/apm-collector/apm-collector-cluster/src/test/resources/collector.config
+++ /dev/null
@@ -1,11 +0,0 @@
-cluster.current.hostname=127.0.0.1
-cluster.current.port=1000
-cluster.current.roles=WorkersListener
-cluster.seed_nodes=127.0.0.1:1000
-
-es.cluster.name=CollectorDBCluster
-es.cluster.nodes=127.0.0.1:9300
-es.cluster.transport.sniffer=true
-
-es.index.shards.number=2
-es.index.replicas.number=0
\ No newline at end of file
diff --git a/apm-collector/apm-collector-cluster/src/test/resources/log4j2.xml b/apm-collector/apm-collector-cluster/src/test/resources/log4j2.xml
deleted file mode 100644
index cc2fe771e..000000000
--- a/apm-collector/apm-collector-cluster/src/test/resources/log4j2.xml
+++ /dev/null
@@ -1,13 +0,0 @@
-
-
-
-
-
-
-
-
-
-
-
-
-
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerModuleDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java
similarity index 63%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerModuleDefine.java
rename to apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java
index 2e67e51be..e564c258a 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerModuleDefine.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.core.worker;
+package org.skywalking.apm.collector.core.agentstream;
import java.util.Map;
import org.skywalking.apm.collector.core.client.Client;
@@ -7,17 +7,17 @@ import org.skywalking.apm.collector.core.cluster.ClusterDataInitializer;
import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
import org.skywalking.apm.collector.core.config.ConfigParseException;
import org.skywalking.apm.collector.core.framework.DataInitializer;
+import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
-import org.skywalking.apm.collector.core.module.ModuleException;
import org.skywalking.apm.collector.core.server.Server;
import org.skywalking.apm.collector.core.server.ServerException;
/**
* @author pengys5
*/
-public abstract class WorkerModuleDefine extends ModuleDefine {
+public abstract class AgentStreamModuleDefine extends ModuleDefine {
- @Override public final void initialize(Map config) throws ModuleException, ClientException {
+ @Override public final void initialize(Map config) throws DefineException, ClientException {
try {
configParser().parse(config);
Server server = server();
@@ -26,15 +26,15 @@ public abstract class WorkerModuleDefine extends ModuleDefine {
String key = ClusterDataInitializer.BASE_CATALOG + "." + name();
ClusterModuleContext.writer.write(key, registration().buildValue());
} catch (ConfigParseException | ServerException e) {
- throw new WorkerModuleException(e.getMessage(), e);
+ throw new AgentStreamModuleException(e.getMessage(), e);
}
}
- @Override public final Client createClient() {
- throw new UnsupportedOperationException();
+ @Override protected final DataInitializer dataInitializer() {
+ throw new UnsupportedOperationException("");
}
- @Override public final DataInitializer dataInitializer() {
- throw new UnsupportedOperationException();
+ @Override protected final Client createClient() {
+ throw new UnsupportedOperationException("");
}
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java
new file mode 100644
index 000000000..11380f6b4
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java
@@ -0,0 +1,16 @@
+package org.skywalking.apm.collector.core.agentstream;
+
+import org.skywalking.apm.collector.core.module.ModuleException;
+
+/**
+ * @author pengys5
+ */
+public class AgentStreamModuleException extends ModuleException{
+ public AgentStreamModuleException(String message) {
+ super(message);
+ }
+
+ public AgentStreamModuleException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java
index 6525dbdd0..980873d4e 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java
@@ -5,4 +5,6 @@ package org.skywalking.apm.collector.core.cluster;
*/
public class ClusterModuleContext {
public static ClusterModuleRegistrationWriter writer;
+
+ public static ClusterModuleRegistrationReader reader;
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java
index cd232a76c..42e4016c3 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java
@@ -39,4 +39,6 @@ public abstract class ClusterModuleDefine extends ModuleDefine {
}
protected abstract ClusterModuleRegistrationWriter registrationWriter();
+
+ protected abstract ClusterModuleRegistrationReader registrationReader();
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java
index 9de01cf38..44534172d 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java
@@ -7,7 +7,10 @@ import org.skywalking.apm.collector.core.module.ModuleConfigLoader;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleDefineLoader;
import org.skywalking.apm.collector.core.module.ModuleGroup;
+import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
+import org.skywalking.apm.collector.core.module.ModuleGroupDefineLoader;
import org.skywalking.apm.collector.core.module.ModuleInstallerAdapter;
+import org.skywalking.apm.collector.core.remote.SerializedDefineLoader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -19,12 +22,19 @@ public class CollectorStarter implements Starter {
private final Logger logger = LoggerFactory.getLogger(CollectorStarter.class);
@Override public void start() throws ConfigException, DefineException, ClientException {
+ Context context = new Context();
ModuleConfigLoader configLoader = new ModuleConfigLoader();
Map configuration = configLoader.load();
+ SerializedDefineLoader serializedDefineLoader = new SerializedDefineLoader();
+ serializedDefineLoader.load();
+
ModuleDefineLoader defineLoader = new ModuleDefineLoader();
Map> moduleDefineMap = defineLoader.load();
+ ModuleGroupDefineLoader groupDefineLoader = new ModuleGroupDefineLoader();
+ Map moduleGroupDefineMap = groupDefineLoader.load();
+
ModuleInstallerAdapter moduleInstallerAdapter = new ModuleInstallerAdapter(ModuleGroup.Cluster);
moduleInstallerAdapter.install(configuration.get(ModuleGroup.Cluster.name().toLowerCase()), moduleDefineMap.get(ModuleGroup.Cluster.name().toLowerCase()));
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java
new file mode 100644
index 000000000..db34bb29a
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java
@@ -0,0 +1,7 @@
+package org.skywalking.apm.collector.core.framework;
+
+/**
+ * @author pengys5
+ */
+public class Context {
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java
index 63ef4d819..1ead094ea 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java
@@ -1,7 +1,6 @@
package org.skywalking.apm.collector.core.module;
import org.skywalking.apm.collector.core.client.Client;
-import org.skywalking.apm.collector.core.client.ClientConfig;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.framework.Define;
import org.skywalking.apm.collector.core.server.Server;
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefineLoader.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefineLoader.java
index 031ab022d..407bc13ae 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefineLoader.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefineLoader.java
@@ -19,10 +19,10 @@ public class ModuleDefineLoader implements Loader