#628 - Refactor Clustering - extract ClusterBroadcast implementations (TCP and Multicast)

This commit is contained in:
Robin Bygrave
2016-03-30 16:20:24 +13:00
parent 30a5769973
commit 9a912390d4
40 changed files with 134 additions and 3976 deletions
@@ -1,59 +0,0 @@
package com.avaje.ebeaninternal.server.cluster.mcast;
import java.util.List;
import org.junit.Assert;
import org.junit.Test;
import com.avaje.ebean.BaseTestCase;
import com.avaje.ebeaninternal.server.cluster.mcast.IncomingPacketsProcessed.GotAllPoint;
public class TestMcastMemberPackets extends BaseTestCase {
@Test
public void test() {
GotAllPoint member = new GotAllPoint("129.12.23.12:9089", 3);
Assert.assertTrue(member.processPacket(1234));
Assert.assertTrue(member.processPacket(1235));
Assert.assertTrue(member.processPacket(1236));
Assert.assertEquals(1236l, member.getGotAllPoint());
Assert.assertEquals(0, member.getMissingPackets().size());
Assert.assertFalse(member.processPacket(1234));
Assert.assertTrue(member.processPacket(1239));
List<Long> missingPackets = member.getMissingPackets();
Assert.assertEquals(2, missingPackets.size());
Assert.assertTrue(missingPackets.contains(1237l));
Assert.assertTrue(missingPackets.contains(1238l));
Assert.assertFalse(missingPackets.contains(1239l));
Assert.assertFalse(missingPackets.contains(1236l));
missingPackets = member.getMissingPackets();
Assert.assertEquals(2, missingPackets.size());
Assert.assertTrue(missingPackets.contains(1237l));
Assert.assertTrue(missingPackets.contains(1238l));
Assert.assertEquals(1236l, member.getGotAllPoint());
// get a missing packet
Assert.assertTrue(member.processPacket(1237));
Assert.assertEquals(1237l, member.getGotAllPoint());
missingPackets = member.getMissingPackets();
Assert.assertEquals(1, missingPackets.size());
Assert.assertTrue(missingPackets.contains(1238l));
// but we now hit maxResendIncoming
missingPackets = member.getMissingPackets();
Assert.assertEquals(0, missingPackets.size());
// gave up on 1238 ..
Assert.assertEquals(1239l, member.getGotAllPoint());
}
}
@@ -1,32 +0,0 @@
package com.avaje.ebeaninternal.server.cluster.mcast;
import org.junit.Assert;
import org.junit.Test;
import com.avaje.ebean.BaseTestCase;
public class TestPacketsAcked extends BaseTestCase {
@Test
public void test() {
OutgoingPacketsAcked packetsAcked = new OutgoingPacketsAcked();
Assert.assertEquals(0l, packetsAcked.getMinimumGotAllPacketId());
long receivedAck = packetsAcked.receivedAck("A", new MessageAck("A", 1020l));
Assert.assertEquals(1020l, packetsAcked.getMinimumGotAllPacketId());
Assert.assertEquals(1020l, receivedAck);
receivedAck = packetsAcked.receivedAck("B", new MessageAck("B", 1030l));
Assert.assertEquals(1020l, packetsAcked.getMinimumGotAllPacketId());
Assert.assertEquals(0l, receivedAck);
receivedAck = packetsAcked.receivedAck("C", new MessageAck("C", 1025l));
Assert.assertEquals(0l, receivedAck);
receivedAck = packetsAcked.receivedAck("A", new MessageAck("A", 1040l));
Assert.assertEquals(1025l, receivedAck);
}
}
@@ -1,79 +0,0 @@
package com.avaje.ebeaninternal.server.cluster.socket;
import com.avaje.ebean.config.ContainerConfig;
import com.avaje.ebeaninternal.api.TDSpiEbeanServer;
import com.avaje.ebeaninternal.api.TransactionEventTable;
import com.avaje.ebeaninternal.server.cluster.ClusterManager;
import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent;
import org.junit.Test;
import java.util.Arrays;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
public class SocketClusterBroadcastTest {
class TestServer extends TDSpiEbeanServer {
RemoteTransactionEvent event;
TestServer(String name) {
super(name);
}
@Override
public void remoteTransactionEvent(RemoteTransactionEvent event) {
this.event = event;
}
}
private ContainerConfig createContainerConfig(String local, String threadPoolName) {
ContainerConfig container0 = new ContainerConfig();
container0.setMode(ContainerConfig.ClusterMode.SOCKET);
ContainerConfig.SocketConfig socketConfig = new ContainerConfig.SocketConfig();
socketConfig.setLocalHostPort(local);
socketConfig.setThreadPoolName(threadPoolName);
socketConfig.setMembers(Arrays.asList("127.0.0.1:9876", "127.0.0.1:9866"));
container0.setSocketConfig(socketConfig);
return container0;
}
@Test
public void testStartup() throws Exception {
ContainerConfig container0 = createContainerConfig("127.0.0.1:9876", "pool0");
ClusterManager mgr0 = new ClusterManager(container0);
TestServer server0 = new TestServer("s001");
mgr0.registerServer(server0);
ContainerConfig container1 = createContainerConfig("127.0.0.1:9866", "pool1");
ClusterManager mgr1 = new ClusterManager(container1);
TestServer server1 = new TestServer("s001");
mgr1.registerServer(server1);
Thread.sleep(1000);
RemoteTransactionEvent evt = new RemoteTransactionEvent("s001");
TransactionEventTable.TableIUD tableIUD = new TransactionEventTable.TableIUD("noSuchTable", true, false, false);
evt.addTableIUD(tableIUD);
assertNull(server1.event);
mgr0.broadcast(evt);
Thread.sleep(100);
assertNotNull(server1.event);
Thread.sleep(1000);
mgr0.shutdown();
mgr1.shutdown();
}
}