Merge branch 'bug5991'

This commit is contained in:
Axel Uhl
2024-04-09 23:48:28 +02:00
5 changed files with 171 additions and 5 deletions
@@ -21,6 +21,7 @@ import org.junit.rules.TestRule;
import org.junit.rules.Timeout;
import com.sap.sailing.domain.common.PassingInstruction;
import com.sap.sailing.domain.tractracadapter.DomainFactory;
import com.sap.sailing.domain.tractracadapter.Receiver;
import com.sap.sailing.domain.tractracadapter.TracTracConnectionConstants;
import com.sap.sailing.domain.tractracadapter.TracTracControlPoint;
@@ -110,7 +111,7 @@ public abstract class AbstractTracTracLiveTest extends StoredTrackBasedTest {
eventSubscriber = subscriberFactory.createEventSubscriber(race.getEvent());
raceSubscriber = subscriberFactory.createRaceSubscriber(race);
} else {
eventSubscriber = subscriberFactory.createEventSubscriber(race.getEvent(), liveUri, storedUri);
eventSubscriber = DomainFactory.INSTANCE.getOrCreateEventSubscriber(race.getEvent(), liveUri, storedUri);
raceSubscriber = subscriberFactory.createRaceSubscriber(race, liveUri, storedUri);
}
assertNotNull(race);
@@ -339,4 +339,13 @@ public interface DomainFactory {
*/
ControlPoint getExistingControlWithTwoMarks(Iterable<IControl> candidates, Mark first, Mark second);
/**
* Event subscribers created by this call are cached in this domain factory, using the three parameters as a compound
* caching key. Event subscribers found in the cache are returned by this method. The event subscriber returned will
* be a wrapper around the actual {@link IEventSubscriber}, managing the {@link IEventSubscriber#start()} and {@link IEventSubscriber#stop()}
* calls such that only the first {@link IEventSubscriber#start()} call is actually forwarded to the wrapper subscriber, and only
* the last {@link IEventSubscriber#stop()} call is forwarded. This is managed by an atomic counter that keeps track of the
* start/stop invocations.
*/
IEventSubscriber getOrCreateEventSubscriber(IEvent tractracEvent, URI liveURI, URI storedURI);
}
@@ -20,6 +20,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Consumer;
import java.util.logging.Level;
@@ -100,6 +101,7 @@ import com.sap.sse.common.Duration;
import com.sap.sse.common.TimePoint;
import com.sap.sse.common.Util;
import com.sap.sse.common.Util.Pair;
import com.sap.sse.common.Util.Triple;
import com.sap.sse.common.impl.AbstractColor;
import com.sap.sse.common.impl.DegreeBearingImpl;
import com.sap.sse.common.impl.MillisecondsDurationImpl;
@@ -161,11 +163,17 @@ public class DomainFactoryImpl implements DomainFactory {
* monitor.
*/
private final Set<Competitor> competitorsCurrentlyBeingMigrated;
/**
* The key consists of the {@link IEvent}, the live and the stored URI.
*/
private final ConcurrentMap<Triple<IEvent, URI, URI>, IEventSubscriber> eventSubscriberCache;
public DomainFactoryImpl(com.sap.sailing.domain.base.DomainFactory baseDomainFactory) {
this.baseDomainFactory = baseDomainFactory;
this.metadataParser = new MetadataParserImpl();
competitorsCurrentlyBeingMigrated = Collections.synchronizedSet(new HashSet<>());
eventSubscriberCache = new ConcurrentHashMap<>();
}
@Override
@@ -1082,4 +1090,15 @@ public class DomainFactoryImpl implements DomainFactory {
return new JSONServiceImpl(jsonURL, raceId, loadClientParams);
}
@Override
public IEventSubscriber getOrCreateEventSubscriber(IEvent tractracEvent, URI liveURI, URI storedURI) {
return eventSubscriberCache.computeIfAbsent(new Triple<>(tractracEvent, liveURI, storedURI), key->
{
try {
return new EventSubscriberWrapper(key.getA(), key.getB(), key.getC());
} catch (SubscriberInitializationException e) {
throw new RuntimeException(e);
}
});
}
}
@@ -0,0 +1,139 @@
package com.sap.sailing.domain.tractracadapter.impl;
import java.net.URI;
import java.util.concurrent.atomic.AtomicInteger;
import com.tractrac.model.lib.api.event.IEvent;
import com.tractrac.subscription.lib.api.IEventSubscriber;
import com.tractrac.subscription.lib.api.SubscriberInitializationException;
import com.tractrac.subscription.lib.api.SubscriptionLocator;
import com.tractrac.subscription.lib.api.competitor.ICompetitorsListener;
import com.tractrac.subscription.lib.api.control.IControlsListener;
import com.tractrac.subscription.lib.api.event.IConnectionStatusListener;
import com.tractrac.subscription.lib.api.event.IEventMessageListener;
import com.tractrac.subscription.lib.api.event.IServerTimeListener;
import com.tractrac.subscription.lib.api.race.IRacesListener;
import com.tractrac.subscription.lib.api.race.IStartStopTimesChangeListener;
/**
* A wrapper around a {@link IEventSubscriber} that can be shared across many {@link TracTracRaceTrackerImpl} instances
* that each invoke {@link #start} and {@link #stop()} symmetrically. This wrapper manages a counter (as an
* {@link AtomicInteger}) such that {@link #start} will only delegate to the instance wrapped if the counter is 0;
* likewise, {@link #stop} will delegate only if the counter goes to 0.
*
* @author Axel Uhl (d043530)
*
*/
public class EventSubscriberWrapper implements IEventSubscriber {
private IEventSubscriber delegate;
private final IEvent tractracEvent;
private final URI liveURI;
private final URI storedURI;
private int startCounter;
public EventSubscriberWrapper(IEvent tractracEvent, URI liveURI, URI storedURI) throws SubscriberInitializationException {
this.tractracEvent = tractracEvent;
this.liveURI = liveURI;
this.storedURI = storedURI;
this.startCounter = 0;
this.delegate = createEventSubscriber();
}
private IEventSubscriber createEventSubscriber() throws SubscriberInitializationException {
return SubscriptionLocator.getSusbcriberFactory().createEventSubscriber(tractracEvent, liveURI, storedURI);
}
@Override
public void subscribeConnectionStatus(IConnectionStatusListener listener) {
delegate.subscribeConnectionStatus(listener);
}
@Override
public void unsubscribeConnectionStatus(IConnectionStatusListener listener) {
delegate.unsubscribeConnectionStatus(listener);
}
@Override
public synchronized void start() {
if (startCounter++ == 0) {
delegate.start();
}
}
@Override
public synchronized void stop() {
if (--startCounter == 0) {
delegate.stop();
try {
delegate = createEventSubscriber();
} catch (SubscriberInitializationException e) {
throw new RuntimeException(e);
}
}
}
@Override
public boolean isRunning() {
return delegate.isRunning();
}
@Override
public void subscribeControls(IControlsListener listener) {
delegate.subscribeControls(listener);
}
@Override
public void unsubscribeControls(IControlsListener listener) {
delegate.unsubscribeControls(listener);
}
@Override
public void subscribeEventTimesChanges(IStartStopTimesChangeListener listener) {
delegate.subscribeEventTimesChanges(listener);
}
@Override
public void unsubscribeEventTimesChanges(IStartStopTimesChangeListener listener) {
delegate.unsubscribeEventTimesChanges(listener);
}
@Override
public void subscribeEventMessages(IEventMessageListener listener) {
delegate.subscribeEventMessages(listener);
}
@Override
public void unsubscribeEventMessages(IEventMessageListener listener) {
delegate.unsubscribeEventMessages(listener);
}
@Override
public void subscribeServerTime(IServerTimeListener serverTimeListener) {
delegate.subscribeServerTime(serverTimeListener);
}
@Override
public void unsubscribeServerTime(IServerTimeListener serverTimeListener) {
delegate.unsubscribeServerTime(serverTimeListener);
}
@Override
public void subscribeRaces(IRacesListener listener) {
delegate.subscribeRaces(listener);
}
@Override
public void unsubscribeRaces(IRacesListener listener) {
delegate.unsubscribeRaces(listener);
}
@Override
public void subscribeCompetitors(ICompetitorsListener listener) {
delegate.subscribeCompetitors(listener);
}
@Override
public void unsubscribeCompetitors(ICompetitorsListener listener) {
delegate.unsubscribeCompetitors(listener);
}
}
@@ -65,7 +65,6 @@ import com.tractrac.model.lib.api.event.IRaceCompetitor;
import com.tractrac.model.lib.api.route.IControl;
import com.tractrac.subscription.lib.api.IEventSubscriber;
import com.tractrac.subscription.lib.api.IRaceSubscriber;
import com.tractrac.subscription.lib.api.ISubscriberFactory;
import com.tractrac.subscription.lib.api.SubscriberInitializationException;
import com.tractrac.subscription.lib.api.SubscriptionLocator;
import com.tractrac.subscription.lib.api.event.IConnectionStatusListener;
@@ -350,8 +349,7 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl<RaceTrackin
+ (endOfTracking != null ? endOfTracking.asMillis() : "n/a"));
// Initialize data controller using live and stored data sources
ISubscriberFactory subscriberFactory = SubscriptionLocator.getSusbcriberFactory();
eventSubscriber = subscriberFactory.createEventSubscriber(tractracEvent, liveURI, effectiveStoredURI);
eventSubscriber = domainFactory.getOrCreateEventSubscriber(tractracEvent, liveURI, effectiveStoredURI);
if (useOfficialEventsToUpdateRaceLog) {
reconciler = new RaceAndCompetitorStatusWithRaceLogReconciler(domainFactory, raceLogResolver, tractracRace);
} else {
@@ -401,7 +399,7 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl<RaceTrackin
eventSubscriber.subscribeRaces(racesListener);
// Start live and stored data streams
final Regatta effectiveRegatta;
raceSubscriber = subscriberFactory.createRaceSubscriber(tractracRace, liveURI, effectiveStoredURI);
raceSubscriber = SubscriptionLocator.getSusbcriberFactory().createRaceSubscriber(tractracRace, liveURI, effectiveStoredURI);
raceSubscriber.subscribeConnectionStatus(this);
// Try to find a pre-associated event based on the Race ID
if (regatta == null) {