diff --git a/java/com.sap.sailing.feature.p2build/raceanalysis.product b/java/com.sap.sailing.feature.p2build/raceanalysis.product index 904b3ed6703..c38f8177f84 100644 --- a/java/com.sap.sailing.feature.p2build/raceanalysis.product +++ b/java/com.sap.sailing.feature.p2build/raceanalysis.product @@ -1,54 +1,52 @@ - - - - - - - - - - -os ${target.os} -ws ${target.ws} -arch ${target.arch} -nl ${target.nl} -consoleLog -console -clean + + + + + + + + + + -os ${target.os} -ws ${target.ws} -arch ${target.arch} -nl ${target.nl} -consoleLog -console -clean -Declipse.ignoreApp=true -Dosgi.noShutdown=true --Xmx1024m -Djetty.home=configuration/jetty - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - +-Xmx1024m -Djetty.home=configuration/jetty + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/adminconsole/ReplicationPanel.java b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/adminconsole/ReplicationPanel.java index 43c0f0ae627..3454b5c7f0f 100755 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/adminconsole/ReplicationPanel.java +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/adminconsole/ReplicationPanel.java @@ -12,6 +12,7 @@ import com.google.gwt.user.client.ui.IntegerBox; import com.google.gwt.user.client.ui.Label; import com.google.gwt.user.client.ui.TextBox; import com.google.gwt.user.client.ui.Widget; +import com.sap.sailing.domain.common.impl.Util.Pair; import com.sap.sailing.domain.common.impl.Util.Triple; import com.sap.sailing.gwt.ui.client.DataEntryDialog; import com.sap.sailing.gwt.ui.client.ErrorReporter; @@ -61,16 +62,18 @@ public class ReplicationPanel extends FlowPanel { private void addReplication() { AddReplicationDialog dialog = new AddReplicationDialog(null, - new AsyncCallback>() { + new AsyncCallback, Integer, Integer>>() { @Override - public void onSuccess(final Triple masterNameAndJMSPortNumberAndServletPortNumber) { - sailingService.startReplicatingFromMaster(masterNameAndJMSPortNumberAndServletPortNumber.getA(), - /* TODO servlet port */ masterNameAndJMSPortNumberAndServletPortNumber.getC(), - /* TODO JMS port */ masterNameAndJMSPortNumberAndServletPortNumber.getB(), new AsyncCallback() { + public void onSuccess(final Triple, Integer, Integer> masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber) { + sailingService.startReplicatingFromMaster(masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getA(), + masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getB(), + /* TODO servlet port */ masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getC(), + /* TODO JMS port */ masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getB(), new AsyncCallback() { @Override public void onFailure(Throwable e) { errorReporter.reportError(stringMessages.errorStartingReplication( - masterNameAndJMSPortNumberAndServletPortNumber.getA(), e.getMessage())); + masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getA(), + masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getB(), e.getMessage())); } @Override @@ -131,17 +134,19 @@ public class ReplicationPanel extends FlowPanel { * @author Axel Uhl (d043530) * */ - private class AddReplicationDialog extends DataEntryDialog> { - private final TextBox entryField; - private final IntegerBox jmsPortField; + private class AddReplicationDialog extends DataEntryDialog, Integer, Integer>> { + private final TextBox hostnameEntryField; + private final TextBox exchangenameEntryField; + private final IntegerBox messagingPortField; private final IntegerBox servletPortField; - public AddReplicationDialog(final Validator> validator, - final AsyncCallback> callback) { + public AddReplicationDialog(final Validator, Integer, Integer>> validator, + final AsyncCallback, Integer, Integer>> callback) { super(stringMessages.add(), stringMessages.enterMaster(), stringMessages.ok(), stringMessages.cancel(), validator, callback); - entryField = createTextBox(""); - jmsPortField = createIntegerBox(61616, /* visible length */ 5); + hostnameEntryField = createTextBox(""); + exchangenameEntryField = createTextBox(""); + messagingPortField = createIntegerBox(61616, /* visible length */ 5); servletPortField = createIntegerBox(8888, /* visibleLength */ 5); } @@ -151,25 +156,29 @@ public class ReplicationPanel extends FlowPanel { */ @Override protected Widget getAdditionalWidget() { - Grid grid = new Grid(3, 2); + Grid grid = new Grid(4, 2); grid.setWidget(0, 0, new Label(stringMessages.hostname())); - grid.setWidget(0, 1, entryField); - grid.setWidget(1, 0, new Label(stringMessages.jmsPortNumber())); - grid.setWidget(1, 1, jmsPortField); - grid.setWidget(2, 0, new Label(stringMessages.servletPortNumber())); - grid.setWidget(2, 1, servletPortField); + grid.setWidget(0, 1, hostnameEntryField); + grid.setWidget(0, 0, new Label(stringMessages.exchangeName())); + grid.setWidget(0, 1, exchangenameEntryField); + grid.setWidget(2, 0, new Label(stringMessages.jmsPortNumber())); + grid.setWidget(2, 1, messagingPortField); + grid.setWidget(3, 0, new Label(stringMessages.servletPortNumber())); + grid.setWidget(4, 1, servletPortField); return grid; } @Override public void show() { super.show(); - entryField.setFocus(true); + hostnameEntryField.setFocus(true); } @Override - protected Triple getResult() { - return new Triple(entryField.getText(), jmsPortField.getValue(), servletPortField.getValue()); + protected Triple, Integer, Integer> getResult() { + return new Triple, Integer, Integer>( + new Pair(hostnameEntryField.getText(), exchangenameEntryField.getText()), + messagingPortField.getValue(), servletPortField.getValue()); } } diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingService.java b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingService.java index 3dcb49b0253..8ef1f895124 100644 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingService.java +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingService.java @@ -213,7 +213,7 @@ public interface SailingService extends RemoteService { ReplicationStateDTO getReplicaInfo(); - void startReplicatingFromMaster(String masterName, int servletPort, int jmsPort) throws Exception; + void startReplicatingFromMaster(String masterName, String exchangeName, int servletPort, int messagingPort) throws Exception; void updateRaceDelayToLive(RegattaAndRaceIdentifier regattaAndRaceIdentifier, long delayToLiveInMs); diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingServiceAsync.java b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingServiceAsync.java index 71df563be43..b95388c0e7d 100644 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingServiceAsync.java +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/SailingServiceAsync.java @@ -330,7 +330,8 @@ public interface SailingServiceAsync { void getReplicaInfo(AsyncCallback callback); - void startReplicatingFromMaster(String masterName, int servletPort, int jmsPort, AsyncCallback callback); + void startReplicatingFromMaster(String masterName, String exchangeName, int servletPort, int messagingPort, + AsyncCallback callback); void getEvents(AsyncCallback> callback); diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.java b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.java index 5056851e3ad..c119092cb56 100644 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.java +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.java @@ -291,7 +291,7 @@ public interface StringMessages extends Messages { String errorFetchingReplicaData(String message); String averageCrossTrackErrorInMeters(); String enterMaster(); - String errorStartingReplication(String hostname, String message); + String errorStartingReplication(String hostname, String exchangeName, String message); String helpLines(); String startLine(); String finishLine(); @@ -339,4 +339,5 @@ public interface StringMessages extends Messages { String regattaExistForSelectedBoatClass(); String reload(); String addRegatta(); + String exchangeName(); } \ No newline at end of file diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.properties b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.properties index e3d8e5569ac..b8c9930b8f5 100644 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.properties +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages.properties @@ -289,7 +289,7 @@ replication=Replication errorFetchingReplicaData=Error fetching replica data: {0} averageCrossTrackErrorInMeters=\u2205 XTE enterMaster=Hostname of Master Instance -errorStartingReplication=Error starting replication from host {0}: {1} +errorStartingReplication=Error starting replication from host {0} / exchange {1}: {2} helpLines=Help lines startLine=Start line finishLine=Finish line @@ -337,4 +337,5 @@ Do you really want to use the regatta ''{1}''? regattaExistForSelectedBoatClass=There is at least one regatta for the selected boat classes.\ Do you really want to use the ''no regatta'' selection? reload=Reload -addRegatta=Add Regatta... \ No newline at end of file +addRegatta=Add Regatta... +exchangeName=Exchange Name \ No newline at end of file diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages_de.properties b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages_de.properties index 2e3fb5dffb6..b6c3057916b 100755 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages_de.properties +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/client/StringMessages_de.properties @@ -290,7 +290,7 @@ replication=Replikation errorFetchingReplicaData=Fehler beim Beschaffen der Replikationsdaten: {0} averageCrossTrackErrorInMeters=\u2205 XTE enterMaster=Hostname der Master-Instanz -errorStartingReplication=Fehler beim Starten der Replikation von Host {0}: {1} +errorStartingReplication=Fehler beim Starten der Replikation von Host {0} / Exchange {1}: {2} helpLines=Hilfslinien startLine=Startlinie finishLine=Ziellinie @@ -338,4 +338,5 @@ Willst du wirklich diese Regatta benutzen ''{1}''? regattaExistForSelectedBoatClass=Es gibt mindestens eine Regatta zu den ausgewählten Bootsklassen.\ Willst du wirklich die Auswahl ''keine Regatta'' benutzen? reload=Neu laden -addRegatta=Regatta hinzufügen... \ No newline at end of file +addRegatta=Regatta hinzufügen... +exchangeName=Exchange Name \ No newline at end of file diff --git a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/server/SailingServiceImpl.java b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/server/SailingServiceImpl.java index 72da8835f24..fb8419e4bf1 100755 --- a/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/server/SailingServiceImpl.java +++ b/java/com.sap.sailing.gwt.ui/src/main/java/com/sap/sailing/gwt/ui/server/SailingServiceImpl.java @@ -30,7 +30,6 @@ import java.util.concurrent.FutureTask; import java.util.concurrent.RunnableFuture; import java.util.logging.Logger; -import javax.jms.JMSException; import javax.servlet.ServletContext; import org.osgi.framework.BundleContext; @@ -2217,9 +2216,9 @@ public class SailingServiceImpl extends RemoteServiceServlet implements SailingS } @Override - public void startReplicatingFromMaster(String masterName, int servletPort, int jmsPort) throws IOException, ClassNotFoundException, JMSException { + public void startReplicatingFromMaster(String masterName, String exchangeName, int servletPort, int messagingPort) throws IOException, ClassNotFoundException { getReplicationService().startToReplicateFrom( - ReplicationFactory.INSTANCE.createReplicationMasterDescriptor(masterName, servletPort, jmsPort)); + ReplicationFactory.INSTANCE.createReplicationMasterDescriptor(masterName, exchangeName, servletPort, messagingPort)); } @Override diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java index 5a5ed286b03..695dac45422 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java @@ -1,28 +1,16 @@ package com.sap.sailing.server.replication.test; -import java.io.File; -import java.io.FileNotFoundException; import java.io.IOException; import java.io.ObjectInputStream; import java.io.ObjectOutputStream; import java.io.PipedInputStream; import java.io.PipedOutputStream; import java.net.InetAddress; -import java.net.MalformedURLException; -import java.net.URL; -import java.net.UnknownHostException; -import javax.jms.Connection; -import javax.jms.JMSException; -import javax.jms.Session; -import javax.jms.Topic; -import javax.jms.TopicSubscriber; - -import org.apache.activemq.ActiveMQConnection; -import org.apache.activemq.ActiveMQConnectionFactory; import org.junit.After; import org.junit.Before; +import com.rabbitmq.client.QueueingConsumer; import com.sap.sailing.domain.base.DomainFactory; import com.sap.sailing.domain.common.impl.Util.Pair; import com.sap.sailing.mongodb.MongoDBService; @@ -31,10 +19,8 @@ import com.sap.sailing.server.impl.RacingEventServiceImpl; 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.impl.Activator; -import com.sap.sailing.server.replication.impl.MessageBrokerConfiguration; -import com.sap.sailing.server.replication.impl.MessageBrokerManager; 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.Replicator; @@ -42,8 +28,6 @@ public abstract class AbstractServerReplicationTest { private DomainFactory resolveAgainst; protected RacingEventServiceImpl replica; protected RacingEventServiceImpl master; - private MessageBrokerManager brokerMgr; - private File brokerPersistenceDir; private ReplicaDescriptor replicaDescriptor; private ReplicationServiceImpl masterReplicator; @@ -73,8 +57,8 @@ public abstract class AbstractServerReplicationTest { * service will be created as replica */ protected Pair basicSetUp( - boolean dropDB, RacingEventServiceImpl master, RacingEventServiceImpl replica) throws FileNotFoundException, Exception, - JMSException, UnknownHostException { + boolean dropDB, RacingEventServiceImpl master, RacingEventServiceImpl replica) throws IOException { + final String exchangeName = "test-sapsailinganalytics-exchange"; final MongoDBService mongoDBService = MongoDBService.INSTANCE; if (dropDB) { mongoDBService.getDB().dropDatabase(); @@ -91,53 +75,11 @@ public abstract class AbstractServerReplicationTest { this.replica = new RacingEventServiceImpl(mongoDBService); } ReplicationInstancesManager rim = new ReplicationInstancesManager(); - final String IN_VM_BROKER_URL = "vm://localhost-jms-connection?broker.useJmx=false"; - final String activeMQPersistenceParentDir = System.getProperty("java.io.tmpdir"); - final String brokerName = "local_in-VM_test_broker"; - brokerPersistenceDir = new File(activeMQPersistenceParentDir, brokerName); - Activator.removeTemporaryTestBrokerPersistenceDirectory(brokerPersistenceDir); - brokerMgr = new MessageBrokerManager(new MessageBrokerConfiguration(brokerName, - IN_VM_BROKER_URL, activeMQPersistenceParentDir)); - brokerMgr.startMessageBroker(/* useJmx */ false); - brokerMgr.createAndStartConnection(); - masterReplicator = new ReplicationServiceImpl(rim, brokerMgr, this.master); + masterReplicator = new ReplicationServiceImpl(exchangeName, rim, this.master); replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost()); masterReplicator.registerReplica(replicaDescriptor); - ReplicationMasterDescriptor masterDescriptor = new ReplicationMasterDescriptor() { - @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, IN_VM_BROKER_URL); - 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()); - } - @Override - public int getJMSPort() { - return 0; - } - @Override - public int getServletPort() { - return 0; - } - @Override - public String getHostname() { - return null; - } - }; - ReplicationServiceTestImpl replicaReplicator = new ReplicationServiceTestImpl(resolveAgainst, rim, brokerMgr, + ReplicationMasterDescriptor masterDescriptor = new ReplicationMasterDescriptorImpl(null, exchangeName, 0, 0); + ReplicationServiceTestImpl replicaReplicator = new ReplicationServiceTestImpl(exchangeName, resolveAgainst, rim, replicaDescriptor, this.replica, this.master, masterReplicator); Pair result = new Pair<>(replicaReplicator, masterDescriptor); return result; @@ -145,10 +87,6 @@ public abstract class AbstractServerReplicationTest { @After public void tearDown() throws Exception { - brokerMgr.closeSessions(); - brokerMgr.closeConnections(); - brokerMgr.stopMessageBroker(); - Activator.removeTemporaryTestBrokerPersistenceDirectory(brokerPersistenceDir); masterReplicator.unregisterReplica(replicaDescriptor); } @@ -158,10 +96,11 @@ public abstract class AbstractServerReplicationTest { private final ReplicaDescriptor replicaDescriptor; private final ReplicationService masterReplicationService; - public ReplicationServiceTestImpl(DomainFactory resolveAgainst, - ReplicationInstancesManager replicationInstancesManager, MessageBrokerManager messageBrokerManager, - ReplicaDescriptor replicaDescriptor, RacingEventService replica, RacingEventService master, ReplicationService masterReplicationService) { - super(replicationInstancesManager, messageBrokerManager, replica); + public ReplicationServiceTestImpl(String exchangeName, DomainFactory resolveAgainst, + ReplicationInstancesManager replicationInstancesManager, ReplicaDescriptor replicaDescriptor, + RacingEventService replica, RacingEventService master, ReplicationService masterReplicationService) + throws IOException { + super(exchangeName, replicationInstancesManager, replica); this.resolveAgainst = resolveAgainst; this.replicaDescriptor = replicaDescriptor; this.master = master; @@ -173,19 +112,19 @@ public abstract class AbstractServerReplicationTest { */ @Override public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, - ClassNotFoundException, JMSException { + ClassNotFoundException { Replicator replicator = startToReplicateFromButDontYetFetchInitialLoad(master, /* startReplicatorSuspended */ true); initialLoad(); replicator.setSuspended(false); // resume after initial load } protected Replicator startToReplicateFromButDontYetFetchInitialLoad(ReplicationMasterDescriptor master, boolean startReplicatorSuspended) - throws JMSException, UnknownHostException { + throws IOException { masterReplicationService.registerReplica(replicaDescriptor); registerReplicaUuidForMaster(replicaDescriptor.getUuid().toString(), master); - TopicSubscriber replicationSubscription = master.getTopicSubscriber(replicaDescriptor.getUuid().toString()); - final Replicator replicator = new Replicator(master, this, startReplicatorSuspended); - replicationSubscription.setMessageListener(replicator); + QueueingConsumer consumer = master.getConsumer(); + final Replicator replicator = new Replicator(master, this, startReplicatorSuspended, consumer); + new Thread(replicator).start(); return replicator; } diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/DelayedLeaderboardCorrectionsReplicationTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/DelayedLeaderboardCorrectionsReplicationTest.java index aad57788d01..0d6cd3ae108 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/DelayedLeaderboardCorrectionsReplicationTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/DelayedLeaderboardCorrectionsReplicationTest.java @@ -7,8 +7,6 @@ import static org.junit.Assert.assertNull; import java.io.FileNotFoundException; import java.net.UnknownHostException; -import javax.jms.JMSException; - import org.junit.Before; import org.junit.Test; @@ -43,7 +41,7 @@ public class DelayedLeaderboardCorrectionsReplicationTest extends AbstractServer private ReplicationMasterDescriptor masterDescriptor; @Before - public void setUp() throws FileNotFoundException, UnknownHostException, JMSException, Exception { + public void setUp() throws FileNotFoundException, UnknownHostException { final MongoDBService mongoDBService = MongoDBService.INSTANCE; mongoDBService.getDB().dropDatabase(); } diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/SimpleRabbitMQTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/SimpleRabbitMQTest.java index 0e3a6d4ba43..5e650ab2e4d 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/SimpleRabbitMQTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/SimpleRabbitMQTest.java @@ -96,7 +96,7 @@ public class SimpleRabbitMQTest { @Override public void run() { try { - channel.basicConsume(getQueueName(), true, consumer); + channel.basicConsume(getQueueName(), /* auto-ack */ true, consumer); QueueingConsumer.Delivery delivery = consumer.nextDelivery(); final String s = new String(delivery.getBody()); received.put(this, s); diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationFactory.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationFactory.java index 0b590c533a6..ddf3ab852e4 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationFactory.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationFactory.java @@ -5,5 +5,5 @@ import com.sap.sailing.server.replication.impl.ReplicationFactoryImpl; public interface ReplicationFactory { static ReplicationFactory INSTANCE = new ReplicationFactoryImpl(); - ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, int servletPort, int jmsPort); + ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int jmsPort); } diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationMasterDescriptor.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationMasterDescriptor.java index 6a2bdcc8423..7b50abf4d65 100644 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationMasterDescriptor.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationMasterDescriptor.java @@ -1,11 +1,10 @@ package com.sap.sailing.server.replication; +import java.io.IOException; import java.net.MalformedURLException; import java.net.URL; -import java.net.UnknownHostException; -import javax.jms.JMSException; -import javax.jms.TopicSubscriber; +import com.rabbitmq.client.QueueingConsumer; /** * Identifies a master server instance from which a replica can obtain an initial load and continuous updates. @@ -19,11 +18,16 @@ public interface ReplicationMasterDescriptor { URL getInitialLoadURL() throws MalformedURLException; - TopicSubscriber getTopicSubscriber(String clientID) throws JMSException, UnknownHostException; - int getJMSPort(); int getServletPort(); String getHostname(); + + /** + * Creates a queue, declares the master's fanout exchange on the calling client and binds the queue to the exchange. + * Then, adds a consumer to the queue just created and starts consuming. The caller may keep calling + * {@link QueueingConsumer#nextDelivery()} on the consumer returned in order to obtain the next message. + */ + QueueingConsumer getConsumer() throws IOException; } diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationService.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationService.java index 490c06223eb..c7a68b7a47d 100644 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationService.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/ReplicationService.java @@ -3,8 +3,6 @@ package com.sap.sailing.server.replication; import java.io.IOException; import java.util.Map; -import javax.jms.JMSException; - import com.sap.sailing.server.RacingEventServiceOperation; import com.sap.sailing.server.replication.impl.ReplicationServlet; @@ -27,16 +25,15 @@ public interface ReplicationService { * the JMS replication topic is created, then subscribing for the master's JMS replication topic and asking the servlet * for the stream containing the initial load. */ - void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException, JMSException; + void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException; /** - * Registers a replica with this master instance. If the replication topic hasn't been created in the - * JMS message broker yet, it will be when this method returns. The replica will be considered - * in the result of {@link #getReplicaInfo()} when this call has succeeded. + * Registers a replica with this master instance. The replica will be considered in the result of + * {@link #getReplicaInfo()} when this call has succeeded. */ - void registerReplica(ReplicaDescriptor replica) throws JMSException; + void registerReplica(ReplicaDescriptor replica); - void unregisterReplica(ReplicaDescriptor replica) throws JMSException; + void unregisterReplica(ReplicaDescriptor replica) throws IOException; /** * For a replica replicating off this master, provides statistics in the form of number of operations sent to that diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Activator.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Activator.java index ef23c840cff..ae8231604b6 100644 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Activator.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Activator.java @@ -1,7 +1,5 @@ package com.sap.sailing.server.replication.impl; -import java.io.File; -import java.io.FileNotFoundException; import java.util.logging.Logger; import org.osgi.framework.BundleActivator; @@ -12,72 +10,25 @@ import com.sap.sailing.server.replication.ReplicationService; public class Activator implements BundleActivator { private static final Logger logger = Logger.getLogger(Activator.class.getName()); - private static final String REPLICATION_PERSISTENCE_DIR_PROPERTY = "replication.persistenceDir"; - - private static final String BROKER_URL_PROPERTY = "replication.brokerURL"; - - private static final String REPLICATION_USE_JMX_PROPERTY = "replication.useJMX"; - - private MessageBrokerManager messageBrokerManager; + private static final String PROPERTY_NAME_EXCHANGE_NAME = "replication.exchangeName"; private ReplicationInstancesManager replicationInstancesManager; private static BundleContext defaultContext; - + public void start(BundleContext bundleContext) throws Exception { defaultContext = bundleContext; - String replicationPersistenceDirectory = bundleContext.getProperty(REPLICATION_PERSISTENCE_DIR_PROPERTY); - if (replicationPersistenceDirectory == null) { - replicationPersistenceDirectory = System.getProperty("java.io.tmpdir"); + String exchangeName = bundleContext.getProperty(PROPERTY_NAME_EXCHANGE_NAME); + if (exchangeName == null) { + exchangeName = "sapsailinganalytics"; } - String brokerURL = bundleContext.getProperty(BROKER_URL_PROPERTY); - if (brokerURL == null) { - brokerURL = "tcp://localhost:61616"; - } - final File brokerPersistenceDir = new File(replicationPersistenceDirectory, "kahadb"); - removeTemporaryTestBrokerPersistenceDirectory(brokerPersistenceDir); - MessageBrokerConfiguration brokerConfig = new MessageBrokerConfiguration("SailingServerReplicationBroker", - brokerURL, brokerPersistenceDir.getAbsolutePath()); - messageBrokerManager = new MessageBrokerManager(brokerConfig); - String useJMX = bundleContext.getProperty(REPLICATION_USE_JMX_PROPERTY); - if (useJMX == null || useJMX.length() == 0) { - useJMX = "false"; - } - messageBrokerManager.startMessageBroker(Boolean.valueOf(useJMX)); - messageBrokerManager.createAndStartConnection(); replicationInstancesManager = new ReplicationInstancesManager(); - ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(replicationInstancesManager, messageBrokerManager); + ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(exchangeName, replicationInstancesManager); bundleContext.registerService(ReplicationService.class, serverReplicationMasterService, null); + logger.info("Registered replication service "+serverReplicationMasterService); } - public static void removeTemporaryTestBrokerPersistenceDirectory(File brokerPersistenceDir) throws FileNotFoundException { - if (brokerPersistenceDir.exists() && brokerPersistenceDir.isDirectory()) { - logger.info("Deleting message broker persistence director "+brokerPersistenceDir); - deleteRecursive(brokerPersistenceDir); - } - File failoverStore = new File("activemq-data"); - if (failoverStore.exists() && failoverStore.isDirectory()) { - logger.info("Deleting message broker failover store "+failoverStore); - deleteRecursive(failoverStore); - } - } - - public static boolean deleteRecursive(File path) throws FileNotFoundException{ - if (!path.exists()) throw new FileNotFoundException(path.getAbsolutePath()); - boolean ret = true; - if (path.isDirectory()){ - for (File f : path.listFiles()){ - ret = ret && deleteRecursive(f); - } - } - return ret && path.delete(); - } - - public void stop(BundleContext bundleContext) throws Exception { - messageBrokerManager.closeSessions(); - messageBrokerManager.closeConnections(); - messageBrokerManager.stopMessageBroker(); } public static BundleContext getDefaultContext() { diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/MessageBrokerConfiguration.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/MessageBrokerConfiguration.java deleted file mode 100644 index 73b9a03a5ec..00000000000 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/MessageBrokerConfiguration.java +++ /dev/null @@ -1,32 +0,0 @@ -package com.sap.sailing.server.replication.impl; - -/** - * A simple configuration for the message broker - */ -public class MessageBrokerConfiguration { - private final String brokerName; - - private final String brokerUrl; - - private final String dataStoreDirectory; - - public MessageBrokerConfiguration(String brokerName, String brokerUrl, String dataStoreDirectory) { - super(); - this.brokerName = brokerName; - this.brokerUrl = brokerUrl; - this.dataStoreDirectory = dataStoreDirectory; - } - - public String getBrokerName() { - return brokerName; - } - - public String getDataStoreDirectory() { - return dataStoreDirectory; - } - - public String getBrokerUrl() { - return brokerUrl; - } - -} diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/MessageBrokerManager.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/MessageBrokerManager.java deleted file mode 100644 index 7e472639180..00000000000 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/MessageBrokerManager.java +++ /dev/null @@ -1,73 +0,0 @@ -package com.sap.sailing.server.replication.impl; - -import java.net.URI; - -import javax.jms.Connection; -import javax.jms.JMSException; -import javax.jms.Session; - -import org.apache.activemq.ActiveMQConnection; -import org.apache.activemq.ActiveMQConnectionFactory; -import org.apache.activemq.broker.BrokerService; -import org.apache.activemq.broker.TransportConnector; - -public class MessageBrokerManager { - private final MessageBrokerConfiguration configuration; - - private ActiveMQConnectionFactory connectionFactory; - private Connection connection; - private Session session; - - private BrokerService broker; - - public MessageBrokerManager(final MessageBrokerConfiguration configuration) { - this.configuration = configuration; - } - - public void startMessageBroker(boolean useJmx) throws Exception { - broker = new BrokerService(); - broker.setBrokerName(configuration.getBrokerName()); - if (configuration.getDataStoreDirectory() != null) { - broker.setDataDirectory(configuration.getDataStoreDirectory()); - } - broker.setUseJmx(useJmx); - TransportConnector transportConnector = new TransportConnector(); - transportConnector.setUri(new URI(configuration.getBrokerUrl())); - broker.addConnector(transportConnector); - broker.start(); - } - - public void stopMessageBroker() throws Exception { - if (broker != null) { - broker.stop(); - } - } - - public void createAndStartConnection() throws JMSException { - connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER, - ActiveMQConnection.DEFAULT_PASSWORD, configuration.getBrokerUrl()); - connection = connectionFactory.createConnection(); - connection.start(); - } - - public void closeConnections() throws JMSException { - if (connection != null) { - connection.close(); - } - } - - public Session createSession(boolean transacted) throws JMSException { - session = connection.createSession(transacted, Session.AUTO_ACKNOWLEDGE); - return session; - } - - public void closeSessions() throws JMSException { - if (session != null) { - session.close(); - } - } - - public Session getSession() { - return session; - } -} diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationFactoryImpl.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationFactoryImpl.java index 0fb74745e07..50837ce3b7b 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationFactoryImpl.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationFactoryImpl.java @@ -4,10 +4,9 @@ import com.sap.sailing.server.replication.ReplicationFactory; import com.sap.sailing.server.replication.ReplicationMasterDescriptor; public class ReplicationFactoryImpl implements ReplicationFactory { - @Override - public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, int servletPort, int jmsPort) { - return new ReplicationMasterDescriptorImpl(hostname, servletPort, jmsPort); + public ReplicationMasterDescriptor createReplicationMasterDescriptor(String hostname, String exchangeName, int servletPort, int messagingPort) { + return new ReplicationMasterDescriptorImpl(hostname, exchangeName, servletPort, messagingPort); } } diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java index 3086d48ea40..b4ef987e6af 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationMasterDescriptorImpl.java @@ -1,32 +1,27 @@ package com.sap.sailing.server.replication.impl; -import java.net.InetAddress; +import java.io.IOException; import java.net.MalformedURLException; import java.net.URL; -import java.net.UnknownHostException; - -import javax.jms.Connection; -import javax.jms.JMSException; -import javax.jms.Session; -import javax.jms.Topic; -import javax.jms.TopicSubscriber; - -import org.apache.activemq.ActiveMQConnection; -import org.apache.activemq.ActiveMQConnectionFactory; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.Connection; +import com.rabbitmq.client.ConnectionFactory; +import com.rabbitmq.client.QueueingConsumer; import com.sap.sailing.server.replication.ReplicationMasterDescriptor; -import com.sap.sailing.server.replication.ReplicationService; public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescriptor { private static final String REPLICATION_SERVLET = "/replication/replication"; private final String hostname; + private final String exchangeName; private final int servletPort; private final int jmsPort; - public ReplicationMasterDescriptorImpl(String hostname, int servletPort, int jmsPort) { + public ReplicationMasterDescriptorImpl(String hostname, String exchangeName, int servletPort, int jmsPort) { this.hostname = hostname; this.servletPort = servletPort; this.jmsPort = jmsPort; + this.exchangeName = exchangeName; } @Override @@ -42,15 +37,17 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip } @Override - public TopicSubscriber getTopicSubscriber(String clientID) throws JMSException, UnknownHostException { - ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(ActiveMQConnection.DEFAULT_USER, - ActiveMQConnection.DEFAULT_PASSWORD, "tcp://" + hostname + ":" + jmsPort); - 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()); + public QueueingConsumer getConsumer() throws IOException { + ConnectionFactory connectionFactory = new ConnectionFactory(); + connectionFactory.setHost(getHostname()); + Connection connection = connectionFactory.newConnection(); + Channel channel = connection.createChannel(); + channel.exchangeDeclare(exchangeName, "fanout"); + QueueingConsumer consumer = new QueueingConsumer(channel); + String queueName = channel.queueDeclare().getQueue(); + channel.queueBind(queueName, exchangeName, ""); + channel.basicConsume(queueName, consumer); + return consumer; } @Override diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java index 57825a57263..a4ac435d8e0 100644 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java @@ -11,16 +11,12 @@ import java.net.URLConnection; import java.util.HashMap; import java.util.Map; -import javax.jms.BytesMessage; -import javax.jms.DeliveryMode; -import javax.jms.JMSException; -import javax.jms.MessageProducer; -import javax.jms.Session; -import javax.jms.Topic; -import javax.jms.TopicSubscriber; - import org.osgi.util.tracker.ServiceTracker; +import com.rabbitmq.client.AMQP.Exchange; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.ConnectionFactory; +import com.rabbitmq.client.QueueingConsumer; import com.sap.sailing.server.OperationExecutionListener; import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.RacingEventServiceOperation; @@ -31,8 +27,8 @@ import com.sap.sailing.server.replication.ReplicationService; /** * Can observe a {@link RacingEventService} for the operations it performs that require replication. Only observes as * long as there are replicas registered. If the last replica is de-registered, the service stops observing the - * {@link RacingEventService}. Operations received that require replication are broadcast to the - * {@link #getReplicationTopic() replication topic}. + * {@link RacingEventService}. Operations received that require replication are sent to the {@link Exchange} to which + * replica queues can bind. The exchange name is provided to this service during construction. *

* * This service object {@link RacingEventService#addOperationExecutionListener(OperationExecutionListener) registers} as @@ -46,12 +42,6 @@ import com.sap.sailing.server.replication.ReplicationService; public class ReplicationServiceImpl implements ReplicationService, OperationExecutionListener, HasRacingEventService { private final ReplicationInstancesManager replicationInstancesManager; - private final MessageBrokerManager messageBrokerManager; - - private MessageProducer messageProducer; - - private Topic replicationTopic; - private ServiceTracker racingEventServiceTracker; private final RacingEventService localService; @@ -66,29 +56,43 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec */ private final Map replicaUUIDs; - public ReplicationServiceImpl(final ReplicationInstancesManager replicationInstancesManager, - final MessageBrokerManager messageBrokerManager) throws Exception { + private final Channel channel; + + /** + * The name of the RabbitMQ exchange to which this replication service sends its replication operations in + * serialized form. Clients need to know this name to be able to bind their queues to the exchange. + */ + private final String exchangeName; + + public ReplicationServiceImpl(String exchangeName, final ReplicationInstancesManager replicationInstancesManager) throws IOException { this.replicationInstancesManager = replicationInstancesManager; replicaUUIDs = new HashMap(); - this.messageBrokerManager = messageBrokerManager; racingEventServiceTracker = new ServiceTracker( Activator.getDefaultContext(), RacingEventService.class.getName(), null); racingEventServiceTracker.open(); localService = null; + this.exchangeName = exchangeName; + channel = createChannel(exchangeName); } /** - * Like {@link #ReplicationServiceImpl(ReplicationInstancesManager, MessageBrokerManager)}, only that instead of using + * Like {@link #ReplicationServiceImpl(String, ReplicationInstancesManager)}, only that instead of using * an OSGi service tracker to discover the {@link RacingEventService}, the service to replicate is "injected" here. + * @param exchangeName the name of the exchange to which replicas can bind */ - public ReplicationServiceImpl(final ReplicationInstancesManager replicationInstancesManager, - final MessageBrokerManager messageBrokerManager, RacingEventService localService) { + public ReplicationServiceImpl(String exchangeName, + final ReplicationInstancesManager replicationInstancesManager, RacingEventService localService) throws IOException { this.replicationInstancesManager = replicationInstancesManager; replicaUUIDs = new HashMap(); - this.messageBrokerManager = messageBrokerManager; this.localService = localService; + this.exchangeName = exchangeName; + channel = createChannel(exchangeName); } + private Channel createChannel(String exchangeName) throws IOException { + return new ConnectionFactory().newConnection().createChannel(); + } + @Override public RacingEventService getRacingEventService() { RacingEventService result; @@ -101,13 +105,9 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec } @Override - public void registerReplica(ReplicaDescriptor replica) throws JMSException { - Topic topic = getReplicationTopic(); - assert topic != null; + public void registerReplica(ReplicaDescriptor replica) { if (!replicationInstancesManager.hasReplicas()) { addAsListenerToRacingEventService(); - messageBrokerManager.createAndStartConnection(); - messageBrokerManager.createSession(/* transacted */ false); } replicationInstancesManager.registerReplica(replica); } @@ -117,53 +117,27 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec } @Override - public void unregisterReplica(ReplicaDescriptor replica) throws JMSException { + public void unregisterReplica(ReplicaDescriptor replica) throws IOException { replicationInstancesManager.unregisterReplica(replica); if (!replicationInstancesManager.hasReplicas()) { removeAsListenerFromRacingEventService(); - messageBrokerManager.closeSessions(); - messageBrokerManager.closeConnections(); } } private void removeAsListenerFromRacingEventService() { getRacingEventService().removeOperationExecutionListener(this); - messageProducer = null; } - private Topic getReplicationTopic() throws JMSException{ - if (replicationTopic == null) { - Session session = messageBrokerManager.getSession(); - if (session == null) { - session = messageBrokerManager.createSession(true); - } - replicationTopic = session.createTopic(SAILING_SERVER_REPLICATION_TOPIC); - } - return replicationTopic; - } - private void broadcastOperation(RacingEventServiceOperation operation) throws Exception { - Topic topic = getReplicationTopic(); - Session session = messageBrokerManager.getSession(); - getMessageProducer(topic).setDeliveryMode(DeliveryMode.NON_PERSISTENT); - BytesMessage operationAsMessage = session.createBytesMessage(); // serialize operation into message ByteArrayOutputStream bos = new ByteArrayOutputStream(); ObjectOutputStream oos = new ObjectOutputStream(bos); oos.writeObject(operation); oos.close(); - operationAsMessage.writeBytes(bos.toByteArray()); - messageProducer.send(operationAsMessage); + channel.basicPublish(exchangeName, /* routingKey */ "", /* properties */ null, bos.toByteArray()); replicationInstancesManager.log(operation); } - private MessageProducer getMessageProducer(Topic topic) throws JMSException { - if (messageProducer == null) { - messageProducer = messageBrokerManager.getSession().createProducer(topic); - } - return messageProducer; - } - @Override public Iterable getReplicaInfo() { return replicationInstancesManager.getReplicaDescriptors(); @@ -175,13 +149,12 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec } @Override - public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException, JMSException { + public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException { replicatingFromMaster = master; - String uuid = registerReplicaWithMaster(master); - TopicSubscriber replicationSubscription = master.getTopicSubscriber(uuid); + registerReplicaWithMaster(master); + QueueingConsumer consumer = master.getConsumer(); URL initialLoadURL = master.getInitialLoadURL(); - final Replicator replicator = new Replicator(master, this, /* startSuspended */ true); - replicationSubscription.setMessageListener(replicator); + final Replicator replicator = new Replicator(master, this, /* startSuspended */ true, consumer); InputStream is = initialLoadURL.openStream(); ObjectInputStream ois = new ObjectInputStream(is) { @Override diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServlet.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServlet.java index d8113ae4654..4612618a6e6 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServlet.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/ReplicationServlet.java @@ -7,7 +7,6 @@ import java.net.UnknownHostException; import java.util.Arrays; import java.util.logging.Logger; -import javax.jms.JMSException; import javax.servlet.ServletException; import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; @@ -17,8 +16,8 @@ import org.osgi.util.tracker.ServiceTracker; import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.SailingServerHttpServlet; -import com.sap.sailing.server.replication.ReplicationService; import com.sap.sailing.server.replication.ReplicaDescriptor; +import com.sap.sailing.server.replication.ReplicationService; /** * As the response to any type of GET request, sends a serialized copy of the {@link RacingEventService} to @@ -60,11 +59,7 @@ public class ReplicationServlet extends SailingServerHttpServlet { String action = req.getParameter(ACTION); switch (Action.valueOf(action)) { case REGISTER: - try { - registerClientWithReplicationService(req, resp); - } catch (JMSException e) { - resp.sendError(HttpServletResponse.SC_INTERNAL_SERVER_ERROR, e.getMessage()); - } + registerClientWithReplicationService(req, resp); break; case INITIAL_LOAD: ObjectOutputStream oos = new ObjectOutputStream(resp.getOutputStream()); @@ -84,7 +79,7 @@ public class ReplicationServlet extends SailingServerHttpServlet { } private void registerClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp) - throws JMSException, IOException { + throws IOException { ReplicaDescriptor replica = getReplicaDescriptor(req); getReplicationService().registerReplica(replica); resp.setContentType("text/plain"); diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java index 48f89d11783..294d93be3a6 100755 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java +++ b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/Replicator.java @@ -6,12 +6,12 @@ import java.io.ObjectInputStream; import java.util.ArrayList; import java.util.Iterator; import java.util.List; +import java.util.logging.Logger; -import javax.jms.BytesMessage; -import javax.jms.JMSException; -import javax.jms.Message; -import javax.jms.MessageListener; - +import com.rabbitmq.client.ConsumerCancelledException; +import com.rabbitmq.client.QueueingConsumer; +import com.rabbitmq.client.QueueingConsumer.Delivery; +import com.rabbitmq.client.ShutdownSignalException; import com.sap.sailing.domain.base.DomainFactory; import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.RacingEventServiceOperation; @@ -30,10 +30,13 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor; * @author Axel Uhl (d043530) * */ -public class Replicator implements MessageListener { +public class Replicator implements Runnable { + private final static Logger logger = Logger.getLogger(Replicator.class.getName()); + private final ReplicationMasterDescriptor master; private final HasRacingEventService racingEventServiceTracker; private final List> queue; + private final QueueingConsumer consumer; /** * If the replicator is suspended, messages received are queued. @@ -47,46 +50,52 @@ public class Replicator implements MessageListener { * descriptor of the master server from which this replicator receives messages * @param racingEventServiceTracker * OSGi service tracker for the replica to which to apply the messages received + * @param consumer the RabbitMQ consumer from which to load messages */ - public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker) { - this(master, racingEventServiceTracker, /* startSuspended */ false); + public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, QueueingConsumer consumer) { + this(master, racingEventServiceTracker, /* startSuspended */ false, consumer); } - public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, boolean startSuspended) { + public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, boolean startSuspended, QueueingConsumer consumer) { this.queue = new ArrayList>(); this.master = master; this.racingEventServiceTracker = racingEventServiceTracker; this.suspended = startSuspended; + this.consumer = consumer; + } + + /** + * Starts fetching messages from the {@link #consumer}. After receiving a single message, assumes it's a serialized + * {@link RacingEventServiceOperation}, and applies it to the {@link RacingEventService} which is obtained from the + * service tracker passed to this replicator at construction time. + */ + @Override + public void run() { + ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader(); + while (true) { + try { + Delivery delivery = consumer.nextDelivery(); + byte[] bytesFromMessage = delivery.getBody(); + // Set this object's class's class loader as context for de-serialization so that all exported classes + // of all required bundles/packages can be deserialized at least + Thread.currentThread().setContextClassLoader(getClass().getClassLoader()); + ObjectInputStream ois = DomainFactory.INSTANCE.createObjectInputStreamResolvingAgainstThisFactory( + new ByteArrayInputStream(bytesFromMessage)); + RacingEventServiceOperation operation = (RacingEventServiceOperation) ois.readObject(); + applyOrQueue(operation); + } catch (ShutdownSignalException | ConsumerCancelledException | InterruptedException | IOException | ClassNotFoundException e) { + logger.info("Exception while processing replica: "+e.getMessage()); + logger.throwing(Replicator.class.getName(), "run", e); + } finally { + Thread.currentThread().setContextClassLoader(oldClassLoader); + } + } } public synchronized boolean isQueueEmpty() { return queue.isEmpty(); } - /** - * Receives a single message, assuming it's a {@link RacingEventServiceOperation}, and applies it to the - * {@link RacingEventService} which is obtained from the service tracker passed to this replicator at construction - * time. - */ - @Override - public synchronized void onMessage(Message m) { - ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader(); - try { - byte[] bytesFromMessage = getBytes((BytesMessage) m); - // Set this object's class's class loader as context for de-serialization so that all exported classes - // of all required bundles/packages can be deserialized at least - Thread.currentThread().setContextClassLoader(getClass().getClassLoader()); - ObjectInputStream ois = DomainFactory.INSTANCE.createObjectInputStreamResolvingAgainstThisFactory( - new ByteArrayInputStream(bytesFromMessage)); - RacingEventServiceOperation operation = (RacingEventServiceOperation) ois.readObject(); - applyOrQueue(operation); - } catch (IOException | ClassNotFoundException | JMSException e) { - throw new RuntimeException(e); - } finally { - Thread.currentThread().setContextClassLoader(oldClassLoader); - } - } - /** * If the replicator is currently {@link #suspended}, the operation is queued, otherwise immediately applied to * the receiving replica. @@ -134,12 +143,6 @@ public class Replicator implements MessageListener { return suspended; } - private byte[] getBytes(BytesMessage m) throws JMSException { - byte[] buf = new byte[(int) m.getBodyLength()]; - m.readBytes(buf); - return buf; - } - @Override public String toString() { return "Replicator for master "+master+", queue size: "+queue.size(); diff --git a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/SampleMessageConsumer.java b/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/SampleMessageConsumer.java deleted file mode 100644 index da2c4ad8c72..00000000000 --- a/java/com.sap.sailing.server.replication/src/com/sap/sailing/server/replication/impl/SampleMessageConsumer.java +++ /dev/null @@ -1,27 +0,0 @@ -package com.sap.sailing.server.replication.impl; - -import javax.jms.ExceptionListener; -import javax.jms.JMSException; -import javax.jms.Message; -import javax.jms.MessageListener; -import javax.jms.TextMessage; - -public class SampleMessageConsumer implements MessageListener, ExceptionListener { - - synchronized public void onException(JMSException ex) { - System.out.println("JMS Exception occured: " + ex.getMessage()); - } - - public void onMessage(Message message) { - if (message instanceof TextMessage) { - TextMessage textMessage = (TextMessage) message; - try { - System.out.println("Received message: " + textMessage.getText()); - } catch (JMSException ex) { - System.out.println("Error reading message: " + ex); - } - } else { - System.out.println("Received: " + message); - } - } -}