diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/AbstractTracTracLiveTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/AbstractTracTracLiveTest.java index 7a948df9cbc..9fc3ab449c3 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/AbstractTracTracLiveTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/AbstractTracTracLiveTest.java @@ -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); diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/DomainFactory.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/DomainFactory.java index dccce036d84..0a8993cfe8b 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/DomainFactory.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/DomainFactory.java @@ -339,4 +339,13 @@ public interface DomainFactory { */ ControlPoint getExistingControlWithTwoMarks(Iterable 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); } diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java index d6182d21b8e..c6c2da74ce4 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/DomainFactoryImpl.java @@ -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 competitorsCurrentlyBeingMigrated; + + /** + * The key consists of the {@link IEvent}, the live and the stored URI. + */ + private final ConcurrentMap, 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); + } + }); + } } diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/EventSubscriberWrapper.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/EventSubscriberWrapper.java new file mode 100644 index 00000000000..7de64c6bcac --- /dev/null +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/EventSubscriberWrapper.java @@ -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); + } +} diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/TracTracRaceTrackerImpl.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/TracTracRaceTrackerImpl.java index c81aea5ab41..004a2820889 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/TracTracRaceTrackerImpl.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/TracTracRaceTrackerImpl.java @@ -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