diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java
index 5837feed1..b994c7710 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java
@@ -1,6 +1,5 @@
package com.a.eye.skywalking.collector.actor;
-import akka.actor.ActorRef;
import akka.actor.UntypedActor;
import akka.cluster.ClusterEvent;
import akka.cluster.Member;
@@ -13,18 +12,46 @@ import com.a.eye.skywalking.collector.cluster.WorkersRefCenter;
import java.util.List;
/**
+ * Abstract implementation of the {@link akka.actor.UntypedActor} that represents an
+ * analysis unit. AbstractWorker implementation process the message in
+ * {@link #receive(Object)} method.
+ *
+ *
+ * Subclasses must implement the abstract {@link #receive(Object)} method to process message.
+ * Subclasses forbid to override the {@link #onReceive(Object)} method.
+ *
+ * Here is an example on how to create and use an {@link AbstractWorker}:
+ *
+ * {{{
+ * public class SampleWorker extends AbstractWorker {
+ *
+ * @Override
+ * public void receive(Object message) throws Throwable {
+ * if (message.equals("Tell Next")) {
+ * Object sendMessage = new Object();
+ * tell(new NextSampleWorkerFactory(), RollingSelector.INSTANCE, sendMessage);
+ * }
+ * }
+ * }
+ * }}}
+ *
* @author pengys5
*/
public abstract class AbstractWorker extends UntypedActor {
- final String workerRole;
-
- public AbstractWorker(String workerRole) {
- this.workerRole = workerRole;
- }
-
- public abstract void receive(Object message);
+ /**
+ * Receive the message to analyse.
+ *
+ * @param message is the data send from the forward worker
+ * @throws Throwable is the exception thrown by that worker implementation processing
+ */
+ public abstract void receive(Object message) throws Throwable;
+ /**
+ * 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 Throwable {
if (message instanceof ClusterEvent.CurrentClusterState) {
@@ -42,14 +69,28 @@ public abstract class AbstractWorker extends UntypedActor {
}
}
+ /**
+ * Send analysed data to next Worker.
+ *
+ * @param targetWorkerProvider is the worker provider to create worker instance.
+ * @param selector is the selector to select a same role worker instance form cluster.
+ * @param message is the data used to send to next worker.
+ * @throws Throwable
+ */
public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable {
List availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName());
selector.select(availableWorks, message).tell(message, getSelf());
}
+ /**
+ * When member role is {@link WorkersListener#WorkName} 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.getRoles().equals(WorkersListener.WorkName)) {
- WorkerListenerMessage.RegisterMessage registerMessage = new WorkerListenerMessage.RegisterMessage(workerRole);
+ WorkerListenerMessage.RegisterMessage registerMessage = new WorkerListenerMessage.RegisterMessage(getClass().getSimpleName());
getContext().actorSelection(member.address() + "/user/" + WorkersListener.WorkName).tell(registerMessage, getSelf());
}
}
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java
index 0020097ee..11bd7dfb9 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java
@@ -4,6 +4,24 @@ import akka.actor.ActorSystem;
import akka.actor.Props;
/**
+ * The AbstractWorkerProvider should be implemented by any class whose
+ * instances are intended to provide create instance of the {@link AbstractWorker}.
+ *
+ * Here is an example on how to create and use an {@link AbstractWorkerProvider}:
+ *
+ * {{{
+ * public class SampleWorkerFactory extends AbstractWorkerProvider {
+ *
+ * @Override public Class workerClass() {
+ * return SampleWorker.class;
+ * }
+ *
+ * @Override public int workerNum() {
+ * return Config.SampleWorkerNum;
+ * }
+ * }
+ * }}}
+ *
* @author pengys5
*/
public abstract class AbstractWorkerProvider {
@@ -12,6 +30,11 @@ public abstract class AbstractWorkerProvider {
public abstract int workerNum();
+ /**
+ * Use {@link ActorSystem} to Create worker instance with the {@link #workerClass()} method returned class.
+ *
+ * @param system is a akka {@link ActorSystem} instance.
+ */
public void createWorker(ActorSystem system) {
if (workerClass() == null) {
throw new IllegalArgumentException("cannot createWorker() with nothing obtained from workerClass()");
@@ -25,7 +48,12 @@ public abstract class AbstractWorkerProvider {
}
}
- protected String roleName(){
+ /**
+ * Use {@link #workerClass()} method returned class's simple name as a role name.
+ *
+ * @return is role of Worker
+ */
+ protected String roleName() {
return workerClass().getSimpleName();
}
}
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java
index c82b89470..d53deafca 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java
@@ -3,30 +3,25 @@ package com.a.eye.skywalking.collector.actor;
import akka.actor.ActorRef;
/**
+ * The Worker reference
+ *
* @author pengys5
*/
public class WorkerRef {
final ActorRef actorRef;
- public WorkerRef(ActorRef actorRef) {
+ final String workerRole;
+
+ public WorkerRef(ActorRef actorRef, String workerRole) {
this.actorRef = actorRef;
+ this.workerRole = workerRole;
}
void tell(Object message, ActorRef actorRef) {
actorRef.tell(message, actorRef);
}
- public String path(){
- return actorRef.path().toString();
- }
-
- @Override
- public boolean equals(Object obj) {
- return actorRef.equals(obj);
- }
-
- @Override
- public String toString() {
- return actorRef.toString();
+ public String getWorkerRole() {
+ return workerRole;
}
}
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java
index 40217f654..b48da7c68 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java
@@ -5,11 +5,19 @@ import akka.actor.ActorSystem;
import java.util.ServiceLoader;
/**
+ * WorkersCreator is a util that use Java Spi to create
+ * workers by META-INF config file.
+ *
* @author pengys5
*/
public enum WorkersCreator {
INSTANCE;
+ /**
+ * create worker to use Java Spi.
+ *
+ * @param system is create by akka {@link ActorSystem}
+ */
public void boot(ActorSystem system) {
ServiceLoader serviceLoader = ServiceLoader.load(AbstractWorkerProvider.class);
for (AbstractWorkerProvider provider : serviceLoader) {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfig.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfig.java
index caf92c3c3..42c13016e 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfig.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfig.java
@@ -1,6 +1,16 @@
package com.a.eye.skywalking.collector.cluster;
+import akka.actor.ActorSystem;
+
/**
+ * 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#nodes} is a nodes which cluster have.
+ * {@link Cluster#appname} is a name of {@link ActorSystem} in cluster.
+ *
* @author pengys5
*/
public class ClusterConfig {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfigInitializer.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfigInitializer.java
index 6885d1f0f..85447a89c 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfigInitializer.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/ClusterConfigInitializer.java
@@ -9,6 +9,14 @@ import java.io.InputStream;
import java.util.Properties;
/**
+ * ClusterConfigInitializer Contains static methods for setting
+ * {@link ClusterConfig} attributes value.
+ *
+ *
+ * The priority of value setting is
+ * system property -> collector.config -> {@link ClusterConfig} default value
+ *
+ *
* @author pengys5
*/
public class ClusterConfigInitializer {
@@ -17,6 +25,11 @@ public class ClusterConfigInitializer {
public static final String ConfigFileName = "collector.config";
+ /**
+ * Read config file to setting {@link ClusterConfig} then get system property to overwrite it.
+ *
+ * @param configFileName is the config file name, the file format is key-value pairs
+ */
public static void initialize(String configFileName) {
InputStream configFileStream = ClusterConfigInitializer.class.getResourceAsStream("/" + configFileName);
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkerListenerMessage.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkerListenerMessage.java
index 014964d0b..4bf9ac303 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkerListenerMessage.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkerListenerMessage.java
@@ -3,6 +3,10 @@ package com.a.eye.skywalking.collector.cluster;
import java.io.Serializable;
/**
+ * WorkerListenerMessage is a message just for the worker
+ * implementation of the {@link com.a.eye.skywalking.collector.actor.AbstractWorker}
+ * to register.
+ *
* @author pengys5
*/
public class WorkerListenerMessage {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java
index ded11b15d..29fa694ab 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java
@@ -5,6 +5,15 @@ import akka.actor.Terminated;
import akka.actor.UntypedActor;
/**
+ * WorkersListener listening the register message from workers
+ * implementation of the {@link com.a.eye.skywalking.collector.actor.AbstractWorker}
+ * and terminated message from akka cluster.
+ *
+ * when listened register message then begin to watch the state for this worker
+ * and register to {@link WorkersRefCenter}.
+ *
+ * when listened terminate message then unregister from {@link WorkersRefCenter}.
+ *
* @author pengys5
*/
public class WorkersListener extends UntypedActor {
diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java
index 9856fe194..d9d0007c6 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java
@@ -19,31 +19,25 @@ import java.util.concurrent.ConcurrentHashMap;
public enum WorkersRefCenter {
INSTANCE;
- private Map> roleToActor = new ConcurrentHashMap();
+ private Map> roleToWorkerRef = new ConcurrentHashMap();
- private Map actorToRole = new ConcurrentHashMap();
-
-// private Map pathToWorkerRef = new ConcurrentHashMap();
+ private Map actorRefToWorkerRef = new ConcurrentHashMap<>();
public void register(ActorRef newActorRef, String workerRole) {
- if (!roleToActor.containsKey(workerRole)) {
+ if (!roleToWorkerRef.containsKey(workerRole)) {
List actorList = Collections.synchronizedList(new ArrayList());
- roleToActor.putIfAbsent(workerRole, actorList);
+ roleToWorkerRef.putIfAbsent(workerRole, actorList);
}
- WorkerRef newWorkerRef = new WorkerRef(newActorRef);
- roleToActor.get(workerRole).add(newWorkerRef);
- actorToRole.put(newWorkerRef, workerRole);
-// pathToWorkerRef.put(newWorkerRef.path(), newWorkerRef);
+ WorkerRef newWorkerRef = new WorkerRef(newActorRef, workerRole);
+ roleToWorkerRef.get(workerRole).add(newWorkerRef);
+ actorRefToWorkerRef.put(newActorRef, newWorkerRef);
}
- public void unregister(ActorRef newActorRef) {
- String workerRole = actorToRole.get(newActorRef.path());
-// WorkerRef workerRef = pathToWorkerRef.get(newActorRef.path());
-
- roleToActor.get(workerRole).remove(newActorRef);
- actorToRole.remove(newActorRef);
-// pathToWorkerRef.remove(newActorRef.path());
+ public void unregister(ActorRef oldActorRef) {
+ WorkerRef oldWorkerRef = actorRefToWorkerRef.get(oldActorRef);
+ roleToWorkerRef.get(oldWorkerRef.getWorkerRole()).remove(oldWorkerRef);
+ actorRefToWorkerRef.remove(oldActorRef);
}
/**
@@ -54,7 +48,7 @@ public enum WorkersRefCenter {
* @throws NoAvailableWorkerException , when no available worker.
*/
public List availableWorks(String workerRole) throws NoAvailableWorkerException {
- List refs = roleToActor.get(workerRole);
+ List refs = roleToWorkerRef.get(workerRole);
if (refs == null || refs.size() == 0) {
throw new NoAvailableWorkerException("role=" + workerRole + ", no available worker.");
}
diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java
index df929f59d..0ee6d4345 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorker.java
@@ -1,20 +1,21 @@
package com.a.eye.skywalking.collector.actor;
+import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
+
/**
* @author pengys5
*/
public class SpiTestWorker extends AbstractWorker {
- public SpiTestWorker(String workerRole) {
- super(workerRole);
- }
-
@Override
- public void receive(Object message) {
+ public void receive(Object message) throws Throwable {
if (message.equals("Test1")) {
getSender().tell("Yes", getSelf());
} else if (message.equals("Test2")) {
getSender().tell("No", getSelf());
+ } else if (message.equals("Test3")) {
+ Object sendMessage = new Object();
+ tell(new SpiTestWorkerFactory(), RollingSelector.INSTANCE, sendMessage);
}
}
-}
+}
\ No newline at end of file
diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java
index f68399bce..a76c33adc 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/SpiTestWorkerFactory.java
@@ -5,8 +5,6 @@ package com.a.eye.skywalking.collector.actor;
*/
public class SpiTestWorkerFactory extends AbstractWorkerProvider {
- public static final String WorkerRole = "SpiTestWorker";
-
@Override
public Class workerClass() {
return SpiTestWorker.class;
diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/cluster/WorkerListenerTestCase.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/cluster/WorkerListenerTestCase.java
index e0ad8e1f4..af55dc36a 100644
--- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/cluster/WorkerListenerTestCase.java
+++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/cluster/WorkerListenerTestCase.java
@@ -6,6 +6,7 @@ import akka.actor.Props;
import akka.actor.Terminated;
import akka.pattern.Patterns;
import akka.testkit.TestActorRef;
+import com.a.eye.skywalking.collector.actor.WorkerRef;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
@@ -23,10 +24,19 @@ import java.util.concurrent.ConcurrentHashMap;
public class WorkerListenerTestCase {
ActorSystem system;
+ TestActorRef senderActorRef;
+ TestActorRef receiveactorRef;
@Before
- public void createSystem() {
+ public void initData() {
system = ActorSystem.create();
+
+ final Props props = Props.create(WorkersListener.class);
+ senderActorRef = TestActorRef.create(system, props, "WorkersListenerSender");
+ receiveactorRef = TestActorRef.create(system, props, "WorkersListenerReceive");
+
+ WorkerListenerMessage.RegisterMessage message = new WorkerListenerMessage.RegisterMessage("WorkersListener");
+ receiveactorRef.tell(message, senderActorRef);
}
@After
@@ -34,44 +44,32 @@ public class WorkerListenerTestCase {
system.terminate();
system.awaitTermination();
system = null;
- MemberModifier.field(WorkersRefCenter.class, "actorToRole").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
- MemberModifier.field(WorkersRefCenter.class, "roleToActor").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
+ MemberModifier.field(WorkersRefCenter.class, "roleToWorkerRef").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
+ MemberModifier.field(WorkersRefCenter.class, "actorRefToWorkerRef").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
}
@Test
public void testRegister() throws IllegalAccessException {
- final Props props = Props.create(WorkersListener.class);
- final TestActorRef senderActorRef = TestActorRef.create(system, props, "WorkersListenerSender");
- final TestActorRef receiveactorRef = TestActorRef.create(system, props, "WorkersListenerReceive");
+ Map actorRefToWorkerRef = (Map) MemberModifier.field(WorkersRefCenter.class, "actorRefToWorkerRef").get(WorkersRefCenter.INSTANCE);
+ ActorRef senderRefInWorkerRef = (ActorRef) MemberModifier.field(WorkerRef.class, "actorRef").get(actorRefToWorkerRef.get(senderActorRef));
+ Assert.assertEquals(senderActorRef, senderRefInWorkerRef);
- WorkerListenerMessage.RegisterMessage message = new WorkerListenerMessage.RegisterMessage("WorkersListener");
- receiveactorRef.tell(message, senderActorRef);
-
- Map actorToRole = (Map) MemberModifier.field(WorkersRefCenter.class, "actorToRole").get(WorkersRefCenter.INSTANCE);
- Assert.assertEquals("WorkersListener", actorToRole.get(senderActorRef));
-
- Map> roleToActor = (Map>) MemberModifier.field(WorkersRefCenter.class, "roleToActor").get(WorkersRefCenter.INSTANCE);
- ActorRef[] actorRefs = {senderActorRef};
- Assert.assertArrayEquals(actorRefs, roleToActor.get("WorkersListener").toArray());
+ Map> roleToWorkerRef = (Map>) MemberModifier.field(WorkersRefCenter.class, "roleToWorkerRef").get(WorkersRefCenter.INSTANCE);
+ WorkerRef workerRef = roleToWorkerRef.get("WorkersListener").get(0);
+ senderRefInWorkerRef = (ActorRef) MemberModifier.field(WorkerRef.class, "actorRef").get(workerRef);
+ Assert.assertEquals(senderActorRef, senderRefInWorkerRef);
}
@Test
public void testTerminated() throws IllegalAccessException {
- final Props props = Props.create(WorkersListener.class);
- final TestActorRef senderActorRef = TestActorRef.create(system, props, "WorkersListenerSender");
- final TestActorRef receiveactorRef = TestActorRef.create(system, props, "WorkersListenerReceive");
-
- WorkerListenerMessage.RegisterMessage message = new WorkerListenerMessage.RegisterMessage("WorkersListener");
- receiveactorRef.tell(message, senderActorRef);
-
senderActorRef.stop();
- Map actorToRole = (Map) MemberModifier.field(WorkersRefCenter.class, "actorToRole").get(WorkersRefCenter.INSTANCE);
- Assert.assertEquals(null, actorToRole.get(senderActorRef));
+ Map actorRefToWorkerRef = (Map) MemberModifier.field(WorkersRefCenter.class, "actorRefToWorkerRef").get(WorkersRefCenter.INSTANCE);
+ Assert.assertEquals(null, actorRefToWorkerRef.get(senderActorRef));
- Map> roleToActor = (Map>) MemberModifier.field(WorkersRefCenter.class, "roleToActor").get(WorkersRefCenter.INSTANCE);
+ Map> roleToWorkerRef = (Map>) MemberModifier.field(WorkersRefCenter.class, "roleToWorkerRef").get(WorkersRefCenter.INSTANCE);
ActorRef[] actorRefs = {};
- Assert.assertArrayEquals(actorRefs, roleToActor.get("WorkersListener").toArray());
+ Assert.assertArrayEquals(actorRefs, roleToWorkerRef.get("WorkersListener").toArray());
}
@Test