bug4006: in RacingEventService wait for fully-initialized SecurityService and SharedSailingData;

introduced an observer pattern that is used in case the
ReplicationService is still in the process of starting, waiting for
initial loads to complete. Only when the initial loads have completed,
the service references will be returned.
This commit is contained in:
Axel Uhl
2020-03-06 18:48:30 +01:00
parent b87a5fb02b
commit 39608e3f6c
8 changed files with 99 additions and 38 deletions
@@ -40,7 +40,6 @@ import com.sap.sailing.domain.persistence.racelog.tracking.impl.GPSFixMongoHandl
import com.sap.sailing.domain.persistence.racelog.tracking.impl.GPSFixMovingMongoHandlerImpl;
import com.sap.sailing.domain.polars.PolarDataService;
import com.sap.sailing.domain.racelog.tracking.SensorFixStoreSupplier;
import com.sap.sailing.domain.sharedsailingdata.SharedSailingData;
import com.sap.sailing.domain.tracking.TrackedRegattaListener;
import com.sap.sailing.domain.windestimation.WindEstimationFactoryService;
import com.sap.sailing.server.RacingEventServiceMXBean;
@@ -61,6 +60,7 @@ import com.sap.sse.mail.MailService;
import com.sap.sse.mail.queue.MailQueue;
import com.sap.sse.mail.queue.impl.ExecutorMailQueue;
import com.sap.sse.osgi.CachedOsgiTypeBasedServiceFinderFactory;
import com.sap.sse.replication.FullyInitializedReplicableTracker;
import com.sap.sse.replication.Replicable;
import com.sap.sse.replication.ReplicationService;
import com.sap.sse.security.SecurityInitializationCustomizer;
@@ -110,9 +110,9 @@ public class Activator implements BundleActivator {
private ServiceTracker<MailService, MailService> mailServiceTracker;
private ServiceTracker<SecurityService, SecurityService> securityServiceTracker;
private FullyInitializedReplicableTracker<SecurityService> securityServiceTracker;
private ServiceTracker<SharedSailingData, SharedSailingData> sharedSailingDataTracker;
private FullyInitializedReplicableTracker<ReplicatingSharedSailingData> sharedSailingDataTracker;
private ServiceTracker<ReplicationService, ReplicationService> replicationServiceTracker;
@@ -131,7 +131,13 @@ public class Activator implements BundleActivator {
extenderBundleTracker = new ExtenderBundleTracker(context);
extenderBundleTracker.open();
mailServiceTracker = ServiceTrackerFactory.createAndOpen(context, MailService.class);
securityServiceTracker = ServiceTrackerFactory.createAndOpen(context, SecurityService.class);
replicationServiceTracker = ServiceTrackerFactory.createAndOpen(context, ReplicationService.class);
sharedSailingDataTracker = new FullyInitializedReplicableTracker<>(context, ReplicatingSharedSailingData.class,
/* customizer */ null, replicationServiceTracker);
sharedSailingDataTracker.open();
securityServiceTracker = new FullyInitializedReplicableTracker<>(context, SecurityService.class,
/* customizer */ null, replicationServiceTracker);
securityServiceTracker.open();
if (securityServiceTracker != null) {
new Thread("Racingevent wait for securityservice for migration thread") {
public void run() {
@@ -239,8 +245,6 @@ public class Activator implements BundleActivator {
// this code block is not run, and the test case can inject some other type of finder
// instead.
serviceFinderFactory = new CachedOsgiTypeBasedServiceFinderFactory(context);
sharedSailingDataTracker = ServiceTrackerFactory.createAndOpen(context, SharedSailingData.class);
replicationServiceTracker = ServiceTrackerFactory.createAndOpen(context, ReplicationService.class);
racingEventService = new RacingEventServiceImpl(clearPersistentCompetitors,
serviceFinderFactory, trackedRegattaListener, notificationService,
trackedRaceStatisticsCache, restoreTrackedRaces, securityServiceTracker,
@@ -18,8 +18,6 @@ import java.util.logging.Level;
import java.util.logging.Logger;
import java.util.stream.Collectors;
import org.osgi.util.tracker.ServiceTracker;
import com.sap.sailing.domain.abstractlog.AbstractLogEventAuthor;
import com.sap.sailing.domain.abstractlog.race.analyzing.impl.RaceLogResolver;
import com.sap.sailing.domain.abstractlog.race.state.ReadonlyRaceState;
@@ -97,13 +95,14 @@ import com.sap.sse.common.TimePoint;
import com.sap.sse.common.Timed;
import com.sap.sse.common.Util;
import com.sap.sse.common.Util.Pair;
import com.sap.sse.replication.FullyInitializedReplicableTracker;
import com.sap.sse.security.shared.impl.UserGroup;
public class CourseAndMarkConfigurationFactoryImpl implements CourseAndMarkConfigurationFactory {
private static final Logger logger = Logger.getLogger(CourseAndMarkConfigurationFactoryImpl.class.getName());
private final ServiceTracker<SharedSailingData, SharedSailingData> sharedSailingDataTracker;
private final FullyInitializedReplicableTracker<ReplicatingSharedSailingData> sharedSailingDataTracker;
private final SensorFixStore sensorFixStore;
/**
@@ -117,7 +116,7 @@ public class CourseAndMarkConfigurationFactoryImpl implements CourseAndMarkConfi
private final DomainFactory domainFactory;
public CourseAndMarkConfigurationFactoryImpl(
ServiceTracker<SharedSailingData, SharedSailingData> sharedSailingDataTracker,
FullyInitializedReplicableTracker<ReplicatingSharedSailingData> sharedSailingDataTracker,
SensorFixStore sensorFixStore, RaceLogResolver raceLogResolver, DomainFactory domainFactory) {
this.sharedSailingDataTracker = sharedSailingDataTracker;
this.domainFactory = domainFactory;
@@ -139,7 +138,15 @@ public class CourseAndMarkConfigurationFactoryImpl implements CourseAndMarkConfi
}
private SharedSailingData getSharedSailingData() {
return sharedSailingDataTracker.getService();
ReplicatingSharedSailingData result;
try {
result = sharedSailingDataTracker.getInitializedService(0);
} catch (InterruptedException e) {
logger.log(Level.SEVERE, "Interrupted while waiting for a fully initialized SharedSailingData service; "
+ "continuing with null, probably causing a NullPointerException along the way", e);
result = null;
}
return result;
}
private CourseTemplate resolveCourseTemplateSafe(CourseBase course) {
@@ -186,7 +186,6 @@ import com.sap.sailing.domain.regattalike.HasRegattaLike;
import com.sap.sailing.domain.regattalike.IsRegattaLike;
import com.sap.sailing.domain.regattalike.LeaderboardThatHasRegattaLike;
import com.sap.sailing.domain.regattalog.RegattaLogStore;
import com.sap.sailing.domain.sharedsailingdata.SharedSailingData;
import com.sap.sailing.domain.statistics.Statistics;
import com.sap.sailing.domain.tracking.AddResult;
import com.sap.sailing.domain.tracking.DynamicRaceDefinitionSet;
@@ -305,6 +304,7 @@ import com.sap.sse.pairinglist.PairingFrameProvider;
import com.sap.sse.pairinglist.PairingList;
import com.sap.sse.pairinglist.PairingListTemplate;
import com.sap.sse.pairinglist.PairingListTemplateFactory;
import com.sap.sse.replication.FullyInitializedReplicableTracker;
import com.sap.sse.replication.OperationExecutionListener;
import com.sap.sse.replication.OperationWithResult;
import com.sap.sse.replication.OperationWithResultWithIdWrapper;
@@ -556,10 +556,8 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
*/
private OperationsToMasterSendingQueue unsentOperationsToMasterSender;
private final ServiceTracker<SecurityService, SecurityService> securityServiceTracker;
private final FullyInitializedReplicableTracker<SecurityService> securityServiceTracker;
private final ServiceTracker<ReplicationService, ReplicationService> replicationServiceTracker;
private final CourseAndMarkConfigurationFactory courseAndMarkConfigurationFactory;
/**
@@ -628,8 +626,8 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
final TypeBasedServiceFinderFactory serviceFinderFactory, TrackedRegattaListenerManager trackedRegattaListener,
SailingNotificationService sailingNotificationService,
TrackedRaceStatisticsCache trackedRaceStatisticsCache, boolean restoreTrackedRaces,
ServiceTracker<SecurityService, SecurityService> securityServiceTracker,
ServiceTracker<SharedSailingData, SharedSailingData> sharedSailingDataTracker, ServiceTracker<ReplicationService, ReplicationService> replicationServiceTracker) {
FullyInitializedReplicableTracker<SecurityService> securityServiceTracker,
FullyInitializedReplicableTracker<ReplicatingSharedSailingData> sharedSailingDataTracker, ServiceTracker<ReplicationService, ReplicationService> replicationServiceTracker) {
this((final RaceLogAndTrackedRaceResolver raceLogResolver) -> {
return new ConstructorParameters() {
private final MongoObjectFactory mongoObjectFactory = PersistenceFactory.INSTANCE
@@ -761,13 +759,12 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
TypeBasedServiceFinderFactory serviceFinderFactory, TrackedRegattaListenerManager trackedRegattaListener,
SailingNotificationService sailingNotificationService,
TrackedRaceStatisticsCache trackedRaceStatisticsCache, boolean restoreTrackedRaces,
ServiceTracker<SecurityService, SecurityService> securityServiceTracker,
ServiceTracker<SharedSailingData, SharedSailingData> sharedSailingDataTracker,
FullyInitializedReplicableTracker<SecurityService> securityServiceTracker,
FullyInitializedReplicableTracker<ReplicatingSharedSailingData> sharedSailingDataTracker,
ServiceTracker<ReplicationService, ReplicationService> replicationServiceTracker) {
logger.info("Created " + this);
this.currentlyFillingFromInitialLoad = false;
this.securityServiceTracker = securityServiceTracker;
this.replicationServiceTracker = replicationServiceTracker;
this.numberOfTrackedRacesRestored = new AtomicInteger();
this.scoreCorrectionListenersByLeaderboard = new ConcurrentHashMap<>();
this.connectivityParametersByRace = new ConcurrentHashMap<>();
@@ -4686,7 +4683,7 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes
@Override
public SecurityService getSecurityService() {
try {
return securityServiceTracker.waitForService(0);
return securityServiceTracker.getInitializedService(0);
} catch (InterruptedException e) {
logger.severe("Interrupted while waiting for security service; returning null");
return null;
@@ -38,6 +38,7 @@ public class Activator implements BundleActivator {
PersistenceFactory.INSTANCE.getDefaultMongoObjectFactory(serviceFinderFactory), serviceFinderFactory,
securityServiceTracker);
registrations.add(context.registerService(SharedSailingData.class, sharedSailingData, /* properties */ null));
registrations.add(context.registerService(ReplicatingSharedSailingData.class, sharedSailingData, /* properties */ null));
registrations
.add(context.registerService(ClearStateTestSupport.class, sharedSailingData, /* properties */ null));
final Dictionary<String, String> replicableServiceProperties = new Hashtable<>();
@@ -21,6 +21,14 @@
<div class="mainContent">
<h2 class="releaseHeadline">Release Notes - Administration Console</h2>
<div class="innerContent">
<h2 class="articleSubheadline">March 2020</h2>
<ul class="bulletList">
<li>The "Replication" panel in the "Advanced" category now has improved information about the
replicables shown. An advanced replica information string may reveal more than just the address,
and the address is now resolved through any <tt>X-Forwarded-For</tt> header fields present
during replica registration, hence also working through load balancers and reverse proxies.
</li>
</ul>
<h2 class="articleSubheadline">February 2020</h2>
<ul class="bulletList">
<li>The "Local Server" tab now refreshes automatically when it comes into view.</li>
@@ -9,7 +9,9 @@ Bundle-RequiredExecutionEnvironment: JavaSE-1.8
Export-Package: com.sap.sse.replication,
com.sap.sse.replication.impl
Import-Package: com.rabbitmq.client;version="2.8.4",
org.json.simple
org.json.simple,
org.osgi.framework;version="1.8.0",
org.osgi.util.tracker;version="1.5.1"
Require-Bundle: com.sap.sse.operationaltransformation;bundle-version="1.0.0",
com.sap.sse.common;bundle-version="1.0.0",
com.sap.sse;bundle-version="1.0.0"
@@ -1,4 +1,7 @@
package com.sap.sse.replication.impl;
package com.sap.sse.replication;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import org.osgi.framework.BundleContext;
import org.osgi.framework.Filter;
@@ -6,8 +9,7 @@ import org.osgi.framework.ServiceReference;
import org.osgi.util.tracker.ServiceTracker;
import org.osgi.util.tracker.ServiceTrackerCustomizer;
import com.sap.sse.replication.Replicable;
import com.sap.sse.replication.ReplicationService;
import com.sap.sse.replication.ReplicationService.ReplicationStartingListener;
/**
* While a regular OSGi {@link ServiceTracker} would {@link ServiceTracker#waitForService(long) wait} for the service's
@@ -68,19 +70,60 @@ public class FullyInitializedReplicableTracker<R extends Replicable<?, ?>> exten
super(context, clazz, customizer);
this.replicationServiceTracker = replicationServiceTracker;
}
@Override
public R waitForService(long timeout) throws InterruptedException {
final R replicable = super.waitForService(timeout);
// TODO bug4006: continue here...
// waitForReplicationServiceToBeReady();
// waitForReplicableToBeFullyInitialized(replicable);
return replicable;
/**
* @param timeoutInMillis
* 0 means indefinite waiting time
* @return {@code true} if the {@link #replicationServiceTracker} is {@code null} or the {@link ReplicationService}
* was obtained successfully and has been in or has reached the state of not
* {@link ReplicationService#isReplicationStarting()} within the timeout provided. {@code false} otherwise.
*/
private boolean waitForReplicationToBeInitialized(long timeoutInMillis) throws InterruptedException {
final boolean result;
if (replicationServiceTracker != null) {
final CountDownLatch latch = new CountDownLatch(1); // counted down either by the direct check or the listener
final ReplicationService replicationService = replicationServiceTracker.waitForService(timeoutInMillis);
final ReplicationStartingListener replicationStartingListener = newIsReplicationStarting->{
if (!newIsReplicationStarting) {
latch.countDown();
}
};
replicationService.addReplicationStartingListener(replicationStartingListener);
if (!replicationService.isReplicationStarting()) {
latch.countDown();
}
result = latch.await(timeoutInMillis, TimeUnit.MILLISECONDS);
replicationService.removeReplicationStartingListener(replicationStartingListener);
} else {
result = true;
}
return result;
}
@Override
public R getService(ServiceReference<R> reference) {
// TODO Implement FullyInitializedReplicableTracker.getService(...)
return super.getService(reference);
/**
* Waits for one service object tracked to appear for {@code timeoutInMillis} milliseconds (see
* {@link #waitForService(long)}). If no such service object can be found before timing out, {@code null}
* is returned. Once a service object has been retrieved and a non-{@code null} {@link #replicationServiceTracker}
* has been provided at construction time, the {@link ReplicationService} is obtained from that tracker by
* waiting for it at least {@code timeoutInMillis} milliseconds and then
*
* @param timeoutInMillis
* 0 means indefinite wait time
* @return {@code null} if no service was obtained from the registry in the timeout specified or the replication did
* not reach a fully initialized state in the timeout specified.
*/
public R getInitializedService(long timeoutInMillis) throws InterruptedException {
final R service = waitForService(timeoutInMillis);
final R result;
if (service != null) {
if (waitForReplicationToBeInitialized(timeoutInMillis)) {
result = service;
} else {
result = null;
}
} else {
result = null;
}
return result;
}
}
@@ -48,7 +48,6 @@ import com.sap.sse.replication.ReplicablesProvider.ReplicableLifeCycleListener;
import com.sap.sse.replication.ReplicationMasterDescriptor;
import com.sap.sse.replication.ReplicationReceiver;
import com.sap.sse.replication.ReplicationService;
import com.sap.sse.replication.ReplicationService.ReplicationStartingListener;
import com.sap.sse.replication.ReplicationStatus;
import com.sap.sse.replication.persistence.MongoObjectFactory;
import com.sap.sse.util.HttpUrlConnectionHelper;