mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-04 03:13:49 +00:00
adding set of replicables or their IDs to master and replica descriptors and passing through to ReplicationServlet
Change-Id: I51a70718d94b161e747f2a887028fdeeb996f21a
This commit is contained in:
1 parent
a04ea1df4e
commit
7126902cf2
14 files changed
+134
-53
No files matched your search
+3
-4
@@ -3450,11 +3450,11 @@ public class SailingServiceImpl extends ProxiedRemoteServiceServlet implements S
|
||||
|
||||
@Override
|
||||
public void startReplicatingFromMaster(String messagingHost, String masterHost, String exchangeName, int servletPort, int messagingPort) throws IOException, ClassNotFoundException, InterruptedException {
|
||||
// the queue name must be always the same for this server. in order to achieve
|
||||
// The queue name must always be the same for this server. In order to achieve
|
||||
// this we're using the unique server identifier
|
||||
getReplicationService().startToReplicateFrom(
|
||||
ReplicationFactory.INSTANCE.createReplicationMasterDescriptor(messagingHost, masterHost, exchangeName, servletPort, messagingPort,
|
||||
getReplicationService().getServerIdentifier().toString()));
|
||||
/* use local server identifier as queue name */ getReplicationService().getServerIdentifier().toString(), getReplicationService().getAllReplicables()));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -4565,9 +4565,8 @@ public class SailingServiceImpl extends ProxiedRemoteServiceServlet implements S
|
||||
@Override
|
||||
public void stopSingleReplicaInstance(String identifier) {
|
||||
UUID uuid = UUID.fromString(identifier);
|
||||
ReplicaDescriptor replicaDescriptor = new ReplicaDescriptor(null, uuid, "");
|
||||
try {
|
||||
getReplicationService().unregisterReplica(replicaDescriptor);
|
||||
getReplicationService().unregisterReplica(uuid);
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
throw new RuntimeException(e);
|
||||
|
||||
+4
-3
@@ -26,6 +26,7 @@ import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.common.Util.Pair;
|
||||
import com.sap.sse.common.impl.MillisecondsTimePoint;
|
||||
import com.sap.sse.replication.Replicable;
|
||||
import com.sap.sse.replication.ReplicationMasterDescriptor;
|
||||
import com.sap.sse.replication.impl.ReplicationMasterDescriptorImpl;
|
||||
|
||||
@@ -56,12 +57,12 @@ public class ConnectionResetAndReconnectTest extends AbstractServerReplicationTe
|
||||
|
||||
static class MasterReplicationDescriptorMock extends ReplicationMasterDescriptorImpl {
|
||||
|
||||
public MasterReplicationDescriptorMock(String messagingHost, String hostname, String exchangeName, int servletPort, int messagingPort) {
|
||||
super(messagingHost, exchangeName, messagingPort, UUID.randomUUID().toString(), hostname, servletPort);
|
||||
public MasterReplicationDescriptorMock(String messagingHost, String hostname, String exchangeName, int servletPort, int messagingPort, Iterable<Replicable<?, ?>> replicables) {
|
||||
super(messagingHost, exchangeName, messagingPort, UUID.randomUUID().toString(), hostname, servletPort, replicables);
|
||||
}
|
||||
|
||||
public static MasterReplicationDescriptorMock from(ReplicationMasterDescriptor obj) {
|
||||
return new MasterReplicationDescriptorMock(obj.getMessagingHostname(), obj.getHostname(), obj.getExchangeName(), obj.getServletPort(), obj.getMessagingPort());
|
||||
return new MasterReplicationDescriptorMock(obj.getMessagingHostname(), obj.getHostname(), obj.getExchangeName(), obj.getServletPort(), obj.getMessagingPort(), obj.getReplicables());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+1
-1
@@ -23,7 +23,7 @@ public class ReplicationInstancesManagerLoggingPerformanceTest {
|
||||
@Before
|
||||
public void setUp() throws UnknownHostException {
|
||||
replicationInstanceManager = new ReplicationInstancesManager();
|
||||
replica = new ReplicaDescriptor(InetAddress.getLocalHost(), UUID.randomUUID(), "");
|
||||
replica = new ReplicaDescriptor(InetAddress.getLocalHost(), UUID.randomUUID(), "", /* replicableIds */ new String[] { "Humba" });
|
||||
replicationInstanceManager.registerReplica(replica);
|
||||
operation = new CreateLeaderboardGroup("Test Leaderboard Group", "Description of Test Leaderboard Group", /* displayName */ null,
|
||||
/* displayGroupsInReverseOrder */ false,
|
||||
|
||||
+12
-6
@@ -17,12 +17,12 @@ import java.net.ServerSocket;
|
||||
import java.net.Socket;
|
||||
import java.net.URL;
|
||||
import java.net.URLConnection;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
import java.util.logging.Level;
|
||||
import java.util.logging.Logger;
|
||||
|
||||
import net.jpountz.lz4.LZ4BlockOutputStream;
|
||||
|
||||
import org.junit.Rule;
|
||||
import org.junit.rules.Timeout;
|
||||
|
||||
@@ -41,6 +41,8 @@ import com.sap.sse.replication.impl.ReplicationReceiver;
|
||||
import com.sap.sse.replication.impl.ReplicationServiceImpl;
|
||||
import com.sap.sse.replication.impl.SingletonReplicablesProvider;
|
||||
|
||||
import net.jpountz.lz4.LZ4BlockOutputStream;
|
||||
|
||||
public abstract class AbstractServerReplicationTestSetUp<ReplicableInterface extends Replicable<?, ?>, ReplicableImpl extends ReplicableInterface> {
|
||||
private static final Logger logger = Logger.getLogger(AbstractServerReplicationTestSetUp.class.getName());
|
||||
|
||||
@@ -147,12 +149,12 @@ public abstract class AbstractServerReplicationTestSetUp<ReplicableInterface ext
|
||||
}
|
||||
ReplicationInstancesManager rim = new ReplicationInstancesManager();
|
||||
masterReplicator = new ReplicationServiceImpl(exchangeName, exchangeHost, 0, rim, new SingletonReplicablesProvider(this.master));
|
||||
replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost(), serverUuid, "");
|
||||
replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost(), serverUuid, "", new String[] { this.master.getId().toString() });
|
||||
|
||||
// connect to exchange host and local server running as master
|
||||
// master server and exchange host can be two different hosts
|
||||
ReplicationServiceTestImpl<ReplicableInterface> replicaReplicator = new ReplicationServiceTestImpl<ReplicableInterface>(exchangeName, exchangeHost, rim, replicaDescriptor,
|
||||
this.replica, this.master, masterReplicator, masterDescriptor);
|
||||
this.replica, this.master, masterReplicator);
|
||||
masterDescriptor = replicaReplicator.getMasterDescriptor();
|
||||
servletPort = masterDescriptor.getServletPort();
|
||||
Pair<ReplicationServiceTestImpl<ReplicableInterface>, ReplicationMasterDescriptor> result = new Pair<>(replicaReplicator, masterDescriptor);
|
||||
@@ -236,15 +238,19 @@ public abstract class AbstractServerReplicationTestSetUp<ReplicableInterface ext
|
||||
|
||||
public ReplicationServiceTestImpl(String exchangeName, String exchangeHost, ReplicationInstancesManager replicationInstancesManager,
|
||||
ReplicaDescriptor replicaDescriptor, ReplicableInterface replica,
|
||||
ReplicableInterface master, ReplicationService masterReplicationService, ReplicationMasterDescriptor masterDescriptor)
|
||||
ReplicableInterface master, ReplicationService masterReplicationService)
|
||||
throws IOException {
|
||||
super(exchangeName, exchangeHost, 0, replicationInstancesManager, new SingletonReplicablesProvider(replica));
|
||||
this.replicaDescriptor = replicaDescriptor;
|
||||
this.master = master;
|
||||
ss = new ServerSocket(0); // bind to any free port
|
||||
this.masterReplicationService = masterReplicationService;
|
||||
final List<Replicable<?, ?>> replicablesToReplicate = new ArrayList<>();
|
||||
for (final String replicableIdAsString : replicaDescriptor.getReplicableIdsAsStrings()) {
|
||||
replicablesToReplicate.add(getReplicablesProvider().getReplicable(replicableIdAsString, /* wait */ false));
|
||||
}
|
||||
this.masterDescriptor = new ReplicationMasterDescriptorImpl(exchangeHost, exchangeName, /* messagingPort */ 0,
|
||||
UUID.randomUUID().toString(), "localhost", ss.getLocalPort());
|
||||
UUID.randomUUID().toString(), "localhost", ss.getLocalPort(), replicablesToReplicate);
|
||||
}
|
||||
|
||||
ReplicationMasterDescriptor getMasterDescriptor() {
|
||||
|
||||
@@ -5,5 +5,5 @@ import com.sap.sse.replication.impl.ReplicationFactoryImpl;
|
||||
public interface ReplicationFactory {
|
||||
static ReplicationFactory INSTANCE = new ReplicationFactoryImpl();
|
||||
|
||||
ReplicationMasterDescriptor createReplicationMasterDescriptor(String messagingHostname, String hostname, String exchangeName, int servletPort, int jmsPort, String jmsQueueName);
|
||||
ReplicationMasterDescriptor createReplicationMasterDescriptor(String messagingHostname, String hostname, String exchangeName, int servletPort, int jmsPort, String jmsQueueName, Iterable<Replicable<?, ?>> replicables);
|
||||
}
|
||||
+2
@@ -68,4 +68,6 @@ public interface ReplicationMasterDescriptor {
|
||||
*/
|
||||
Channel createChannel() throws IOException;
|
||||
|
||||
Iterable<Replicable<?, ?>> getReplicables();
|
||||
|
||||
}
|
||||
@@ -62,6 +62,11 @@ public interface ReplicationService {
|
||||
* {@link #registerReplica(ReplicaDescriptor) registered} again.
|
||||
*/
|
||||
void unregisterReplica(ReplicaDescriptor replica) throws IOException;
|
||||
|
||||
/**
|
||||
* Same as {@link #unregisterReplica(ReplicaDescriptor)}, identifying the replica by its {@link ReplicaDescriptor#getUuid() ID}.
|
||||
*/
|
||||
ReplicaDescriptor unregisterReplica(UUID replicaId) throws IOException;
|
||||
|
||||
/**
|
||||
* For a replica replicating off this master, provides statistics in the form of number of operations sent to that
|
||||
@@ -94,4 +99,6 @@ public interface ReplicationService {
|
||||
long getNumberOfBytesSent(ReplicaDescriptor replica);
|
||||
|
||||
double getAverageNumberOfBytesPerMessage(ReplicaDescriptor replica);
|
||||
|
||||
Iterable<Replicable<?, ?>> getAllReplicables();
|
||||
}
|
||||
@@ -155,7 +155,7 @@ public class Activator implements BundleActivator {
|
||||
Integer.valueOf(System.getProperty(PROPERTY_NAME_REPLICATE_MASTER_QUEUE_PORT).trim()),
|
||||
serverReplicationMasterService.getServerIdentifier().toString(),
|
||||
System.getProperty(PROPERTY_NAME_REPLICATE_MASTER_SERVLET_HOST),
|
||||
Integer.valueOf(System.getProperty(PROPERTY_NAME_REPLICATE_MASTER_SERVLET_PORT).trim()));
|
||||
Integer.valueOf(System.getProperty(PROPERTY_NAME_REPLICATE_MASTER_SERVLET_PORT).trim()), replicables);
|
||||
try {
|
||||
serverReplicationMasterService.startToReplicateFrom(master, replicables);
|
||||
logger.info("Automatic replication has been started.");
|
||||
|
||||
+17
-2
@@ -6,6 +6,7 @@ import java.util.UUID;
|
||||
|
||||
import com.sap.sse.common.TimePoint;
|
||||
import com.sap.sse.common.impl.MillisecondsTimePoint;
|
||||
import com.sap.sse.replication.Replicable;
|
||||
|
||||
/**
|
||||
* Describes a replica by remembering its IP address as well as the replication time and a UUID. Hash code and equality
|
||||
@@ -21,15 +22,18 @@ public class ReplicaDescriptor implements Serializable {
|
||||
private final InetAddress ipAddress;
|
||||
private final TimePoint registrationTime;
|
||||
private final String additionalInformation;
|
||||
private final String[] replicableIdsAsStrings;
|
||||
|
||||
/**
|
||||
* Sets the registration time to now.
|
||||
*/
|
||||
public ReplicaDescriptor(InetAddress ipAddress, UUID serverUuid, String additionalInformation) {
|
||||
public ReplicaDescriptor(InetAddress ipAddress, UUID serverUuid, String additionalInformation, String[] replicableIdsAsStrings) {
|
||||
assert replicableIdsAsStrings != null && replicableIdsAsStrings.length > 0;
|
||||
this.uuid = serverUuid;
|
||||
this.registrationTime = MillisecondsTimePoint.now();
|
||||
this.ipAddress = ipAddress;
|
||||
this.additionalInformation = additionalInformation;
|
||||
this.replicableIdsAsStrings = replicableIdsAsStrings;
|
||||
}
|
||||
|
||||
public UUID getUuid() {
|
||||
@@ -48,6 +52,17 @@ public class ReplicaDescriptor implements Serializable {
|
||||
return additionalInformation;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@link Replicable#getId() IDs} of the replicables that the replica represented by this descriptor
|
||||
* has requested from the master for replication. The master may send operations for a superset of those
|
||||
* replicables in case other replicas have requested replication for other replicables. Therefore, the
|
||||
* replica must filter the operations received for those replicable IDs it has been requesting replication
|
||||
* for.
|
||||
*/
|
||||
public String[] getReplicableIdsAsStrings() {
|
||||
return replicableIdsAsStrings;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int hashCode() {
|
||||
final int prime = 31;
|
||||
@@ -74,7 +89,7 @@ public class ReplicaDescriptor implements Serializable {
|
||||
}
|
||||
|
||||
public String toString() {
|
||||
return ""+uuid+": "+ipAddress+" ("+additionalInformation+")";
|
||||
return ""+uuid+": "+ipAddress+" ("+additionalInformation+") for replicables "+String.join(", ", getReplicableIdsAsStrings());
|
||||
}
|
||||
|
||||
}
|
||||
+5
-2
@@ -1,12 +1,15 @@
|
||||
package com.sap.sse.replication.impl;
|
||||
|
||||
import com.sap.sse.replication.Replicable;
|
||||
import com.sap.sse.replication.ReplicationFactory;
|
||||
import com.sap.sse.replication.ReplicationMasterDescriptor;
|
||||
|
||||
public class ReplicationFactoryImpl implements ReplicationFactory {
|
||||
@Override
|
||||
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String messagingHostname, String hostname, String exchangeName, int servletPort, int messagingPort, String queueName) {
|
||||
return new ReplicationMasterDescriptorImpl(messagingHostname, exchangeName, messagingPort, queueName, hostname, servletPort);
|
||||
public ReplicationMasterDescriptor createReplicationMasterDescriptor(String messagingHostname, String hostname,
|
||||
String exchangeName, int servletPort, int messagingPort, String queueName,
|
||||
Iterable<Replicable<?, ?>> replicables) {
|
||||
return new ReplicationMasterDescriptorImpl(messagingHostname, exchangeName, messagingPort, queueName, hostname, servletPort, replicables);
|
||||
}
|
||||
|
||||
}
|
||||
+14
-9
@@ -2,10 +2,9 @@ package com.sap.sse.replication.impl;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
|
||||
import com.sap.sse.common.Util;
|
||||
import com.sap.sse.replication.OperationWithResult;
|
||||
@@ -22,9 +21,10 @@ public class ReplicationInstancesManager {
|
||||
|
||||
/**
|
||||
* The set of descriptors of all registered slaves. All broadcast operations will send the messages to all
|
||||
* registered slaves, assuming the slaves will have subscribed for the replication topic.
|
||||
* registered slaves, assuming the slaves will have subscribed for the replication topic. Keys are the
|
||||
* {@link ReplicaDescriptor#getUuid() IDs of the corresponding values}.
|
||||
*/
|
||||
private Set<ReplicaDescriptor> replicaDescriptors;
|
||||
private Map<UUID, ReplicaDescriptor> replicaDescriptors;
|
||||
|
||||
private Map<ReplicaDescriptor, Map<Class<? extends OperationWithResult<?, ?>>, Integer>> replicationCounts;
|
||||
|
||||
@@ -53,7 +53,7 @@ public class ReplicationInstancesManager {
|
||||
private ReplicationMasterDescriptor replicationMasterDescriptor;
|
||||
|
||||
public ReplicationInstancesManager() {
|
||||
replicaDescriptors = new HashSet<ReplicaDescriptor>();
|
||||
replicaDescriptors = new HashMap<>();
|
||||
replicationCounts = new HashMap<ReplicaDescriptor, Map<Class<? extends OperationWithResult<?, ?>>,Integer>>();
|
||||
totalMessageCount = new HashMap<>();
|
||||
totalNumberOfOperations = new HashMap<>();
|
||||
@@ -71,7 +71,7 @@ public class ReplicationInstancesManager {
|
||||
}
|
||||
|
||||
public Iterable<ReplicaDescriptor> getReplicaDescriptors() {
|
||||
return Collections.unmodifiableCollection(replicaDescriptors);
|
||||
return Collections.unmodifiableCollection(replicaDescriptors.values());
|
||||
}
|
||||
|
||||
public ReplicationMasterDescriptor getReplicationMasterDescriptor() {
|
||||
@@ -79,13 +79,18 @@ public class ReplicationInstancesManager {
|
||||
}
|
||||
|
||||
public void registerReplica(ReplicaDescriptor replica) {
|
||||
replicaDescriptors.add(replica);
|
||||
replicationCounts.put(replica, new HashMap<Class<? extends OperationWithResult<?, ?>>, Integer>());
|
||||
replicaDescriptors.put(replica.getUuid(), replica);
|
||||
replicationCounts.put(replica, new HashMap<>());
|
||||
}
|
||||
|
||||
public void unregisterReplica(ReplicaDescriptor replica) {
|
||||
replicaDescriptors.remove(replica);
|
||||
unregisterReplica(replica.getUuid());
|
||||
}
|
||||
|
||||
public ReplicaDescriptor unregisterReplica(UUID replicaUuid) {
|
||||
final ReplicaDescriptor replica = replicaDescriptors.remove(replicaUuid);
|
||||
replicationCounts.remove(replica);
|
||||
return replica;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+18
-6
@@ -41,13 +41,10 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
private final String messagingHostname;
|
||||
private final int messagingPort;
|
||||
private final String queueName;
|
||||
private final Iterable<Replicable<?, ?>> replicables;
|
||||
|
||||
private QueueingConsumer consumer;
|
||||
|
||||
// TODO bug 2465: add the set of {@link Replicable}s that are replicated from the master represented by this descriptor,
|
||||
// considering that this may be a subset only of the replicables running on this instance or the master server. Example:
|
||||
//replicating only the SecurityService from some other server but being a master regarding all other Replicables.
|
||||
|
||||
/**
|
||||
* @param messagingHostname
|
||||
* name of the host on which the exchange is hosted to which this replica connects with a queue whose
|
||||
@@ -65,9 +62,15 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
* replica with the master
|
||||
* @param servletPort
|
||||
* port for HTTP requests to the master
|
||||
* @param replicables
|
||||
* the {@link Replicable} objects to replicate from the master described by this object; the master will
|
||||
* send operations only for those replicables that at least one replica is registered for; this, however,
|
||||
* may mean that the master also sends operations for replicables that a particular replica hasn't
|
||||
* registered for. Replicas shall silently drop operations for such replicables that they haven't
|
||||
* requested replication for.
|
||||
*/
|
||||
public ReplicationMasterDescriptorImpl(String messagingHostname, String exchangeName, int messagingPort,
|
||||
String queueName, String masterServletHostname, int servletPort) {
|
||||
String queueName, String masterServletHostname, int servletPort, Iterable<Replicable<?, ?>> replicables) {
|
||||
this.masterServletHostname = masterServletHostname;
|
||||
this.messagingHostname = messagingHostname;
|
||||
this.servletPort = servletPort;
|
||||
@@ -75,16 +78,20 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
this.exchangeName = exchangeName;
|
||||
this.queueName = queueName;
|
||||
this.consumer = null;
|
||||
this.replicables = replicables;
|
||||
}
|
||||
|
||||
@Override
|
||||
public URL getReplicationRegistrationRequestURL(UUID uuid, String additional) throws MalformedURLException,
|
||||
UnsupportedEncodingException {
|
||||
final String[] replicableIdsAsString = StreamSupport.stream(replicables.spliterator(), /* parallel */ false).map(r->r.getId()).toArray(i->new String[i]);
|
||||
return new URL("http", getHostname(), servletPort, REPLICATION_SERVLET + "?" + ReplicationServlet.ACTION + "="
|
||||
+ ReplicationServlet.Action.REGISTER.name() + "&" + ReplicationServlet.SERVER_UUID + "="
|
||||
+ java.net.URLEncoder.encode(uuid.toString(), "UTF-8") + "&"
|
||||
+ ReplicationServlet.ADDITIONAL_INFORMATION + "="
|
||||
+ java.net.URLEncoder.encode(ServerInfo.getBuildVersion(), "UTF-8"));
|
||||
+ java.net.URLEncoder.encode(ServerInfo.getBuildVersion(), "UTF-8") + "&"
|
||||
+ ReplicationServlet.REPLICABLES_IDS_AS_STRINGS_COMMA_SEPARATED + "="
|
||||
+ java.net.URLEncoder.encode(String.join(",", replicableIdsAsString), "UTF-8"));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -94,6 +101,11 @@ public class ReplicationMasterDescriptorImpl implements ReplicationMasterDescrip
|
||||
+ uuid.toString());
|
||||
}
|
||||
|
||||
@Override
|
||||
public Iterable<Replicable<?, ?>> getReplicables() {
|
||||
return replicables;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int hashCode() {
|
||||
final int prime = 31;
|
||||
|
||||
+22
-5
@@ -371,10 +371,10 @@ public class ReplicationServiceImpl implements ReplicationService {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void unregisterReplica(ReplicaDescriptor replica) throws IOException {
|
||||
logger.info("Unregistering replica " + replica);
|
||||
public ReplicaDescriptor unregisterReplica(UUID replicaUuid) throws IOException {
|
||||
logger.info("Unregistering replica with ID " + replicaUuid);
|
||||
synchronized (replicationInstancesManager) {
|
||||
replicationInstancesManager.unregisterReplica(replica);
|
||||
final ReplicaDescriptor unregisteredReplica = replicationInstancesManager.unregisterReplica(replicaUuid);
|
||||
if (!replicationInstancesManager.hasReplicas()) {
|
||||
removeAsListenerFromReplicables();
|
||||
synchronized (this) {
|
||||
@@ -384,8 +384,15 @@ public class ReplicationServiceImpl implements ReplicationService {
|
||||
}
|
||||
}
|
||||
}
|
||||
return unregisteredReplica;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void unregisterReplica(ReplicaDescriptor replica) throws IOException {
|
||||
logger.info("Unregistering replica " + replica);
|
||||
unregisterReplica(replica.getUuid());
|
||||
}
|
||||
|
||||
private void removeAsListenerFromReplicables() {
|
||||
for (ReplicationServiceExecutionListener<?> listener : executionListenersByReplicableIdAsString.values()) {
|
||||
@@ -556,7 +563,7 @@ public class ReplicationServiceImpl implements ReplicationService {
|
||||
} else {
|
||||
logger.info("Starting to replicate from " + master);
|
||||
try {
|
||||
registerReplicaWithMaster(master);
|
||||
registerReplicaWithMaster(master, replicables);
|
||||
} catch (Exception ex) {
|
||||
logger.log(Level.SEVERE, "ERROR", ex);
|
||||
throw ex;
|
||||
@@ -628,9 +635,14 @@ public class ReplicationServiceImpl implements ReplicationService {
|
||||
}
|
||||
|
||||
/**
|
||||
* @param replicables
|
||||
* the replica is registered for these {@link Replicable}s. The master sends operations only for
|
||||
* replicables that at least one replica has registered for. This may mean that operations are received
|
||||
* for replicables for which no replicable was requested. Replicas shall drop such operations silently.
|
||||
*
|
||||
* @return the UUID that the master generated for this client which is also entered into {@link #replicaUUIDs}
|
||||
*/
|
||||
private String registerReplicaWithMaster(ReplicationMasterDescriptor master) throws IOException,
|
||||
private String registerReplicaWithMaster(ReplicationMasterDescriptor master, Iterable<Replicable<?, ?>> replicables) throws IOException,
|
||||
ClassNotFoundException {
|
||||
URL replicationRegistrationRequestURL = master.getReplicationRegistrationRequestURL(getServerIdentifier(),
|
||||
ServerInfo.getBuildVersion());
|
||||
@@ -703,6 +715,11 @@ public class ReplicationServiceImpl implements ReplicationService {
|
||||
return replicationInstancesManager.getAverageNumberOfBytesPerMessage(replica);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Iterable<Replicable<?, ?>> getAllReplicables() {
|
||||
return getReplicablesProvider().getReplicables();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stopToReplicateFromMaster() throws IOException {
|
||||
ReplicationMasterDescriptor descriptor = getReplicatingFromMaster();
|
||||
|
||||
+27
-13
@@ -17,6 +17,7 @@ import javax.servlet.ServletException;
|
||||
import javax.servlet.http.HttpServletRequest;
|
||||
import javax.servlet.http.HttpServletResponse;
|
||||
|
||||
import net.jpountz.lz4.LZ4BlockInputStream;
|
||||
import net.jpountz.lz4.LZ4BlockOutputStream;
|
||||
|
||||
import org.apache.commons.lang.StringEscapeUtils;
|
||||
@@ -77,15 +78,23 @@ public class ReplicationServlet extends AbstractHttpServlet {
|
||||
|
||||
/**
|
||||
* The client identifies itself in the request. Two servlet operations are supported currently: registering the
|
||||
* client with the replication service (if not already created, the JMS replication topic will be created by this
|
||||
* registration); and obtaining an initial load stream that the replica can use to initialize its
|
||||
* {@link RacingEventService}. The operation performed is selected by passing one of the {@link Action} enumeration
|
||||
* values for the URL parameter named {@link #ACTION}.
|
||||
* client with the replication service (if not already created, the message exchange will be created by this
|
||||
* registration); and triggering sending an initial load stream through RabbitMQ that the replica can use to
|
||||
* initialize its {@link Replicable} s. The IDs of the replicables of which the initial load is to be sent is
|
||||
* expected as the {@link #REPLICABLES_IDS_AS_STRINGS_COMMA_SEPARATED} parameter value. The servlet response
|
||||
* consists of the name of the RabbitMQ queue name that the client can connect to and through which to receive
|
||||
* the LZ4-compressed initial load stream per replicable, using a {@link RabbitInputStreamProvider} and
|
||||
* an {@link LZ4BlockInputStream}..
|
||||
* <p>
|
||||
*
|
||||
* The operation performed is selected by passing one of the {@link Action} enumeration values for the URL parameter
|
||||
* named {@link #ACTION}.
|
||||
*/
|
||||
@Override
|
||||
protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
|
||||
String action = req.getParameter(ACTION);
|
||||
logger.info("Received replication request, action is "+action);
|
||||
String[] replicableIdsAsStrings;
|
||||
switch (Action.valueOf(action)) {
|
||||
case REGISTER:
|
||||
registerClientWithReplicationService(req, resp);
|
||||
@@ -94,7 +103,7 @@ public class ReplicationServlet extends AbstractHttpServlet {
|
||||
deregisterClientWithReplicationService(req, resp);
|
||||
break;
|
||||
case INITIAL_LOAD:
|
||||
String[] replicableIdsAsStrings = req.getParameter(REPLICABLES_IDS_AS_STRINGS_COMMA_SEPARATED).split(",");
|
||||
replicableIdsAsStrings = req.getParameter(REPLICABLES_IDS_AS_STRINGS_COMMA_SEPARATED).split(",");
|
||||
Channel channel = getReplicationService().createMasterChannel();
|
||||
try {
|
||||
RabbitOutputStream ros = new RabbitOutputStream(INITIAL_LOAD_PACKAGE_SIZE, channel,
|
||||
@@ -174,17 +183,22 @@ public class ReplicationServlet extends AbstractHttpServlet {
|
||||
}
|
||||
|
||||
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());
|
||||
final UUID replicaUuid = UUID.fromString(req.getParameter(SERVER_UUID));
|
||||
final ReplicaDescriptor replica = getReplicationService().unregisterReplica(replicaUuid);
|
||||
if (replica != null) {
|
||||
logger.info("Deregistered replication client with this server " + replica.getIpAddress());
|
||||
resp.setContentType("text/plain");
|
||||
resp.getWriter().print(replica.getUuid());
|
||||
} else {
|
||||
logger.warning("Couldn't find replica to de-register with ID "+replicaUuid);
|
||||
}
|
||||
}
|
||||
|
||||
private void registerClientWithReplicationService(HttpServletRequest req, HttpServletResponse resp)
|
||||
throws IOException {
|
||||
ReplicaDescriptor replica = getReplicaDescriptor(req);
|
||||
final ReplicaDescriptor replica = getReplicaDescriptor(req);
|
||||
getReplicationService().registerReplica(replica);
|
||||
logger.info("Registered new replica " + replica);
|
||||
resp.setContentType("text/plain");
|
||||
resp.getWriter().print(replica.getUuid());
|
||||
}
|
||||
@@ -193,7 +207,7 @@ public class ReplicationServlet extends AbstractHttpServlet {
|
||||
InetAddress ipAddress = InetAddress.getByName(req.getRemoteAddr());
|
||||
UUID uuid = UUID.fromString(req.getParameter(SERVER_UUID));
|
||||
String additional = req.getParameter(ADDITIONAL_INFORMATION);
|
||||
logger.info("Registered new replica " + ipAddress + " " + uuid.toString() + " " + additional);
|
||||
return new ReplicaDescriptor(ipAddress, uuid, additional);
|
||||
final String[] replicableIdsAsStrings = req.getParameter(REPLICABLES_IDS_AS_STRINGS_COMMA_SEPARATED).split(",");
|
||||
return new ReplicaDescriptor(ipAddress, uuid, additional, replicableIdsAsStrings);
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user