Created test that proves that Replicator can recover from connection problems

This commit is contained in:
Simon Pamies committed 2013-06-19 15:13:38 +02:00
1 parent 1e8d1ec0ac
commit d337e7206e
6 files changed
+165 -19

No files matched your search

@@ -1,15 +0,0 @@
<?xml version="1.0" encoding="UTF-8" standalone="no"?>
<launchConfiguration type="org.eclipse.jdt.junit.launchconfig">
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_PATHS">
<listEntry value="/com.sap.sailing.gwt.ui.test/src/com/sap/sailing/gwt/ui/test/TestColumnSwapping.java"/>
</listAttribute>
<listAttribute key="org.eclipse.debug.core.MAPPED_RESOURCE_TYPES">
<listEntry value="1"/>
</listAttribute>
<stringAttribute key="org.eclipse.jdt.junit.CONTAINER" value=""/>
<booleanAttribute key="org.eclipse.jdt.junit.KEEPRUNNING_ATTR" value="false"/>
<stringAttribute key="org.eclipse.jdt.junit.TESTNAME" value=""/>
<stringAttribute key="org.eclipse.jdt.junit.TEST_KIND" value="org.eclipse.jdt.junit.loader.junit4"/>
<stringAttribute key="org.eclipse.jdt.launching.MAIN_TYPE" value="com.sap.sailing.gwt.ui.test.TestColumnSwapping"/>
<stringAttribute key="org.eclipse.jdt.launching.PROJECT_ATTR" value="com.sap.sailing.gwt.ui.test"/>
</launchConfiguration>
@@ -38,6 +38,7 @@ public abstract class AbstractServerReplicationTest {
private DomainFactory resolveAgainst;
protected RacingEventServiceImpl replica;
protected RacingEventServiceImpl master;
protected ReplicationServiceTestImpl replicaReplicator;
private ReplicaDescriptor replicaDescriptor;
private ReplicationServiceImpl masterReplicator;
@@ -94,6 +95,7 @@ public abstract class AbstractServerReplicationTest {
replicaDescriptor, this.replica, this.master, masterReplicator, masterDescriptor);
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> result = new Pair<>(replicaReplicator, masterDescriptor);
replicaReplicator.startInitialLoadTransmissionServlet();
this.replicaReplicator = replicaReplicator;
return result;
}
@@ -103,7 +105,7 @@ public abstract class AbstractServerReplicationTest {
URLConnection urlConnection = new URL("http://localhost:"+SERVLET_PORT+"/STOP").openConnection(); // stop the initial load test server thread
urlConnection.getInputStream().close();
}
static class ReplicationServiceTestImpl extends ReplicationServiceImpl {
private final DomainFactory resolveAgainst;
private final RacingEventService master;
@@ -0,0 +1,133 @@
package com.sap.sailing.server.replication.test;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNotSame;
import static org.junit.Assert.assertNull;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.logging.Logger;
import org.junit.Before;
import org.junit.Test;
import com.rabbitmq.client.AlreadyClosedException;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.ConsumerCancelledException;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.ShutdownSignalException;
import com.sap.sailing.domain.base.Event;
import com.sap.sailing.domain.common.impl.Util;
import com.sap.sailing.domain.common.impl.Util.Pair;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.impl.ReplicationMasterDescriptorImpl;
public class ConnectionResetAndReconnectTest extends AbstractServerReplicationTest {
static final Logger logger = Logger.getLogger(ConnectionResetAndReconnectTest.class.getName());
public static boolean forceStopDelivery = false;
static class QueuingConsumerTest extends QueueingConsumer {
public QueuingConsumerTest(Channel ch) {
super(ch);
}
@Override
public Delivery nextDelivery() throws ShutdownSignalException, ConsumerCancelledException, InterruptedException {
if (forceStopDelivery) {
throw new AlreadyClosedException("Could not connect", this);
}
return super.nextDelivery();
}
}
static class MasterReplicationDescriptorMock extends ReplicationMasterDescriptorImpl {
public MasterReplicationDescriptorMock(String hostname, String exchangeName, int servletPort, int messagingPort) {
super(hostname, exchangeName, servletPort, messagingPort);
}
public static MasterReplicationDescriptorMock from(ReplicationMasterDescriptor obj) {
return new MasterReplicationDescriptorMock(obj.getHostname(), obj.getExchangeName(), obj.getServletPort(), obj.getMessagingPort());
}
@Override
public QueueingConsumer getConsumer() throws IOException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost(getHostname());
int port = getMessagingPort();
if (port != 0) {
connectionFactory.setPort(port);
}
Connection connection = connectionFactory.newConnection();
Channel channel = connection.createChannel();
channel.exchangeDeclare(getExchangeName(), "fanout");
QueueingConsumer consumer = new QueuingConsumerTest(channel);
String queueName = channel.queueDeclare().getQueue();
channel.queueBind(queueName, getExchangeName(), "");
channel.basicConsume(queueName, /* auto-ack */ true, consumer);
return consumer;
}
}
private MasterReplicationDescriptorMock masterReplicationDescriptor;
private ReplicationServiceTestImpl replicaReplicationDescriptor;
@Before
public void setUp() throws Exception {
try {
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> result = basicSetUp(
true, /* master=null means create a new one */ null,
/* replica=null means create a new one */null);
masterReplicationDescriptor = MasterReplicationDescriptorMock.from(result.getB());
replicaReplicationDescriptor = result.getA();
} catch (Exception e) {
e.printStackTrace();
tearDown();
}
}
@Test
public void testReplicaLoosingConnectionToExchangeQueue() throws Exception {
assertNotSame(master, replica);
assertEquals(Util.size(master.getAllRegattas()), Util.size(replica.getAllRegattas()));
/* until here both instances should have the same in-memory state.
* now lets add an event on master and stop the messaging queue. */
stopMessagingExchange();
replicaReplicationDescriptor.startToReplicateFrom(masterReplicationDescriptor);
Event event = addEventOnMaster();
Thread.sleep(1000);
assertNull(replica.getEvent(event.getId()));
startMessagingExchange();
Thread.sleep(2000); // wait for connection to recover
assertNotNull(replica.getEvent(event.getId()));
}
private Event addEventOnMaster() {
final String eventName = "ESS Masquat";
final String venueName = "Masquat, Oman";
final String publicationUrl = "http://ess40.sapsailing.com";
final boolean isPublic = false;
List<String> regattas = new ArrayList<String>();
regattas.add("Day1");
regattas.add("Day2");
return master.addEvent(eventName, venueName, publicationUrl, isPublic, UUID.randomUUID(), regattas);
}
private void stopMessagingExchange() {
forceStopDelivery = true;
}
private void startMessagingExchange() {
forceStopDelivery = false;
}
}
@@ -0,0 +1,7 @@
package com.sap.sailing.server.replication.test;
import com.sap.sailing.server.impl.RacingEventServiceImpl;
public class MockedRacingEventServiceImpl extends RacingEventServiceImpl {
}
@@ -5,6 +5,7 @@ import java.io.IOException;
import java.io.InputStream;
import java.io.ObjectInputStream;
import java.io.ObjectOutputStream;
import java.net.ConnectException;
import java.net.URL;
import java.net.URLConnection;
import java.util.HashMap;
@@ -91,10 +92,20 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
this.exchangeName = exchangeName;
}
private Channel createChannel(String exchangeName) throws IOException {
private Channel createMasterChannel(String exchangeName) throws IOException {
final ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("localhost"); // ...and use default port
Channel result = connectionFactory.newConnection().createChannel();
Channel result = null;
try {
result = connectionFactory.newConnection().createChannel();
} catch (ConnectException ex) {
// make sure to log something meaningful
logger.severe("Could not connect to messaging queue on " + connectionFactory.getHost() + ":" + connectionFactory.getPort() + "/" + exchangeName);
throw ex;
}
logger.info("Connected to " + connectionFactory.getHost() + ":" + connectionFactory.getPort() + "/" + exchangeName);
result.exchangeDeclare(exchangeName, "fanout");
return result;
}
@@ -116,7 +127,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
addAsListenerToRacingEventService();
synchronized (this) {
if (masterChannel == null) {
masterChannel = createChannel(exchangeName);
masterChannel = createMasterChannel(exchangeName);
}
}
}
@@ -8,6 +8,7 @@ import java.util.List;
import java.util.logging.Level;
import java.util.logging.Logger;
import com.rabbitmq.client.AlreadyClosedException;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.QueueingConsumer.Delivery;
import com.rabbitmq.client.ShutdownSignalException;
@@ -82,6 +83,13 @@ public class Replicator implements Runnable {
new ByteArrayInputStream(bytesFromMessage));
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
applyOrQueue(operation);
} catch (AlreadyClosedException ace) {
// ignore this exception and try again after some time
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
} catch (ShutdownSignalException sse) {
logger.info("Received "+sse.getMessage()+". Terminating "+this);
break;