diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java index 20623c7efdc..cd22465b890 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/AbstractServerReplicationTest.java @@ -35,6 +35,7 @@ import com.sap.sailing.domain.persistence.media.MediaDBFactory; import com.sap.sailing.domain.racelog.tracking.EmptyGPSFixStore; import com.sap.sailing.domain.tracking.impl.EmptyWindStore; import com.sap.sailing.server.RacingEventService; +import com.sap.sailing.server.RacingEventServiceOperation; import com.sap.sailing.server.impl.RacingEventServiceImpl; import com.sap.sailing.server.replication.ReplicationMasterDescriptor; import com.sap.sailing.server.replication.ReplicationService; @@ -55,7 +56,7 @@ public abstract class AbstractServerReplicationTest { protected RacingEventServiceImpl master; protected ReplicationServiceTestImpl replicaReplicator; private ReplicaDescriptor replicaDescriptor; - private ReplicationServiceImpl masterReplicator; + private ReplicationServiceImpl> masterReplicator; private ReplicationMasterDescriptor masterDescriptor; @Rule public Timeout AbstractTracTracLiveTestTimeout = new Timeout(5 * 60 * 1000); // timeout after 5 minutes @@ -115,7 +116,7 @@ public abstract class AbstractServerReplicationTest { new DomainFactoryImpl()), mongoObjectFactory, MediaDBFactory.INSTANCE.getMediaDB(mongoDBService), EmptyWindStore.INSTANCE, EmptyGPSFixStore.INSTANCE); } ReplicationInstancesManager rim = new ReplicationInstancesManager(); - masterReplicator = new ReplicationServiceImpl(exchangeName, exchangeHost, rim, this.master); + masterReplicator = new ReplicationServiceImpl>(exchangeName, exchangeHost, rim, this.master); replicaDescriptor = new ReplicaDescriptor(InetAddress.getLocalHost(), serverUuid, ""); masterReplicator.registerReplica(replicaDescriptor); // connect to exchange host and local server running as master @@ -166,7 +167,7 @@ public abstract class AbstractServerReplicationTest { replicaReplicator.stopToReplicateFromMaster(); } - static class ReplicationServiceTestImpl extends ReplicationServiceImpl { + static class ReplicationServiceTestImpl extends ReplicationServiceImpl> { protected static final int INITIAL_LOAD_PACKAGE_SIZE = 1024*1024; private final DomainFactory resolveAgainst; private final RacingEventService master; @@ -258,12 +259,13 @@ public abstract class AbstractServerReplicationTest { // replicator.setSuspended(false); // resume after initial load // } - protected Replicator startToReplicateFromButDontYetFetchInitialLoad(ReplicationMasterDescriptor master, boolean startReplicatorSuspended) + protected Replicator> startToReplicateFromButDontYetFetchInitialLoad(ReplicationMasterDescriptor master, boolean startReplicatorSuspended) throws IOException { masterReplicationService.registerReplica(replicaDescriptor); registerReplicaUuidForMaster(replicaDescriptor.getUuid().toString(), master); QueueingConsumer consumer = master.getConsumer(); - final Replicator replicator = new Replicator(master, this, startReplicatorSuspended, consumer); + final Replicator> replicator = new Replicator>( + master, this, startReplicatorSuspended, consumer); new Thread(replicator).start(); return replicator; } diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/InitialLoadReplicationObjectIdentityTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/InitialLoadReplicationObjectIdentityTest.java index 77b893b5fd4..53e6f316bb9 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/InitialLoadReplicationObjectIdentityTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/InitialLoadReplicationObjectIdentityTest.java @@ -42,6 +42,8 @@ import com.sap.sailing.domain.persistence.media.MediaDBFactory; import com.sap.sailing.domain.racelog.tracking.EmptyGPSFixStore; import com.sap.sailing.domain.test.TrackBasedTest; import com.sap.sailing.domain.tracking.impl.EmptyWindStore; +import com.sap.sailing.server.RacingEventService; +import com.sap.sailing.server.RacingEventServiceOperation; import com.sap.sailing.server.impl.RacingEventServiceImpl; import com.sap.sailing.server.operationaltransformation.AddDefaultRegatta; import com.sap.sailing.server.operationaltransformation.AddRaceDefinition; @@ -140,7 +142,7 @@ public class InitialLoadReplicationObjectIdentityTest extends AbstractServerRepl /* fire up replication */ performReplicationSetup(); ReplicationMasterDescriptor the_master = replicationDescriptorPair.getB(); /* master descriptor */ - Replicator replicator = replicationDescriptorPair.getA().startToReplicateFromButDontYetFetchInitialLoad(the_master, /* startReplicatorSuspended */ true); + Replicator> replicator = replicationDescriptorPair.getA().startToReplicateFromButDontYetFetchInitialLoad(the_master, /* startReplicatorSuspended */ true); replicationDescriptorPair.getA().initialLoad(); replicator.setSuspended(false); synchronized (replicator) { diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/PrematureOperationReceiptTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/PrematureOperationReceiptTest.java index 794f70a6d2f..383b5acfd72 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/PrematureOperationReceiptTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/PrematureOperationReceiptTest.java @@ -10,13 +10,15 @@ import org.junit.Test; import com.sap.sailing.domain.leaderboard.Leaderboard; import com.sap.sailing.domain.leaderboard.impl.LowPoint; +import com.sap.sailing.server.RacingEventService; +import com.sap.sailing.server.RacingEventServiceOperation; import com.sap.sailing.server.operationaltransformation.CreateFlexibleLeaderboard; import com.sap.sailing.server.replication.ReplicationMasterDescriptor; import com.sap.sailing.server.replication.impl.Replicator; import com.sap.sse.common.Util; public class PrematureOperationReceiptTest extends AbstractServerReplicationTest { - private Replicator replicator; + private Replicator> replicator; private ReplicationServiceTestImpl replicationService; /** diff --git a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/TrackRaceReplicationTest.java b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/TrackRaceReplicationTest.java index 58023d7e10a..5742642d9e5 100755 --- a/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/TrackRaceReplicationTest.java +++ b/java/com.sap.sailing.server.replication.test/src/com/sap/sailing/server/replication/test/TrackRaceReplicationTest.java @@ -32,6 +32,7 @@ import com.sap.sailing.domain.tracking.DynamicTrackedRace; import com.sap.sailing.domain.tracking.RaceHandle; import com.sap.sailing.domain.tracking.RaceTrackingConnectivityParameters; import com.sap.sailing.domain.tracking.TrackedRace; +import com.sap.sailing.server.RacingEventService; import com.sap.sailing.server.operationaltransformation.AddColumnToLeaderboard; import com.sap.sailing.server.operationaltransformation.ConnectTrackedRaceToLeaderboardColumn; import com.sap.sailing.server.operationaltransformation.CreateFlexibleLeaderboard; @@ -40,7 +41,6 @@ import com.sap.sailing.server.replication.OperationExecutionListener; import com.sap.sse.common.Util; import com.sap.sse.common.impl.MillisecondsTimePoint; import com.sap.sse.operationaltransformation.Operation; -import com.sap.sse.operationaltransformation.OperationWithTransformationSupport; public class TrackRaceReplicationTest extends AbstractServerReplicationTest { private TrackedRace masterTrackedRace; @@ -71,9 +71,9 @@ public class TrackRaceReplicationTest extends AbstractServerReplicationTest { MillisecondsTimePoint startOfTracking = new MillisecondsTimePoint(cal.getTimeInMillis()); cal.set(2011, 05, 23, 15, 14, 31); MillisecondsTimePoint endOfTracking = new MillisecondsTimePoint(cal.getTimeInMillis()); - master.addOperationExecutionListener(new OperationExecutionListener() { + master.addOperationExecutionListener(new OperationExecutionListener() { @Override - public void executed(OperationWithTransformationSupport> operation) { + public void executed(Operation operation) { if (operation instanceof CreateTrackedRace) { synchronized (notifier) { notifier[0] = true; diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java index b00e1d23665..791a452251e 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventService.java @@ -66,7 +66,6 @@ import com.sap.sailing.domain.tracking.TrackedRegattaRegistry; import com.sap.sailing.domain.tracking.TrackerManager; import com.sap.sailing.domain.tracking.WindStore; import com.sap.sailing.server.masterdata.DataImportLockWithProgress; -import com.sap.sailing.server.replication.OperationExecutionListener; import com.sap.sailing.server.replication.Replicable; import com.sap.sse.common.TimePoint; import com.sap.sse.common.Util; diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceOperation.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceOperation.java index b0afcd28c6c..29433430867 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceOperation.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/RacingEventServiceOperation.java @@ -11,15 +11,9 @@ import com.sap.sailing.server.operationaltransformation.MoveLeaderboardColumnUp; import com.sap.sailing.server.operationaltransformation.RemoveLeaderboard; import com.sap.sailing.server.operationaltransformation.RemoveLeaderboardColumn; import com.sap.sailing.server.operationaltransformation.RenameLeaderboardColumn; -import com.sap.sse.operationaltransformation.Operation; +import com.sap.sailing.server.replication.OperationWithResult; -public interface RacingEventServiceOperation extends Operation, Serializable { - /** - * Performs the actual operation, applying it to the toState service. The operation's result is - * returned. - */ - ResultType internalApplyTo(RacingEventService toState) throws Exception; - +public interface RacingEventServiceOperation extends OperationWithResult, Serializable { /** * Assumes this is the "server" operation and transforms the client's removeColumnFromLeaderboardClientOp according to this * operation. The default implementation will probably pass on the untransformed client operation. However, if this diff --git a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java index a516770dfc1..b234c2e63e3 100755 --- a/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java +++ b/java/com.sap.sailing.server/src/com/sap/sailing/server/impl/RacingEventServiceImpl.java @@ -296,7 +296,7 @@ public class RacingEventServiceImpl implements RacingEventServiceWithTestSupport */ private final NamedReentrantReadWriteLock regattaTrackingCacheLock; - private final ConcurrentHashMap operationExecutionListeners; + private final ConcurrentHashMap, OperationExecutionListener> operationExecutionListeners; /** * Keys are the toString() representation of the {@link RaceDefinition#getId() IDs} of races passed to @@ -2084,7 +2084,7 @@ public class RacingEventServiceImpl implements RacingEventServiceWithTestSupport * this service known. */ @Override - public T apply(Operation operation) { + public T apply(Operation operation) { RacingEventServiceOperation reso = (RacingEventServiceOperation) operation; try { T result = reso.internalApplyTo(this); @@ -2108,7 +2108,7 @@ public class RacingEventServiceImpl implements RacingEventServiceWithTestSupport @Override public void replicate(RacingEventServiceOperation operation) { - for (OperationExecutionListener listener : operationExecutionListeners.keySet()) { + for (OperationExecutionListener listener : operationExecutionListeners.keySet()) { try { listener.executed(operation); } catch (Exception e) { @@ -2120,12 +2120,12 @@ public class RacingEventServiceImpl implements RacingEventServiceWithTestSupport } @Override - public void addOperationExecutionListener(OperationExecutionListener listener) { + public void addOperationExecutionListener(OperationExecutionListener listener) { operationExecutionListeners.put(listener, listener); } @Override - public void removeOperationExecutionListener(OperationExecutionListener listener) { + public void removeOperationExecutionListener(OperationExecutionListener listener) { operationExecutionListeners.remove(listener); } diff --git a/java/com.sap.sse.operationaltransformation.test/META-INF/MANIFEST.MF b/java/com.sap.sse.operationaltransformation.test/META-INF/MANIFEST.MF index 911559cc367..75cb42f3baf 100644 --- a/java/com.sap.sse.operationaltransformation.test/META-INF/MANIFEST.MF +++ b/java/com.sap.sse.operationaltransformation.test/META-INF/MANIFEST.MF @@ -5,5 +5,6 @@ Bundle-SymbolicName: com.sap.sse.operationaltransformation.test Bundle-Version: 1.0.0.qualifier Bundle-Vendor: SAP Bundle-RequiredExecutionEnvironment: JavaSE-1.8 -Require-Bundle: org.junit4;bundle-version="4.8.2" +Require-Bundle: org.junit4;bundle-version="4.8.2", + com.sap.sse.replication Fragment-Host: com.sap.sse.operationaltransformation diff --git a/java/com.sap.sse.operationaltransformation.test/src/com/sap/sse/operationaltransformation/test/OperationalTransformationTest.java b/java/com.sap.sse.operationaltransformation.test/src/com/sap/sse/operationaltransformation/test/OperationalTransformationTest.java index 8711f61e787..b6da96b44bd 100755 --- a/java/com.sap.sse.operationaltransformation.test/src/com/sap/sse/operationaltransformation/test/OperationalTransformationTest.java +++ b/java/com.sap.sse.operationaltransformation.test/src/com/sap/sse/operationaltransformation/test/OperationalTransformationTest.java @@ -12,9 +12,9 @@ import org.junit.Test; import com.sap.sse.operationaltransformation.ClientServerOperationPair; import com.sap.sse.operationaltransformation.Operation; import com.sap.sse.operationaltransformation.Peer; +import com.sap.sse.operationaltransformation.Peer.Role; import com.sap.sse.operationaltransformation.PeerImpl; import com.sap.sse.operationaltransformation.Transformer; -import com.sap.sse.operationaltransformation.Peer.Role; import com.sap.sse.operationaltransformation.test.util.Base64; public class OperationalTransformationTest { diff --git a/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationWithTransformationSupport.java b/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationWithTransformationSupport.java index 75c1ce29a5c..d820f393066 100755 --- a/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationWithTransformationSupport.java +++ b/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationWithTransformationSupport.java @@ -1,7 +1,13 @@ package com.sap.sse.operationaltransformation; - - +/** + * Returns the same state object as was passed in. This is required for the Operational Transformation algorithm to work. + * + * @author Axel Uhl (D043530) + * + * @param + * @param + */ public interface OperationWithTransformationSupport> extends Operation { /** * Implements the specific transformation rule for the implementing subclass for the set of possible peer operations diff --git a/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationalTransformer.java b/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationalTransformer.java index 11f2820c269..74099633359 100755 --- a/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationalTransformer.java +++ b/java/com.sap.sse.operationaltransformation/src/com/sap/sse/operationaltransformation/OperationalTransformer.java @@ -1,8 +1,15 @@ package com.sap.sse.operationaltransformation; - +/** + * A default implementation for the {@link Transformer} interface, assuming operations implement the + * {@link OperationWithTransformationSupport} interface which standardizes how operations are being transformed. + * + * @author Axel Uhl (D043530) + * + * @param + * @param + */ public class OperationalTransformer> implements Transformer { - @Override public ClientServerOperationPair transform(O clientOp, O serverOp) { ClientServerOperationPair result = new ClientServerOperationPair( @@ -10,5 +17,4 @@ public class OperationalTransformer void executed(Operation operation); +/** + * Can be registered on a {@link Replicable} and will receive notifications about the execution of + * operations. + * + * @author Axel Uhl (D043530) + * + * @param the type of the state to which the operations are applied + */ +public interface OperationExecutionListener { + void executed(OperationWithResult operation); } \ No newline at end of file diff --git a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/OperationWithResult.java b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/OperationWithResult.java new file mode 100755 index 00000000000..aa6015813cf --- /dev/null +++ b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/OperationWithResult.java @@ -0,0 +1,20 @@ +package com.sap.sailing.server.replication; + +import com.sap.sse.operationaltransformation.Operation; + +/** + * An operational transformation {@link Operation} is expected to return the target state after applying the operation. + * This operation type offers the possibility to let an operation return a result value from its + * {@link #internalApplyTo} method while in the context of operational transformation the operation will still return + * the target state. + * + * @author Axel Uhl (D043530) + * + */ +public interface OperationWithResult extends Operation { + /** + * Performs the actual operation, applying it to the toState service. The operation's result is + * returned. + */ + R internalApplyTo(S toState) throws Exception; +} diff --git a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/Replicable.java b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/Replicable.java index 8c95db2bdd4..22156165a39 100755 --- a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/Replicable.java +++ b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/Replicable.java @@ -41,7 +41,7 @@ import com.sap.sse.operationaltransformation.OperationWithTransformationSupport; * @author Axel Uhl (D043530) * */ -public interface Replicable> extends WithID { +public interface Replicable> extends WithID { /** * The name of the property to use in the properties dictionary in a call to * {@link BundleContext#registerService(Class, Object, java.util.Dictionary)} when registering a {@link Replicable}. @@ -61,11 +61,11 @@ public interface Replicable> extends WithID { * {@link OperationExecutionListener#executed(OperationWithTransformationSupport) notifies} all registered * operation execution listeners about the execution of the operation. */ - T apply(Operation operation); + T apply(OperationWithResult operation); - void addOperationExecutionListener(OperationExecutionListener listener); + void addOperationExecutionListener(OperationExecutionListener listener); - void removeOperationExecutionListener(OperationExecutionListener listener); + void removeOperationExecutionListener(OperationExecutionListener listener); void clearReplicaState() throws MalformedURLException, IOException, InterruptedException; diff --git a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Activator.java b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Activator.java index fdaca49715a..6d4ab4f25e1 100644 --- a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Activator.java +++ b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Activator.java @@ -111,7 +111,8 @@ public class Activator implements BundleActivator { logger.severe("Couldn't parse the replication port specification \""+exchangePortAsString+"\". Using default."); } replicationInstancesManager = new ReplicationInstancesManager(); - ReplicationService serverReplicationMasterService = new ReplicationServiceImpl(exchangeName, exchangeHost, exchangePort, replicationInstancesManager); + ReplicationService serverReplicationMasterService = new ReplicationServiceImpl<>( + exchangeName, exchangeHost, exchangePort, replicationInstancesManager); bundleContext.registerService(ReplicationService.class, serverReplicationMasterService, null); logger.info("Registered replication service "+serverReplicationMasterService+" using exchange name "+exchangeName+" on host "+exchangeHost); checkIfAutomaticReplicationShouldStart(serverReplicationMasterService, exchangeName); diff --git a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/HasReplicable.java b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/HasReplicable.java index e0cdc58425c..9ca17091f09 100755 --- a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/HasReplicable.java +++ b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/HasReplicable.java @@ -1,7 +1,8 @@ package com.sap.sailing.server.replication.impl; +import com.sap.sailing.server.replication.OperationWithResult; import com.sap.sailing.server.replication.Replicable; -public interface HasReplicable { - Replicable getReplicable(); +public interface HasReplicable> { + Replicable getReplicable(); } diff --git a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java index 6f933fcfba9..73019c160d3 100755 --- a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java +++ b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/ReplicationServiceImpl.java @@ -30,11 +30,11 @@ import com.rabbitmq.client.Channel; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.QueueingConsumer; import com.sap.sailing.server.replication.OperationExecutionListener; +import com.sap.sailing.server.replication.OperationWithResult; import com.sap.sailing.server.replication.Replicable; import com.sap.sailing.server.replication.ReplicationMasterDescriptor; import com.sap.sailing.server.replication.ReplicationService; import com.sap.sse.BuildVersion; -import com.sap.sse.operationaltransformation.Operation; import com.sap.sse.operationaltransformation.OperationWithTransformationSupport; /** @@ -52,14 +52,14 @@ import com.sap.sse.operationaltransformation.OperationWithTransformationSupport; * @author Frank Mittag, Axel Uhl (d043530) * */ -public class ReplicationServiceImpl implements ReplicationService, OperationExecutionListener, HasReplicable { +public class ReplicationServiceImpl> implements ReplicationService, OperationExecutionListener, HasReplicable { private static final Logger logger = Logger.getLogger(ReplicationServiceImpl.class.getName()); private final ReplicationInstancesManager replicationInstancesManager; - private final ServiceTracker, Replicable> racingEventServiceTracker; + private final ServiceTracker, Replicable> racingEventServiceTracker; - private final Replicable localService; + private final Replicable localService; /** * null, if this instance is not currently replicating from some master; the master's descriptor otherwise @@ -97,7 +97,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec */ private final UUID serverUUID; - private Replicator replicator; + private Replicator replicator; private Thread replicatorThread; /** @@ -178,7 +178,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec private ReplicationServiceImpl(String exchangeName, String exchangeHost, int exchangePort, final ReplicationInstancesManager replicationInstancesManager, - Replicable localService, boolean createRacingEventServiceTracker) throws IOException { + Replicable localService, boolean createRacingEventServiceTracker) throws IOException { timer = new Timer("ReplicationServiceImpl timer for delayed task sending"); this.replicationInstancesManager = replicationInstancesManager; replicaUUIDs = new HashMap(); @@ -206,12 +206,12 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec * the name of the exchange to which replicas can bind */ public ReplicationServiceImpl(String exchangeName, String exchangeHost, - final ReplicationInstancesManager replicationInstancesManager, Replicable localService) throws IOException { + final ReplicationInstancesManager replicationInstancesManager, Replicable localService) throws IOException { this(exchangeName, exchangeHost, 0, replicationInstancesManager, localService, /* create RacingEventServiceTracker */ false); } - protected ServiceTracker, Replicable> getRacingEventServiceTracker() { - return new ServiceTracker, Replicable>( + protected ServiceTracker, Replicable> getRacingEventServiceTracker() { + return new ServiceTracker, Replicable>( Activator.getDefaultContext(), Replicable.class.getName(), null); } @@ -242,8 +242,8 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec } @Override - public Replicable getReplicable() { - Replicable result; + public Replicable getReplicable() { + Replicable result; if (localService != null) { result = localService; } else { @@ -293,7 +293,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec * Schedules a single operation for broadcast. The operation is added to {@link #outboundBuffer}, and if not already scheduled, * a {@link #timer} is created and scheduled to send in {@link #TRANSMISSION_DELAY_MILLIS} milliseconds. */ - private void broadcastOperation(Operation operation) throws IOException { + private void broadcastOperation(OperationWithResult operation) throws IOException { // need to write the operations one by one, making sure the ObjectOutputStream always writes // identical objects again if required because they may have changed state in between ByteArrayOutputStream bos = new ByteArrayOutputStream(); @@ -426,7 +426,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec logger.info("Connection to exchange successful."); URL initialLoadURL = master.getInitialLoadURL(); logger.info("Initial load URL is "+initialLoadURL); - 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 replicatorThread = new Thread(replicator, "Replicator receiving from "+master.getMessagingHostname()+"/"+master.getExchangeName()); final Replicable replicable = getReplicable(); @@ -497,7 +497,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec * replicas by publishing it to the fan-out exchange. */ @Override - public void executed(Operation operation) { + public void executed(OperationWithResult operation) { try { broadcastOperation(operation); } catch (Exception e) { diff --git a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Replicator.java b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Replicator.java index c4e17ed35ed..3c3faa7541a 100755 --- a/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Replicator.java +++ b/java/com.sap.sse.replication/src/com/sap/sailing/server/replication/impl/Replicator.java @@ -21,9 +21,8 @@ import com.rabbitmq.client.ConsumerCancelledException; import com.rabbitmq.client.QueueingConsumer; import com.rabbitmq.client.QueueingConsumer.Delivery; import com.rabbitmq.client.ShutdownSignalException; +import com.sap.sailing.server.replication.OperationWithResult; import com.sap.sailing.server.replication.ReplicationMasterDescriptor; -import com.sap.sse.operationaltransformation.Operation; -import com.sap.sse.operationaltransformation.OperationWithTransformationSupport; /** * Receives {@link RacingEventServiceOperation}s through JMS and @@ -38,15 +37,15 @@ import com.sap.sse.operationaltransformation.OperationWithTransformationSupport; * @author Axel Uhl (d043530) * */ -public class Replicator implements Runnable { +public class Replicator> implements Runnable { private final static Logger logger = Logger.getLogger(Replicator.class.getName()); private static final long CHECK_INTERVAL_MILLIS = 2000; // how long (milliseconds) to pause before checking connection again private static final int CHECK_COUNT = 150; // how long to check, value is CHECK_INTERVAL second steps private final ReplicationMasterDescriptor master; - private final HasReplicable replicableProvider; - private final List> queue; + private final HasReplicable replicableProvider; + private final List> queue; private QueueingConsumer consumer; @@ -89,8 +88,8 @@ public class Replicator implements Runnable { * @param consumer * the RabbitMQ consumer from which to load messages */ - public Replicator(ReplicationMasterDescriptor master, HasReplicable racingEventServiceTracker, boolean startSuspended, QueueingConsumer consumer) { - this.queue = new ArrayList>(); + public Replicator(ReplicationMasterDescriptor master, HasReplicable racingEventServiceTracker, boolean startSuspended, QueueingConsumer consumer) { + this.queue = new ArrayList>(); this.master = master; this.replicableProvider = racingEventServiceTracker; this.suspended = startSuspended; @@ -147,7 +146,8 @@ public class Replicator implements Runnable { byte[] serializedOperation = (byte[]) ois.readObject(); ObjectInputStream operationOIS = replicableProvider.getReplicable() .createObjectInputStreamResolvingAgainstCache(new ByteArrayInputStream(serializedOperation)); - OperationWithTransformationSupport operation = (OperationWithTransformationSupport) operationOIS.readObject(); + @SuppressWarnings("unchecked") + OperationWithResult operation = (OperationWithResult) operationOIS.readObject(); operationCount++; operationsInMessage++; if (operationCount % 10000l == 0) { @@ -222,7 +222,7 @@ public class Replicator implements Runnable { * If the replicator is currently {@link #suspended}, the operation is queued, otherwise immediately applied to * the receiving replica. */ - private synchronized void applyOrQueue(OperationWithTransformationSupport operation) { + private synchronized void applyOrQueue(OperationWithResult operation) { if (suspended) { queue(operation); } else { @@ -230,7 +230,7 @@ public class Replicator implements Runnable { } } - private synchronized void apply(final Operation operation) { + private synchronized void apply(final OperationWithResult operation) { Runnable runnable = new Runnable() { @Override public void run() { @@ -244,7 +244,7 @@ public class Replicator implements Runnable { } } - private synchronized void queue(Operation operation) { + private synchronized void queue(OperationWithResult operation) { if (queue.isEmpty()) { notifyAll(); } @@ -262,8 +262,8 @@ public class Replicator implements Runnable { } private synchronized void applyQueue() { - for (Iterator> i=queue.iterator(); i.hasNext(); ) { - Operation operation = i.next(); + for (Iterator> i=queue.iterator(); i.hasNext(); ) { + OperationWithResult operation = i.next(); i.remove(); try { apply(operation);