fixing bug 723 by queuing Operations until initial load has completed

This commit is contained in:
Axel Uhl committed 2012-06-07 18:20:51 +02:00
1 parent 348c66e5cf
commit f4928da853
5 files changed
+113 -16

No files matched your search

@@ -52,7 +52,7 @@ public abstract class AbstractServerReplicationTest {
*/
@Before
public void setUp() throws Exception {
Pair<ReplicationService, ReplicationMasterDescriptor> result = basicSetUp(true, /* master=null means create a new one */ null,
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> 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 <code>null</code>, the value will be used for {@link #replica}; otherwise, a new racing event
* service will be created as replica
*/
protected Pair<ReplicationService, ReplicationMasterDescriptor> basicSetUp(
protected Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> 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<ReplicationService, ReplicationMasterDescriptor> result = new Pair<>(replicaReplicator, masterDescriptor);
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> 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);
@@ -99,7 +99,7 @@ public class DelayedLeaderboardCorrectionsReplicationTest extends AbstractServer
masterLeaderboardReloaded.getRaceColumnByName(Q2), MillisecondsTimePoint.now()));
// replicate the re-loaded environment
Pair<ReplicationService, ReplicationMasterDescriptor> descriptors = basicSetUp(/* dropDB */ false, master, replica);
Pair<ReplicationServiceTestImpl, ReplicationMasterDescriptor> descriptors = basicSetUp(/* dropDB */ false, master, replica);
replicaReplicator = descriptors.getA();
masterDescriptor = descriptors.getB();
replicaReplicator.startToReplicateFrom(masterDescriptor);
@@ -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
}
/**
@@ -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.<p>
*
* 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<RacingEventServiceOperation<?>> 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<RacingEventServiceOperation<?>>();
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,7 +79,7 @@ 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 {
@@ -55,6 +87,53 @@ public class Replicator implements MessageListener {
}
}
/**
* If the replicator is currently {@link #suspended}, the <code>operation</code> 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<RacingEventServiceOperation<?>> 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();
}
}
@@ -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<FlexibleLeaderboard> {
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<Flex
@Override
public FlexibleLeaderboard internalApplyTo(RacingEventService toState) {
return toState.addFlexibleLeaderboard(getLeaderboardName(), discardThresholds);
FlexibleLeaderboard result = null;
if (toState.getLeaderboardByName(getLeaderboardName()) == null) {
result = toState.addFlexibleLeaderboard(getLeaderboardName(), discardThresholds);
} else {
logger.warning("Cannot replicate creation of flexible leaderboard "+getLeaderboardName()+" because it already exists in the replica");
}
return result;
}
@Override