Merge branch 'master' into raceboard-simulation

This commit is contained in:
Axel Uhl committed 2015-03-30 22:05:27 +02:00
commit 34169f8e77
6 files changed
+33 -18

No files matched your search

@@ -190,7 +190,7 @@ public class InitialLoadReplicationObjectIdentityTest extends AbstractServerRepl
replicationDescriptorPair.getA().initialLoad();
replicator.setSuspended(false);
synchronized (replicator) {
while (!replicator.isQueueEmpty()) {
while (!replicator.isQueueEmptyOrStopped()) {
replicator.wait();
}
}
@@ -60,7 +60,7 @@ public class PrematureOperationReceiptTest extends AbstractServerReplicationTest
}
replicator.setSuspended(false);
synchronized (replicator) {
while (!replicator.isQueueEmpty()) {
while (!replicator.isQueueEmptyOrStopped()) {
replicator.wait();
}
}
@@ -383,6 +383,8 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
.withInitial(() -> false);
private final Set<ClassLoader> masterDataClassLoaders = new HashSet<ClassLoader>();
private final JoinedClassLoader joinedClassLoader;
/**
* Constructs a {@link DomainFactory base domain factory} that uses this object's {@link #competitorStore competitor
@@ -461,6 +463,7 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
TypeBasedServiceFinderFactory serviceFinderFactory) {
logger.info("Created " + this);
this.masterDataClassLoaders.add(this.getClass().getClassLoader());
joinedClassLoader = new JoinedClassLoader(masterDataClassLoaders);
this.operationsSentToMasterForReplication = new HashSet<>();
if (windStore == null) {
try {
@@ -2307,7 +2310,7 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
@Override
public ClassLoader getDeserializationClassLoader() {
return getCombinedMasterDataClassLoader();
return joinedClassLoader;
}
@Override
@@ -179,10 +179,13 @@ public abstract class AbstractServerReplicationTestSetUp<ReplicableInterface ext
* original exception propagate */
logger.log(Level.SEVERE, "Exception trying to connect to initial load test servlet to STOP it", ex);
}
synchronized (replicaReplicator.getReplicator()) {
while (!replicaReplicator.getReplicator().isQueueEmpty()) {
logger.info("Waiting for replication queue to drain...");
replicaReplicator.getReplicator().wait();
final ReplicationReceiver replicaReplicatorReplicator = replicaReplicator.getReplicator();
if (replicaReplicatorReplicator != null) {
synchronized (replicaReplicatorReplicator) {
while (!replicaReplicatorReplicator.isQueueEmptyOrStopped()) {
logger.info("Waiting for replication queue to drain...");
replicaReplicatorReplicator.wait();
}
}
}
logger.info("Replication queue has been drained...");
@@ -334,7 +337,7 @@ public abstract class AbstractServerReplicationTestSetUp<ReplicableInterface ext
public void waitUntilQueueIsEmpty() throws InterruptedException, IllegalAccessException {
synchronized (getReplicator()) {
while (!getReplicator().isQueueEmpty()) {
while (!getReplicator().isQueueEmptyOrStopped()) {
getReplicator().wait();
}
}
@@ -147,7 +147,6 @@ public class ReplicationReceiver implements Runnable {
*/
@Override
public void run() {
ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader();
long messageCount = 0;
long operationCount = 0;
final boolean logsFine = logger.isLoggable(Level.FINE);
@@ -156,6 +155,11 @@ public class ReplicationReceiver implements Runnable {
Delivery delivery = consumer.nextDelivery();
messageCount++;
if (_queue != null) {
synchronized (this) {
if (getInboundMessageQueue().isEmpty()) {
notifyAll(); // wake up anyone waiting for isQueueEmpty()
}
}
if (logsFine || messageCount % 10l == 0) {
try {
logger.log(messageCount%10l==0 ? Level.INFO : Level.FINE,
@@ -174,7 +178,6 @@ public class ReplicationReceiver implements Runnable {
String replicableIdAsString = new DataInputStream(uncompressedInputStream).readUTF();
Replicable<?, ?> replicable = replicableProvider.getReplicable(replicableIdAsString, /* wait */ false);
if (replicable != null) {
Thread.currentThread().setContextClassLoader(replicable.getClass().getClassLoader());
ObjectInputStream ois = new ObjectInputStream(uncompressedInputStream); // no special stream required; only reading a generic byte[]
int operationsInMessage = 0;
try {
@@ -248,11 +251,13 @@ public class ReplicationReceiver implements Runnable {
} catch (Exception e) {
logger.info("Exception while processing replica: "+e.getMessage());
logger.log(Level.SEVERE, "run", e);
} finally {
Thread.currentThread().setContextClassLoader(oldClassLoader);
}
}
logger.info("Stopped replicator thread. This server will no longer receive events from a master.");
synchronized (this) {
stopped = true;
notifyAll();
}
}
/**
@@ -313,9 +318,6 @@ public class ReplicationReceiver implements Runnable {
queue = new ArrayList<>();
queueByReplicableIdAsString.put(replicableIdAsString, queue);
}
if (queue.isEmpty()) {
notifyAll();
}
queue.add(new Pair<String, OperationWithResult<?, ?>>(replicable.getId().toString(), operation));
assert !queue.isEmpty();
}
@@ -359,6 +361,7 @@ public class ReplicationReceiver implements Runnable {
logger.info("Signaled Replicator thread to stop asap.");
stopped = true;
master.stopConnection();
notifyAll(); // notify those waiting for stopped
}
public synchronized boolean isBeingStopped() {
@@ -374,8 +377,9 @@ public class ReplicationReceiver implements Runnable {
/**
* @return <code>true</code> if all queues for all replicables are empty
*/
public boolean isQueueEmpty() throws IllegalAccessException {
return (_queue == null || getInboundMessageQueue().isEmpty()) && !queueByReplicableIdAsString.values().stream().anyMatch(q->!q.isEmpty());
public boolean isQueueEmptyOrStopped() throws IllegalAccessException {
return isBeingStopped() ||
(_queue == null || getInboundMessageQueue().isEmpty()) && !queueByReplicableIdAsString.values().stream().anyMatch(q->!q.isEmpty());
}
}
@@ -246,7 +246,12 @@ public class ReplicationServiceImpl implements ReplicationService {
replicator = null;
serverUUID = UUID.randomUUID();
logger.info("Setting " + serverUUID.toString() + " as unique replication identifier.");
}
@Override
protected void finalize() {
logger.info("terminating timer "+timer);
timer.cancel();
}
protected ServiceTracker<Replicable<?, ?>, Replicable<?, ?>> getReplicableTracker() {