From c283d9276a6fd417a9577ed4c925a07bdbf4c65d Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 10:51:47 +0100 Subject: [PATCH 01/12] Introduces oldStatus parameter for state change listener polar service can now listen for state changes (loading -> x) --- .../MarkPassingUpdateListener.java | 2 +- .../domain/tracking/RaceChangeListener.java | 2 +- .../impl/AbstractRaceChangeListener.java | 2 +- .../tracking/impl/DynamicTrackedRaceImpl.java | 7 +++-- .../TrackBasedEstimationWindTrackImpl.java | 2 +- .../domain/tracking/impl/TrackedLegImpl.java | 2 +- .../sap/sailing/polars/PolarDataService.java | 2 +- .../polars/impl/PolarDataServiceImpl.java | 2 +- .../sailing/polars/mining/PolarDataMiner.java | 31 ++++++------------- .../server/impl/RacingEventServiceImpl.java | 9 +++--- 10 files changed, 26 insertions(+), 35 deletions(-) diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/markpassingcalculation/MarkPassingUpdateListener.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/markpassingcalculation/MarkPassingUpdateListener.java index c1636f24013..eb9e97dbffa 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/markpassingcalculation/MarkPassingUpdateListener.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/markpassingcalculation/MarkPassingUpdateListener.java @@ -45,7 +45,7 @@ public class MarkPassingUpdateListener extends AbstractRaceChangeListener { } @Override - public void statusChanged(TrackedRaceStatus newStatus) { + public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) { if (newStatus.getStatus() == TrackedRaceStatusEnum.FINISHED) { queue.add(end); } diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/RaceChangeListener.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/RaceChangeListener.java index 6dbefd5a3ca..df09d90ece0 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/RaceChangeListener.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/RaceChangeListener.java @@ -50,5 +50,5 @@ public interface RaceChangeListener extends CourseListener { void windSourcesToExcludeChanged(Iterable windSourcesToExclude); - void statusChanged(TrackedRaceStatus newStatus); + void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus); } diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/AbstractRaceChangeListener.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/AbstractRaceChangeListener.java index 30feacbaf26..4dbf64a1b70 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/AbstractRaceChangeListener.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/AbstractRaceChangeListener.java @@ -40,7 +40,7 @@ public abstract class AbstractRaceChangeListener implements RaceChangeListener { } @Override - public void statusChanged(TrackedRaceStatus newStatus) { + public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) { defaultAction(); } diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRaceImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRaceImpl.java index 675f70a917b..f8def06f984 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRaceImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRaceImpl.java @@ -147,8 +147,9 @@ DynamicTrackedRace, GPSTrackListener { @Override public void setStatus(TrackedRaceStatus newStatus) { + TrackedRaceStatus oldStatus = getStatus(); super.setStatus(newStatus); - notifyListeners(newStatus); + notifyListeners(newStatus, oldStatus); } @Override @@ -415,8 +416,8 @@ DynamicTrackedRace, GPSTrackListener { notifyListeners(listener -> listener.competitorPositionChanged(fix, competitor)); } - private void notifyListeners(TrackedRaceStatus status) { - notifyListeners(listener -> listener.statusChanged(status)); + private void notifyListeners(TrackedRaceStatus status, TrackedRaceStatus oldStatus) { + notifyListeners(listener -> listener.statusChanged(status, oldStatus)); } private void notifyListeners(Wind wind, WindSource windSource) { diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackBasedEstimationWindTrackImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackBasedEstimationWindTrackImpl.java index 0de72df2936..8ff7fcd2c57 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackBasedEstimationWindTrackImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackBasedEstimationWindTrackImpl.java @@ -667,7 +667,7 @@ public class TrackBasedEstimationWindTrackImpl extends VirtualWindTrackImpl { } @Override - public void statusChanged(TrackedRaceStatus newStatus) { + public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) { // This virtual wind track's cache can cope with an empty cache after the LOADING phase and populates the // cache // upon request. Invalidation happens also during the LOADING phase, preserving the cache's invariant. diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedLegImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedLegImpl.java index 9a06c8452ab..37e27570806 100755 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedLegImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/TrackedLegImpl.java @@ -269,7 +269,7 @@ public class TrackedLegImpl implements TrackedLeg { * no-op; the leg doesn't mind the tracked race's status being updated */ @Override - public void statusChanged(TrackedRaceStatus newStatus) { + public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) { } /** diff --git a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/PolarDataService.java b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/PolarDataService.java index 66e493c184b..4f230462a71 100755 --- a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/PolarDataService.java +++ b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/PolarDataService.java @@ -194,5 +194,5 @@ public interface PolarDataService { */ PolynomialFunction getAngleRegressionFunction(BoatClass boatClass, LegType legType, Tack tack) throws NotEnoughDataHasBeenAddedException; - void raceFinishedTracking(TrackedRace race); + void raceFinishedLoading(TrackedRace race); } diff --git a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/impl/PolarDataServiceImpl.java b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/impl/PolarDataServiceImpl.java index 8ae96c222b0..716524b4e05 100755 --- a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/impl/PolarDataServiceImpl.java +++ b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/impl/PolarDataServiceImpl.java @@ -148,7 +148,7 @@ public class PolarDataServiceImpl implements PolarDataService { } @Override - public void raceFinishedTracking(TrackedRace race) { + public void raceFinishedLoading(TrackedRace race) { polarDataMiner.raceFinishedTracking(race); } } diff --git a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java index c2c1ceaac29..12f4e2e4497 100755 --- a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java +++ b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java @@ -43,7 +43,6 @@ import com.sap.sailing.domain.tracking.TrackedRace; import com.sap.sailing.polars.impl.CubicEquation; import com.sap.sailing.polars.regression.NotEnoughDataHasBeenAddedException; import com.sap.sse.common.Util.Pair; -import com.sap.sse.common.impl.MillisecondsTimePoint; import com.sap.sse.datamining.components.Processor; import com.sap.sse.datamining.data.ClusterGroup; import com.sap.sse.datamining.functions.Function; @@ -58,7 +57,7 @@ public class PolarDataMiner { private static final int THREAD_POOL_SIZE = Math.max((int) (Runtime.getRuntime().availableProcessors() * (3.0/4.0)), 3); private final ThreadPoolExecutor executor = createExecutor(); private final PolarSheetGenerationSettings backendPolarSheetGenerationSettings; - private final Map> fixesForReplayRacesWhichAreStillLoading = new HashMap<>(); + private final Map> fixesForRacesWhichAreStillLoading = new HashMap<>(); private final Queue fixQueue = new ConcurrentLinkedQueue(); @@ -188,22 +187,22 @@ public class PolarDataMiner { public void addFix(GPSFixMoving fix, Competitor competitor, TrackedRace trackedRace) { GPSFixMovingWithOriginInfo fixWithOriginInfo = new GPSFixMovingWithOriginInfo(fix, trackedRace, competitor); - if (isReplayAndFinishedOrLive(trackedRace)) { - processFix(trackedRace, fixWithOriginInfo); - } else { + if (trackedRace.getStatus().getStatus() == TrackedRaceStatusEnum.LOADING) { /* * logger.info("Received fix for replay race which has not finished loading. Queuing. " + * (trackedRace.getRace() != null ? trackedRace.getRace().getName() : trackedRace.getRaceIdentifier() * .getRaceName())); */ - synchronized (fixesForReplayRacesWhichAreStillLoading) { - Set fixes = fixesForReplayRacesWhichAreStillLoading.get(trackedRace); + synchronized (fixesForRacesWhichAreStillLoading) { + Set fixes = fixesForRacesWhichAreStillLoading.get(trackedRace); if (fixes == null) { fixes = new HashSet<>(); - fixesForReplayRacesWhichAreStillLoading.put(trackedRace, fixes); + fixesForRacesWhichAreStillLoading.put(trackedRace, fixes); } fixes.add(fixWithOriginInfo); } + } else { + processFix(trackedRace, fixWithOriginInfo); } } @@ -221,16 +220,6 @@ public class PolarDataMiner { } } - private boolean isReplayAndFinishedOrLive(TrackedRace trackedRace) { - boolean isLive = trackedRace.isLive(new MillisecondsTimePoint(System.currentTimeMillis())); - boolean loadingFinished = false; - if (trackedRace.getStatus().getStatus() == TrackedRaceStatusEnum.FINISHED) { - loadingFinished = true; - } - boolean replayRaceAndFinished = !isLive && loadingFinished; - return replayRaceAndFinished || isLive; - } - public boolean isCurrentlyActiveAndOrHasQueue() { boolean isActive = executor.getActiveCount() > 0; boolean hasQueue = executor.getQueue().size() > 0; @@ -399,11 +388,11 @@ public class PolarDataMiner { public void raceFinishedTracking(TrackedRace race) { Set fixes = null; - synchronized (fixesForReplayRacesWhichAreStillLoading) { - fixes = fixesForReplayRacesWhichAreStillLoading.remove(race); + synchronized (fixesForRacesWhichAreStillLoading) { + fixes = fixesForRacesWhichAreStillLoading.remove(race); } if (fixes != null) { - logger.info("All queued fixes for newly completed race will process now. " + (race.getRace() != null ? race + logger.info("All queued fixes for newly loaded race will process now. " + (race.getRace() != null ? race .getRace().getName() : race.getRaceIdentifier().getRaceName())); for (GPSFixMovingWithOriginInfo fix : fixes) { processFix(race, fix); 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 b5379a259e6..fa8bc632793 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 @@ -1540,9 +1540,10 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes } @Override - public void statusChanged(TrackedRaceStatus newStatus) { - if (newStatus.getStatus() == TrackedRaceStatusEnum.FINISHED) { - polarDataService.raceFinishedTracking(race); + public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) { + if (oldStatus.getStatus() == TrackedRaceStatusEnum.LOADING + && newStatus.getStatus() != TrackedRaceStatusEnum.LOADING) { + polarDataService.raceFinishedLoading(race); } } @@ -1617,7 +1618,7 @@ public class RacingEventServiceImpl implements RacingEventService, ClearStateTes } @Override - public void statusChanged(TrackedRaceStatus newStatus) { + public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) { replicate(new UpdateTrackedRaceStatus(getRaceIdentifier(), newStatus)); } From 6340e59026f821733b6616b960e4d9bbf4fbc076 Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 12:58:53 +0100 Subject: [PATCH 02/12] Introduced callbacks for tractrac receivers for true loading completion --- .../domain/test/PositionConversionTest.java | 5 ++ .../test/ReceiveMarkPassingDataTest.java | 5 ++ .../domain/test/RouteAssemblyTest.java | 5 ++ .../LoadingQueueDoneCallBack.java | 7 +++ .../domain/tractracadapter/Receiver.java | 2 + .../impl/AbstractReceiverWithQueue.java | 25 ++++++-- .../impl/TracTracRaceTrackerImpl.java | 58 +++++++++++++++---- 7 files changed, 91 insertions(+), 16 deletions(-) create mode 100644 java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/PositionConversionTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/PositionConversionTest.java index 7fc54a77afe..4e02d9a158a 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/PositionConversionTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/PositionConversionTest.java @@ -11,6 +11,7 @@ import org.junit.Test; import com.sap.sailing.domain.common.Position; import com.sap.sailing.domain.tractracadapter.DomainFactory; +import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; import com.tractrac.model.lib.api.data.IPosition; import com.tractrac.model.lib.api.route.IControl; @@ -73,6 +74,10 @@ public class PositionConversionTest extends AbstractTracTracLiveTest { } }); } + + @Override + public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { + } }; addListenersForStoredDataAndStartController(Collections.singleton(receiver)); synchronized (semaphor) { diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveMarkPassingDataTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveMarkPassingDataTest.java index 05e5feb33cc..fa78cfaab45 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveMarkPassingDataTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/ReceiveMarkPassingDataTest.java @@ -21,6 +21,7 @@ import com.sap.sailing.domain.tracking.impl.DynamicTrackedRegattaImpl; import com.sap.sailing.domain.tracking.impl.EmptyWindStore; import com.sap.sailing.domain.tracking.impl.GPSFixMovingImpl; import com.sap.sailing.domain.tractracadapter.DomainFactory; +import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; import com.sap.sailing.domain.tractracadapter.ReceiverType; import com.sap.sailing.domain.tractracadapter.impl.ControlPointAdapter; @@ -88,6 +89,10 @@ public class ReceiveMarkPassingDataTest extends AbstractTracTracLiveTest { } }); } + + @Override + public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { + } }; List receivers = new ArrayList(); receivers.add(receiver); diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/RouteAssemblyTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/RouteAssemblyTest.java index f2c16669e08..4776abbc85a 100755 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/RouteAssemblyTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/test/RouteAssemblyTest.java @@ -11,6 +11,7 @@ import org.junit.Test; import com.sap.sailing.domain.base.Course; import com.sap.sailing.domain.tractracadapter.DomainFactory; +import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; import com.sap.sse.common.Util; import com.tractrac.model.lib.api.route.IControlRoute; @@ -67,6 +68,10 @@ public class RouteAssemblyTest extends AbstractTracTracLiveTest { } }); } + + @Override + public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { + } }; addListenersForStoredDataAndStartController(Collections.singleton(receiver)); synchronized (semaphor) { diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java new file mode 100644 index 00000000000..6c78125e07a --- /dev/null +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java @@ -0,0 +1,7 @@ +package com.sap.sailing.domain.tractracadapter; + +public interface LoadingQueueDoneCallBack { + + void loadingQueueDone(Receiver receiver); + +} diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java index 9bfa996cec7..5fa47b124ef 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java @@ -34,4 +34,6 @@ public interface Receiver { * Waits until this received has stopped, but no longer than timeout milliseconds */ void join(long timeoutInMilliseconds) throws InterruptedException; + + void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback); } diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java index aaadc77479f..29edd8e5f99 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java @@ -1,6 +1,7 @@ package com.sap.sailing.domain.tractracadapter.impl; -import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.LinkedBlockingDeque; import java.util.concurrent.TimeUnit; import java.util.logging.Logger; @@ -11,9 +12,10 @@ import com.sap.sailing.domain.tracking.RaceTracker; import com.sap.sailing.domain.tracking.TrackedRace; import com.sap.sailing.domain.tracking.TrackedRegatta; import com.sap.sailing.domain.tractracadapter.DomainFactory; +import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; -import com.tractrac.model.lib.api.event.IEvent; import com.sap.sse.common.Util; +import com.tractrac.model.lib.api.event.IEvent; import com.tractrac.model.lib.api.event.IRace; import com.tractrac.subscription.lib.api.IEventSubscriber; import com.tractrac.subscription.lib.api.IRaceSubscriber; @@ -29,7 +31,7 @@ import com.tractrac.subscription.lib.api.IRaceSubscriber; public abstract class AbstractReceiverWithQueue implements Runnable, Receiver { private static Logger logger = Logger.getLogger(AbstractReceiverWithQueue.class.getName()); - private final LinkedBlockingQueue> queue; + private final LinkedBlockingDeque> queue; private final DomainFactory domainFactory; private final IEvent tractracEvent; private final IEventSubscriber eventSubscriber; @@ -37,6 +39,7 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re private final DynamicTrackedRegatta trackedRegatta; private final Simulator simulator; private final Thread thread; + private final ConcurrentHashMap, LoadingQueueDoneCallBack> loadingQueueDoneCallBacks; /** * used by {@link #stopAfterNotReceivingEventsForSomeTime(long)} and {@link #run()} to check if an event was received @@ -53,8 +56,9 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re this.trackedRegatta = trackedRegatta; this.domainFactory = domainFactory; this.simulator = simulator; - this.queue = new LinkedBlockingQueue>(); + this.queue = new LinkedBlockingDeque>(); this.thread = new Thread(this, getClass().getName()); + this.loadingQueueDoneCallBacks = new ConcurrentHashMap<>(); } protected IEventSubscriber getEventSubscriber() { @@ -136,6 +140,10 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re if (!isStopEvent(event)) { handleEvent(event); } + LoadingQueueDoneCallBack callBack = loadingQueueDoneCallBacks.remove(event); + if (callBack != null) { + callBack.loadingQueueDone(this); + } } catch (InterruptedException e) { e.printStackTrace(); } @@ -180,4 +188,13 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re } return result; } + + @Override + public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { + if (queue.isEmpty()) { + callback.loadingQueueDone(this); + } else { + loadingQueueDoneCallBacks.put(queue.getLast(), callback); + } + } } 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 84f579a59c9..9fb6347369e 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 @@ -43,6 +43,7 @@ import com.sap.sailing.domain.tracking.WindTrack; import com.sap.sailing.domain.tracking.impl.EmptyWindStore; import com.sap.sailing.domain.tracking.impl.TrackedRaceStatusImpl; import com.sap.sailing.domain.tractracadapter.DomainFactory; +import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; import com.sap.sailing.domain.tractracadapter.TracTracConnectionConstants; import com.sap.sailing.domain.tractracadapter.TracTracRaceTracker; @@ -577,7 +578,25 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements } } logger.info("Stored data progress in tracker "+getID()+" for race(s) "+getRaces()+": "+progress); - lastStatus = new TrackedRaceStatusImpl(progress==1.0 ? TrackedRaceStatusEnum.TRACKING : TrackedRaceStatusEnum.LOADING, progress); + if (progress==1.0) { + LoadingQueueDoneCallBack callBackHandler = new LoadingQueueDoneCallBack() { + + private Set receivers = new HashSet<>(TracTracRaceTrackerImpl.this.receivers); + + @Override + public void loadingQueueDone(Receiver receiver) { + receivers.remove(receiver); + if (this.receivers.isEmpty()) { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.TRACKING, progress); + updateStatusOfTrackedRaces(); + } + } + }; + for (Receiver receiver : receivers) { + receiver.callBackWhenLoadingQueueIsDone(callBackHandler); + } + } + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, progress); lastProgressPerID.put(getID(), new Util.Pair(counter, progress)); updateStatusOfTrackedRaces(); } @@ -621,20 +640,35 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements @Override public void stopped(Object o) { logger.info("stopped TracTrac tracking in tracker "+getID()+" for "+getRaces()+" while in status "+lastStatus); - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, 1.0); - updateStatusOfTrackedRaces(); - if (!stopped) { - try { - for (RaceDefinition race : getRaces()) { - // See also bug 1517; with TracAPI we assume that when stopped(IEvent) is called by the TracAPI then - // all subscriptions have received all their data and it's therefore safe to stop all subscriptions - // at this point without missing any data. - trackedRegattaRegistry.stopTracking(regatta, race); + LoadingQueueDoneCallBack callBackHandler = new LoadingQueueDoneCallBack() { + + private Set receivers = new HashSet<>(TracTracRaceTrackerImpl.this.receivers); + + @Override + public void loadingQueueDone(Receiver receiver) { + receivers.remove(receiver); + if (this.receivers.isEmpty()) { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, 1.0); + updateStatusOfTrackedRaces(); + if (!stopped) { + try { + for (RaceDefinition race : getRaces()) { + // See also bug 1517; with TracAPI we assume that when stopped(IEvent) is called by the TracAPI then + // all subscriptions have received all their data and it's therefore safe to stop all subscriptions + // at this point without missing any data. + trackedRegattaRegistry.stopTracking(regatta, race); + } + } catch (InterruptedException | IOException e) { + logger.log(Level.INFO, "Interrupted while trying to stop tracker "+this, e); + } + } } - } catch (InterruptedException | IOException e) { - logger.log(Level.INFO, "Interrupted while trying to stop tracker "+this, e); } + }; + for (Receiver receiver : receivers) { + receiver.callBackWhenLoadingQueueIsDone(callBackHandler); } + } @Override From 14a403ccb69d8d81dd30399e7e0b4785e6d87ecd Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 13:50:46 +0100 Subject: [PATCH 03/12] Moved prefilter into polar mining processor chain --- .../impl/TracTracRaceTrackerImpl.java | 38 +++++++------ .../sailing/polars/mining/PolarDataMiner.java | 55 ++++++++++++++----- 2 files changed, 61 insertions(+), 32 deletions(-) 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 9fb6347369e..44317ae7b03 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 @@ -585,10 +585,12 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements @Override public void loadingQueueDone(Receiver receiver) { - receivers.remove(receiver); - if (this.receivers.isEmpty()) { - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.TRACKING, progress); - updateStatusOfTrackedRaces(); + synchronized (this.receivers) { + receivers.remove(receiver); + if (this.receivers.isEmpty()) { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.TRACKING, progress); + updateStatusOfTrackedRaces(); + } } } }; @@ -646,20 +648,22 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements @Override public void loadingQueueDone(Receiver receiver) { - receivers.remove(receiver); - if (this.receivers.isEmpty()) { - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, 1.0); - updateStatusOfTrackedRaces(); - if (!stopped) { - try { - for (RaceDefinition race : getRaces()) { - // See also bug 1517; with TracAPI we assume that when stopped(IEvent) is called by the TracAPI then - // all subscriptions have received all their data and it's therefore safe to stop all subscriptions - // at this point without missing any data. - trackedRegattaRegistry.stopTracking(regatta, race); + synchronized (this.receivers) { + receivers.remove(receiver); + if (this.receivers.isEmpty()) { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, 1.0); + updateStatusOfTrackedRaces(); + if (!stopped) { + try { + for (RaceDefinition race : getRaces()) { + // See also bug 1517; with TracAPI we assume that when stopped(IEvent) is called by the TracAPI then + // all subscriptions have received all their data and it's therefore safe to stop all subscriptions + // at this point without missing any data. + trackedRegattaRegistry.stopTracking(regatta, race); + } + } catch (InterruptedException | IOException e) { + logger.log(Level.INFO, "Interrupted while trying to stop tracker "+this, e); } - } catch (InterruptedException | IOException e) { - logger.log(Level.INFO, "Interrupted while trying to stop tracker "+this, e); } } } diff --git a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java index 12f4e2e4497..d23ad8855bb 100755 --- a/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java +++ b/java/com.sap.sailing.polars/src/com/sap/sailing/polars/mining/PolarDataMiner.java @@ -43,6 +43,7 @@ import com.sap.sailing.domain.tracking.TrackedRace; import com.sap.sailing.polars.impl.CubicEquation; import com.sap.sailing.polars.regression.NotEnoughDataHasBeenAddedException; import com.sap.sse.common.Util.Pair; +import com.sap.sse.datamining.components.FilterCriterion; import com.sap.sse.datamining.components.Processor; import com.sap.sse.datamining.data.ClusterGroup; import com.sap.sse.datamining.functions.Function; @@ -78,14 +79,14 @@ public class PolarDataMiner { synchronized (fixQueue) { while (this.getQueue().size() < (EXECUTOR_QUEUE_SIZE / 10) && !fixQueue.isEmpty()) { GPSFixMovingWithOriginInfo fix = fixQueue.poll(); - enrichingProcessor.processElement(fix); + preFilteringProcessor.processElement(fix); } } } }; } - private AbstractEnrichingProcessor enrichingProcessor; + /** * This processor keeps track of the moving average of the speed values and the average angle for each course (legtype tack combination) @@ -99,6 +100,7 @@ public class PolarDataMiner { private CubicRegressionPerCourseProcessor cubicRegressionPerCourseProcessor; private SpeedRegressionPerAngleClusterProcessor speedRegressionPerAngleClusterProcessor; + private ParallelFilteringProcessor preFilteringProcessor; public PolarDataMiner() { this(PolarSheetGenerationSettingsImpl.createBackendPolarSettings()); @@ -169,15 +171,42 @@ public class PolarDataMiner { .asList(filteringProcessor); - enrichingProcessor = new AbstractEnrichingProcessor( - GPSFixMovingWithOriginInfo.class, GPSFixMovingWithPolarContext.class, executor, enrichingResultReceivers) { + AbstractEnrichingProcessor enrichingProcessor = new AbstractEnrichingProcessor( + GPSFixMovingWithOriginInfo.class, GPSFixMovingWithPolarContext.class, executor, + enrichingResultReceivers) { @Override protected GPSFixMovingWithPolarContext enrich(GPSFixMovingWithOriginInfo element) { - return new GPSFixMovingWithPolarContext(element.getFix(), element.getTrackedRace(), element.getCompetitor(), - speedClusterGroup, angleClusterGroup); + GPSFixMovingWithPolarContext result = null; + result = new GPSFixMovingWithPolarContext(element.getFix(), element.getTrackedRace(), + element.getCompetitor(), speedClusterGroup, angleClusterGroup); + return result; } }; + + Collection> preFilterResultReceivers = Arrays + .asList(enrichingProcessor); + + preFilteringProcessor = new ParallelFilteringProcessor( + GPSFixMovingWithOriginInfo.class, executor, preFilterResultReceivers, new FilterCriterion() { + + @Override + public boolean matches(GPSFixMovingWithOriginInfo element) { + boolean result = false; + if (PolarFixFilterCriteria.isInLeadingCompetitors(element.getTrackedRace(), element.getCompetitor(), + backendPolarSheetGenerationSettings.getPctOfLeadingCompetitorsToInclude())) { + result = true; + } + return result; + } + + @Override + public Class getElementType() { + return GPSFixMovingWithOriginInfo.class; + } + }); + + } @@ -207,15 +236,11 @@ public class PolarDataMiner { } private void processFix(TrackedRace trackedRace, GPSFixMovingWithOriginInfo fixWithOriginInfo) { - // Pre Filter for performance. We don't need to run wind estimation for this filter - if (PolarFixFilterCriteria.isInLeadingCompetitors(trackedRace, fixWithOriginInfo.getCompetitor(), - backendPolarSheetGenerationSettings.getPctOfLeadingCompetitorsToInclude())) { - synchronized (fixQueue) { - if (executor.getQueue().size() >= EXECUTOR_QUEUE_SIZE / 10) { - fixQueue.add(fixWithOriginInfo); - } else { - enrichingProcessor.processElement(fixWithOriginInfo); - } + synchronized (fixQueue) { + if (executor.getQueue().size() >= EXECUTOR_QUEUE_SIZE / 10) { + fixQueue.add(fixWithOriginInfo); + } else { + preFilteringProcessor.processElement(fixWithOriginInfo); } } } From 7cb3e8d4ef263e63ef01edde97e1fd1ac4348778 Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 15:00:14 +0100 Subject: [PATCH 04/12] Added comments for new receiver callback --- .../tractracadapter/LoadingQueueDoneCallBack.java | 9 +++++++++ .../sailing/domain/tractracadapter/Receiver.java | 15 +++++++++++++++ 2 files changed, 24 insertions(+) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java index 6c78125e07a..60b5df94a02 100644 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java @@ -1,7 +1,16 @@ package com.sap.sailing.domain.tractracadapter; +import com.sap.sailing.domain.tractracadapter.impl.AbstractReceiverWithQueue; +import com.sap.sailing.domain.tractracadapter.impl.TracTracRaceTrackerImpl; + public interface LoadingQueueDoneCallBack { + /** + * Callback method + * See {@link AbstractReceiverWithQueue} and {@link TracTracRaceTrackerImpl} for example usage. + * + * @param receiver Receiver that handled the marked event. + */ void loadingQueueDone(Receiver receiver); } diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java index 5fa47b124ef..a23cb51888c 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java @@ -1,5 +1,9 @@ package com.sap.sailing.domain.tractracadapter; +import com.sap.sailing.domain.tracking.TrackedRace; +import com.sap.sailing.domain.tracking.TrackedRaceStatus; +import com.sap.sailing.domain.tractracadapter.impl.TracTracRaceTrackerImpl; + public interface Receiver { /** @@ -35,5 +39,16 @@ public interface Receiver { */ void join(long timeoutInMilliseconds) throws InterruptedException; + /** + * Allows to "mark" the currently last event in the queue and provides callback as soon as this event has been handled. + * + * This is used in {@link TracTracRaceTrackerImpl} to ensure that queued events during loading phase will be processed + * before the new {@link TrackedRaceStatus} is propagated to the {@link TrackedRace}. + * + * @param callback + * {@link LoadingQueueDoneCallBack#loadingQueueDone(Receiver)} Will be called as soon as the event that + * was last in the queue at the time of calling + * {@link #callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack)} has been handled. + */ void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback); } From 0d360ace6bb8678bee883a98812ef110388c72fb Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 15:05:09 +0100 Subject: [PATCH 05/12] More comments concerning new tractrac receiver callback --- .../tractracadapter/LoadingQueueDoneCallBack.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java index 60b5df94a02..502d6e9f955 100644 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/LoadingQueueDoneCallBack.java @@ -1,8 +1,20 @@ package com.sap.sailing.domain.tractracadapter; +import com.sap.sailing.domain.tracking.TrackedRace; import com.sap.sailing.domain.tractracadapter.impl.AbstractReceiverWithQueue; import com.sap.sailing.domain.tractracadapter.impl.TracTracRaceTrackerImpl; +/** + * Callback interface which is currently used, so that {@link Receiver}s can be told to callback when current queue + * content is worked through. + * + * When loading -> tracking or loading -> finished would happen, we want to wait for all the current events in the queue + * to be handled before notifying the {@link TrackedRace}. This will ensure that caches, polar miner, etc. will not be + * resumed/started, as long as loading events are handled. + * + * @author Frederik Petersen + * + */ public interface LoadingQueueDoneCallBack { /** From 4f06a96d2c063c0e929518be2d96ce6e6fc34c46 Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 15:24:35 +0100 Subject: [PATCH 06/12] Refactored LoadingQueueDone callback for tractrac receivers --- .../AbstractLoadingQueueDoneCallBack.java | 49 +++++++++++++++ .../impl/TracTracRaceTrackerImpl.java | 63 +++++++------------ 2 files changed, 70 insertions(+), 42 deletions(-) create mode 100644 java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractLoadingQueueDoneCallBack.java diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractLoadingQueueDoneCallBack.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractLoadingQueueDoneCallBack.java new file mode 100644 index 00000000000..1417f405052 --- /dev/null +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractLoadingQueueDoneCallBack.java @@ -0,0 +1,49 @@ +package com.sap.sailing.domain.tractracadapter.impl; + +import java.util.Collection; +import java.util.HashSet; +import java.util.Set; + +import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; +import com.sap.sailing.domain.tractracadapter.Receiver; + +/** + * Performs {@link #executeWhenAllReceiversAreDoneLoading()}, when all Receivers have called back. + * + * @author Frederik Petersen + * + */ +public abstract class AbstractLoadingQueueDoneCallBack implements LoadingQueueDoneCallBack { + + /** + * Set keeping track of all receivers that still need to callback. When it's empty, desired + * action should be performed. + */ + private final Set receiversToCallback; + + /** + * + * @param receivers + * All of these receivers will be queried to call back when they are done handling currently queued + * events + */ + public AbstractLoadingQueueDoneCallBack(Collection receivers) { + this.receiversToCallback = new HashSet<>(receivers); + for (Receiver receiver : receivers) { + receiver.callBackWhenLoadingQueueIsDone(this); + } + } + + @Override + public void loadingQueueDone(Receiver receiver) { + synchronized (this.receiversToCallback) { + receiversToCallback.remove(receiver); + if (this.receiversToCallback.isEmpty()) { + executeWhenAllReceiversAreDoneLoading(); + } + } + } + + protected abstract void executeWhenAllReceiversAreDoneLoading(); + +} 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 44317ae7b03..3120ad3a6f7 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 @@ -43,7 +43,6 @@ import com.sap.sailing.domain.tracking.WindTrack; import com.sap.sailing.domain.tracking.impl.EmptyWindStore; import com.sap.sailing.domain.tracking.impl.TrackedRaceStatusImpl; import com.sap.sailing.domain.tractracadapter.DomainFactory; -import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; import com.sap.sailing.domain.tractracadapter.TracTracConnectionConstants; import com.sap.sailing.domain.tractracadapter.TracTracRaceTracker; @@ -579,24 +578,13 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements } logger.info("Stored data progress in tracker "+getID()+" for race(s) "+getRaces()+": "+progress); if (progress==1.0) { - LoadingQueueDoneCallBack callBackHandler = new LoadingQueueDoneCallBack() { - - private Set receivers = new HashSet<>(TracTracRaceTrackerImpl.this.receivers); - + new AbstractLoadingQueueDoneCallBack(receivers) { @Override - public void loadingQueueDone(Receiver receiver) { - synchronized (this.receivers) { - receivers.remove(receiver); - if (this.receivers.isEmpty()) { - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.TRACKING, progress); - updateStatusOfTrackedRaces(); - } - } + protected void executeWhenAllReceiversAreDoneLoading() { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.TRACKING, progress); + updateStatusOfTrackedRaces(); } }; - for (Receiver receiver : receivers) { - receiver.callBackWhenLoadingQueueIsDone(callBackHandler); - } } lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, progress); lastProgressPerID.put(getID(), new Util.Pair(counter, progress)); @@ -641,38 +629,29 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements @Override public void stopped(Object o) { - logger.info("stopped TracTrac tracking in tracker "+getID()+" for "+getRaces()+" while in status "+lastStatus); - LoadingQueueDoneCallBack callBackHandler = new LoadingQueueDoneCallBack() { - - private Set receivers = new HashSet<>(TracTracRaceTrackerImpl.this.receivers); - + logger.info("stopped TracTrac tracking in tracker " + getID() + " for " + getRaces() + " while in status " + + lastStatus); + new AbstractLoadingQueueDoneCallBack(receivers) { @Override - public void loadingQueueDone(Receiver receiver) { - synchronized (this.receivers) { - receivers.remove(receiver); - if (this.receivers.isEmpty()) { - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, 1.0); - updateStatusOfTrackedRaces(); - if (!stopped) { - try { - for (RaceDefinition race : getRaces()) { - // See also bug 1517; with TracAPI we assume that when stopped(IEvent) is called by the TracAPI then - // all subscriptions have received all their data and it's therefore safe to stop all subscriptions - // at this point without missing any data. - trackedRegattaRegistry.stopTracking(regatta, race); - } - } catch (InterruptedException | IOException e) { - logger.log(Level.INFO, "Interrupted while trying to stop tracker "+this, e); - } + protected void executeWhenAllReceiversAreDoneLoading() { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, 1.0); + updateStatusOfTrackedRaces(); + if (!stopped) { + try { + for (RaceDefinition race : getRaces()) { + // See also bug 1517; with TracAPI we assume that when stopped(IEvent) is called by the + // TracAPI then + // all subscriptions have received all their data and it's therefore safe to stop all + // subscriptions + // at this point without missing any data. + trackedRegattaRegistry.stopTracking(regatta, race); } + } catch (InterruptedException | IOException e) { + logger.log(Level.INFO, "Interrupted while trying to stop tracker " + this, e); } } } }; - for (Receiver receiver : receivers) { - receiver.callBackWhenLoadingQueueIsDone(callBackHandler); - } - } @Override From 83fa47585ef9fdea6d58f741618bd2ee0dc5bc5f Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 15:45:02 +0100 Subject: [PATCH 07/12] Improved thread safety of new receiver callback There was a minimal chance that queues isEmpty returned false and until the getLast call emptied leading to a nullpointer. This was solved using peekLast instead of isEmpty and getLast --- .../impl/AbstractReceiverWithQueue.java | 23 ++++++++++++------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java index 29edd8e5f99..c7c8ac2f0bc 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java @@ -1,6 +1,7 @@ package com.sap.sailing.domain.tractracadapter.impl; -import java.util.concurrent.ConcurrentHashMap; +import java.util.HashMap; +import java.util.Map; import java.util.concurrent.LinkedBlockingDeque; import java.util.concurrent.TimeUnit; import java.util.logging.Logger; @@ -39,7 +40,7 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re private final DynamicTrackedRegatta trackedRegatta; private final Simulator simulator; private final Thread thread; - private final ConcurrentHashMap, LoadingQueueDoneCallBack> loadingQueueDoneCallBacks; + private final Map, LoadingQueueDoneCallBack> loadingQueueDoneCallBacks; /** * used by {@link #stopAfterNotReceivingEventsForSomeTime(long)} and {@link #run()} to check if an event was received @@ -58,7 +59,7 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re this.simulator = simulator; this.queue = new LinkedBlockingDeque>(); this.thread = new Thread(this, getClass().getName()); - this.loadingQueueDoneCallBacks = new ConcurrentHashMap<>(); + this.loadingQueueDoneCallBacks = new HashMap<>(); } protected IEventSubscriber getEventSubscriber() { @@ -140,10 +141,14 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re if (!isStopEvent(event)) { handleEvent(event); } - LoadingQueueDoneCallBack callBack = loadingQueueDoneCallBacks.remove(event); + LoadingQueueDoneCallBack callBack; + synchronized (loadingQueueDoneCallBacks) { + callBack = loadingQueueDoneCallBacks.remove(event); + } if (callBack != null) { callBack.loadingQueueDone(this); } + } catch (InterruptedException e) { e.printStackTrace(); } @@ -191,10 +196,12 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re @Override public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { - if (queue.isEmpty()) { - callback.loadingQueueDone(this); - } else { - loadingQueueDoneCallBacks.put(queue.getLast(), callback); + synchronized (loadingQueueDoneCallBacks) { + if (queue.isEmpty()) { + callback.loadingQueueDone(this); + } else { + loadingQueueDoneCallBacks.put(queue.getLast(), callback); + } } } } From 0b8048f8aa25cc3c5089fc2808f5ec59a735d33b Mon Sep 17 00:00:00 2001 From: Frederik Petersen Date: Wed, 25 Mar 2015 16:21:08 +0100 Subject: [PATCH 08/12] Thread safety improvements which are explained in last commit just forgot to add some changes --- .../tractracadapter/impl/AbstractReceiverWithQueue.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java index c7c8ac2f0bc..0c6543b53b4 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java @@ -16,6 +16,7 @@ import com.sap.sailing.domain.tractracadapter.DomainFactory; import com.sap.sailing.domain.tractracadapter.LoadingQueueDoneCallBack; import com.sap.sailing.domain.tractracadapter.Receiver; import com.sap.sse.common.Util; +import com.sap.sse.common.Util.Triple; import com.tractrac.model.lib.api.event.IEvent; import com.tractrac.model.lib.api.event.IRace; import com.tractrac.subscription.lib.api.IEventSubscriber; @@ -148,7 +149,7 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re if (callBack != null) { callBack.loadingQueueDone(this); } - + } catch (InterruptedException e) { e.printStackTrace(); } @@ -197,10 +198,11 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re @Override public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { synchronized (loadingQueueDoneCallBacks) { - if (queue.isEmpty()) { + Triple lastInQueue = queue.peekLast(); + if (lastInQueue == null) { callback.loadingQueueDone(this); } else { - loadingQueueDoneCallBacks.put(queue.getLast(), callback); + loadingQueueDoneCallBacks.put(lastInQueue, callback); } } } From 878b90eb89f8e7548f1a1c475aeb8bf153a1f3c9 Mon Sep 17 00:00:00 2001 From: Axel Uhl Date: Thu, 26 Mar 2015 16:51:18 +0100 Subject: [PATCH 09/12] perform check for completeness and status change from LOADING to TRACKING *after* setting the default LOADING state --- .../com/sap/sailing/domain/tractracadapter/Receiver.java | 8 +++++--- .../tractracadapter/impl/TracTracRaceTrackerImpl.java | 2 +- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java index a23cb51888c..c7861e5819f 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/Receiver.java @@ -40,10 +40,12 @@ public interface Receiver { void join(long timeoutInMilliseconds) throws InterruptedException; /** - * Allows to "mark" the currently last event in the queue and provides callback as soon as this event has been handled. + * Allows to "mark" the currently last event in the queue and provides callback as soon as this event has been + * handled. + *

* - * This is used in {@link TracTracRaceTrackerImpl} to ensure that queued events during loading phase will be processed - * before the new {@link TrackedRaceStatus} is propagated to the {@link TrackedRace}. + * This is used in {@link TracTracRaceTrackerImpl} to ensure that events queued during loading phase will be + * processed before the new {@link TrackedRaceStatus} is propagated to the {@link TrackedRace}. * * @param callback * {@link LoadingQueueDoneCallBack#loadingQueueDone(Receiver)} Will be called as soon as the event that 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 3120ad3a6f7..b6097d6e1a7 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 @@ -577,6 +577,7 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements } } logger.info("Stored data progress in tracker "+getID()+" for race(s) "+getRaces()+": "+progress); + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, progress); if (progress==1.0) { new AbstractLoadingQueueDoneCallBack(receivers) { @Override @@ -586,7 +587,6 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements } }; } - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, progress); lastProgressPerID.put(getID(), new Util.Pair(counter, progress)); updateStatusOfTrackedRaces(); } From 1d4d9640a0e0d853d5c13e31fa3c6262dbec1d1e Mon Sep 17 00:00:00 2001 From: Axel Uhl Date: Thu, 26 Mar 2015 16:54:11 +0100 Subject: [PATCH 10/12] also wait for queues to drain in TracTrac connector in non-preemptive stop(...) call --- .../impl/TracTracRaceTrackerImpl.java | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) 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 b6097d6e1a7..31258110815 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 @@ -510,8 +510,20 @@ public class TracTracRaceTrackerImpl extends AbstractRaceTrackerImpl implements receiver.stopAfterProcessingQueuedEvents(); } } - lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, /* will be ignored */1.0); - updateStatusOfTrackedRaces(); + if (!stopReceiversPreemtively) { + // wait for their queues to be worked down before signalling the FINISHED state. + new AbstractLoadingQueueDoneCallBack(receivers) { + @Override + protected void executeWhenAllReceiversAreDoneLoading() { + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, /* will be ignored */1.0); + updateStatusOfTrackedRaces(); + } + }; + } else { + // queues contents were cleared preemptively; this means we're done with loading immediately + lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.FINISHED, /* will be ignored */1.0); + updateStatusOfTrackedRaces(); + } } } From 9570a5d3f03701dd79ba9de5076c610278d248e3 Mon Sep 17 00:00:00 2001 From: Axel Uhl Date: Thu, 26 Mar 2015 17:41:47 +0100 Subject: [PATCH 11/12] allow multiple loading completion listeners to subscribe at the same last event in queue --- .../impl/AbstractReceiverWithQueue.java | 24 +++++++++++++------ 1 file changed, 17 insertions(+), 7 deletions(-) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java index 0c6543b53b4..094305259fd 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java @@ -1,7 +1,9 @@ package com.sap.sailing.domain.tractracadapter.impl; import java.util.HashMap; +import java.util.HashSet; import java.util.Map; +import java.util.Set; import java.util.concurrent.LinkedBlockingDeque; import java.util.concurrent.TimeUnit; import java.util.logging.Logger; @@ -41,7 +43,7 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re private final DynamicTrackedRegatta trackedRegatta; private final Simulator simulator; private final Thread thread; - private final Map, LoadingQueueDoneCallBack> loadingQueueDoneCallBacks; + private final Map, Set> loadingQueueDoneCallBacks; /** * used by {@link #stopAfterNotReceivingEventsForSomeTime(long)} and {@link #run()} to check if an event was received @@ -142,12 +144,14 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re if (!isStopEvent(event)) { handleEvent(event); } - LoadingQueueDoneCallBack callBack; + Set callBacks; synchronized (loadingQueueDoneCallBacks) { - callBack = loadingQueueDoneCallBacks.remove(event); + callBacks = loadingQueueDoneCallBacks.remove(event); } - if (callBack != null) { - callBack.loadingQueueDone(this); + if (callBacks != null) { + for (LoadingQueueDoneCallBack callback : callBacks) { + callback.loadingQueueDone(this); + } } } catch (InterruptedException e) { @@ -199,10 +203,16 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) { synchronized (loadingQueueDoneCallBacks) { Triple lastInQueue = queue.peekLast(); - if (lastInQueue == null) { + // when simulator is attached, consider loading already done; the simulator simulates "live" tracking + if (lastInQueue == null || getSimulator() != null) { callback.loadingQueueDone(this); } else { - loadingQueueDoneCallBacks.put(lastInQueue, callback); + Set set = loadingQueueDoneCallBacks.get(lastInQueue); + if (set == null) { + set = new HashSet<>(); + loadingQueueDoneCallBacks.put(lastInQueue, set); + } + set.add(callback); } } } From 26ad184a9f8648b861e2442ece6f3674fb6800f9 Mon Sep 17 00:00:00 2001 From: Axel Uhl Date: Thu, 26 Mar 2015 17:50:19 +0100 Subject: [PATCH 12/12] fixing simulator issue with loading finished callbacks; loading finished callbacks are now notified immediately when simulator is running --- .../impl/AbstractReceiverWithQueue.java | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java index 094305259fd..5a5ed361b94 100755 --- a/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java +++ b/java/com.sap.sailing.domain.tractracadapter/src/com/sap/sailing/domain/tractracadapter/impl/AbstractReceiverWithQueue.java @@ -144,9 +144,21 @@ public abstract class AbstractReceiverWithQueue implements Runnable, Re if (!isStopEvent(event)) { handleEvent(event); } - Set callBacks; + final Set callBacks; synchronized (loadingQueueDoneCallBacks) { - callBacks = loadingQueueDoneCallBacks.remove(event); + if (getSimulator() != null) { + // when simulator is running, loading is considered finished and all callbacks will + // be satisfied instantly + callBacks = new HashSet<>(); + for (Set set : loadingQueueDoneCallBacks.values()) { + callBacks.addAll(set); + } + loadingQueueDoneCallBacks.clear(); + } else { + // otherwise, check only if there are callbacks that registered at the event + // currently consumed and notify if any are found + callBacks = loadingQueueDoneCallBacks.remove(event); + } } if (callBacks != null) { for (LoadingQueueDoneCallBack callback : callBacks) {