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 4bdb809e28c..e6b7999d8c7 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 @@ -52,7 +52,7 @@ public abstract class AbstractServerReplicationTest { */ @Before public void setUp() throws Exception { - Pair result = basicSetUp(true, /* master=null means create a new one */ null, + Pair result = basicSetUp(true, /* master=null means create a new one */ null, /* replica=null means create a new one */ null); result.getA().startToReplicateFrom(result.getB()); } @@ -67,7 +67,7 @@ public abstract class AbstractServerReplicationTest { * if not null, the value will be used for {@link #replica}; otherwise, a new racing event * service will be created as replica */ - protected Pair basicSetUp( + protected Pair basicSetUp( boolean dropDB, RacingEventServiceImpl master, RacingEventServiceImpl replica) throws FileNotFoundException, Exception, JMSException, UnknownHostException { final MongoDBService mongoDBService = MongoDBService.INSTANCE; @@ -132,9 +132,9 @@ public abstract class AbstractServerReplicationTest { return null; } }; - ReplicationService replicaReplicator = new ReplicationServiceTestImpl(resolveAgainst, rim, brokerMgr, + ReplicationServiceTestImpl replicaReplicator = new ReplicationServiceTestImpl(resolveAgainst, rim, brokerMgr, replicaDescriptor, this.replica, this.master, masterReplicator); - Pair result = new Pair<>(replicaReplicator, masterDescriptor); + Pair result = new Pair<>(replicaReplicator, masterDescriptor); return result; } @@ -147,7 +147,7 @@ public abstract class AbstractServerReplicationTest { Activator.removeTemporaryTestBrokerPersistenceDirectory(brokerPersistenceDir); } - private static class ReplicationServiceTestImpl extends ReplicationServiceImpl { + static class ReplicationServiceTestImpl extends ReplicationServiceImpl { private final DomainFactory resolveAgainst; private final RacingEventService master; private final ReplicaDescriptor replicaDescriptor; @@ -169,11 +169,19 @@ public abstract class AbstractServerReplicationTest { @Override public void startToReplicateFrom(ReplicationMasterDescriptor master) throws IOException, ClassNotFoundException, JMSException { + Replicator replicator = startToReplicateFromButDontYetFetchInitialLoad(master, /* startReplicatorSuspended */ true); + initialLoad(); + replicator.setSuspended(false); // resume after initial load + } + + protected Replicator startToReplicateFromButDontYetFetchInitialLoad(ReplicationMasterDescriptor master, boolean startReplicatorSuspended) + throws JMSException, UnknownHostException { masterReplicationService.registerReplica(replicaDescriptor); registerReplicaUuidForMaster(replicaDescriptor.getUuid().toString(), master); TopicSubscriber replicationSubscription = master.getTopicSubscriber(replicaDescriptor.getUuid().toString()); - replicationSubscription.setMessageListener(new Replicator(master, this)); - initialLoad(); + final Replicator replicator = new Replicator(master, this, startReplicatorSuspended); + replicationSubscription.setMessageListener(replicator); + return replicator; } /** @@ -181,7 +189,7 @@ public abstract class AbstractServerReplicationTest { * {@link RacingEventServiceImpl#serializeForInitialReplication(ObjectOutputStream)} and * {@link RacingEventServiceImpl#initiallyFillFrom(ObjectInputStream)} through a piped input/output stream. */ - private void initialLoad() throws IOException, ClassNotFoundException { + protected void initialLoad() throws IOException, ClassNotFoundException { PipedOutputStream pos = new PipedOutputStream(); PipedInputStream pis = new PipedInputStream(pos); final ObjectOutputStream oos = new ObjectOutputStream(pos); 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 8a44bb81353..aad57788d01 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 @@ -99,7 +99,7 @@ public class DelayedLeaderboardCorrectionsReplicationTest extends AbstractServer masterLeaderboardReloaded.getRaceColumnByName(Q2), MillisecondsTimePoint.now())); // replicate the re-loaded environment - Pair descriptors = basicSetUp(/* dropDB */ false, master, replica); + Pair descriptors = basicSetUp(/* dropDB */ false, master, replica); replicaReplicator = descriptors.getA(); masterDescriptor = descriptors.getB(); replicaReplicator.startToReplicateFrom(masterDescriptor); 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 243a7f2d273..d583e60e07e 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 @@ -179,11 +179,12 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec String uuid = registerReplicaWithMaster(master); TopicSubscriber replicationSubscription = master.getTopicSubscriber(uuid); URL initialLoadURL = master.getInitialLoadURL(); - // TODO bug 723 hold back Replicator from performing operations received until initial load has completed - replicationSubscription.setMessageListener(new Replicator(master, this)); + final Replicator replicator = new Replicator(master, this, /* startSuspended */ true); + replicationSubscription.setMessageListener(replicator); InputStream is = initialLoadURL.openStream(); ObjectInputStream ois = new ObjectInputStream(is); getRacingEventService().initiallyFillFrom(ois); + replicator.setSuspended(false); // apply queued operations } /** 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 719de502002..48f89d11783 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 @@ -3,6 +3,9 @@ package com.sap.sailing.server.replication.impl; import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.ObjectInputStream; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; import javax.jms.BytesMessage; import javax.jms.JMSException; @@ -17,7 +20,12 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor; /** * Receives {@link RacingEventServiceOperation}s through JMS and * {@link RacingEventService#apply(RacingEventServiceOperation) applies} them to the {@link RacingEventService} passed - * to this replicator at construction. + * to this replicator at construction. When started in suspended mode, messages received will be turned into + * {@link RacingEventServiceOperation}s and then queued until {@link #setSuspended(boolean) setSuspended(false)} is invoked + * which applies all queued operations before applying the ones received later.

+ * + * The receiver takes care of synchronizing receiving, suspending/resuming and queuing. Waiters are notified + * whenever the result of {@link #isQueueEmpty} changes. * * @author Axel Uhl (d043530) * @@ -25,10 +33,34 @@ import com.sap.sailing.server.replication.ReplicationMasterDescriptor; public class Replicator implements MessageListener { private final ReplicationMasterDescriptor master; private final HasRacingEventService racingEventServiceTracker; + private final List> queue; + /** + * If the replicator is suspended, messages received are queued. + */ + private boolean suspended; + + /** + * Starts the replicator immediately, not holding back messages received but forwarding them directly. + * + * @param master + * 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 + */ public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker) { + this(master, racingEventServiceTracker, /* startSuspended */ false); + } + + public Replicator(ReplicationMasterDescriptor master, HasRacingEventService racingEventServiceTracker, boolean startSuspended) { + this.queue = new ArrayList>(); this.master = master; this.racingEventServiceTracker = racingEventServiceTracker; + this.suspended = startSuspended; + } + + public synchronized boolean isQueueEmpty() { + return queue.isEmpty(); } /** @@ -37,7 +69,7 @@ public class Replicator implements MessageListener { * time. */ @Override - public void onMessage(Message m) { + public synchronized void onMessage(Message m) { ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader(); try { byte[] bytesFromMessage = getBytes((BytesMessage) m); @@ -47,14 +79,61 @@ public class Replicator implements MessageListener { ObjectInputStream ois = DomainFactory.INSTANCE.createObjectInputStreamResolvingAgainstThisFactory( new ByteArrayInputStream(bytesFromMessage)); RacingEventServiceOperation operation = (RacingEventServiceOperation) ois.readObject(); - racingEventServiceTracker.getRacingEventService().apply(operation); + 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. + */ + private synchronized void applyOrQueue(RacingEventServiceOperation operation) { + if (suspended) { + queue(operation); + } else { + apply(operation); + } + } + + private synchronized void apply(RacingEventServiceOperation operation) { + racingEventServiceTracker.getRacingEventService().apply(operation); + } + private synchronized void queue(RacingEventServiceOperation operation) { + if (queue.isEmpty()) { + notifyAll(); + } + queue.add(operation); + assert !queue.isEmpty(); + } + + public synchronized void setSuspended(final boolean suspended) { + if (this.suspended != suspended) { + this.suspended = suspended; + if (!this.suspended) { + applyQueue(); + } + } + } + + private synchronized void applyQueue() { + for (Iterator> i=queue.iterator(); i.hasNext(); ) { + RacingEventServiceOperation operation = i.next(); + i.remove(); + apply(operation); + } + assert queue.isEmpty(); + notifyAll(); + } + + public synchronized boolean isSuspended() { + return suspended; + } + private byte[] getBytes(BytesMessage m) throws JMSException { byte[] buf = new byte[(int) m.getBodyLength()]; m.readBytes(buf); @@ -63,7 +142,7 @@ public class Replicator implements MessageListener { @Override public String toString() { - return "Replicator for master "+master; + return "Replicator for master "+master+", queue size: "+queue.size(); } } diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/CreateFlexibleLeaderboard.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/CreateFlexibleLeaderboard.java index f250ce5408f..008fdc55171 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/CreateFlexibleLeaderboard.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/operationaltransformation/CreateFlexibleLeaderboard.java @@ -1,10 +1,13 @@ package com.sap.sailing.server.operationaltransformation; +import java.util.logging.Logger; + import com.sap.sailing.domain.leaderboard.FlexibleLeaderboard; import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.RacingEventServiceOperation; public class CreateFlexibleLeaderboard extends AbstractLeaderboardOperation { + private static final Logger logger = Logger.getLogger(CreateFlexibleLeaderboard.class.getName()); private static final long serialVersionUID = 891352705068098580L; private final int[] discardThresholds; @@ -15,7 +18,13 @@ public class CreateFlexibleLeaderboard extends AbstractLeaderboardOperation