mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-11 06:40:53 +00:00
Adding UI and API support for stopping replication.
This commit is contained in:
1 parent
b23ae0defb
commit
3e0f694e65
16 files changed
+214
-17
No files matched your search
+2
@@ -15,6 +15,8 @@ import com.rabbitmq.client.QueueingConsumer;
|
||||
public interface ReplicationMasterDescriptor {
|
||||
|
||||
URL getReplicationRegistrationRequestURL() throws MalformedURLException;
|
||||
|
||||
URL getReplicationDeRegistrationRequestURL() throws MalformedURLException;
|
||||
|
||||
URL getInitialLoadURL() throws MalformedURLException;
|
||||
|
||||
|
||||
+7
@@ -42,4 +42,11 @@ public interface ReplicationService {
|
||||
* replica by type, where the operation type is the key, represented as the operation's class name
|
||||
*/
|
||||
Map<Class<? extends RacingEventServiceOperation<?>>, Integer> getStatistics(ReplicaDescriptor replicaDescriptor);
|
||||
|
||||
/**
|
||||
* Stops the currently running replication. As there can be only one replication running
|
||||
* this method needs no parameters.
|
||||
* @throws IOException
|
||||
*/
|
||||
void stopToReplicateFromMaster() throws IOException;
|
||||
}
|
||||
+7
@@ -33,6 +33,12 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
+ ReplicationServlet.Action.REGISTER.name());
|
||||
}
|
||||
|
||||
@Override
|
||||
public URL getReplicationDeRegistrationRequestURL() throws MalformedURLException {
|
||||
return new URL("http", hostname, servletPort, REPLICATION_SERVLET + "?" + ReplicationServlet.ACTION + "="
|
||||
+ ReplicationServlet.Action.DEREGISTER.name());
|
||||
}
|
||||
|
||||
@Override
|
||||
public URL getInitialLoadURL() throws MalformedURLException {
|
||||
return new URL("http", hostname, servletPort, REPLICATION_SERVLET + "?" + ReplicationServlet.ACTION + "="
|
||||
@@ -79,4 +85,5 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
public String getExchangeName() {
|
||||
return exchangeName;
|
||||
}
|
||||
|
||||
}
|
||||
+41
-4
@@ -69,6 +69,8 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
*/
|
||||
private final String exchangeName;
|
||||
|
||||
private Replicator replicator;
|
||||
|
||||
public ReplicationServiceImpl(String exchangeName, final ReplicationInstancesManager replicationInstancesManager) throws IOException {
|
||||
this.replicationInstancesManager = replicationInstancesManager;
|
||||
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
|
||||
@@ -77,6 +79,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
racingEventServiceTracker.open();
|
||||
localService = null;
|
||||
this.exchangeName = exchangeName;
|
||||
replicator = null;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -87,9 +90,10 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
public ReplicationServiceImpl(String exchangeName,
|
||||
final ReplicationInstancesManager replicationInstancesManager, RacingEventService localService) throws IOException {
|
||||
this.replicationInstancesManager = replicationInstancesManager;
|
||||
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>();
|
||||
replicaUUIDs = new HashMap<ReplicationMasterDescriptor, String>(); // XXX why is this a map? there should be only one connection to a master
|
||||
this.localService = localService;
|
||||
this.exchangeName = exchangeName;
|
||||
replicator = null;
|
||||
}
|
||||
|
||||
private Channel createMasterChannel(String exchangeName) throws IOException {
|
||||
@@ -144,9 +148,12 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
if (!replicationInstancesManager.hasReplicas()) {
|
||||
removeAsListenerFromRacingEventService();
|
||||
synchronized (this) {
|
||||
masterChannel.close();
|
||||
masterChannel = null;
|
||||
if (masterChannel != null) {
|
||||
masterChannel.close();
|
||||
masterChannel = null;
|
||||
}
|
||||
}
|
||||
logger.info("Unregistered replica " + replica.getIpAddress().toString());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -183,6 +190,8 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
registerReplicaWithMaster(master);
|
||||
logger.info("Registered replica with master");
|
||||
QueueingConsumer consumer = null;
|
||||
// logging exception here because it will not propagate
|
||||
// thru the client with all details
|
||||
try {
|
||||
consumer = master.getConsumer();
|
||||
} catch (Exception ex) {
|
||||
@@ -192,7 +201,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
logger.info("Connection to exchange successful.");
|
||||
URL initialLoadURL = master.getInitialLoadURL();
|
||||
logger.info("Initial load URL is "+initialLoadURL);
|
||||
final Replicator replicator = new Replicator(master, this, /* startSuspended */ true, consumer);
|
||||
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();
|
||||
logger.info("Started replicator thread");
|
||||
@@ -224,6 +233,21 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
registerReplicaUuidForMaster(replicaUUID, master);
|
||||
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);
|
||||
}
|
||||
content.close();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> void executed(RacingEventServiceOperation<T> operation) {
|
||||
@@ -243,4 +267,17 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
|
||||
return replicationInstancesManager.getStatistics(replicaDescriptor);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stopToReplicateFromMaster() throws IOException {
|
||||
ReplicationMasterDescriptor descriptor = isReplicatingFromMaster();
|
||||
if (descriptor != null) {
|
||||
synchronized(replicaUUIDs) {
|
||||
replicator.stop();
|
||||
deregisterReplicaWithMaster(descriptor);
|
||||
descriptor.getConsumer().getChannel().close();
|
||||
replicaUUIDs.clear();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
+13
-1
@@ -31,7 +31,7 @@ public class ReplicationServlet extends SailingServerHttpServlet {
|
||||
|
||||
private static final long serialVersionUID = 4835516998934433846L;
|
||||
|
||||
public enum Action { REGISTER, INITIAL_LOAD }
|
||||
public enum Action { REGISTER, INITIAL_LOAD, DEREGISTER }
|
||||
|
||||
public static final String ACTION = "action";
|
||||
|
||||
@@ -61,6 +61,9 @@ public class ReplicationServlet extends SailingServerHttpServlet {
|
||||
case REGISTER:
|
||||
registerClientWithReplicationService(req, resp);
|
||||
break;
|
||||
case DEREGISTER:
|
||||
deregisterClientWithReplicationService(req, resp);
|
||||
break;
|
||||
case INITIAL_LOAD:
|
||||
ObjectOutputStream oos = new ObjectOutputStream(resp.getOutputStream());
|
||||
try {
|
||||
@@ -78,6 +81,14 @@ public class ReplicationServlet extends SailingServerHttpServlet {
|
||||
}
|
||||
}
|
||||
|
||||
private void deregisterClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp) throws IOException {
|
||||
ReplicaDescriptor replica = getReplicaDescriptor(req);
|
||||
getReplicationService().unregisterReplica(replica);
|
||||
logger.info("Deregistered replication client with this server " + replica.getIpAddress());
|
||||
resp.setContentType("text/plain");
|
||||
resp.getWriter().print(replica.getUuid());
|
||||
}
|
||||
|
||||
private void registerClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp)
|
||||
throws IOException {
|
||||
ReplicaDescriptor replica = getReplicaDescriptor(req);
|
||||
@@ -87,6 +98,7 @@ public class ReplicationServlet extends SailingServerHttpServlet {
|
||||
}
|
||||
|
||||
private ReplicaDescriptor getReplicaDescriptor(HttpServletRequest req) throws UnknownHostException {
|
||||
// XXX: this can lead to problems if there are multiple replicas on one server
|
||||
InetAddress ipAddress = InetAddress.getByName(req.getRemoteAddr());
|
||||
return new ReplicaDescriptor(ipAddress);
|
||||
}
|
||||
|
||||
+39
-1
@@ -52,6 +52,8 @@ public class Replicator implements Runnable {
|
||||
*/
|
||||
private boolean suspended;
|
||||
|
||||
private boolean stopped = false;
|
||||
|
||||
/**
|
||||
* Starts the replicator immediately, not holding back messages received but forwarding them directly.
|
||||
*
|
||||
@@ -83,8 +85,22 @@ public class Replicator implements Runnable {
|
||||
ClassLoader oldClassLoader = Thread.currentThread().getContextClassLoader();
|
||||
|
||||
while (true) {
|
||||
if (isBeingStopped()) {
|
||||
break;
|
||||
}
|
||||
try {
|
||||
Delivery delivery = consumer.nextDelivery();
|
||||
|
||||
/* Delivery is blocking, upon unblock we need to check if there
|
||||
* we has been stopped. If this is the case we assume that
|
||||
* we do not handle any new deliveries. It's a bit odd to have so
|
||||
* many checks for stopped event but we need to make sure that
|
||||
* we check this at every stage.
|
||||
*/
|
||||
if (isBeingStopped()) {
|
||||
break;
|
||||
}
|
||||
|
||||
byte[] bytesFromMessage = delivery.getBody();
|
||||
checksPerformed = 0;
|
||||
// Set this object's class's class loader as context for de-serialization so that all exported classes
|
||||
@@ -95,6 +111,11 @@ public class Replicator implements Runnable {
|
||||
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) ois.readObject();
|
||||
applyOrQueue(operation);
|
||||
} catch (ShutdownSignalException sse) {
|
||||
/* make sure to respond to a stop event without waiting */
|
||||
if (isBeingStopped()) {
|
||||
break;
|
||||
}
|
||||
|
||||
if (sse.isInitiatedByApplication()) {
|
||||
logger.severe("Application shut down messaging queue for " + this.toString());
|
||||
break;
|
||||
@@ -103,12 +124,15 @@ public class Replicator implements Runnable {
|
||||
logger.info(sse.getMessage());
|
||||
if (checksPerformed <= CHECK_COUNT) {
|
||||
try {
|
||||
logger.info("Replication reciever is sleeping because of " + sse.getLocalizedMessage());
|
||||
Thread.sleep(CHECK_INTERVAL);
|
||||
|
||||
/* isOpen() will return false if the channel has been closed. This
|
||||
* does not hold when the connection is dropped.
|
||||
*/
|
||||
if (!this.consumer.getChannel().isOpen()) {
|
||||
/* for a reconnection we need to instantiate a new consumer */
|
||||
try {
|
||||
logger.info("Channel seems to be closed. Trying to reconnect consumer queue...");
|
||||
this.consumer = master.getConsumer();
|
||||
Thread.sleep(CHECK_INTERVAL);
|
||||
checksPerformed += 1;
|
||||
@@ -133,6 +157,7 @@ public class Replicator implements Runnable {
|
||||
Thread.currentThread().setContextClassLoader(oldClassLoader);
|
||||
}
|
||||
}
|
||||
logger.info("Stopped replicator thread. This server will no longer receive events from a master.");
|
||||
}
|
||||
|
||||
public synchronized boolean isQueueEmpty() {
|
||||
@@ -185,6 +210,19 @@ public class Replicator implements Runnable {
|
||||
public synchronized boolean isSuspended() {
|
||||
return suspended;
|
||||
}
|
||||
|
||||
public synchronized void stop() {
|
||||
if (isSuspended()) {
|
||||
/* make sure to apply everything in queue before stopping this thread */
|
||||
applyQueue();
|
||||
}
|
||||
stopped = true;
|
||||
logger.info("Signaled Replicator thread to stop asap.");
|
||||
}
|
||||
|
||||
public synchronized boolean isBeingStopped() {
|
||||
return stopped;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
|
||||
Reference in new issue
Block a user