made JMS test work locally with mocked in-VM connection

This commit is contained in:
Axel Uhl
2012-05-02 21:50:30 +02:00
parent fd71724cec
commit 9d74329444
2 changed files with 86 additions and 33 deletions
@@ -8,8 +8,12 @@ import java.io.ObjectOutputStream;
import java.io.PipedInputStream; import java.io.PipedInputStream;
import java.io.PipedOutputStream; import java.io.PipedOutputStream;
import java.net.InetAddress; import java.net.InetAddress;
import java.net.MalformedURLException;
import java.net.URL;
import java.net.UnknownHostException;
import javax.jms.Connection; import javax.jms.Connection;
import javax.jms.JMSException;
import javax.jms.Session; import javax.jms.Session;
import javax.jms.Topic; import javax.jms.Topic;
import javax.jms.TopicSubscriber; import javax.jms.TopicSubscriber;
@@ -24,12 +28,12 @@ import com.sap.sailing.domain.common.impl.Util;
import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.RacingEventService;
import com.sap.sailing.server.impl.RacingEventServiceImpl; import com.sap.sailing.server.impl.RacingEventServiceImpl;
import com.sap.sailing.server.replication.ReplicaDescriptor; import com.sap.sailing.server.replication.ReplicaDescriptor;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.ReplicationService; import com.sap.sailing.server.replication.ReplicationService;
import com.sap.sailing.server.replication.impl.HasRacingEventService; import com.sap.sailing.server.replication.impl.HasRacingEventService;
import com.sap.sailing.server.replication.impl.MessageBrokerConfiguration; import com.sap.sailing.server.replication.impl.MessageBrokerConfiguration;
import com.sap.sailing.server.replication.impl.MessageBrokerManager; import com.sap.sailing.server.replication.impl.MessageBrokerManager;
import com.sap.sailing.server.replication.impl.ReplicationInstancesManager; import com.sap.sailing.server.replication.impl.ReplicationInstancesManager;
import com.sap.sailing.server.replication.impl.ReplicationMasterDescriptorImpl;
import com.sap.sailing.server.replication.impl.ReplicationServiceImpl; import com.sap.sailing.server.replication.impl.ReplicationServiceImpl;
import com.sap.sailing.server.replication.impl.Replicator; import com.sap.sailing.server.replication.impl.Replicator;
@@ -54,8 +58,29 @@ public class ServerReplicationTest {
ReplicationService masterReplicator = new ReplicationServiceImpl(rim, brokerMgr, master); ReplicationService masterReplicator = new ReplicationServiceImpl(rim, brokerMgr, master);
ReplicaDescriptor replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost()); ReplicaDescriptor replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost());
masterReplicator.registerReplica(replicaDescriptor); masterReplicator.registerReplica(replicaDescriptor);
ReplicationMasterDescriptorImpl masterDescriptor = new ReplicationMasterDescriptorImpl("localhost", 8888, 61616); ReplicationMasterDescriptor masterDescriptor = new ReplicationMasterDescriptor() {
ReplicationService replicaReplicator = new ReplicationServiceImpl(rim, brokerMgr, replica); @Override
public URL getReplicationRegistrationRequestURL() throws MalformedURLException {
throw new UnsupportedOperationException();
}
@Override
public URL getInitialLoadURL() throws MalformedURLException {
throw new UnsupportedOperationException();
}
@Override
public TopicSubscriber getTopicSubscriber(String clientID) throws JMSException, UnknownHostException {
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER,
ActiveMQConnection.DEFAULT_PASSWORD, "vm://localhost-jms-connection");
connectionFactory.setClientID(clientID);
Connection connection = connectionFactory.createConnection();
connection.start();
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Topic topic = session.createTopic(ReplicationService.SAILING_SERVER_REPLICATION_TOPIC);
return session.createDurableSubscriber(topic, InetAddress.getLocalHost().getHostAddress());
}
};
ReplicationService replicaReplicator = new ReplicationServiceTestImpl(resolveAgainst, rim, brokerMgr, replicaDescriptor, replica, master, masterReplicator);
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER, ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER,
ActiveMQConnection.DEFAULT_PASSWORD, "vm://localhost-jms-connection"); ActiveMQConnection.DEFAULT_PASSWORD, "vm://localhost-jms-connection");
connectionFactory.setClientID("Test Client"); connectionFactory.setClientID("Test Client");
@@ -68,7 +93,7 @@ public class ServerReplicationTest {
replicaReplicator.startToReplicateFrom(masterDescriptor); replicaReplicator.startToReplicateFrom(masterDescriptor);
} }
private Replicator createReplicator(ReplicationMasterDescriptorImpl masterDescriptor, final RacingEventService master) { private Replicator createReplicator(ReplicationMasterDescriptor masterDescriptor, final RacingEventService master) {
return new Replicator(masterDescriptor, new HasRacingEventService() { return new Replicator(masterDescriptor, new HasRacingEventService() {
@Override @Override
public RacingEventService getRacingEventService() { public RacingEventService getRacingEventService() {
@@ -79,35 +104,63 @@ public class ServerReplicationTest {
@Test @Test
public void testBasicInitialLoad() throws Exception { public void testBasicInitialLoad() throws Exception {
initialLoad();
assertEquals(Util.size(master.getAllEvents()), Util.size(replica.getAllEvents())); assertEquals(Util.size(master.getAllEvents()), Util.size(replica.getAllEvents()));
assertEquals(master.getLeaderboardGroups().size(), replica.getLeaderboardGroups().size()); assertEquals(master.getLeaderboardGroups().size(), replica.getLeaderboardGroups().size());
assertEquals(master.getLeaderboards().size(), replica.getLeaderboards().size()); assertEquals(master.getLeaderboards().size(), replica.getLeaderboards().size());
assertEquals(master.getLeaderboards().keySet(), replica.getLeaderboards().keySet()); assertEquals(master.getLeaderboards().keySet(), replica.getLeaderboards().keySet());
} }
/** private static class ReplicationServiceTestImpl extends ReplicationServiceImpl {
* Clones the {@link #master}'s state to the {@link #replica} using private final DomainFactory resolveAgainst;
* {@link RacingEventServiceImpl#serializeForInitialReplication(ObjectOutputStream)} and private final RacingEventService master;
* {@link RacingEventServiceImpl#initiallyFillFrom(ObjectInputStream)} through a piped input/output stream. private final ReplicaDescriptor replicaDescriptor;
*/ private final ReplicationService masterReplicationService;
private void initialLoad() throws IOException, ClassNotFoundException {
PipedOutputStream pos = new PipedOutputStream(); public ReplicationServiceTestImpl(DomainFactory resolveAgainst,
PipedInputStream pis = new PipedInputStream(pos); ReplicationInstancesManager replicationInstancesManager, MessageBrokerManager messageBrokerManager,
final ObjectOutputStream oos = new ObjectOutputStream(pos); ReplicaDescriptor replicaDescriptor, RacingEventService replica, RacingEventService master, ReplicationService masterReplicationService) {
new Thread("clone writer") { super(replicationInstancesManager, messageBrokerManager, replica);
public void run() { this.resolveAgainst = resolveAgainst;
try { this.replicaDescriptor = replicaDescriptor;
master.serializeForInitialReplication(oos); this.master = master;
oos.close(); this.masterReplicationService = masterReplicationService;
} catch (IOException e) { }
e.printStackTrace();
throw new RuntimeException(e); /**
* Ignore the master descriptor and replicate from the local master passed to the constructor instead.
*/
@Override
public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException,
ClassNotFoundException, JMSException {
masterReplicationService.registerReplica(replicaDescriptor);
TopicSubscriber replicationSubscription = master.getTopicSubscriber(replicaDescriptor.getUuid().toString());
replicationSubscription.setMessageListener(new Replicator(master, this));
initialLoad();
}
/**
* Clones the {@link #master}'s state to the {@link #replica} using
* {@link RacingEventServiceImpl#serializeForInitialReplication(ObjectOutputStream)} and
* {@link RacingEventServiceImpl#initiallyFillFrom(ObjectInputStream)} through a piped input/output stream.
*/
private void initialLoad() throws IOException, ClassNotFoundException {
PipedOutputStream pos = new PipedOutputStream();
PipedInputStream pis = new PipedInputStream(pos);
final ObjectOutputStream oos = new ObjectOutputStream(pos);
new Thread("clone writer") {
public void run() {
try {
master.serializeForInitialReplication(oos);
oos.close();
} catch (IOException e) {
e.printStackTrace();
throw new RuntimeException(e);
}
} }
} }.start();
}.start(); ObjectInputStream dis = resolveAgainst.createObjectInputStreamResolvingAgainstThisFactory(pis);
ObjectInputStream dis = resolveAgainst.createObjectInputStreamResolvingAgainstThisFactory(pis); getRacingEventService().initiallyFillFrom(dis);
replica.initiallyFillFrom(dis); dis.close();
dis.close(); }
} }
} }
@@ -54,7 +54,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
private ServiceTracker<RacingEventService, RacingEventService> racingEventServiceTracker; private ServiceTracker<RacingEventService, RacingEventService> racingEventServiceTracker;
private final RacingEventService master; private final RacingEventService localService;
/** /**
* The UUIDs with which this replica is registered by the master identified by the corresponding key * The UUIDs with which this replica is registered by the master identified by the corresponding key
@@ -69,7 +69,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
racingEventServiceTracker = new ServiceTracker<RacingEventService, RacingEventService>( racingEventServiceTracker = new ServiceTracker<RacingEventService, RacingEventService>(
Activator.getDefaultContext(), RacingEventService.class.getName(), null); Activator.getDefaultContext(), RacingEventService.class.getName(), null);
racingEventServiceTracker.open(); racingEventServiceTracker.open();
master = null; localService = null;
} }
/** /**
@@ -77,18 +77,18 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
* an OSGi service tracker to discover the {@link RacingEventService}, the service to replicate is "injected" here. * an OSGi service tracker to discover the {@link RacingEventService}, the service to replicate is "injected" here.
*/ */
public ReplicationServiceImpl(final ReplicationInstancesManager replicationInstancesManager, public ReplicationServiceImpl(final ReplicationInstancesManager replicationInstancesManager,
final MessageBrokerManager messageBrokerManager, RacingEventService master) { final MessageBrokerManager messageBrokerManager, RacingEventService localService) {
this.replicationInstancesManager = replicationInstancesManager; this.replicationInstancesManager = replicationInstancesManager;
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>(); replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
this.messageBrokerManager = messageBrokerManager; this.messageBrokerManager = messageBrokerManager;
this.master = master; this.localService = localService;
} }
@Override @Override
public RacingEventService getRacingEventService() { public RacingEventService getRacingEventService() {
RacingEventService result; RacingEventService result;
if (master != null) { if (localService != null) {
result = master; result = localService;
} else { } else {
result = racingEventServiceTracker.getService(); result = racingEventServiceTracker.getService();
} }