Make sure that replicator thread is interrupted when application

requests shutdown of the connection. Unfortunately the RabbitMQ API does
not support closing a connection while unblocking the nextDelivery call.
There fore we need to interrupt the thread manually.
This commit is contained in:
Simon Pamies committed 2013-06-20 12:36:54 +02:00
1 parent 4956ee0a0c
commit 671e68e53f
3 files changed
+43 -19

No files matched your search

@@ -85,13 +85,17 @@ public class ReplicationPanel extends FlowPanel {
}
private void stopReplication() {
stopReplicationButton.setEnabled(false);
sailingService.stopReplicatingFromMaster(new AsyncCallback<Void>() {
@Override
public void onFailure(Throwable caught) {
errorReporter.reportError(caught.getMessage());
stopReplicationButton.setEnabled(true);
}
@Override
public void onSuccess(Void result) {
addButton.setEnabled(true);
stopReplicationButton.setEnabled(false);
updateReplicaList();
}
});
@@ -105,12 +109,15 @@ public class ReplicationPanel extends FlowPanel {
registeredMasters.removeRow(0);
registeredMasters.insertRow(0);
registeredMasters.setWidget(0, 0, new Label(stringMessages.loading()));
addButton.setEnabled(false);
stopReplicationButton.setEnabled(false);
sailingService.startReplicatingFromMaster(masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getA(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getB(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getC(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getB(), new AsyncCallback<Void>() {
@Override
public void onFailure(Throwable e) {
addButton.setEnabled(true);
errorReporter.reportError(stringMessages.errorStartingReplication(
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getA(),
masterNameAndExchangeNameAndMessagingPortNumberAndServletPortNumber.getA().getB(), e.getMessage()));
@@ -118,8 +125,9 @@ public class ReplicationPanel extends FlowPanel {
@Override
public void onSuccess(Void arg0) {
updateReplicaList();
addButton.setEnabled(false);
stopReplicationButton.setEnabled(true);
updateReplicaList();
}
});
}
@@ -70,6 +70,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
private final String exchangeName;
private Replicator replicator;
private Thread replicatorThread;
public ReplicationServiceImpl(String exchangeName, final ReplicationInstancesManager replicationInstancesManager) throws IOException {
this.replicationInstancesManager = replicationInstancesManager;
@@ -203,7 +204,8 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
logger.info("Initial load URL is "+initialLoadURL);
replicator = new Replicator(master, this, /* startSuspended */ true, consumer);
// start receiving messages already now, but start in suspended mode
new Thread(replicator, "Replicator receiving from "+master.getHostname()+"/"+master.getExchangeName()).start();
replicatorThread = new Thread(replicator, "Replicator receiving from "+master.getHostname()+"/"+master.getExchangeName());
replicatorThread.start();
logger.info("Started replicator thread");
InputStream is = initialLoadURL.openStream();
final RacingEventService racingEventService = getRacingEventService();
@@ -234,19 +236,25 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
return replicaUUID;
}
protected void deregisterReplicaWithMaster(ReplicationMasterDescriptor master) throws IOException {
URL replicationDeRegistrationRequestURL = master.getReplicationDeRegistrationRequestURL();
final URLConnection deregistrationRequestConnection = replicationDeRegistrationRequestURL.openConnection();
deregistrationRequestConnection.connect();
StringBuilder uuid = new StringBuilder();
InputStream content = (InputStream) deregistrationRequestConnection.getContent();
byte[] buf = new byte[256];
int read = content.read(buf);
while (read != -1) {
uuid.append(new String(buf, 0, read));
read = content.read(buf);
protected void deregisterReplicaWithMaster(ReplicationMasterDescriptor master) {
try {
URL replicationDeRegistrationRequestURL = master.getReplicationDeRegistrationRequestURL();
final URLConnection deregistrationRequestConnection = replicationDeRegistrationRequestURL.openConnection();
deregistrationRequestConnection.connect();
StringBuilder uuid = new StringBuilder();
InputStream content = (InputStream) deregistrationRequestConnection.getContent();
byte[] buf = new byte[256];
int read = content.read(buf);
while (read != -1) {
uuid.append(new String(buf, 0, read));
read = content.read(buf);
}
content.close();
} catch (Exception ex) {
// ignore exceptions here - they will mostly be caused by an incompatible server
// it is also not problematic if the server does not get this deregistration
// a new registration will overwrite the current one
}
content.close();
}
@Override
@@ -275,7 +283,12 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
replicator.stop();
deregisterReplicaWithMaster(descriptor);
descriptor.getConsumer().getChannel().close();
replicatingFromMaster = null;
replicaUUIDs.clear();
// this is needed because QueuingConsumer.nextDelivery() wont unblock
// if the connection is closed by application.
replicatorThread.interrupt();
}
}
}
@@ -90,10 +90,6 @@ public class Replicator implements Runnable {
}
try {
Delivery delivery = consumer.nextDelivery();
/* Delivery is blocking, upon unblock we will not check for
* stopping because we want to receive at least the last event.
*/
byte[] bytesFromMessage = delivery.getBody();
checksPerformed = 0;
// Set this object's class's class loader as context for de-serialization so that all exported classes
@@ -103,6 +99,8 @@ public class Replicator implements Runnable {
new ByteArrayInputStream(bytesFromMessage));
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
applyOrQueue(operation);
} catch (InterruptedException irr) {
logger.info("Application requested shutdown.");
} catch (ShutdownSignalException sse) {
/* make sure to respond to a stop event without waiting */
if (isBeingStopped()) {
@@ -209,8 +207,13 @@ public class Replicator implements Runnable {
/* make sure to apply everything in queue before stopping this thread */
applyQueue();
}
stopped = true;
logger.info("Signaled Replicator thread to stop asap.");
stopped = true;
try {
master.getConsumer().getChannel().getConnection().close(1);
} catch (Exception ex) {
// ignore any exception during abort.
}
}
public synchronized boolean isBeingStopped() {