mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-10 22:30:56 +00:00
Merge branch 'master' into mark-rounding-inference
Conflicts: java/com.sap.sailing.domain/src/com/sap/sailing/domain/markpassingcalculation/MarkPassingUpdateListener.java java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/DynamicTrackedRaceImpl.java
This commit is contained in:
commit
ba47883a59
17 files changed
+273
-69
No files matched your search
+5
@@ -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) {
|
||||
|
||||
+5
@@ -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<Receiver> receivers = new ArrayList<Receiver>();
|
||||
receivers.add(receiver);
|
||||
|
||||
+5
@@ -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) {
|
||||
|
||||
+28
@@ -0,0 +1,28 @@
|
||||
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 {
|
||||
|
||||
/**
|
||||
* Callback method
|
||||
* See {@link AbstractReceiverWithQueue} and {@link TracTracRaceTrackerImpl} for example usage.
|
||||
*
|
||||
* @param receiver Receiver that handled the marked event.
|
||||
*/
|
||||
void loadingQueueDone(Receiver receiver);
|
||||
|
||||
}
|
||||
+19
@@ -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 {
|
||||
/**
|
||||
@@ -34,4 +38,19 @@ public interface Receiver {
|
||||
* Waits until this received has stopped, but no longer than <code>timeout</code> milliseconds
|
||||
*/
|
||||
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.
|
||||
* <p>
|
||||
*
|
||||
* 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
|
||||
* was last in the queue at the time of calling
|
||||
* {@link #callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack)} has been handled.
|
||||
*/
|
||||
void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback);
|
||||
}
|
||||
+49
@@ -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<Receiver> receiversToCallback;
|
||||
|
||||
/**
|
||||
*
|
||||
* @param receivers
|
||||
* All of these receivers will be queried to call back when they are done handling currently queued
|
||||
* events
|
||||
*/
|
||||
public AbstractLoadingQueueDoneCallBack(Collection<Receiver> 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();
|
||||
|
||||
}
|
||||
+52
-4
@@ -1,6 +1,10 @@
|
||||
package com.sap.sailing.domain.tractracadapter.impl;
|
||||
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
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;
|
||||
|
||||
@@ -11,9 +15,11 @@ 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.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;
|
||||
import com.tractrac.subscription.lib.api.IRaceSubscriber;
|
||||
@@ -29,7 +35,7 @@ import com.tractrac.subscription.lib.api.IRaceSubscriber;
|
||||
public abstract class AbstractReceiverWithQueue<A, B, C> implements Runnable, Receiver {
|
||||
private static Logger logger = Logger.getLogger(AbstractReceiverWithQueue.class.getName());
|
||||
|
||||
private final LinkedBlockingQueue<Util.Triple<A, B, C>> queue;
|
||||
private final LinkedBlockingDeque<Util.Triple<A, B, C>> queue;
|
||||
private final DomainFactory domainFactory;
|
||||
private final IEvent tractracEvent;
|
||||
private final IEventSubscriber eventSubscriber;
|
||||
@@ -37,6 +43,7 @@ public abstract class AbstractReceiverWithQueue<A, B, C> implements Runnable, Re
|
||||
private final DynamicTrackedRegatta trackedRegatta;
|
||||
private final Simulator simulator;
|
||||
private final Thread thread;
|
||||
private final Map<Util.Triple<A, B, C>, Set<LoadingQueueDoneCallBack>> loadingQueueDoneCallBacks;
|
||||
|
||||
/**
|
||||
* used by {@link #stopAfterNotReceivingEventsForSomeTime(long)} and {@link #run()} to check if an event was received
|
||||
@@ -53,8 +60,9 @@ public abstract class AbstractReceiverWithQueue<A, B, C> implements Runnable, Re
|
||||
this.trackedRegatta = trackedRegatta;
|
||||
this.domainFactory = domainFactory;
|
||||
this.simulator = simulator;
|
||||
this.queue = new LinkedBlockingQueue<Util.Triple<A, B, C>>();
|
||||
this.queue = new LinkedBlockingDeque<Util.Triple<A, B, C>>();
|
||||
this.thread = new Thread(this, getClass().getName());
|
||||
this.loadingQueueDoneCallBacks = new HashMap<>();
|
||||
}
|
||||
|
||||
protected IEventSubscriber getEventSubscriber() {
|
||||
@@ -136,6 +144,28 @@ public abstract class AbstractReceiverWithQueue<A, B, C> implements Runnable, Re
|
||||
if (!isStopEvent(event)) {
|
||||
handleEvent(event);
|
||||
}
|
||||
final Set<LoadingQueueDoneCallBack> callBacks;
|
||||
synchronized (loadingQueueDoneCallBacks) {
|
||||
if (getSimulator() != null) {
|
||||
// when simulator is running, loading is considered finished and all callbacks will
|
||||
// be satisfied instantly
|
||||
callBacks = new HashSet<>();
|
||||
for (Set<LoadingQueueDoneCallBack> 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) {
|
||||
callback.loadingQueueDone(this);
|
||||
}
|
||||
}
|
||||
|
||||
} catch (InterruptedException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
@@ -180,4 +210,22 @@ public abstract class AbstractReceiverWithQueue<A, B, C> implements Runnable, Re
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void callBackWhenLoadingQueueIsDone(LoadingQueueDoneCallBack callback) {
|
||||
synchronized (loadingQueueDoneCallBacks) {
|
||||
Triple<A, B, C> lastInQueue = queue.peekLast();
|
||||
// when simulator is attached, consider loading already done; the simulator simulates "live" tracking
|
||||
if (lastInQueue == null || getSimulator() != null) {
|
||||
callback.loadingQueueDone(this);
|
||||
} else {
|
||||
Set<LoadingQueueDoneCallBack> set = loadingQueueDoneCallBacks.get(lastInQueue);
|
||||
if (set == null) {
|
||||
set = new HashSet<>();
|
||||
loadingQueueDoneCallBacks.put(lastInQueue, set);
|
||||
}
|
||||
set.add(callback);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+45
-16
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -577,7 +589,16 @@ 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);
|
||||
lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.LOADING, progress);
|
||||
if (progress==1.0) {
|
||||
new AbstractLoadingQueueDoneCallBack(receivers) {
|
||||
@Override
|
||||
protected void executeWhenAllReceiversAreDoneLoading() {
|
||||
lastStatus = new TrackedRaceStatusImpl(TrackedRaceStatusEnum.TRACKING, progress);
|
||||
updateStatusOfTrackedRaces();
|
||||
}
|
||||
};
|
||||
}
|
||||
lastProgressPerID.put(getID(), new Util.Pair<Integer, Float>(counter, progress));
|
||||
updateStatusOfTrackedRaces();
|
||||
}
|
||||
@@ -620,21 +641,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);
|
||||
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);
|
||||
logger.info("stopped TracTrac tracking in tracker " + getID() + " for " + getRaces() + " while in status "
|
||||
+ lastStatus);
|
||||
new AbstractLoadingQueueDoneCallBack(receivers) {
|
||||
@Override
|
||||
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);
|
||||
}
|
||||
}
|
||||
} catch (InterruptedException | IOException e) {
|
||||
logger.log(Level.INFO, "Interrupted while trying to stop tracker "+this, e);
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+1
-1
@@ -50,5 +50,5 @@ public interface RaceChangeListener extends CourseListener {
|
||||
|
||||
void windSourcesToExcludeChanged(Iterable<? extends WindSource> windSourcesToExclude);
|
||||
|
||||
void statusChanged(TrackedRaceStatus newStatus);
|
||||
void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus);
|
||||
}
|
||||
+1
-1
@@ -40,7 +40,7 @@ public abstract class AbstractRaceChangeListener implements RaceChangeListener {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void statusChanged(TrackedRaceStatus newStatus) {
|
||||
public void statusChanged(TrackedRaceStatus newStatus, TrackedRaceStatus oldStatus) {
|
||||
defaultAction();
|
||||
}
|
||||
|
||||
|
||||
+4
-3
@@ -150,8 +150,9 @@ DynamicTrackedRace, GPSTrackListener<Competitor, GPSFixMoving> {
|
||||
|
||||
@Override
|
||||
public void setStatus(TrackedRaceStatus newStatus) {
|
||||
TrackedRaceStatus oldStatus = getStatus();
|
||||
super.setStatus(newStatus);
|
||||
notifyListeners(newStatus);
|
||||
notifyListeners(newStatus, oldStatus);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -418,8 +419,8 @@ DynamicTrackedRace, GPSTrackListener<Competitor, GPSFixMoving> {
|
||||
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) {
|
||||
|
||||
+1
-1
@@ -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.
|
||||
|
||||
+1
-1
@@ -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) {
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
+1
-1
@@ -148,7 +148,7 @@ public class PolarDataServiceImpl implements PolarDataService {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void raceFinishedTracking(TrackedRace race) {
|
||||
public void raceFinishedLoading(TrackedRace race) {
|
||||
polarDataMiner.raceFinishedTracking(race);
|
||||
}
|
||||
}
|
||||
@@ -43,7 +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.common.impl.MillisecondsTimePoint;
|
||||
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;
|
||||
@@ -58,7 +58,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<TrackedRace, Set<GPSFixMovingWithOriginInfo>> fixesForReplayRacesWhichAreStillLoading = new HashMap<>();
|
||||
private final Map<TrackedRace, Set<GPSFixMovingWithOriginInfo>> fixesForRacesWhichAreStillLoading = new HashMap<>();
|
||||
|
||||
private final Queue<GPSFixMovingWithOriginInfo> fixQueue = new ConcurrentLinkedQueue<GPSFixMovingWithOriginInfo>();
|
||||
|
||||
@@ -79,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<GPSFixMovingWithOriginInfo, GPSFixMovingWithPolarContext> enrichingProcessor;
|
||||
|
||||
|
||||
/**
|
||||
* This processor keeps track of the moving average of the speed values and the average angle for each course (legtype tack combination)
|
||||
@@ -100,6 +100,7 @@ public class PolarDataMiner {
|
||||
private CubicRegressionPerCourseProcessor cubicRegressionPerCourseProcessor;
|
||||
|
||||
private SpeedRegressionPerAngleClusterProcessor speedRegressionPerAngleClusterProcessor;
|
||||
private ParallelFilteringProcessor<GPSFixMovingWithOriginInfo> preFilteringProcessor;
|
||||
|
||||
public PolarDataMiner() {
|
||||
this(PolarSheetGenerationSettingsImpl.createBackendPolarSettings());
|
||||
@@ -170,15 +171,42 @@ public class PolarDataMiner {
|
||||
.asList(filteringProcessor);
|
||||
|
||||
|
||||
enrichingProcessor = new AbstractEnrichingProcessor<GPSFixMovingWithOriginInfo, GPSFixMovingWithPolarContext>(
|
||||
GPSFixMovingWithOriginInfo.class, GPSFixMovingWithPolarContext.class, executor, enrichingResultReceivers) {
|
||||
AbstractEnrichingProcessor<GPSFixMovingWithOriginInfo, GPSFixMovingWithPolarContext> enrichingProcessor = new AbstractEnrichingProcessor<GPSFixMovingWithOriginInfo, GPSFixMovingWithPolarContext>(
|
||||
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<Processor<GPSFixMovingWithOriginInfo, ?>> preFilterResultReceivers = Arrays
|
||||
.asList(enrichingProcessor);
|
||||
|
||||
preFilteringProcessor = new ParallelFilteringProcessor<GPSFixMovingWithOriginInfo>(
|
||||
GPSFixMovingWithOriginInfo.class, executor, preFilterResultReceivers, new FilterCriterion<GPSFixMovingWithOriginInfo>() {
|
||||
|
||||
@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<GPSFixMovingWithOriginInfo> getElementType() {
|
||||
return GPSFixMovingWithOriginInfo.class;
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -188,49 +216,35 @@ 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<GPSFixMovingWithOriginInfo> fixes = fixesForReplayRacesWhichAreStillLoading.get(trackedRace);
|
||||
synchronized (fixesForRacesWhichAreStillLoading) {
|
||||
Set<GPSFixMovingWithOriginInfo> 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);
|
||||
}
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
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 +413,11 @@ public class PolarDataMiner {
|
||||
|
||||
public void raceFinishedTracking(TrackedRace race) {
|
||||
Set<GPSFixMovingWithOriginInfo> 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);
|
||||
|
||||
+5
-4
@@ -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));
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user