|
|
@ -40,19 +40,19 @@ import java.util.concurrent.atomic.AtomicLong; |
|
|
public class ActorSystemTest { |
|
|
public class ActorSystemTest { |
|
|
|
|
|
|
|
|
public static final String ROOT_DISPATCHER = "root-dispatcher"; |
|
|
public static final String ROOT_DISPATCHER = "root-dispatcher"; |
|
|
private static final int _1M = 1024 * 1024; |
|
|
private static final int _100K = 100 * 1024; |
|
|
|
|
|
|
|
|
private TbActorSystem actorSystem; |
|
|
private volatile TbActorSystem actorSystem; |
|
|
private ExecutorService submitPool; |
|
|
private volatile ExecutorService submitPool; |
|
|
|
|
|
private int parallelism; |
|
|
|
|
|
|
|
|
@Before |
|
|
@Before |
|
|
public void initActorSystem() { |
|
|
public void initActorSystem() { |
|
|
int cores = Runtime.getRuntime().availableProcessors(); |
|
|
int cores = Runtime.getRuntime().availableProcessors(); |
|
|
int parallelism = Math.max(1, cores / 2); |
|
|
parallelism = Math.max(2, cores / 2); |
|
|
TbActorSystemSettings settings = new TbActorSystemSettings(5, parallelism, 42); |
|
|
TbActorSystemSettings settings = new TbActorSystemSettings(5, parallelism, 42); |
|
|
actorSystem = new DefaultTbActorSystem(settings); |
|
|
actorSystem = new DefaultTbActorSystem(settings); |
|
|
submitPool = Executors.newWorkStealingPool(parallelism); |
|
|
submitPool = Executors.newWorkStealingPool(parallelism); |
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@After |
|
|
@After |
|
|
@ -62,22 +62,44 @@ public class ActorSystemTest { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void test10actorsAnd1MMessages() throws InterruptedException { |
|
|
public void test1actorsAnd100KMessages() throws InterruptedException { |
|
|
testActorsAndMessages(10, _1M); |
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
|
|
|
testActorsAndMessages(1, _100K, 1); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Test |
|
|
|
|
|
public void test10actorsAnd100KMessages() throws InterruptedException { |
|
|
|
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
|
|
|
testActorsAndMessages(10, _100K, 1); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Test |
|
|
|
|
|
public void test100KActorsAnd1Messages5timesSingleThread() throws InterruptedException { |
|
|
|
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newSingleThreadExecutor()); |
|
|
|
|
|
testActorsAndMessages(_100K, 1, 5); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Test |
|
|
|
|
|
public void test100KActorsAnd1Messages5times() throws InterruptedException { |
|
|
|
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
|
|
|
testActorsAndMessages(_100K, 1, 5); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void test1MActorsAnd10Messages() throws InterruptedException { |
|
|
public void test100KActorsAnd10Messages() throws InterruptedException { |
|
|
testActorsAndMessages(_1M, 10); |
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
|
|
|
testActorsAndMessages(_100K, 10, 1); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void test1KActorsAnd1KMessages() throws InterruptedException { |
|
|
public void test1KActorsAnd1KMessages() throws InterruptedException { |
|
|
testActorsAndMessages(1000, 1000); |
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
|
|
|
testActorsAndMessages(1000, 1000, 10); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void testNoMessagesAfterDestroy() throws InterruptedException { |
|
|
public void testNoMessagesAfterDestroy() throws InterruptedException { |
|
|
|
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
ActorTestCtx testCtx1 = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx1 = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx2 = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx2 = getActorTestCtx(1); |
|
|
|
|
|
|
|
|
@ -86,16 +108,17 @@ public class ActorSystemTest { |
|
|
TbActorRef actorId2 = actorSystem.createRootActor(ROOT_DISPATCHER, new SlowInitActor.SlowInitActorCreator( |
|
|
TbActorRef actorId2 = actorSystem.createRootActor(ROOT_DISPATCHER, new SlowInitActor.SlowInitActorCreator( |
|
|
new TbEntityActorId(new DeviceId(UUID.randomUUID())), testCtx2)); |
|
|
new TbEntityActorId(new DeviceId(UUID.randomUUID())), testCtx2)); |
|
|
|
|
|
|
|
|
actorSystem.tell(actorId1, new IntTbActorMsg(42)); |
|
|
actorId1.tell(new IntTbActorMsg(42)); |
|
|
actorSystem.tell(actorId2, new IntTbActorMsg(42)); |
|
|
actorId2.tell(new IntTbActorMsg(42)); |
|
|
actorSystem.stop(actorId1); |
|
|
actorSystem.stop(actorId1); |
|
|
|
|
|
|
|
|
Assert.assertTrue(testCtx2.getLatch().await(1, TimeUnit.SECONDS)); |
|
|
Assert.assertTrue(testCtx2.getLatch().await(1, TimeUnit.SECONDS)); |
|
|
Assert.assertFalse(testCtx1.getLatch().await(2, TimeUnit.SECONDS)); |
|
|
Assert.assertFalse(testCtx1.getLatch().await(1, TimeUnit.SECONDS)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void testOneActorCreated() throws InterruptedException { |
|
|
public void testOneActorCreated() throws InterruptedException { |
|
|
|
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
ActorTestCtx testCtx1 = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx1 = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx2 = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx2 = getActorTestCtx(1); |
|
|
TbActorId actorId = new TbEntityActorId(new DeviceId(UUID.randomUUID())); |
|
|
TbActorId actorId = new TbEntityActorId(new DeviceId(UUID.randomUUID())); |
|
|
@ -105,15 +128,16 @@ public class ActorSystemTest { |
|
|
Thread.sleep(1000); |
|
|
Thread.sleep(1000); |
|
|
actorSystem.tell(actorId, new IntTbActorMsg(42)); |
|
|
actorSystem.tell(actorId, new IntTbActorMsg(42)); |
|
|
|
|
|
|
|
|
Assert.assertTrue(testCtx1.getLatch().await(3, TimeUnit.SECONDS)); |
|
|
Assert.assertTrue(testCtx1.getLatch().await(1, TimeUnit.SECONDS)); |
|
|
Assert.assertFalse(testCtx2.getLatch().await(3, TimeUnit.SECONDS)); |
|
|
Assert.assertFalse(testCtx2.getLatch().await(1, TimeUnit.SECONDS)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void testActorCreatorCalledOnce() throws InterruptedException { |
|
|
public void testActorCreatorCalledOnce() throws InterruptedException { |
|
|
|
|
|
actorSystem.createDispatcher(ROOT_DISPATCHER, Executors.newWorkStealingPool(parallelism)); |
|
|
ActorTestCtx testCtx = getActorTestCtx(1); |
|
|
ActorTestCtx testCtx = getActorTestCtx(1); |
|
|
TbActorId actorId = new TbEntityActorId(new DeviceId(UUID.randomUUID())); |
|
|
TbActorId actorId = new TbEntityActorId(new DeviceId(UUID.randomUUID())); |
|
|
for(int i =0; i < 1000; i++) { |
|
|
for (int i = 0; i < 1000; i++) { |
|
|
submitPool.submit(() -> actorSystem.createRootActor(ROOT_DISPATCHER, new SlowCreateActor.SlowCreateActorCreator(actorId, testCtx))); |
|
|
submitPool.submit(() -> actorSystem.createRootActor(ROOT_DISPATCHER, new SlowCreateActor.SlowCreateActorCreator(actorId, testCtx))); |
|
|
} |
|
|
} |
|
|
Thread.sleep(1000); |
|
|
Thread.sleep(1000); |
|
|
@ -125,7 +149,7 @@ public class ActorSystemTest { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public void testActorsAndMessages(int actorsCount, int msgNumber) throws InterruptedException { |
|
|
public void testActorsAndMessages(int actorsCount, int msgNumber, int times) throws InterruptedException { |
|
|
Random random = new Random(); |
|
|
Random random = new Random(); |
|
|
int[] randomIntegers = new int[msgNumber]; |
|
|
int[] randomIntegers = new int[msgNumber]; |
|
|
long sumTmp = 0; |
|
|
long sumTmp = 0; |
|
|
@ -141,32 +165,35 @@ public class ActorSystemTest { |
|
|
List<TbActorRef> actorRefs = new ArrayList<>(); |
|
|
List<TbActorRef> actorRefs = new ArrayList<>(); |
|
|
for (int actorIdx = 0; actorIdx < actorsCount; actorIdx++) { |
|
|
for (int actorIdx = 0; actorIdx < actorsCount; actorIdx++) { |
|
|
ActorTestCtx testCtx = getActorTestCtx(msgNumber); |
|
|
ActorTestCtx testCtx = getActorTestCtx(msgNumber); |
|
|
|
|
|
|
|
|
actorRefs.add(actorSystem.createRootActor(ROOT_DISPATCHER, new TestRootActor.TestRootActorCreator( |
|
|
actorRefs.add(actorSystem.createRootActor(ROOT_DISPATCHER, new TestRootActor.TestRootActorCreator( |
|
|
new TbEntityActorId(new DeviceId(UUID.randomUUID())), testCtx))); |
|
|
new TbEntityActorId(new DeviceId(UUID.randomUUID())), testCtx))); |
|
|
testCtxes.add(testCtx); |
|
|
testCtxes.add(testCtx); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
long start = System.nanoTime(); |
|
|
for (int t = 0; t < times; t++) { |
|
|
|
|
|
long start = System.nanoTime(); |
|
|
for (int i = 0; i < msgNumber; i++) { |
|
|
for (int i = 0; i < msgNumber; i++) { |
|
|
int tmp = randomIntegers[i]; |
|
|
int tmp = randomIntegers[i]; |
|
|
submitPool.execute(() -> actorRefs.forEach(actorId -> actorSystem.tell(actorId, new IntTbActorMsg(tmp)))); |
|
|
submitPool.execute(() -> actorRefs.forEach(actorId -> actorId.tell(new IntTbActorMsg(tmp)))); |
|
|
} |
|
|
|
|
|
log.info("Submitted all messages"); |
|
|
|
|
|
|
|
|
|
|
|
testCtxes.forEach(ctx -> { |
|
|
|
|
|
try { |
|
|
|
|
|
Assert.assertTrue(ctx.getLatch().await(1, TimeUnit.MINUTES)); |
|
|
|
|
|
Assert.assertEquals(expected, ctx.getActual().get()); |
|
|
|
|
|
Assert.assertEquals(msgNumber, ctx.getInvocationCount().get()); |
|
|
|
|
|
} catch (InterruptedException e) { |
|
|
|
|
|
e.printStackTrace(); |
|
|
|
|
|
} |
|
|
} |
|
|
}); |
|
|
log.info("Submitted all messages"); |
|
|
|
|
|
testCtxes.forEach(ctx -> { |
|
|
long duration = System.nanoTime() - start; |
|
|
try { |
|
|
log.info("Time spend: {}ns ({} ms)", duration, TimeUnit.NANOSECONDS.toMillis(duration)); |
|
|
boolean success = ctx.getLatch().await(1, TimeUnit.MINUTES); |
|
|
|
|
|
if (!success) { |
|
|
|
|
|
log.warn("Failed: {}, {}", ctx.getActual().get(), ctx.getInvocationCount().get()); |
|
|
|
|
|
} |
|
|
|
|
|
Assert.assertTrue(success); |
|
|
|
|
|
Assert.assertEquals(expected, ctx.getActual().get()); |
|
|
|
|
|
Assert.assertEquals(msgNumber, ctx.getInvocationCount().get()); |
|
|
|
|
|
ctx.clear(); |
|
|
|
|
|
} catch (InterruptedException e) { |
|
|
|
|
|
e.printStackTrace(); |
|
|
|
|
|
} |
|
|
|
|
|
}); |
|
|
|
|
|
long duration = System.nanoTime() - start; |
|
|
|
|
|
log.info("Time spend: {}ns ({} ms)", duration, TimeUnit.NANOSECONDS.toMillis(duration)); |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ActorTestCtx getActorTestCtx(int i) { |
|
|
private ActorTestCtx getActorTestCtx(int i) { |
|
|
|